primus.command_worker package

Portable command service and worker foundation.

class primus.command_worker.AuditEvent(event_type: AuditEventType, occurred_at: datetime, tenant_id: str, principal_id: str, principal_type: str, authentication_method: str, correlation_id: str, outcome: str, policy_version: str | None = None, command_id: str | None = None, target: str | None = None)

Bases: object

Structured payload-free security and command audit event.

__init__(event_type: AuditEventType, occurred_at: datetime, tenant_id: str, principal_id: str, principal_type: str, authentication_method: str, correlation_id: str, outcome: str, policy_version: str | None = None, command_id: str | None = None, target: str | None = None) → None
authentication_method: str
command_id: str | None
correlation_id: str
event_type: AuditEventType
occurred_at: datetime
outcome: str
policy_version: str | None
principal_id: str
principal_type: str
target: str | None
tenant_id: str
class primus.command_worker.AuditEventType(value)

Bases: str, Enum

AUTHORIZATION_DENIED = 'AUTHORIZATION_DENIED'
CANCELLATION = 'CANCELLATION'
IDEMPOTENCY_CONFLICT = 'IDEMPOTENCY_CONFLICT'
READ = 'READ'
SUBMIT_CREATED = 'SUBMIT_CREATED'
SUBMIT_EXISTING = 'SUBMIT_EXISTING'
class primus.command_worker.AuditSink(*args, **kwargs)

Bases: Protocol

__init__(*args, **kwargs)
record(event: AuditEvent) → None
class primus.command_worker.AuthenticatedCommandContext(tenant_id: str, principal_id: str, principal_type: str, authentication_method: str, correlation_id: str)

Bases: object

Identity asserted only by a future authentication adapter.

__init__(tenant_id: str, principal_id: str, principal_type: str, authentication_method: str, correlation_id: str) → None
authentication_method: str
correlation_id: str
principal_id: str
principal_type: str
tenant_id: str
class primus.command_worker.AuthorizationDecision(allowed: 'bool', policy_version: 'str')

Bases: object

__init__(allowed: bool, policy_version: str) → None
allowed: bool
policy_version: str
exception primus.command_worker.AuthorizationDenied

Bases: PermissionError

Submission was denied or its target is not registered.

__init__() → None
exception primus.command_worker.CancellationRequestedError

Bases: RuntimeError

Raised at a cooperative boundary after durable cancellation is requested.

class primus.command_worker.CancellationResult(command: 'CommandRecord', changed: 'bool')

Bases: object

__init__(command: CommandRecord, changed: bool) → None
changed: bool
command: CommandRecord
class primus.command_worker.Claim(token: str, owner: str, expires_at: datetime, cancellation_requested: bool = False)

Bases: object

A lifecycle-store-issued fencing lease.

__init__(token: str, owner: str, expires_at: datetime, cancellation_requested: bool = False) → None
cancellation_requested: bool
expires_at: datetime
owner: str
token: str
class primus.command_worker.ClaimStatus(value)

Bases: str, Enum

ACTIVE = 'active'
INTEGRITY_MISMATCH = 'integrity_mismatch'
TERMINAL = 'terminal'
class primus.command_worker.Clock(*args, **kwargs)

Bases: Protocol

__init__(*args, **kwargs)
now() → datetime
class primus.command_worker.CommandAuthorizer(*args, **kwargs)

Bases: Protocol

Fail-closed policy for authenticated contexts and registered targets.

__init__(*args, **kwargs)
can_cancel(context: AuthenticatedCommandContext, command: CommandRecord) → AuthorizationDecision
can_read(context: AuthenticatedCommandContext, command: CommandRecord) → AuthorizationDecision
can_submit(context: AuthenticatedCommandContext, request: CommandRequest) → AuthorizationDecision
class primus.command_worker.CommandEnvelope(schema_version: int, command_id: str, tenant_id: str, target: str, idempotency_key: str, created_at: datetime, payload: Mapping[str, str | int | float | bool | None | Mapping[str, str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...]] | tuple[str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...], ...]])

