primus.command_worker.ports module
Provider ports for the portable command-worker runtime.
- class primus.command_worker.ports.AuditSink(*args, **kwargs)
Bases:
Protocol- __init__(*args, **kwargs)
- record(event: AuditEvent) None
- class primus.command_worker.ports.ClaimStatus(value)
Bases:
str,Enum- ACTIVE = 'active'
- INTEGRITY_MISMATCH = 'integrity_mismatch'
- TERMINAL = 'terminal'
- class primus.command_worker.ports.Clock(*args, **kwargs)
Bases:
Protocol- __init__(*args, **kwargs)
- now() datetime
- class primus.command_worker.ports.CommandAuthorizer(*args, **kwargs)
Bases:
ProtocolFail-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.ports.CommandIdGenerator(*args, **kwargs)
Bases:
Protocol- __init__(*args, **kwargs)
- new_id() str
- class primus.command_worker.ports.CommandRepository(*args, **kwargs)
Bases:
ProtocolAtomic persistence boundary for portable command state.
announcemust 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.getandrequest_cancelare tenant scoped. Implementations must returnNoneboth when a command is absent and when it belongs to another tenant.request_cancelatomically 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.ports.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.ports.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.ports.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
- 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.ports.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.ports.HeartbeatHandle(*args, **kwargs)
Bases:
Protocol- __init__(*args, **kwargs)
- stop() bool
Stop scheduling and synchronously settle activity when possible.
- class primus.command_worker.ports.HeartbeatScheduler(*args, **kwargs)
Bases:
Protocol- __init__(*args, **kwargs)
- start(interval: timedelta, callback: Callable[[], None]) HeartbeatHandle
- class primus.command_worker.ports.LifecycleStore(*args, **kwargs)
Bases:
ProtocolDurable lifecycle and integrity boundary for untrusted envelopes.
Before issuing a lease,
claimmust 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.
- report_progress(command_id: str, token: str, progress: ProgressUpdate, now: datetime) bool
- class primus.command_worker.ports.TaskScaleInProtection(*args, **kwargs)
Bases:
ProtocolProtect the current worker task while a claimed command is executing.
This is deliberately a small runtime port: the portable worker does not need to know whether the deployment is ECS, Kubernetes, or a local test. Implementations must make enablement durable before work begins and clear protection after the command reaches a terminal lifecycle state.
- __init__(*args, **kwargs)
- clear() bool
- enable() bool