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:
objectStructured 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:
objectIdentity 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:
PermissionErrorSubmission was denied or its target is not registered.
- __init__() None
- exception primus.command_worker.CancellationRequestedError
Bases:
RuntimeErrorRaised 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:
objectA 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:
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.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:
objectImmutable, 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:
objectProvider-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:
LookupErrorA 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:
objectImmutable 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:
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.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:
objectWork 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:
objectApplication 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:
objectReceive 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:
objectThread-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
- 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:
ValueErrorAn actor reused an idempotency key for different canonical work.
- __init__() None
- class primus.command_worker.InMemoryCommandRepository
Bases:
objectThread-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:
RuntimeErrorRaised when the lifecycle store rejects a fenced mutation.
- class primus.command_worker.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.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:
objectNormal 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:
RuntimeErrorRaised 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:
objectRun 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
- 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
- primus.command_worker.adapters package
- primus.command_worker.executors package
- primus.command_worker.runtime package
CommandWorkerRuntimeConfigCommandWorkerRuntimeConfig.__init__()CommandWorkerRuntimeConfig.api_urlCommandWorkerRuntimeConfig.executor_factoryCommandWorkerRuntimeConfig.from_environment()CommandWorkerRuntimeConfig.heartbeat_intervalCommandWorkerRuntimeConfig.lease_durationCommandWorkerRuntimeConfig.log_levelCommandWorkerRuntimeConfig.queue_nameCommandWorkerRuntimeConfig.queue_urlCommandWorkerRuntimeConfig.regionCommandWorkerRuntimeConfig.task_nameCommandWorkerRuntimeConfig.visibility_timeout
main()- Submodules
Submodules
- primus.command_worker.application module
- primus.command_worker.models module
AnnouncementDispositionAnnouncementResultAuditEventAuditEventTypeAuthenticatedCommandContextAuthorizationDecisionCancellationResultClaimCommandEnvelopeCommandEnvelope.CURRENT_SCHEMA_VERSIONCommandEnvelope.__init__()CommandEnvelope.command_idCommandEnvelope.created_atCommandEnvelope.from_message()CommandEnvelope.idempotency_keyCommandEnvelope.payloadCommandEnvelope.schema_versionCommandEnvelope.targetCommandEnvelope.tenant_idCommandEnvelope.to_message()
CommandLimitsCommandRecordCommandRecord.__init__()CommandRecord.command_idCommandRecord.created_atCommandRecord.envelopeCommandRecord.idempotency_keyCommandRecord.idempotency_namespaceCommandRecord.payloadCommandRecord.request_digestCommandRecord.statusCommandRecord.submitted_byCommandRecord.targetCommandRecord.tenant_idCommandRecord.updated_at
CommandRequestCommandStatusProgressUpdateRequestDigestSubmissionDispositionSubmissionResultfreeze_json()request_digest()
- primus.command_worker.ports module
- primus.command_worker.repository module
- primus.command_worker.scheduler module
- primus.command_worker.service module
- primus.command_worker.smoke module
- primus.command_worker.worker module