primus.command_worker.adapters.task_store module
Task-backed portable command repository and lifecycle adapter.
TaskStoreGateway is the boundary implemented by the dashboard Task data
service. Its mutations are conditional mutations of one Task record (and its
authoritative lifecycle), never a projection of another command record.
TaskStages are separate UI detail written by command executors. The in-memory
gateway is deliberately a local proof double; production gateways must provide
the same atomic operations against the Task store.
- class primus.command_worker.adapters.task_store.GraphQLTaskStoreGateway(client: Any)
Bases:
objectTask-only lifecycle gateway using AppSync generated model mutations.
The task id is the idempotency identity. Conditional failures are normal contention signals; every such path reloads the single Task record rather than consulting or creating a secondary command index.
- __init__(client: Any) None
- announce_task(command: CommandRecord, task_fields: Mapping[str, object]) AnnouncementResult
- claim_task(envelope: CommandEnvelope, owner: str, now: datetime, lease_duration: timedelta) Claim | ClaimStatus
- complete_task(command_id: str, token: str, result: str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...], now: datetime) bool
- fail_task(command_id: str, token: str, error: str, now: datetime) bool
- finalize_task_cancel(command_id: str, token: str, now: datetime) bool
- get_command(tenant_id: str, command_id: str) CommandRecord | None
- get_task(command_id: str) Any | None
- progress_task(command_id: str, token: str, progress: ProgressUpdate, now: datetime) bool
- request_task_cancel(tenant_id: str, command_id: str, now: datetime) CancellationResult | None
- class primus.command_worker.adapters.task_store.TaskBackedCommandStore(gateway: TaskStoreGateway, task_fields: Mapping[str, object])
Bases:
CommandRepository,LifecycleStoreOne-store adapter: Task is both command record and lifecycle record.
- __init__(gateway: TaskStoreGateway, task_fields: Mapping[str, object]) None
- announce(command: CommandRecord) AnnouncementResult
- 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.
- get(tenant_id: str, command_id: str) CommandRecord | None
- report_progress(command_id: str, token: str, progress: ProgressUpdate, now: datetime) bool
- request_cancel(tenant_id: str, command_id: str, now: datetime) CancellationResult | None
- class primus.command_worker.adapters.task_store.TaskStoreGateway(*args, **kwargs)
Bases:
ProtocolAtomic Task-store operations required by the portable lifecycle.
- __init__(*args, **kwargs)
- announce_task(command: CommandRecord, task_fields: Mapping[str, object]) AnnouncementResult
- claim_task(envelope: CommandEnvelope, owner: str, now: datetime, lease_duration: timedelta) Claim | ClaimStatus
- complete_task(command_id: str, token: str, result: str | int | float | bool | None | Mapping[str, JSONValue] | tuple[JSONValue, ...], now: datetime) bool
- fail_task(command_id: str, token: str, error: str, now: datetime) bool
- finalize_task_cancel(command_id: str, token: str, now: datetime) bool
- get_command(tenant_id: str, command_id: str) CommandRecord | None
- get_task(command_id: str) Any | None
- progress_task(command_id: str, token: str, progress: ProgressUpdate, now: datetime) bool
- request_task_cancel(tenant_id: str, command_id: str, now: datetime) CancellationResult | None