Skip to main content

operations API

SDK 0.4.0. This reference describes the supported authoring interface.

OperationStatus​

Observation envelope of one operation; opaque references, safe errors.

Fields​

operation_id: Identifier
state: OperationState
submission: Submission
progress: float | None
stage: str | None
queue_position: int | None
submitted_at: int
updated_at: int
finished_at: int | None
stale: bool
input_digest: Digest
backend_reference: str | None
input: dict[str, JsonValue] | None
result: dict[str, JsonValue] | None
error: str | None
FieldDescriptionDeclared defaults and constraints
operation_id——
state——
submission——
progress—default=null; ge=0; le=1
stage—default=null; max_length=64
queue_position—default=null; ge=0
submitted_at—ge=0
updated_at—ge=0
finished_at—default=null; ge=0
stale——
input_digest——
backend_reference—default=null
inputThe recorded start input, so a caller can run the operation again or vary it. Carried by status when it fits beside the result (MAX_INPUT_ECHO_BYTES); listings omit it.default=null
result—default=null
error—default=null; max_length=256

OperationCall​

One verb of the operation contract; the dispatcher checks authority.

OperationCall.verb_shape​

Require exactly the fields each verb needs.

def verb_shape(self) -> Self:
...

Fields​

op: Verb
operation_id: Identifier | None
input: dict[str, JsonValue] | None
offset: int | None
FieldDescriptionDeclared defaults and constraints
op——
operation_id—default=null
input—default=null
offset—default=null; ge=0

OperationPlan​

Immutable estimate a caller reviews before start.

Fields​

plan_digest: Digest
summary: str
estimated_seconds: int | None
estimated_cost_microunits: int | None
expires_at: int
input: dict[str, JsonValue] | None
FieldDescriptionDeclared defaults and constraints
plan_digest——
summary—min_length=1; max_length=512
estimated_seconds—default=null; ge=0
estimated_cost_microunits—default=null; ge=0
expires_at—ge=0
inputThe input to start with: the caller's, completed with what the plan chose (a seed, for one) and carrying plan_digest, so the reviewed plan is the one that starts. Absent when the operation chooses nothing beyond the caller's input.default=null

OperationPage​

One page of retained operations, newest first, compacted.

Results and backend references can each reach their bounds, so a listing carries statuses without them (status returns the full record) and at most PAGE_SIZE entries; next_offset continues the listing.

Fields​

operations: tuple[OperationStatus, ...]
next_offset: int | None
active: int
outstanding: int
FieldDescriptionDeclared defaults and constraints
operations——
next_offset—default=null; ge=0
active—ge=0
outstandingInterrupted rows: the backend may still hold their jobs.default=0; ge=0

OperationLog​

One offset-addressed slice of an operation's bounded log.

Fields​

operation_id: Identifier
offset: int
next_offset: int
truncated: bool
text: str
FieldDescriptionDeclared defaults and constraints
operation_id——
offset—ge=0
next_offset—ge=0
truncated——
text——

operation_descriptor​

Build the public descriptor of an operation capability.

The public input is :class:OperationCall; input_schema describes the start and plan input and result_schema the result object of a succeeded status. Both are embedded so a caller that has never heard of the operation can validate what it sends and receives.

def operation_descriptor(*, id: str, version: str, title: str, description: str, input_schema: dict[str, JsonValue], result_schema: dict[str, JsonValue]) -> Descriptor:
...

input_digest​

Digest canonical input under the installed node and descriptor scope.

def input_digest(scope: str, payload: dict[str, JsonValue]) -> str:
...

SubmissionRefused​

submit refused the job before any request left the child.

The message is published as the operation's error, so it must be operator-safe. Raising it records a definite failure; any other exception from submit records an uncertain submission, because the request may have reached the backend.

SubmissionRefused.__init__​

def __init__(self, reason: str) -> None:
...

JobAdapter​

Backend-specific submit, observe and cancel of one job.

submit returns the backend reference once the backend has accepted; raising :class:SubmissionRefused before anything is sent records a failure with its reason, and raising anything else after the request went out leaves the submission uncertain and the journal records that instead of retrying. observe returns (state, progress, stage, result, error) for a reference; cancel asks the backend to stop the job it owns.

stage and error are published verbatim in the operation's status, so the adapter must return operator-safe text: a fixed reason or a redacted message, never a provider response body, URL, or request detail. The same rule already governs every log line the adapter's owner appends; the runner cannot redact what it does not understand.

JobAdapter.submit​

async def submit(self, operation_id: str, payload: dict[str, JsonValue]) -> str:
...

JobAdapter.observe​

async def observe(self, reference: str) -> tuple[OperationState, float | None, str | None, dict[str, JsonValue] | None, str | None]:
...

JobAdapter.cancel​

Ask the backend to stop the job; return once the request is accepted.

Acceptance is not termination: the runner keeps observing until observe reports a terminal state, so a backend that stops jobs asynchronously is never recorded as stopped before it is.

async def cancel(self, reference: str) -> None:
...

ResultRefused​

A terminal result the journal will not persist; the message is safe to show.

Raised for an oversized result and for one carrying a non-finite number: the wire encodes canonically without NaN or infinity, so a record holding one could never be served again.

journal_key​

Row key: the descriptor scope and the caller's operation id together.

Two descriptors sharing a journal may receive the same caller id; keying by both keeps their reservations, records, and logs apart.

def journal_key(scope: str, operation_id: str) -> str:
...

OperationJournal​

