Class: Kaal::Internal::ActiveRecord::DelayedJobRegistry

Inherits:
DelayedJob::Registry show all
Defined in:
lib/kaal/internal/active_record/delayed_job_registry.rb,
sig/kaal/internal/active_record/delayed_job_registry.rbs

Overview

Active Record-backed store for delayed jobs.

Class Method Summary collapse

Instance Method Summary collapse

Methods inherited from DelayedJob::Registry

#requires_dispatch_lock?

Constructor Details

#initialize(connection: nil, model: DelayedJobRecord, use_skip_locked: false) ⇒ DelayedJobRegistry

Returns a new instance of DelayedJobRegistry.

Parameters:

  • connection: (Kaal::rbs_any, nil) (defaults to: nil)
  • model: (Kaal::rbs_any) (defaults to: DelayedJobRecord)
  • use_skip_locked: (Boolean) (defaults to: false)


15
16
17
18
19
20
# File 'lib/kaal/internal/active_record/delayed_job_registry.rb', line 15

def initialize(connection: nil, model: DelayedJobRecord, use_skip_locked: false)
  super()
  ConnectionSupport.configure!(connection)
  @model = model
  @use_skip_locked = use_skip_locked
end

Class Method Details

.normalize(record) ⇒ Kaal::rbs_any

Parameters:

  • record (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


84
85
86
87
88
89
90
91
92
93
94
95
96
97
# File 'lib/kaal/internal/active_record/delayed_job_registry.rb', line 84

def self.normalize(record)
  return nil unless record

  {
    job_id: record.job_id,
    run_at: record.run_at,
    job_class: record.job_class,
    args: parse_args(record.args),
    queue: record.queue,
    created_at: record.created_at
  }
rescue JSON::ParserError
  nil
end

Instance Method Details

#all_jobs ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


76
77
78
# File 'lib/kaal/internal/active_record/delayed_job_registry.rb', line 76

def all_jobs
  @model.order(:run_at, :job_id).filter_map { |record| self.class.normalize(record) }
end

#claim_strategy ⇒ :skip_locked, :delete_confirmation

Returns:

  • (:skip_locked, :delete_confirmation)


80
81
82
# File 'lib/kaal/internal/active_record/delayed_job_registry.rb', line 80

def claim_strategy
  @use_skip_locked ? :skip_locked : :delete_confirmation
end

#enqueue(job_id:, run_at:, job_class:, args:, queue: nil, connection: nil) ⇒ Kaal::rbs_any

Parameters:

  • job_id: (Kaal::rbs_any)
  • run_at: (Kaal::rbs_any)
  • job_class: (Kaal::rbs_any)
  • args: (Kaal::rbs_any)
  • queue: (Kaal::rbs_any, nil) (defaults to: nil)
  • connection: (Kaal::rbs_any, nil) (defaults to: nil)

Returns:

  • (Kaal::rbs_any)


22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
# File 'lib/kaal/internal/active_record/delayed_job_registry.rb', line 22

def enqueue(job_id:, run_at:, job_class:, args:, queue: nil, connection: nil)
  now = Time.now.utc
  attributes = {
    job_id: job_id,
    run_at: run_at,
    job_class: job_class,
    args: JSON.generate(args),
    queue: queue,
    created_at: now
  }

  if connection
    insert_with_connection(connection, attributes)
  else
    @model.create!(attributes)
  end

  self.class.normalize(@model.new(attributes))
rescue ::ActiveRecord::RecordNotUnique
  raise Kaal::DelayedJob::DuplicateJobError, "Delayed job #{job_id.inspect} already exists"
end

#find_job(job_id) ⇒ Kaal::rbs_any

Parameters:

  • job_id (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


72
73
74
# File 'lib/kaal/internal/active_record/delayed_job_registry.rb', line 72

def find_job(job_id)
  self.class.normalize(@model.find_by(job_id: job_id))
end

#insert_with_connection(connection, attributes) ⇒ Kaal::rbs_any

Parameters:

  • connection (Kaal::rbs_any)
  • attributes (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


101
102
103
104
105
106
107
108
109
110
# File 'lib/kaal/internal/active_record/delayed_job_registry.rb', line 101

def insert_with_connection(connection, attributes)
  table_name = @model.table_name
  columns = attributes.keys
  quoted_pairs = columns.map do |column|
    [connection.quote_column_name(column), connection.quote(attributes.fetch(column))]
  end
  quoted_columns = quoted_pairs.map(&:first).join(', ')
  quoted_values = quoted_pairs.map(&:last).join(', ')
  connection.execute("INSERT INTO #{connection.quote_table_name(table_name)} (#{quoted_columns}) VALUES (#{quoted_values})")
end

#pop_due(now:, limit:) ⇒ Kaal::rbs_any

Parameters:

  • now: (Kaal::rbs_any)
  • limit: (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


44
45
46
47
48
# File 'lib/kaal/internal/active_record/delayed_job_registry.rb', line 44

def pop_due(now:, limit:)
  return pop_due_with_skip_locked(now:, limit:) if @use_skip_locked

  pop_due_with_delete_confirmation(now:, limit:)
end

#pop_due_with_delete_confirmation(now:, limit:) ⇒ Kaal::rbs_any

Parameters:

  • now: (Kaal::rbs_any)
  • limit: (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


61
62
63
64
65
66
67
68
# File 'lib/kaal/internal/active_record/delayed_job_registry.rb', line 61

def pop_due_with_delete_confirmation(now:, limit:)
  @model.transaction do
    @model.where('run_at <= ?', now).order(:run_at, :job_id).limit(limit).each_with_object([]) do |record, jobs|
      normalized_job = self.class.normalize(record)
      jobs << normalized_job if @model.where(job_id: record.job_id).delete_all.positive? && normalized_job
    end
  end
end

#pop_due_with_skip_locked(now:, limit:) ⇒ Kaal::rbs_any

Parameters:

  • now: (Kaal::rbs_any)
  • limit: (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


52
53
54
55
56
57
58
59
# File 'lib/kaal/internal/active_record/delayed_job_registry.rb', line 52

def pop_due_with_skip_locked(now:, limit:)
  @model.transaction do
    due_records = @model.where('run_at <= ?', now).order(:run_at, :job_id).lock('FOR UPDATE SKIP LOCKED').limit(limit).to_a
    job_ids = due_records.map(&:job_id)
    @model.where(job_id: job_ids).delete_all unless job_ids.empty?
    due_records.filter_map { |record| self.class.normalize(record) }
  end
end