Class: SolidQueue::Semaphore::Proxy
- Inherits:
-
Object
- Object
- SolidQueue::Semaphore::Proxy
- Defined in:
- lib/solid_queue_mongoid/models/semaphore.rb
Overview
── Proxy inner class ─────────────────────────────────────────────────────
Class Method Summary collapse
-
.signal_all(jobs) ⇒ Object
Decrement value for every job's semaphore key (signal = release a used slot).
Instance Method Summary collapse
-
#initialize(job) ⇒ Proxy
constructor
A new instance of Proxy.
-
#signal ⇒ Object
Release a slot: decrement value (marks one used slot as freed).
-
#wait ⇒ Object
Acquire a slot: succeeds when value < limit.
Constructor Details
#initialize(job) ⇒ Proxy
Returns a new instance of Proxy.
106 107 108 |
# File 'lib/solid_queue_mongoid/models/semaphore.rb', line 106 def initialize(job) @job = job end |
Class Method Details
.signal_all(jobs) ⇒ Object
Decrement value for every job's semaphore key (signal = release a used slot).
94 95 96 97 98 99 100 101 102 103 104 |
# File 'lib/solid_queue_mongoid/models/semaphore.rb', line 94 def self.signal_all(jobs) keys = jobs.map(&:concurrency_key) return if keys.empty? Semaphore.in(key: keys).each do |sem| Semaphore.collection.find_one_and_update( { _id: sem.id, "value" => { "$gt" => 0 } }, { "$inc" => { "value" => -1 } } ) end end |
Instance Method Details
#signal ⇒ Object
Release a slot: decrement value (marks one used slot as freed).
124 125 126 |
# File 'lib/solid_queue_mongoid/models/semaphore.rb', line 124 def signal attempt_release_slot end |
#wait ⇒ Object
Acquire a slot: succeeds when value < limit. Creates the semaphore document on first use.
112 113 114 115 116 117 118 119 120 121 |
# File 'lib/solid_queue_mongoid/models/semaphore.rb', line 112 def wait semaphore = Semaphore.where(key: key).first if semaphore # Atomically increment if value < limit attempt_acquire(semaphore.id, semaphore.limit) else attempt_creation end end |