class Scheduler
Matches pending requests to concrete backend permits.
Nested Classes and Modules
class ReservationA reserved processing permit and response exchange.
class RejectionA queue admission rejection.
Definitions
def initialize(registry, configuration = Configuration.default)
Signature
-
parameter
registryRegistry The available backend registry.
-
parameter
configurationConfiguration Request and scheduling policy.
Implementation
def initialize(registry, configuration = Configuration.default)
@registry = registry
@configuration = configuration
@guard = Thread::Mutex.new
@pending = configuration.queues.to_h{|name, queue| [name, []]}
@pending_count = 0
@closed = false
@registry.on_available{schedule}
end
def acquire(request)
Admit a request, wait for a matching permit, or return a rejection.
Implementation
def acquire(request)
queue = @configuration.classify(request)
entry = nil
result = @guard.synchronize do
if @closed
nil
elsif reservation = reserve(queue, request)
reservation
elsif @registry.empty?
nil
elsif reject?(queue, request)
Rejection.new(queue)
else
entry = Entry.new(request, queue, now, Async::Queue.new, true, nil)
@pending.fetch(queue.name) << entry
@pending_count += 1
schedule_locked
entry.assignment
end
end
return result unless entry
if assignment = entry.assignment
entry = nil
return assignment
end
if wait_limit = queue.wait_limit
remaining = wait_limit - (now - entry.enqueued_at)
result = entry.result.dequeue(timeout: remaining) if remaining.positive?
if result
entry = nil
return result
end
result = cancel(entry)
entry = nil
return result
else
result = entry.result.dequeue
entry = nil
return result
end
ensure
if entry && assignment = cancel(entry, rejection: false)
# The request was assigned concurrently but its waiting task was
# interrupted before receiving the reservation. No upstream request
# was started, so release both the permit and response exchange.
assignment.failed
end
end
def schedule
Try to dispatch pending requests after capacity changes.
Implementation
def schedule
@guard.synchronize{schedule_locked unless @closed}
end
def close
Stop accepting requests and wake all tasks waiting for a permit.
Implementation
def close
entries = @guard.synchronize do
return if @closed
@closed = true
entries = @pending.values.flatten(1)
@pending.each_value(&:clear)
@pending_count = 0
entries.each{|entry| entry.pending = false}
entries
end
entries.each{|entry| entry.result.close}
end
def pending_count(queue_name = nil)
Signature
-
parameter
queue_nameSymbol | String | Nil An optional queue to inspect.
-
returns
Integer The number of requests waiting for a permit.
Implementation
def pending_count(queue_name = nil)
@guard.synchronize do
if queue_name
@pending.fetch(queue_name.to_sym).size
else
@pending_count
end
end
end