class Scheduler

Matches pending requests to concrete backend permits.

Nested Classes and Modules

class Reservation

A reserved processing permit and response exchange.

class Rejection

A queue admission rejection.

Definitions

def initialize(registry, configuration = Configuration.default)

Signature

parameter registry Registry

The available backend registry.

parameter configuration Configuration

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_name Symbol | 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