primus.command_worker.worker module

Synchronous orchestration for a single command delivery.

exception primus.command_worker.worker.CancellationRequestedError

Bases: RuntimeError

Raised at a cooperative boundary after durable cancellation is requested.

class primus.command_worker.worker.CommandWorker(lifecycle: LifecycleStore, executor: Executor, clock: Clock, lease_duration: timedelta, heartbeat_interval: timedelta | None = None, delivery_lease_duration: timedelta | None = None, heartbeat_scheduler: HeartbeatScheduler | None = None, task_scale_in_protection: TaskScaleInProtection | None = None)

Bases: object

__init__(lifecycle: LifecycleStore, executor: Executor, clock: Clock, lease_duration: timedelta, heartbeat_interval: timedelta | None = None, delivery_lease_duration: timedelta | None = None, heartbeat_scheduler: HeartbeatScheduler | None = None, task_scale_in_protection: TaskScaleInProtection | None = None) → None
process(delivery: Delivery, owner: str) → ProcessOutcome
exception primus.command_worker.worker.LeaseLostError

Bases: RuntimeError

Raised when the lifecycle store rejects a fenced mutation.

class primus.command_worker.worker.ProcessOutcome(value)

Bases: str, Enum

ACTIVE_DUPLICATE = 'active_duplicate'
CANCELLED = 'cancelled'
COMPLETED = 'completed'
FAILED = 'failed'
INTEGRITY_MISMATCH = 'integrity_mismatch'
LEASE_LOST = 'lease_lost'
TERMINAL_DUPLICATE = 'terminal_duplicate'