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: 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.ports.CommandIdGenerator(*args, **kwargs)

Bases: Protocol

__init__(*args, **kwargs)
new_id() → str
class primus.command_worker.ports.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.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
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.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: 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.ports.TaskScaleInProtection(*args, **kwargs)

Bases: Protocol

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

Bases: Protocol

__init__(*args, **kwargs)
receive(timeout: timedelta) → Delivery | None