Class: SolidQueue::Semaphore

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

Overview

Semaphore for concurrency control.

Convention (matches spec expectations):

value = number of USED (acquired) slots   (0 = none in use)
limit = maximum concurrent slots

wait:   acquire a slot — succeeds when value < limit (increments value)
signal: release a slot — increments value (marks another slot returned)

NOTE: This intentionally differs from the ActiveRecord original which uses value = remaining available slots. The specs and BlockedExecution logic here both expect the "used slots" convention.

Defined Under Namespace

Classes: Proxy

Constant Summary

Constants inherited from Record

Record::INDEX_HINTS

Class Method Summary collapse

Instance Method Summary collapse

Methods inherited from Record

create_or_find_by!, find_by!, index, index_specifications, index_specifications=, inherited, non_blocking_lock, supports_insert_conflict_target?, transaction, use_index

Class Method Details

.create_unique_by(attributes) ⇒ Object

Requires a unique index on key. Returns true if created/inserted; false on duplicate.



46
47
48
49
50
51
52
53
# File 'lib/solid_queue_mongoid/models/semaphore.rb', line 46

def create_unique_by(attributes)
  create!(attributes)
  true
rescue Mongoid::Errors::Validations, Mongo::Error::OperationFailure => e
  raise unless duplicate_key_error?(e)

  false
end

.signal(job) ⇒ Object



36
37
38
# File 'lib/solid_queue_mongoid/models/semaphore.rb', line 36

def signal(job)
  Proxy.new(job).signal
end

.signal_all(jobs) ⇒ Object



40
41
42
# File 'lib/solid_queue_mongoid/models/semaphore.rb', line 40

def signal_all(jobs)
  Proxy.signal_all(jobs)
end

.wait(job) ⇒ Object



32
33
34
# File 'lib/solid_queue_mongoid/models/semaphore.rb', line 32

def wait(job)
  Proxy.new(job).wait
end

Instance Method Details

#acquire ⇒ Object

Atomically acquire one slot (increment value if value < limit). Returns true on success, false when at limit.



66
67
68
69
70
71
72
73
# File 'lib/solid_queue_mongoid/models/semaphore.rb', line 66

def acquire
  result = self.class.collection.find_one_and_update(
    { _id: id, "value" => { "$lt" => limit } },
    { "$inc" => { "value" => 1 } },
    return_document: :after
  )
  result.present?
end

#available? ⇒ Boolean

True when there is still room to acquire (value < limit).

Returns:

  • (Boolean)


86
87
88
89
# File 'lib/solid_queue_mongoid/models/semaphore.rb', line 86

def available?
  reload
  value < limit
end

#release ⇒ Object

Release one slot (decrement value). No-op if already at 0.



76
77
78
79
80
81
82
83
# File 'lib/solid_queue_mongoid/models/semaphore.rb', line 76

def release
  self.class.collection.find_one_and_update(
    { _id: id, "value" => { "$gt" => 0 } },
    { "$inc" => { "value" => -1 } },
    return_document: :after
  )
  true
end