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: object

Task-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
renew_task(command_id: str, token: str, now: datetime, lease_duration: timedelta) → Claim | None
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, LifecycleStore

One-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
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
request_cancel(tenant_id: str, command_id: str, now: datetime) → CancellationResult | None
class primus.command_worker.adapters.task_store.TaskStoreGateway(*args, **kwargs)

Bases: Protocol

Atomic 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
renew_task(command_id: str, token: str, now: datetime, lease_duration: timedelta) → Claim | None
request_task_cancel(tenant_id: str, command_id: str, now: datetime) → CancellationResult | None