Bases: object

Immutable, versioned command contract shared by all providers.

CURRENT_SCHEMA_VERSION: ClassVar[int] = 2
__init__(schema_version: int, command_id: str, tenant_id: str, target: str, idempotency_key: str, created_at: datetime, payload: Mapping[str, str | int | float | bool | None | Mapping[str, str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...]] | tuple[str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...], ...]]) → None
command_id: str
created_at: datetime
classmethod from_message(message: Mapping[str, Any]) → CommandEnvelope

Validate and deserialize a transport message.

idempotency_key: str
payload: Mapping[str, str | int | float | bool | None | Mapping[str, str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...]] | tuple[str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...], ...]]
schema_version: int
target: str
tenant_id: str
to_message() → dict[str, Any]

Create a strict JSON-serializable transport message.

class primus.command_worker.CommandIdGenerator(*args, **kwargs)

Bases: Protocol

__init__(*args, **kwargs)
new_id() → str
class primus.command_worker.CommandLimits(max_identifier_length: int = 256, max_idempotency_key_length: int = 512, max_json_depth: int = 32, max_json_containers: int = 10000, max_json_bytes: int = 1048576)

Bases: object

Provider-neutral safety limits; adapters may impose stricter values.

__init__(max_identifier_length: int = 256, max_idempotency_key_length: int = 512, max_json_depth: int = 32, max_json_containers: int = 10000, max_json_bytes: int = 1048576) → None
max_idempotency_key_length: int
max_identifier_length: int
max_json_bytes: int
max_json_containers: int
max_json_depth: int
exception primus.command_worker.CommandNotFound(command_id: str)

Bases: LookupError

A tenant-scoped command lookup did not resolve or was not authorized.

__init__(command_id: str) → None
class primus.command_worker.CommandRecord(command_id: str, tenant_id: str, target: str, idempotency_key: str, idempotency_namespace: str, created_at: datetime, updated_at: datetime, submitted_by: str, payload: Mapping[str, str | int | float | bool | None | Mapping[str, str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...]] | tuple[str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...], ...]], status: CommandStatus, request_digest: RequestDigest)

Bases: object

Immutable durable command snapshot.

__init__(command_id: str, tenant_id: str, target: str, idempotency_key: str, idempotency_namespace: str, created_at: datetime, updated_at: datetime, submitted_by: str, payload: Mapping[str, str | int | float | bool | None | Mapping[str, str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...]] | tuple[str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...], ...]], status: CommandStatus, request_digest: RequestDigest) → None
command_id: str
created_at: datetime
property envelope: CommandEnvelope
idempotency_key: str
idempotency_namespace: str
payload: Mapping[str, str | int | float | bool | None | Mapping[str, str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...]] | tuple[str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...], ...]]
request_digest: RequestDigest
status: CommandStatus
submitted_by: str
target: str
tenant_id: str
updated_at: datetime
class primus.command_worker.CommandRepository(*args, **kwargs)

Bases: Protocol

Atomic persistence boundary for portable command state.

announce must atomically resolve the tenant/principal/operation-scoped idempotency key, durably store a new command, and make that command’s dispatch discoverable. A caller must never observe durable state without its corresponding dispatch eligibility, or dispatch eligibility without durable state.

get and request_cancel are tenant scoped. Implementations must return None both when a command is absent and when it belongs to another tenant. request_cancel atomically applies the cancellation transition and, for an ANNOUNCED command, removes dispatch eligibility in the same operation.

__init__(*args, **kwargs)
announce(command: CommandRecord) → AnnouncementResult
get(tenant_id: str, command_id: str) → CommandRecord | None
request_cancel(tenant_id: str, command_id: str, now: datetime) → CancellationResult | None
class primus.command_worker.CommandRequest(target: str, idempotency_key: str, payload: Mapping[str, str | int | float | bool | None | Mapping[str, str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...]] | tuple[str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...], ...]])

