Class: Pgbus::ActiveJob::Adapter
- Inherits:
-
Object
- Object
- Pgbus::ActiveJob::Adapter
- Defined in:
- lib/pgbus/active_job/adapter.rb
Instance Method Summary collapse
- #enqueue(active_job) ⇒ Object
- #enqueue_all(active_jobs) ⇒ Object
- #enqueue_at(active_job, timestamp) ⇒ Object
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, ) 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 = [( - 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 |