Durable operation records with reservation-scoped idempotency.

Rows are keyed by scope plus operation id. reservations is the replay fence and is bounded only by its capacity, which refuses new starts when full rather than forgetting ids. operations is the display history bounded by retain; pruning it cannot revive a key because reserve consults the fence first, and rows that may still own a backend job (uncertain submissions, interrupted rows) are never pruned.

OperationJournal.__init__​

def __init__(self, path: Path, bounds: OperationBounds, clock: Callable[[], int]=...) -> None:
...

OperationJournal.reserve​

Reserve an id atomically or return the existing matching record.

Raises ValueError when the id was reserved under a different digest, when the queue or the fence is full, or when the id is reserved but its display row has been pruned (the fence still holds).

def reserve(self, operation_id: str, digest: str, payload: dict[str, JsonValue], scope: str='') -> tuple[Literal['created', 'existing'], OperationStatus]:
...

OperationJournal.update​

Record one observation; terminal states stamp finished_at.

Raises ValueError for an unknown row, a result past the bound, a reference past the bound, or a value the status contract refuses.

def update(self, operation_id: str, *, scope: str='', state: OperationState | None=None, submission: Submission | None=None, progress: float | None=None, stage: str | None=None, backend_reference: str | None=None, result: dict[str, JsonValue] | None=None, error: str | None=None) -> OperationStatus:
...

OperationJournal.status​

Return one retained record, or None.

def status(self, operation_id: str, scope: str='') -> OperationStatus | None:
...

OperationJournal.started_at​

Return the persisted first running transition, if any.

def started_at(self, operation_id: str, scope: str='') -> int | None:
...

OperationJournal.payload​

Return the recorded input of one operation for the worker.

def payload(self, operation_id: str, scope: str='') -> dict[str, JsonValue] | None:
...

OperationJournal.page​

Return one page of a scope's records, newest first, compacted.

Results and backend references are omitted so the page provably fits one frame; status returns the full record.

active and outstanding count the whole scope, not the page.

def page(self, scope: str='', offset: int=0) -> OperationPage:
...

OperationJournal.request_cancel​

Persist a cancel intent so it survives a restart before it is honored.

def request_cancel(self, operation_id: str, scope: str='') -> None:
...

OperationJournal.cancel_requested​

Whether a persisted cancel intent is outstanding for the row.

def cancel_requested(self, operation_id: str, scope: str='') -> bool:
...

OperationJournal.next_queued​

Return the oldest queued operation id in scope.

def next_queued(self, scope: str='') -> str | None:
...

OperationJournal.interrupted_accepted​

Return interrupted rows in scope the backend accepted: id, reference.

def interrupted_accepted(self, scope: str='') -> list[tuple[str, str]]:
...

OperationJournal.drain_retired​

Return and clear the (scope, operation_id) pairs pruned so far.

def drain_retired(self) -> list[tuple[str, str]]:
...

OperationLogs​

Bounded per-operation log files the child owns, keyed by scope and id.

OperationLogs.__init__​

def __init__(self, directory: Path, bounds: OperationBounds) -> None:
...

OperationLogs.append​

Append one already-redacted line; drop lines past the byte bound.

def append(self, operation_id: str, line: str, scope: str='') -> None:
...

OperationLogs.retire​

Delete the log of an operation pruned from history.

def retire(self, operation_id: str, scope: str='') -> None:
...

OperationLogs.read​

Return one bounded slice starting at offset.

def read(self, operation_id: str, offset: int, scope: str='') -> OperationLog:
...

OperationRunner​

One worker draining the queue through a backend adapter.

The worker submits, observes until terminal, and records every step in the journal. A runner owns one descriptor scope: it dequeues and reconciles only that scope's rows, so descriptors sharing a journal each run through their own adapter and result schema. Cancellation asks the adapter; the runner owns no process group and kills nothing. max_runtime_seconds cancels through the same adapter path, anchored to the persisted first running transition. Every adapter await is bounded by adapter_seconds.

OperationRunner.__init__​

def __init__(self, journal: OperationJournal, logs: OperationLogs, adapter: JobAdapter, *, scope: str='', poll_seconds: float=1.0, adapter_seconds: float=30.0, result_schema: dict[str, JsonValue], clock: Callable[[], int]=...) -> None:
...

OperationRunner.start​

Start the single worker task.

def start(self) -> None:
...

OperationRunner.stop​

Stop the worker; active rows stay as they are for reconciliation.

async def stop(self) -> None:
...

OperationRunner.enqueue​

Wake the worker after a reservation.

def enqueue(self) -> None:
...

OperationRunner.request_cancel​

Cancel a queued operation now or ask the adapter for a running one.

def request_cancel(self, operation_id: str) -> OperationStatus:
...

OperationService​

Verb dispatcher for one operation descriptor under explicit authority.

authority is decided by the caller from its own grant: read may observe (status, list, log); effect may also start, cancel and plan. The dispatcher validates the verb against the authority before touching the journal, so exposing reads never exposes starts. scope names the descriptor (installed node plus identity); every row is keyed under it, so services sharing one journal cannot see or stop each other's operations.

OperationService.__init__​

def __init__(self, *, scope: str, journal: OperationJournal, logs: OperationLogs, runner: OperationRunner, plan: Callable[[dict[str, JsonValue]], Awaitable[OperationPlan]]) -> None:
...

OperationService.dispatch​

Run one verb and return its typed result.

async def dispatch(self, call: OperationCall, authority: Authority) -> Contract:
...