Bases: object

Work requested by a trusted adapter; tenant identity is intentionally absent.

__init__(target: str, idempotency_key: str, payload: Mapping[str, str | int | float | bool | None | Mapping[str, str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...]] | tuple[str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...], ...]]) → None
idempotency_key: str
payload: Mapping[str, str | int | float | bool | None | Mapping[str, str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...]] | tuple[str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...], ...]]
target: str
class primus.command_worker.CommandService(repository: CommandRepository, clock: Clock, command_ids: CommandIdGenerator, authorizer: CommandAuthorizer, audit: AuditSink, limits: CommandLimits = CommandLimits(max_identifier_length=256, max_idempotency_key_length=512, max_json_depth=32, max_json_containers=10000, max_json_bytes=1048576))

Bases: object

Application boundary for context verified by an authentication adapter.

__init__(repository: CommandRepository, clock: Clock, command_ids: CommandIdGenerator, authorizer: CommandAuthorizer, audit: AuditSink, limits: CommandLimits = CommandLimits(max_identifier_length=256, max_idempotency_key_length=512, max_json_depth=32, max_json_containers=10000, max_json_bytes=1048576)) → None
cancel(context: AuthenticatedCommandContext, command_id: str) → CancellationResult
get(context: AuthenticatedCommandContext, command_id: str) → CommandRecord
submit(context: AuthenticatedCommandContext, request: CommandRequest) → SubmissionResult
class primus.command_worker.CommandStatus(value)

Bases: str, Enum

ANNOUNCED = 'ANNOUNCED'
CANCELLED = 'CANCELLED'
CANCEL_REQUESTED = 'CANCEL_REQUESTED'
FAILED = 'FAILED'
RUNNING = 'RUNNING'
SUCCEEDED = 'SUCCEEDED'
property is_terminal: bool
class primus.command_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
class primus.command_worker.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.Delivery(*args, **kwargs)

Bases: Protocol

__init__(*args, **kwargs)
acknowledge() → None
envelope: CommandEnvelope
extend_lease(duration: timedelta) → bool
quarantine(reason: str) → None
release() → None
class primus.command_worker.DrainSignal(*args, **kwargs)

Bases: Protocol

__init__(*args, **kwargs)
is_requested() → bool
wait(timeout: timedelta) → bool

Wait until draining is requested or the timeout elapses.

class primus.command_worker.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.ExecutionContext(*args, **kwargs)

Bases: Protocol

__init__(*args, **kwargs)
property cancellation_requested: bool
property ownership_lost: bool
raise_if_cancellation_requested() → None
raise_if_lease_lost() → None
renew_lease() → Claim
report_progress(fraction: float, message: str | None = None, details: dict[str, str | int | float | bool | None | Mapping[str, str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...]] | tuple[str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...], ...]] | None = None) → None
class primus.command_worker.Executor(*args, **kwargs)

Bases: Protocol

__init__(*args, **kwargs)
execute(envelope: CommandEnvelope, context: ExecutionContext) → str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...]
class primus.command_worker.HeartbeatHandle(*args, **kwargs)

Bases: Protocol

__init__(*args, **kwargs)
stop() → bool

Stop scheduling and synchronously settle activity when possible.

class primus.command_worker.HeartbeatScheduler(*args, **kwargs)

Bases: Protocol

__init__(*args, **kwargs)
start(interval: timedelta, callback: Callable[[], None]) → HeartbeatHandle
exception primus.command_worker.IdempotencyConflict

Bases: ValueError

An actor reused an idempotency key for different canonical work.

__init__() → None
class primus.command_worker.InMemoryCommandRepository

Bases: object

Thread-safe executable reference for the CommandRepository contract.

