Class: SolidQueue::Semaphore::Proxy

Inherits:
Object
  • Object
show all
Defined in:
lib/solid_queue_mongoid/models/semaphore.rb

Overview

── Proxy inner class ─────────────────────────────────────────────────────

Class Method Summary collapse

Instance Method Summary collapse

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