Class: Kaal::DelayedJob::DatabaseEngine

Inherits:
Registry
  • Object
show all
Defined in:
lib/kaal/delayed_job/database_engine.rb,
sig/kaal/delayed_job/database_engine.rbs

Overview

Sequel-backed delayed-job store persisted in kaal_delayed_jobs.

Class Method Summary collapse

Instance Method Summary collapse

Methods inherited from Registry

#requires_dispatch_lock?

Constructor Details

#initialize(database:, use_skip_locked: false) ⇒ DatabaseEngine

Returns a new instance of DatabaseEngine.

Parameters:

  • database: (Kaal::rbs_any)
  • use_skip_locked: (Boolean) (defaults to: false)


15
16
17
18
19
# File 'lib/kaal/delayed_job/database_engine.rb', line 15

def initialize(database:, use_skip_locked: false)
  super()
  @database = Kaal::Persistence::Database.new(database)
  @use_skip_locked = use_skip_locked
end

Class Method Details

.normalize_row(row) ⇒ Kaal::rbs_any

Parameters:

  • row (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


56
57
58
59
60
61
62
63
64
65
66
67
68
69
# File 'lib/kaal/delayed_job/database_engine.rb', line 56

def self.normalize_row(row)
  return nil unless row

  {
    job_id: row[:job_id],
    run_at: row[:run_at],
    job_class: row[:job_class],
    args: parse_args(row[:args]),
    queue: row[:queue],
    created_at: row[:created_at]
  }
rescue JSON::ParserError
  nil
end

Instance Method Details

#all_jobs ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


48
49
50
# File 'lib/kaal/delayed_job/database_engine.rb', line 48

def all_jobs
  connection[:kaal_delayed_jobs].order(:run_at, :job_id).filter_map { |row| self.class.normalize_row(row) }
end

#claim_strategy ⇒ :skip_locked, :delete_confirmation

Returns:

  • (:skip_locked, :delete_confirmation)


52
53
54
# File 'lib/kaal/delayed_job/database_engine.rb', line 52

def claim_strategy
  @use_skip_locked ? :skip_locked : :delete_confirmation
end

#connection ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


111
112
113
# File 'lib/kaal/delayed_job/database_engine.rb', line 111

def connection
  @database.connection
end

#dataset ⇒ Kaal::rbs_any

Returns:

  • (Kaal::rbs_any)


107
108
109
# File 'lib/kaal/delayed_job/database_engine.rb', line 107

def dataset
  @database.delayed_jobs_dataset
end

#dataset_for(connection) ⇒ Kaal::rbs_any

Parameters:

  • connection (Kaal::rbs_any)

Returns:

  • (Kaal::rbs_any)


101
102
103
104
105
# File 'lib/kaal/delayed_job/database_engine.rb', line 101

def dataset_for(connection)
  return dataset unless connection

  Kaal::Persistence::Database.new(connection).delayed_jobs_dataset
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)


21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
# File 'lib/kaal/delayed_job/database_engine.rb', line 21

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

  dataset_for(connection).insert(payload)
  self.class.normalize_row(payload)
rescue ::Sequel::UniqueConstraintViolation
  raise DuplicateJobError, "Delayed job #{job_id.inspect} already exists"
end

#find_job(job_id, connection: @database.connection) ⇒ Kaal::rbs_any

Parameters:

  • job_id (Kaal::rbs_any)
  • connection: (Kaal::rbs_any) (defaults to: @database.connection)

Returns:

  • (Kaal::rbs_any)


44
45
46
# File 'lib/kaal/delayed_job/database_engine.rb', line 44

def find_job(job_id, connection: @database.connection)
  self.class.normalize_row(connection[:kaal_delayed_jobs].where(job_id: job_id).first)
end

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

Parameters:

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

Returns:

  • (Kaal::rbs_any)


38
39
40
41
42
# File 'lib/kaal/delayed_job/database_engine.rb', line 38

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)


84
85
86
87
88
89
90
91
92
93
94
# File 'lib/kaal/delayed_job/database_engine.rb', line 84

def pop_due_with_delete_confirmation(now:, limit:)
  connection.transaction do
    delayed_jobs_dataset = connection[:kaal_delayed_jobs]
    due_rows = delayed_jobs_dataset.where { run_at <= now }.order(:run_at, :job_id).limit(limit).all
    due_rows.each_with_object([]) do |row, jobs|
      deleted = delayed_jobs_dataset.where(job_id: row[:job_id]).delete
      normalized_job = self.class.normalize_row(row)
      jobs << normalized_job if deleted.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)


73
74
75
76
77
78
79
80
81
82
# File 'lib/kaal/delayed_job/database_engine.rb', line 73

def pop_due_with_skip_locked(now:, limit:)
  connection.transaction do
    delayed_jobs_dataset = connection[:kaal_delayed_jobs]
    due_rows = delayed_jobs_dataset.where { run_at <= now }.order(:run_at, :job_id).for_update.skip_locked.limit(limit).all
    job_ids = due_rows.map { |row| row[:job_id] }
    normalized_jobs = due_rows.filter_map { |row| self.class.normalize_row(row) }
    delayed_jobs_dataset.where(job_id: job_ids).delete unless job_ids.empty?
    normalized_jobs
  end
end