__init__() → None
announce(command: CommandRecord) → AnnouncementResult
discoverable_dispatches() → tuple[CommandEnvelope, ...]
get(tenant_id: str, command_id: str) → CommandRecord | None
request_cancel(tenant_id: str, command_id: str, now: datetime) → CancellationResult | None
set_status(tenant_id: str, command_id: str, status: CommandStatus) → CommandRecord

Advance status in reference integration tests and future adapter tests.

exception primus.command_worker.LeaseLostError

Bases: RuntimeError

Raised when the lifecycle store rejects a fenced mutation.

class primus.command_worker.LifecycleStore(*args, **kwargs)

Bases: Protocol

Durable lifecycle and integrity boundary for untrusted envelopes.

Before issuing a lease, claim must atomically load durable command state and verify command identity, tenant, target, idempotency identity, and the canonical payload digest. CANCELLED is terminal even when a broker message was published before cancellation. Any mismatch returns INTEGRITY_MISMATCH.

__init__(*args, **kwargs)
claim(envelope: CommandEnvelope, owner: str, now: datetime, lease_duration: timedelta) → Claim | ClaimStatus
complete(command_id: str, token: str, result: str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...], now: datetime) → bool
fail(command_id: str, token: str, error: str, now: datetime) → bool
finalize_cancel(command_id: str, token: str, now: datetime) → bool

Atomically settle a cancellation request held by this lease.

renew(command_id: str, token: str, now: datetime, lease_duration: timedelta) → Claim | None
report_progress(command_id: str, token: str, progress: ProgressUpdate, now: datetime) → bool
class primus.command_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'
class primus.command_worker.ProgressUpdate(fraction: 'float', message: 'str | None' = None, details: 'Mapping[str, JSONValue] | None' = None)

Bases: object

__init__(fraction: float, message: str | None = None, details: Mapping[str, str | int | float | bool | None | Mapping[str, str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...]] | tuple[str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...], ...]] | None = None) → None
details: Mapping[str, str | int | float | bool | None | Mapping[str, str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...]] | tuple[str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...], ...]] | None
fraction: float
message: str | None
class primus.command_worker.RequestDigest(algorithm: 'str', canonicalization_version: 'int', value: 'str')

Bases: object

__init__(algorithm: str, canonicalization_version: int, value: str) → None
algorithm: str
canonicalization_version: int
value: str
class primus.command_worker.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.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.ServiceStopReason(value)

Bases: str, Enum

DRAIN_REQUESTED = 'drain_requested'
class primus.command_worker.SubmissionDisposition(value)

Bases: str, Enum

EXISTING = 'EXISTING'
NEW = 'NEW'
class primus.command_worker.SubmissionResult(command: 'CommandRecord', disposition: 'SubmissionDisposition')

Bases: object

__init__(command: CommandRecord, disposition: SubmissionDisposition) → None
command: CommandRecord
disposition: SubmissionDisposition
class primus.command_worker.ThreadHeartbeatScheduler(shutdown_timeout: timedelta = datetime.timedelta(seconds=5))

Bases: object

Run recurring callbacks on a daemon thread with bounded shutdown.

__init__(shutdown_timeout: timedelta = datetime.timedelta(seconds=5)) → None
start(interval: timedelta, callback: Callable[[], None]) → _ThreadHeartbeatHandle
class primus.command_worker.Transport(*args, **kwargs)

Bases: Protocol

__init__(*args, **kwargs)
receive(timeout: timedelta) → Delivery | None
class primus.command_worker.UUIDCommandIdGenerator

Bases: object

new_id() → str
primus.command_worker.request_digest(target: str, payload: Mapping[str, Any], limits: CommandLimits = CommandLimits(max_identifier_length=256, max_idempotency_key_length=512, max_json_depth=32, max_json_containers=10000, max_json_bytes=1048576)) → RequestDigest

Digest strict, bounded, canonically ordered request content.

Subpackages

Submodules