Class: SolidQueue::RecurringTask

Inherits:
Record
  • Object
show all
Defined in:
app/models/solid_queue/recurring_task.rb

Defined Under Namespace

Classes: Arguments

Class Method Summary collapse

Instance Method Summary collapse

Methods inherited from Record

non_blocking_lock, supports_insert_conflict_target?, use_index

Class Method Details

.create_dynamic_task(key, **options) ⇒ Object



39
40
41
# File 'app/models/solid_queue/recurring_task.rb', line 39

def create_dynamic_task(key, **options)
  from_configuration(key, **options.merge(static: false)).save!
end

.create_or_update_all(tasks) ⇒ Object



47
48
49
50
51
52
53
54
55
56
# File 'app/models/solid_queue/recurring_task.rb', line 47

def create_or_update_all(tasks)
  if supports_insert_conflict_target?
    # PostgreSQL fails and aborts the current transaction when it hits a duplicate key conflict
    # during two concurrent INSERTs for the same value of an unique index. We need to explicitly
    # indicate unique_by to ignore duplicate rows by this value when inserting
    upsert_all tasks.map(&:attributes_for_upsert), unique_by: :key
  else
    upsert_all tasks.map(&:attributes_for_upsert)
  end
end

.delete_dynamic_task(key) ⇒ Object



43
44
45
# File 'app/models/solid_queue/recurring_task.rb', line 43

def delete_dynamic_task(key)
  RecurringTask.dynamic.find_by!(key: key).destroy
end

.from_configuration(key, **options) ⇒ Object



26
27
28
29
30
31
32
33
34
35
36
37
# File 'app/models/solid_queue/recurring_task.rb', line 26

def from_configuration(key, **options)
  new \
    key: key,
    class_name: options[:class],
    command: options[:command],
    arguments: options[:args],
    schedule: options[:schedule],
    queue_name: options[:queue].presence,
    priority: options[:priority].presence,
    description: options[:description],
    static: options.fetch(:static, true)
end

.wrap(args) ⇒ Object



22
23
24
# File 'app/models/solid_queue/recurring_task.rb', line 22

def wrap(args)
  args.is_a?(self) ? args : from_configuration(args.first, **args.second)
end

Instance Method Details

#attributes_for_upsertObject



112
113
114
# File 'app/models/solid_queue/recurring_task.rb', line 112

def attributes_for_upsert
  attributes.without("id", "created_at", "updated_at")
end

#enqueue(at:) ⇒ Object



80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
# File 'app/models/solid_queue/recurring_task.rb', line 80

def enqueue(at:)
  SolidQueue.instrument(:enqueue_recurring_task, task: key, at: at) do |payload|
    active_job = if using_solid_queue_adapter?
      enqueue_and_record(run_at: at)
    else
      payload[:other_adapter] = true

      perform_later.tap do |job|
        unless job.successfully_enqueued?
          report_enqueue_error(job.enqueue_error, at: at)
          payload[:enqueue_error] = job.enqueue_error&.message
        end
      end
    end

    active_job.tap do |enqueued_job|
      payload[:active_job_id] = enqueued_job.job_id
    end
  rescue RecurringExecution::AlreadyRecorded
    payload[:skipped] = true
    false
  rescue Job::EnqueueError => error
    report_enqueue_error(error, at: at)
    payload[:enqueue_error] = error.message
    false
  end
end

#last_enqueued_timeObject



72
73
74
75
76
77
78
# File 'app/models/solid_queue/recurring_task.rb', line 72

def last_enqueued_time
  if recurring_executions.loaded?
    recurring_executions.map(&:run_at).max
  else
    recurring_executions.maximum(:run_at)
  end
end

#next_timeObject



64
65
66
# File 'app/models/solid_queue/recurring_task.rb', line 64

def next_time
  parsed_schedule_with_time_zone.next_time.utc
end

#next_time_after(time) ⇒ Object



60
61
62
# File 'app/models/solid_queue/recurring_task.rb', line 60

def next_time_after(time)
  parsed_schedule_with_time_zone.next_time(time).utc
end

#previous_timeObject



68
69
70
# File 'app/models/solid_queue/recurring_task.rb', line 68

def previous_time
  parsed_schedule_with_time_zone.previous_time.utc
end

#to_sObject



108
109
110
# File 'app/models/solid_queue/recurring_task.rb', line 108

def to_s
  "#{class_name}.perform_later(#{arguments.map(&:inspect).join(",")}) [ #{parsed_schedule.original} ]"
end