Class: Pgbus::ActiveJob::Adapter

Inherits:
Object
  • Object
show all
Defined in:
lib/pgbus/active_job/adapter.rb

Instance Method Summary collapse

Instance Method Details

#enqueue(active_job) ⇒ Object



8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
# File 'lib/pgbus/active_job/adapter.rb', line 8

def enqueue(active_job)
  queue = active_job.queue_name || Pgbus.configuration.default_queue
  payload_hash = Serializer.serialize_job_hash(active_job)
  payload_hash = Concurrency.(active_job, payload_hash)
  payload_hash = Uniqueness.(active_job, payload_hash)
  payload_hash = FairShare.(active_job, payload_hash)
  payload_hash = (payload_hash, active_job: active_job)

  if uniqueness_rejected?(active_job, payload_hash, queue: queue)
    uncount_batch_job(payload_hash)
    return active_job
  end

  enqueue_with_concurrency(active_job, queue, payload_hash)
end

#enqueue_all(active_jobs) ⇒ Object



41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
# File 'lib/pgbus/active_job/adapter.rb', line 41

def enqueue_all(active_jobs)
  # Jobs with uniqueness or concurrency must go through individual enqueue
  # to acquire locks/semaphores — the bulk path cannot (issue #413)
  individual, bulk = active_jobs.partition { |j| Uniqueness.uniqueness_config(j) || concurrency_config(j) }
  individual.each do |j|
    if scheduled_in_future?(j)
      enqueue_at(j, j.scheduled_at.to_f)
    else
      enqueue(j)
    end
  end

  # Group by priority too: send_batch routes through the queue strategy,
  # so a mixed-priority bulk send needs one produce_batch per level.
  bulk.group_by { |j| [j.queue_name || Pgbus.configuration.default_queue, j.try(:priority)] }
      .each do |(queue, priority), jobs|
    immediate, scheduled = jobs.partition { |j| !scheduled_in_future?(j) }
    enqueue_immediate(queue, immediate, priority: priority)
    scheduled.each { |j| enqueue_at(j, j.scheduled_at.to_f) }
  end

  active_jobs.count
end

#enqueue_at(active_job, timestamp) ⇒ Object



24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
# File 'lib/pgbus/active_job/adapter.rb', line 24

def enqueue_at(active_job, timestamp)
  queue = active_job.queue_name || Pgbus.configuration.default_queue
  payload_hash = Serializer.serialize_job_hash(active_job)
  payload_hash = Concurrency.(active_job, payload_hash)
  payload_hash = Uniqueness.(active_job, payload_hash)
  payload_hash = FairShare.(active_job, payload_hash)
  payload_hash = (payload_hash, active_job: active_job)
  delay = [(timestamp - Time.current.to_f).ceil, 0].max

  if uniqueness_rejected?(active_job, payload_hash, queue: queue)
    uncount_batch_job(payload_hash)
    return active_job
  end

  enqueue_with_concurrency(active_job, queue, payload_hash, delay: delay)
end