primus.command_worker.service module

Provider-neutral service loop for a warm command worker container.

class primus.command_worker.service.CommandWorkerService(transport: Transport, worker: CommandWorker, owner: str, drain: DrainSignal, receive_timeout: timedelta, idle_wait: timedelta)

Bases: object

Receive and synchronously settle one delivery at a time until draining.

__init__(transport: Transport, worker: CommandWorker, owner: str, drain: DrainSignal, receive_timeout: timedelta, idle_wait: timedelta) → None
run() → ServiceOutcome
class primus.command_worker.service.EventDrainSignal

Bases: object

Thread-safe drain signal backed by a standard-library event.

__init__() → None
is_requested() → bool
request() → None
wait(timeout: timedelta) → bool
class primus.command_worker.service.ServiceOutcome(stop_reason: ServiceStopReason, process_outcomes: tuple[ProcessOutcome, ...])

Bases: object

Normal service termination and the work settled before it stopped.

__init__(stop_reason: ServiceStopReason, process_outcomes: tuple[ProcessOutcome, ...]) → None
count(outcome: ProcessOutcome) → int
property outcome_counts: dict[ProcessOutcome, int]
process_outcomes: tuple[ProcessOutcome, ...]
property processed_count: int
stop_reason: ServiceStopReason
exception primus.command_worker.service.ServiceReceiveError(error: Exception, process_outcomes: tuple[ProcessOutcome, ...])

Bases: RuntimeError

Raised when the transport cannot complete a bounded receive.

__init__(error: Exception, process_outcomes: tuple[ProcessOutcome, ...]) → None
property processed_count: int
class primus.command_worker.service.ServiceStopReason(value)

Bases: str, Enum

DRAIN_REQUESTED = 'drain_requested'