m3.storage
Signatures use ... for factory-backed or opaque defaults. Model field tables show required status, defaults, constraints, and descriptions.
ACPProbeDimension
m3.storage.ACPProbeDimension(
*,
profile_id: str,
revision_id: str,
probe_type: m3.services.acp_probes.ACPProbeKind,
transport: str = 'stdio',
agent_mode_id: str | None = None,
session_config: collections.abc.Mapping[str, JsonValue] = ...,
) -> NoneExact profile/revision and ACP request dimensions.
Model fields:
| Field | Type | Required | Default | Constraints | Description |
|---|---|---|---|---|---|
profile_id | str | Yes | — | min_length=1, max_length=256 | — |
revision_id | str | Yes | — | min_length=1, max_length=256 | — |
probe_type | m3.services.acp_probes.ACPProbeKind | Yes | — | — | — |
transport | str | No | 'stdio' | min_length=1, max_length=32 | — |
agent_mode_id | str | None | No | None | max_length=256 | — |
session_config | collections.abc.Mapping[str, JsonValue] | No | factory builtins.dict() | — | — |
kind(property)mode_id(property)stable_session_config(property)stable_key(property)
ACPProbeResult
m3.storage.ACPProbeResult(
*,
profile_id: str,
revision_id: str,
probe_type: m3.services.acp_probes.ACPProbeKind,
transport: str = 'stdio',
agent_mode_id: str | None = None,
session_config: collections.abc.Mapping[str, JsonValue] = ...,
id: str = ...,
status: m3.services.acp_probes.ACPProbeStatus,
agent_capabilities: collections.abc.Mapping[str, JsonValue] = ...,
config_options: tuple[collections.abc.Mapping[str, JsonValue], ...] = (),
evidence: collections.abc.Mapping[str, JsonValue] = ...,
diagnostics: str | None = None,
error: str | None = None,
created_at: datetime.datetime = ...,
started_at: datetime.datetime | None = None,
finished_at: datetime.datetime | None = None,
duration_ms: float | None = None,
agent_identity: m3.services.acp_probes.ACPAgentIdentity | None = None,
agent_modes: tuple[m3.services.acp_probes.ACPAgentMode, ...] = (),
current_agent_mode_id: str | None = None,
) -> NoneOne terminal or in-flight probe observation.
Model fields:
| Field | Type | Required | Default | Constraints | Description |
|---|---|---|---|---|---|
profile_id | str | Yes | — | min_length=1, max_length=256 | — |
revision_id | str | Yes | — | min_length=1, max_length=256 | — |
probe_type | m3.services.acp_probes.ACPProbeKind | Yes | — | — | — |
transport | str | No | 'stdio' | min_length=1, max_length=32 | — |
agent_mode_id | str | None | No | None | max_length=256 | — |
session_config | collections.abc.Mapping[str, JsonValue] | No | factory builtins.dict() | — | — |
id | str | No | factory m3.services.acp_probes.ACPProbeResult.<lambda>() | min_length=1, max_length=256 | — |
status | m3.services.acp_probes.ACPProbeStatus | Yes | — | — | — |
agent_capabilities | collections.abc.Mapping[str, JsonValue] | No | factory builtins.dict() | — | — |
config_options | tuple[collections.abc.Mapping[str, JsonValue], ...] | No | () | — | — |
evidence | collections.abc.Mapping[str, JsonValue] | No | factory builtins.dict() | — | — |
diagnostics | str | None | No | None | max_length=65536 | — |
error | str | None | No | None | max_length=2048 | — |
created_at | datetime.datetime | No | factory m3.services.acp_probes.ACPProbeResult.<lambda>() | — | — |
started_at | datetime.datetime | None | No | None | — | — |
finished_at | datetime.datetime | None | No | None | — | — |
duration_ms | float | None | No | None | ge=0 | — |
agent_identity | m3.services.acp_probes.ACPAgentIdentity | None | No | None | — | — |
agent_modes | tuple[m3.services.acp_probes.ACPAgentMode, ...] | No | () | — | — |
current_agent_mode_id | str | None | No | None | max_length=256 | — |
kind(property)mode_id(property)stable_key(property)stable_session_config(property)
ACPProbeStore
m3.storage.ACPProbeStore(
*args,
**kwargs,
)save_acp_probe(
self,
result: ACPProbeResult,
) -> ACPProbeResultget_acp_probe(
self,
probe_id: str,
) -> ACPProbeResult | Nonelist_acp_probes(
self,
dimension: ACPProbeDimension | None = None,
*,
include_inflight: bool = True,
) -> tuple[ACPProbeResult, ...]latest_acp_probe(
self,
dimension: ACPProbeDimension,
) -> ACPProbeResult | NoneArtifactNotFound
m3.storage.ArtifactNotFoundAn artifact or content-addressed blob is not present.
ArtifactStore
m3.storage.ArtifactStore(
*args,
**kwargs,
)Content-addressed artifact/blob store contract.
put(
self,
execution_id: ExecutionId | str,
name: str,
content: bytes,
*,
media_type: str | None = None,
) -> ArtifactRefget(
self,
artifact: ArtifactRef | ArtifactId | str,
) -> bytesget_ref(
self,
artifact_id: ArtifactId | str,
) -> ArtifactRefiter_refs(
self,
execution_id: ExecutionId | str | None = None,
) -> Iterator[ArtifactRef]delete(
self,
artifact: ArtifactRef | ArtifactId | str,
) -> Nonecleanup(
self,
) -> NoneBlobIntegrityError
m3.storage.BlobIntegrityErrorA compressed blob does not match its recorded hash or length.
BlobRecord
m3.storage.BlobRecord(
sha256: str,
size_bytes: int,
compressed_size_bytes: int,
path: Path,
) -> NoneVerified metadata for one compressed content-addressed blob.
Command
m3.storage.Command(
id: str,
execution_id: ExecutionId,
kind: str,
status: str,
payload: Mapping[str, Any],
session_id: SessionId | None = None,
turn_id: TurnId | None = None,
) -> NoneDurableSerializationError
m3.storage.DurableSerializationError(
reason: str = 'durable value is invalid',
) -> NoneA value cannot safely cross a durable specification boundary.
EventCallback
m3.storage.EventCallback(
*args,
**kwargs,
)ExecutionStore
m3.storage.ExecutionStore(
*args,
**kwargs,
)Store contract for immutable execution snapshots and event streams.
create(
self,
snapshot: ExecutionState,
*,
specification: Mapping[str, object] | None = None,
provenance: Mapping[str, object] | None = None,
server_bindings: Sequence[Mapping[str, object]] = (),
harness_binding: Mapping[str, object] | None = None,
parent_execution_id: ExecutionId | str | None = None,
run_id: RunId | str | None = None,
) -> Noneget_snapshot(
self,
execution_id: ExecutionId | str,
) -> ExecutionState | Noneget_execution_spec(
self,
execution_id: ExecutionId | str,
) -> ExecutionSpec | Nonelist_executions(
self,
*,
limit: int = 50,
offset: int = 0,
lifecycle: ExecutionStatus | str | None = None,
outcome: ExecutionOutcome | str | None = None,
run_id: RunId | str | None = None,
suite_id: int | None = None,
project_id: str | None = None,
) -> ExecutionPageget_report(
self,
execution_id: ExecutionId | str,
*,
after_sequence: int = -1,
event_limit: int | None = None,
artifact_limit: int | None = None,
) -> ExecutionReport | Noneget_trace(
self,
execution_id: ExecutionId | str,
) -> TraceResult | Noneget_trace_view(
self,
execution_id: ExecutionId | str,
) -> TraceView | Nonesave_evaluation(
self,
execution_id: ExecutionId | str,
result: EvaluationResult | Mapping[str, object],
*,
evaluation_id: str | None = None,
turn_id: TurnId | str | None = None,
) -> strevaluations(
self,
execution_id: ExecutionId | str,
*,
turn_id: TurnId | str | None = None,
) -> tuple[EvaluationRecord, ...]aggregate_evaluations(
self,
query: EvaluationQuery,
) -> EvaluationReportturns(
self,
execution_id: ExecutionId | str,
) -> tuple[tuple[TurnState, TurnResult | None], ...]save_snapshot(
self,
snapshot: ExecutionState,
) -> Noneappend_events(
self,
events: Sequence[Event],
) -> Noneappend_event(
self,
event: Event,
content: bytes,
*,
media_type: str,
) -> Eventiter_events(
self,
execution_id: ExecutionId | str,
*,
after_sequence: int = -1,
) -> Iterator[Event]transaction(
self,
execution_id: ExecutionId | str,
) -> ExecutionTransactionrelease(
self,
execution_id: ExecutionId | str,
sequences: Sequence[int],
) -> Nonesubscribe(
self,
execution_id: ExecutionId | str,
callback: EventCallback,
) -> Callable[[], None]put_raw_evidence(
self,
event_id: EventId | str,
content: bytes,
*,
media_type: str,
) -> EvidenceCaptureread_raw_evidence(
self,
reference: EvidenceRef,
*,
max_bytes: int = 1048576,
) -> RawEvidencedelete_execution(
self,
execution_id: ExecutionId | str,
) -> NoneExecutionTransaction
m3.storage.ExecutionTransaction(
*args,
**kwargs,
)Uncommitted event batch used by :class:ExecutionStore.
append(
self,
events: Sequence[Event],
) -> Nonecommit(
self,
) -> Nonerollback(
self,
) -> NoneFilesystemBlobStore
m3.storage.FilesystemBlobStore(
root: str | os.PathLike[str],
*,
max_read_bytes: int = 536870912,
) -> NoneAtomic, compressed, content-addressed filesystem storage.
root(property)
path_for(
self,
sha256: str,
) -> Pathput(
self,
content: bytes,
*,
sha256: str | None = None,
size_bytes: int | None = None,
) -> BlobRecordAtomically persist content and return verified metadata.
put_blob(
self,
content: bytes,
*,
sha256: str | None = None,
size_bytes: int | None = None,
) -> BlobRecordAtomically persist content and return verified metadata.
read(
self,
sha256: str,
*,
size_bytes: int | None = None,
) -> bytesget(
self,
sha256: str,
*,
size_bytes: int | None = None,
) -> bytesread_blob(
self,
sha256: str,
*,
size_bytes: int | None = None,
) -> bytesverify(
self,
sha256: str,
size_bytes: int,
) -> BlobRecorditer_records(
self,
) -> tuple[BlobRecord, ...]garbage_collect(
self,
references: Mapping[str, int] | Iterable[str],
) -> tuple[str, ...]Delete only explicitly unreferenced, valid blobs.
collect_garbage(
self,
references: Mapping[str, int] | Iterable[str],
) -> tuple[str, ...]Delete only explicitly unreferenced, valid blobs.
cleanup_temporary_files(
self,
) -> tuple[Path, ...]Remove incomplete temp files left by interrupted writers.
InMemoryArtifactStore
m3.storage.InMemoryArtifactStore(
*,
config: RedactionConfig | None = None,
) -> NoneIn-memory compressed, content-addressed artifact store.
put(
self,
execution_id: ExecutionId | str,
name: str,
content: bytes,
*,
media_type: str | None = None,
) -> ArtifactRefget(
self,
artifact: ArtifactRef | ArtifactId | str,
) -> bytesget_ref(
self,
artifact_id: ArtifactId | str,
) -> ArtifactRefiter_refs(
self,
execution_id: ExecutionId | str | None = None,
) -> Iterator[ArtifactRef]delete(
self,
artifact: ArtifactRef | ArtifactId | str,
) -> Nonecleanup(
self,
) -> Noneclose(
self,
) -> NoneInMemoryExecutionStore
m3.storage.InMemoryExecutionStore(
*,
config: RedactionConfig | None = None,
capture_config: CaptureOptions | None = None,
) -> NoneThread-safe execution metadata store with commit-gated visibility.
save_acp_probe(
self,
result: ACPProbeResult,
) -> ACPProbeResultget_acp_probe(
self,
probe_id: str,
) -> ACPProbeResult | Nonelist_acp_probes(
self,
dimension: ACPProbeDimension | None = None,
*,
include_inflight: bool = True,
) -> tuple[ACPProbeResult, ...]latest_acp_probe(
self,
dimension: ACPProbeDimension,
) -> ACPProbeResult | Nonecreate(
self,
snapshot: ExecutionState,
*,
specification: Mapping[str, object] | None = None,
provenance: Mapping[str, object] | None = None,
server_bindings: Sequence[Mapping[str, object]] = (),
harness_binding: Mapping[str, object] | None = None,
parent_execution_id: ExecutionId | str | None = None,
run_id: RunId | str | None = None,
) -> Noneensure_project(
self,
project_id: str,
project_name: str,
) -> tuple[str, str]get_project(
self,
project_id: str,
) -> tuple[str, str] | Noneensure_suite(
self,
suite_name: str,
project_id: str | None = None,
) -> Suiteget_suite(
self,
suite_id: SuiteId | str,
) -> Suite | Noneget_suite_by_name(
self,
suite_name: str,
project_id: str | None = None,
) -> Suite | Nonelist_suites(
self,
) -> tuple[Suite, ...]create_execution(
self,
snapshot: ExecutionState,
*,
specification: Mapping[str, object] | None = None,
provenance: Mapping[str, object] | None = None,
server_bindings: Sequence[Mapping[str, object]] = (),
harness_binding: Mapping[str, object] | None = None,
parent_execution_id: ExecutionId | str | None = None,
run_id: RunId | str | None = None,
) -> Nonesave_test_run(
self,
run_id: str,
value: Mapping[str, object],
) -> Noneget_test_run(
self,
run_id: str,
) -> Mapping[str, object] | Nonelist_test_runs(
self,
) -> tuple[Mapping[str, object], ...]list_test_run_page(
self,
*,
limit: int | None = None,
offset: int = 0,
suite_id: int | None = None,
project_id: str | None = None,
q: str | None = None,
) -> tuple[tuple[Mapping[str, object], ...], int]Return newest-first run manifests with their suites, and the total.
save_test_result(
self,
run_id: str,
attempt_id: str,
value: Mapping[str, object],
) -> Nonelist_test_results(
self,
run_id: str,
) -> tuple[Mapping[str, object], ...]get_snapshot(
self,
execution_id: ExecutionId | str,
) -> ExecutionState | Noneget_execution_spec(
self,
execution_id: ExecutionId | str,
) -> ExecutionSpec | Nonelist_executions(
self,
*,
limit: int = 50,
offset: int = 0,
lifecycle: ExecutionStatus | str | None = None,
outcome: ExecutionOutcome | str | None = None,
run_id: RunId | str | None = None,
suite_id: int | None = None,
project_id: str | None = None,
) -> ExecutionPageget_report(
self,
execution_id: ExecutionId | str,
*,
after_sequence: int = -1,
event_limit: int | None = None,
artifact_limit: int | None = None,
) -> ExecutionReport | Nonesave_evaluation(
self,
execution_id: ExecutionId | str,
result: EvaluationResult | Mapping[str, object],
*,
evaluation_id: str | None = None,
turn_id: object | None = None,
) -> strsave_turn(
self,
snapshot: TurnState,
result: TurnResult | None = None,
) -> Noneappend_turn(
self,
snapshot: TurnState,
result: TurnResult | None = None,
) -> Noneturns(
self,
execution_id: ExecutionId | str,
) -> tuple[tuple[TurnState, TurnResult | None], ...]evaluations(
self,
execution_id: ExecutionId | str,
*,
turn_id: object | None = None,
) -> tuple[EvaluationRecord, ...]aggregate_evaluations(
self,
query: EvaluationQuery,
) -> EvaluationReportCalculate summaries from the evaluations currently in memory.
get_trace(
self,
execution_id: ExecutionId | str,
) -> TraceResult | Noneget_trace_view(
self,
execution_id: ExecutionId | str,
) -> TraceView | Nonesave_snapshot(
self,
snapshot: ExecutionState,
) -> Noneupdate_snapshot(
self,
snapshot: ExecutionState,
) -> Noneappend_events(
self,
events: Sequence[Event],
) -> Noneappend(
self,
events: Sequence[Event],
) -> Noneappend_event(
self,
event: Event,
content: bytes,
*,
media_type: str,
) -> EventCommit an event and its raw blob as one in-memory operation.
iter_events(
self,
execution_id: ExecutionId | str,
*,
after_sequence: int = -1,
) -> Iterator[Event]events(
self,
execution_id: ExecutionId | str,
*,
after_sequence: int = -1,
) -> tuple[Event, ...]transaction(
self,
execution_id: ExecutionId | str,
) -> ExecutionTransactionsubscribe(
self,
execution_id: ExecutionId | str,
callback: EventCallback,
) -> Callable[[], None]allocate(
self,
execution_id: ExecutionId | str,
*,
count: int = 1,
) -> tuple[int, ...]Reserve the next per-execution sequence numbers.
allocate_sequence(
self,
execution_id: ExecutionId | str,
) -> intallocate_sequences(
self,
execution_id: ExecutionId | str,
*,
count: int = 1,
) -> tuple[int, ...]Reserve the next per-execution sequence numbers.
release(
self,
execution_id: ExecutionId | str,
sequences: Sequence[int],
) -> NoneRelease uncommitted reservations after producer cancellation.
release_sequences(
self,
execution_id: ExecutionId | str,
sequences: Sequence[int],
) -> NoneRelease uncommitted reservations after producer cancellation.
close(
self,
) -> Noneput_raw_evidence(
self,
event_id: EventId | str,
content: bytes,
*,
media_type: str,
) -> EvidenceCaptureread_raw_evidence(
self,
reference: EvidenceRef,
*,
max_bytes: int = 1048576,
) -> RawEvidencedelete_execution(
self,
execution_id: ExecutionId | str,
) -> Nonedelete(
self,
execution_id: ExecutionId | str,
) -> NoneLease
m3.storage.Lease(
execution_id: ExecutionId,
owner_id: str,
lease_token: str,
expires_at: datetime,
) -> NoneManagedInputLease
m3.storage.ManagedInputLease(
*,
owner_id: str,
lease_token: str,
expires_at: datetime.datetime,
) -> NoneThe compare-and-set token currently allowed to mutate a round.
Model fields:
| Field | Type | Required | Default | Constraints | Description |
|---|---|---|---|---|---|
owner_id | str | Yes | — | min_length=1 | — |
lease_token | str | Yes | — | min_length=1 | — |
expires_at | datetime.datetime | Yes | — | — | — |
ManagedInputRecord
m3.storage.ManagedInputRecord(
*,
pending: m3.elicitation.PendingElicitationRound,
round_index: int,
round_limit: int,
status: Literal['pending', 'response_validated', 'delivery_started', 'delivered', 'resolved', 'failed'],
lease: m3.storage.managed_input.ManagedInputLease,
responses: collections.abc.Mapping[str, m3.elicitation.ElicitationResponse] | None = None,
response_idempotency_key: str | None = None,
harness_session_id: str | None = None,
native_resume_token: str | None = None,
delivery_idempotency_key: str | None = None,
session_id: str | None = None,
turn_id: str | None = None,
operation_parameters: collections.abc.Mapping[str, object] = ...,
delivery_attempts: int = 0,
created_at: datetime.datetime,
updated_at: datetime.datetime,
response_validated_at: datetime.datetime | None = None,
delivery_started_at: datetime.datetime | None = None,
delivered_at: datetime.datetime | None = None,
resolved_at: datetime.datetime | None = None,
failed_at: datetime.datetime | None = None,
failure_code: str | None = None,
failure_message: str | None = None,
) -> NoneImmutable view of one durable managed-input round.
Model fields:
| Field | Type | Required | Default | Constraints | Description |
|---|---|---|---|---|---|
pending | m3.elicitation.PendingElicitationRound | Yes | — | — | — |
round_index | int | Yes | — | ge=0 | — |
round_limit | int | Yes | — | gt=0 | — |
status | Literal['pending', 'response_validated', 'delivery_started', 'delivered', 'resolved', 'failed'] | Yes | — | — | — |
lease | m3.storage.managed_input.ManagedInputLease | Yes | — | — | — |
responses | collections.abc.Mapping[str, m3.elicitation.ElicitationResponse] | None | No | None | — | — |
response_idempotency_key | str | None | No | None | — | — |
harness_session_id | str | None | No | None | — | — |
native_resume_token | str | None | No | None | — | — |
delivery_idempotency_key | str | None | No | None | — | — |
session_id | str | None | No | None | — | — |
turn_id | str | None | No | None | — | — |
operation_parameters | collections.abc.Mapping[str, object] | No | factory builtins.dict() | — | — |
delivery_attempts | int | No | 0 | ge=0 | — |
created_at | datetime.datetime | Yes | — | — | — |
updated_at | datetime.datetime | Yes | — | — | — |
response_validated_at | datetime.datetime | None | No | None | — | — |
delivery_started_at | datetime.datetime | None | No | None | — | — |
delivered_at | datetime.datetime | None | No | None | — | — |
resolved_at | datetime.datetime | None | No | None | — | — |
failed_at | datetime.datetime | None | No | None | — | — |
failure_code | str | None | No | None | — | — |
failure_message | str | None | No | None | — | — |
execution_id(property)round_id(property)request_state(property)lease_token(property)owner_id(property)
ManagedInputStatus
m3.storage.ManagedInputStatus(
*args,
**kwargs,
)ManagedInputStore
m3.storage.ManagedInputStore(
*args,
**kwargs,
)Storage contract consumed by future execution handles.
create_round(
self,
pending: PendingElicitationRound,
*,
round_index: int,
round_limit: int,
owner_id: str,
lease_seconds: float,
lease_token: str | None = None,
harness_session_id: str | None = None,
native_resume_token: str | None = None,
delivery_idempotency_key: str | None = None,
session_id: str | None = None,
turn_id: str | None = None,
operation_parameters: Mapping[str, object] | None = None,
) -> ManagedInputRecordget_round(
self,
execution_id: str,
round_id: str,
) -> ManagedInputRecord | Nonelist_rounds(
self,
execution_id: str,
) -> tuple[ManagedInputRecord, ...]claim_round(
self,
execution_id: str,
round_id: str,
*,
owner_id: str,
lease_seconds: float,
expected_lease_token: str | None = None,
delivery_state: Literal['not_started', 'not_delivered', 'delivered'] | None = None,
idempotent_delivery: bool = False,
) -> ManagedInputRecordrenew(
self,
execution_id: str,
round_id: str,
*,
owner_id: str,
lease_token: str,
lease_seconds: float,
) -> ManagedInputRecordsubmit_responses(
self,
execution_id: str,
round_id: str,
responses: Mapping[str, ElicitationResponse],
*,
owner_id: str,
lease_token: str,
response_idempotency_key: str,
) -> ManagedInputRecordstart_delivery(
self,
execution_id: str,
round_id: str,
*,
owner_id: str,
lease_token: str,
) -> ManagedInputRecordmark_delivered(
self,
execution_id: str,
round_id: str,
*,
owner_id: str,
lease_token: str,
) -> ManagedInputRecordresolve(
self,
execution_id: str,
round_id: str,
*,
owner_id: str,
lease_token: str,
) -> ManagedInputRecordfail(
self,
execution_id: str,
round_id: str,
*,
owner_id: str,
lease_token: str,
code: str,
message: str,
) -> ManagedInputRecordfail_recovery(
self,
execution_id: str,
round_id: str,
*,
expected_lease_token: str,
message: str,
) -> ManagedInputRecordredacted_responses(
self,
execution_id: str,
round_id: str,
) -> Mapping[str, object] | NonePersistentExecutionStore
m3.storage.PersistentExecutionStore(
database: str | Path,
*,
blob_root: str | Path | None = None,
config: RedactionConfig | None = None,
capture_config: CaptureOptions | None = None,
payload_blob_threshold: int = 65536,
**kwargs: Any,
) -> NoneSQLite implementation of the public :class:ExecutionStore protocol.
managed_input_store(property): Lazily open the managed-input tables for opted-in executions.
resolve_managed_input_store(
self,
) -> SQLiteManagedInputStoreResolve managed-input storage lazily for opted-in executions.
ensure_project(
self,
project_id: str,
project_name: str,
) -> tuple[str, str]Register a stable project identity and refresh its display name.
get_project(
self,
project_id: str,
) -> tuple[str, str] | Noneensure_suite(
self,
suite_name: str,
project_id: str | None = None,
) -> Suiteget_suite(
self,
suite_id: SuiteId | str,
) -> Suite | Noneget_suite_by_name(
self,
suite_name: str,
project_id: str | None = None,
) -> Suite | Nonelist_suites(
self,
) -> tuple[Suite, ...]create(
self,
snapshot: ExecutionState,
*,
specification: Mapping[str, Any] | None = None,
provenance: Mapping[str, Any] | None = None,
server_bindings: Sequence[Mapping[str, Any]] = (),
harness_binding: Mapping[str, Any] | None = None,
parent_execution_id: ExecutionId | str | None = None,
run_id: RunId | str | None = None,
) -> Nonecreate_execution(
self,
snapshot: ExecutionState,
*,
specification: Mapping[str, Any] | None = None,
provenance: Mapping[str, Any] | None = None,
server_bindings: Sequence[Mapping[str, Any]] = (),
harness_binding: Mapping[str, Any] | None = None,
parent_execution_id: ExecutionId | str | None = None,
run_id: RunId | str | None = None,
) -> Nonesave_test_run(
self,
run_id: str,
value: Mapping[str, object],
) -> Noneget_test_run(
self,
run_id: str,
) -> Mapping[str, object] | Nonelist_test_runs(
self,
) -> tuple[Mapping[str, object], ...]list_test_run_page(
self,
*,
limit: int | None = None,
offset: int = 0,
suite_id: int | None = None,
project_id: str | None = None,
q: str | None = None,
) -> tuple[tuple[Mapping[str, object], ...], int]Return newest-first run manifests with their suites, and the total.
save_test_result(
self,
run_id: str,
attempt_id: str,
value: Mapping[str, object],
) -> Nonelist_test_results(
self,
run_id: str,
) -> tuple[Mapping[str, object], ...]get_snapshot(
self,
execution_id: ExecutionId | str,
) -> ExecutionState | Noneget_execution_spec(
self,
execution_id: ExecutionId | str,
) -> ExecutionSpec | NoneReturn the immutable typed submission spec, if one was saved.
get_spec(
self,
execution_id: ExecutionId | str,
) -> ExecutionSpec | NoneReturn the immutable typed submission spec, if one was saved.
list_executions(
self,
*,
limit: int = 50,
offset: int = 0,
lifecycle: ExecutionStatus | str | None = None,
outcome: ExecutionOutcome | str | None = None,
run_id: str | None = None,
suite_id: int | None = None,
project_id: str | None = None,
) -> ExecutionPageget_report(
self,
execution_id: ExecutionId | str,
*,
after_sequence: int = -1,
event_limit: int | None = None,
artifact_limit: int | None = None,
) -> ExecutionReport | Noneget_trace(
self,
execution_id: ExecutionId | str,
) -> TraceResult | Noneget_trace_view(
self,
execution_id: ExecutionId | str,
) -> TraceView | Nonesave_snapshot(
self,
snapshot: ExecutionState,
) -> Noneupdate_snapshot(
self,
snapshot: ExecutionState,
) -> Noneappend_events(
self,
events: Sequence[Event],
) -> Noneappend(
self,
events: Sequence[Event],
) -> Noneappend_event(
self,
event: Event,
content: bytes,
*,
media_type: str,
) -> EventCommit an event and its raw blob in one SQLite transaction.
iter_events(
self,
execution_id: ExecutionId | str,
*,
after_sequence: int = -1,
) -> Iterator[Event]events(
self,
execution_id: ExecutionId | str,
*,
after_sequence: int = -1,
) -> tuple[Event, ...]create_session(
self,
execution_id: ExecutionId | str,
session_id: SessionId | str,
*,
state: str = 'open',
) -> SessionIdclose_session(
self,
session_id: SessionId | str,
) -> Nonesave_turn(
self,
snapshot: TurnState,
result: TurnResult | Mapping[str, Any] | None = None,
) -> Noneappend_turn(
self,
snapshot: TurnState,
result: TurnResult | Mapping[str, Any] | None = None,
) -> Noneturns(
self,
execution_id: ExecutionId | str,
) -> tuple[tuple[TurnState, TurnResult | None], ...]save_evaluation(
self,
execution_id: ExecutionId | str,
result: Mapping[str, Any] | Any,
*,
evaluation_id: str | None = None,
turn_id: TurnId | str | None = None,
) -> strreserve_judge_request(
self,
run_id: str,
limit: int,
) -> boolAtomically reserve one judge request for a run across workers.
save(
self,
result: EvaluationResult,
) -> Noneget(
self,
evaluation_id: str,
) -> EvaluationResult | Noneall(
self,
) -> tuple[EvaluationResult, ...]evaluations(
self,
execution_id: ExecutionId | str,
*,
turn_id: TurnId | str | None = None,
) -> tuple[EvaluationRecord, ...]aggregate_evaluations(
self,
query: EvaluationQuery,
) -> EvaluationReportCalculate summaries from persisted evaluations and execution traces.
persisted_evaluations(
self,
execution_id: ExecutionId | str,
*,
turn_id: TurnId | str | None = None,
) -> tuple[EvaluationRecord, ...]evaluation_json(
self,
execution_id: ExecutionId | str,
*,
turn_id: TurnId | str | None = None,
) -> tuple[Mapping[str, Any], ...]transaction(
self,
execution_id: ExecutionId | str,
) -> ExecutionTransactionallocate(
self,
execution_id: ExecutionId | str,
*,
count: int = 1,
) -> tuple[int, ...]allocate_sequence(
self,
execution_id,
)allocate_sequences(
self,
execution_id: ExecutionId | str,
*,
count: int = 1,
) -> tuple[int, ...]release(
self,
execution_id: ExecutionId | str,
sequences: Sequence[int],
) -> Nonerelease_sequences(
self,
execution_id: ExecutionId | str,
sequences: Sequence[int],
) -> Nonesubscribe(
self,
execution_id: ExecutionId | str,
callback: EventCallback,
) -> Callable[[], None]save_acp_probe(
self,
result: ACPProbeResult,
) -> ACPProbeResultPersist one redacted ACP probe result and return its safe copy.
get_acp_probe(
self,
probe_id: str,
) -> ACPProbeResult | Nonelist_acp_probes(
self,
dimension: ACPProbeDimension | None = None,
*,
include_inflight: bool = True,
) -> tuple[ACPProbeResult, ...]latest_acp_probe(
self,
dimension: ACPProbeDimension,
) -> ACPProbeResult | Nonecreate_profile(
self,
kind: str,
name: str,
value: Mapping[str, Any],
*,
description: str = '',
profile_id: str | None = None,
revision_id: str | None = None,
) -> ProfileRecordcreate_server_profile(
self,
name: str,
value: Mapping[str, Any],
**kwargs: Any,
) -> ProfileRecordcreate_harness_profile(
self,
name: str,
value: Mapping[str, Any],
**kwargs: Any,
) -> ProfileRecordget_profile(
self,
profile_id: str,
*,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRecord | Noneresolve_profile(
self,
profile_id: str,
selection: RevisionSelection,
*,
kind: Literal['server', 'harness'],
) -> tuple[ProfileRecord | None, ProfileRevisionRecord | None]Read one profile and its selected immutable revision.
list_profiles(
self,
kind: str,
*,
include_archived: bool = False,
) -> tuple[ProfileRecord, ...]List one profile family in stable name/id order.
list_profile_revisions(
self,
profile_id: str,
*,
kind: Literal['server', 'harness'] | None = None,
) -> tuple[ProfileRevisionRecord, ...]Return every revision ordered by revision number then id.
update_profile(
self,
profile_id: str,
*,
name: str | None = None,
description: str | None = None,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRecordUpdate mutable metadata without changing the immutable revision.
add_revision(
self,
profile_id: str,
value: Mapping[str, Any],
*,
revision_id: str | None = None,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRevisionRecordarchive_profile(
self,
profile_id: str,
*,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRecordrestore_profile(
self,
profile_id: str,
*,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRecordresolve_revision(
self,
profile_id: str,
selection: RevisionSelection | str = 'latest',
*,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRevisionRecordget_revision(
self,
revision_id: RevisionId | str,
*,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRevisionRecord | Noneenqueue_command(
self,
execution_id: ExecutionId | str,
kind: str = 'execution',
payload: Mapping[str, Any] | None = None,
*,
command_id: str | None = None,
session_id: SessionId | str | None = None,
turn_id: TurnId | str | None = None,
) -> Commandenqueue(
self,
execution_id: ExecutionId | str,
kind: str = 'execution',
payload: Mapping[str, Any] | None = None,
*,
command_id: str | None = None,
session_id: SessionId | str | None = None,
turn_id: TurnId | str | None = None,
) -> Commandget_command(
self,
command_id: str,
) -> Command | Noneclaim_next(
self,
owner_id: str,
*,
lease_seconds: float = 30.0,
) -> tuple[Command, Lease] | Noneclaim(
self,
owner_id: str,
*,
lease_seconds: float = 30.0,
) -> tuple[Command, Lease] | Noneheartbeat(
self,
lease: Lease | str,
*,
owner_id: str | None = None,
lease_seconds: float = 30.0,
) -> boolrenew_lease(
self,
lease: Lease | str,
*,
owner_id: str | None = None,
lease_seconds: float = 30.0,
) -> boolmark_interrupted_if_lease_lost(
self,
lease: Lease,
*,
reason: str = 'worker lease lost',
) -> boolAtomically close work whose owner can no longer renew its lease.
mark_managed_recovery_unavailable_if_lease_lost(
self,
lease: Lease,
*,
reason: str = 'managed interaction cannot be resumed safely after worker loss',
) -> boolTerminalize managed input when its owning worker is lost.
release_lease(
self,
lease: Lease | str,
*,
owner_id: str | None = None,
) -> boolcomplete_command(
self,
command_id: str,
*,
owner_id: str,
lease_token: str,
status: str = 'done',
) -> boolMark one claimed command terminal under its current lease.
request_cancel(
self,
execution_id: ExecutionId | str,
reason: str | None = None,
) -> boolcancel(
self,
execution_id: ExecutionId | str,
reason: str | None = None,
) -> boolfinalize_cancelled(
self,
execution_id: ExecutionId | str,
*,
reason: str = 'cancelled',
) -> boolPersist a cancellation terminal event when a worker observed it.
cancellation_requested(
self,
execution_id: ExecutionId | str,
) -> boolmark_stale_interrupted(
self,
*,
now: datetime | None = None,
) -> tuple[ExecutionId, ...]delete_execution(
self,
execution_id: ExecutionId | str,
) -> Nonedelete(
self,
execution_id: ExecutionId | str,
) -> Noneput_raw_evidence(
self,
event_id: EventId | str,
content: bytes,
*,
media_type: str,
) -> EvidenceCaptureRedact, bound, and durably associate evidence with one event.
read_raw_evidence(
self,
reference: EvidenceRef,
*,
max_bytes: int = 1048576,
) -> RawEvidenceclone_execution(
self,
execution_id: ExecutionId | str,
*,
use_latest: bool = False,
) -> ExecutionIdclone(
self,
execution_id: ExecutionId | str,
*,
use_latest: bool = False,
) -> ExecutionIdresolved_bindings(
self,
execution_id: ExecutionId | str,
) -> Mapping[str, Any]Return the immutable, submission-time profile binding snapshot.
close(
self,
) -> NoneProfileRecord
m3.storage.ProfileRecord(
id: str,
kind: str,
name: str,
description: str,
archived: bool,
current_revision_id: RevisionId | None,
created_at: datetime,
updated_at: datetime,
) -> NoneProfileResolver
m3.storage.ProfileResolver(
*args,
**kwargs,
)Read-only saved-profile lookup used by execution runtimes.
resolve_profile(
self,
profile_id: str,
selection: RevisionSelection,
*,
kind: Literal['server', 'harness'],
) -> tuple[Any, Any]ProfileRevisionRecord
m3.storage.ProfileRevisionRecord(
id: RevisionId,
profile_id: str,
revision_number: int,
value: Mapping[str, Any],
created_at: datetime,
) -> NoneSQLiteArtifactStore
m3.storage.SQLiteArtifactStore(
database: str | Path,
blob_root: str | Path | None = None,
**kwargs: Any,
) -> NoneFilesystem-backed content-addressed artifact store with SQLite refs.
put(
self,
execution_id: ExecutionId | str,
name: str,
content: bytes,
*,
media_type: str | None = None,
) -> ArtifactRefget(
self,
artifact: ArtifactRef | ArtifactId | str,
) -> bytesget_ref(
self,
artifact_id: ArtifactId | str,
) -> ArtifactRefiter_refs(
self,
execution_id: ExecutionId | str | None = None,
) -> Iterator[ArtifactRef]delete(
self,
artifact: ArtifactRef | ArtifactId | str,
) -> Nonecleanup(
self,
) -> Noneblob_store(property): The app-owned content-addressed store used by metadata rows.
SQLiteExecutionStore
m3.storage.SQLiteExecutionStore(
database: str | Path,
*,
blob_root: str | Path | None = None,
config: RedactionConfig | None = None,
capture_config: CaptureOptions | None = None,
payload_blob_threshold: int = 65536,
**kwargs: Any,
) -> NoneSQLite implementation of the public :class:ExecutionStore protocol.
managed_input_store(property): Lazily open the managed-input tables for opted-in executions.
resolve_managed_input_store(
self,
) -> SQLiteManagedInputStoreResolve managed-input storage lazily for opted-in executions.
ensure_project(
self,
project_id: str,
project_name: str,
) -> tuple[str, str]Register a stable project identity and refresh its display name.
get_project(
self,
project_id: str,
) -> tuple[str, str] | Noneensure_suite(
self,
suite_name: str,
project_id: str | None = None,
) -> Suiteget_suite(
self,
suite_id: SuiteId | str,
) -> Suite | Noneget_suite_by_name(
self,
suite_name: str,
project_id: str | None = None,
) -> Suite | Nonelist_suites(
self,
) -> tuple[Suite, ...]create(
self,
snapshot: ExecutionState,
*,
specification: Mapping[str, Any] | None = None,
provenance: Mapping[str, Any] | None = None,
server_bindings: Sequence[Mapping[str, Any]] = (),
harness_binding: Mapping[str, Any] | None = None,
parent_execution_id: ExecutionId | str | None = None,
run_id: RunId | str | None = None,
) -> Nonecreate_execution(
self,
snapshot: ExecutionState,
*,
specification: Mapping[str, Any] | None = None,
provenance: Mapping[str, Any] | None = None,
server_bindings: Sequence[Mapping[str, Any]] = (),
harness_binding: Mapping[str, Any] | None = None,
parent_execution_id: ExecutionId | str | None = None,
run_id: RunId | str | None = None,
) -> Nonesave_test_run(
self,
run_id: str,
value: Mapping[str, object],
) -> Noneget_test_run(
self,
run_id: str,
) -> Mapping[str, object] | Nonelist_test_runs(
self,
) -> tuple[Mapping[str, object], ...]list_test_run_page(
self,
*,
limit: int | None = None,
offset: int = 0,
suite_id: int | None = None,
project_id: str | None = None,
q: str | None = None,
) -> tuple[tuple[Mapping[str, object], ...], int]Return newest-first run manifests with their suites, and the total.
save_test_result(
self,
run_id: str,
attempt_id: str,
value: Mapping[str, object],
) -> Nonelist_test_results(
self,
run_id: str,
) -> tuple[Mapping[str, object], ...]get_snapshot(
self,
execution_id: ExecutionId | str,
) -> ExecutionState | Noneget_execution_spec(
self,
execution_id: ExecutionId | str,
) -> ExecutionSpec | NoneReturn the immutable typed submission spec, if one was saved.
get_spec(
self,
execution_id: ExecutionId | str,
) -> ExecutionSpec | NoneReturn the immutable typed submission spec, if one was saved.
list_executions(
self,
*,
limit: int = 50,
offset: int = 0,
lifecycle: ExecutionStatus | str | None = None,
outcome: ExecutionOutcome | str | None = None,
run_id: str | None = None,
suite_id: int | None = None,
project_id: str | None = None,
) -> ExecutionPageget_report(
self,
execution_id: ExecutionId | str,
*,
after_sequence: int = -1,
event_limit: int | None = None,
artifact_limit: int | None = None,
) -> ExecutionReport | Noneget_trace(
self,
execution_id: ExecutionId | str,
) -> TraceResult | Noneget_trace_view(
self,
execution_id: ExecutionId | str,
) -> TraceView | Nonesave_snapshot(
self,
snapshot: ExecutionState,
) -> Noneupdate_snapshot(
self,
snapshot: ExecutionState,
) -> Noneappend_events(
self,
events: Sequence[Event],
) -> Noneappend(
self,
events: Sequence[Event],
) -> Noneappend_event(
self,
event: Event,
content: bytes,
*,
media_type: str,
) -> EventCommit an event and its raw blob in one SQLite transaction.
iter_events(
self,
execution_id: ExecutionId | str,
*,
after_sequence: int = -1,
) -> Iterator[Event]events(
self,
execution_id: ExecutionId | str,
*,
after_sequence: int = -1,
) -> tuple[Event, ...]create_session(
self,
execution_id: ExecutionId | str,
session_id: SessionId | str,
*,
state: str = 'open',
) -> SessionIdclose_session(
self,
session_id: SessionId | str,
) -> Nonesave_turn(
self,
snapshot: TurnState,
result: TurnResult | Mapping[str, Any] | None = None,
) -> Noneappend_turn(
self,
snapshot: TurnState,
result: TurnResult | Mapping[str, Any] | None = None,
) -> Noneturns(
self,
execution_id: ExecutionId | str,
) -> tuple[tuple[TurnState, TurnResult | None], ...]save_evaluation(
self,
execution_id: ExecutionId | str,
result: Mapping[str, Any] | Any,
*,
evaluation_id: str | None = None,
turn_id: TurnId | str | None = None,
) -> strreserve_judge_request(
self,
run_id: str,
limit: int,
) -> boolAtomically reserve one judge request for a run across workers.
save(
self,
result: EvaluationResult,
) -> Noneget(
self,
evaluation_id: str,
) -> EvaluationResult | Noneall(
self,
) -> tuple[EvaluationResult, ...]evaluations(
self,
execution_id: ExecutionId | str,
*,
turn_id: TurnId | str | None = None,
) -> tuple[EvaluationRecord, ...]aggregate_evaluations(
self,
query: EvaluationQuery,
) -> EvaluationReportCalculate summaries from persisted evaluations and execution traces.
persisted_evaluations(
self,
execution_id: ExecutionId | str,
*,
turn_id: TurnId | str | None = None,
) -> tuple[EvaluationRecord, ...]evaluation_json(
self,
execution_id: ExecutionId | str,
*,
turn_id: TurnId | str | None = None,
) -> tuple[Mapping[str, Any], ...]transaction(
self,
execution_id: ExecutionId | str,
) -> ExecutionTransactionallocate(
self,
execution_id: ExecutionId | str,
*,
count: int = 1,
) -> tuple[int, ...]allocate_sequence(
self,
execution_id,
)allocate_sequences(
self,
execution_id: ExecutionId | str,
*,
count: int = 1,
) -> tuple[int, ...]release(
self,
execution_id: ExecutionId | str,
sequences: Sequence[int],
) -> Nonerelease_sequences(
self,
execution_id: ExecutionId | str,
sequences: Sequence[int],
) -> Nonesubscribe(
self,
execution_id: ExecutionId | str,
callback: EventCallback,
) -> Callable[[], None]save_acp_probe(
self,
result: ACPProbeResult,
) -> ACPProbeResultPersist one redacted ACP probe result and return its safe copy.
get_acp_probe(
self,
probe_id: str,
) -> ACPProbeResult | Nonelist_acp_probes(
self,
dimension: ACPProbeDimension | None = None,
*,
include_inflight: bool = True,
) -> tuple[ACPProbeResult, ...]latest_acp_probe(
self,
dimension: ACPProbeDimension,
) -> ACPProbeResult | Nonecreate_profile(
self,
kind: str,
name: str,
value: Mapping[str, Any],
*,
description: str = '',
profile_id: str | None = None,
revision_id: str | None = None,
) -> ProfileRecordcreate_server_profile(
self,
name: str,
value: Mapping[str, Any],
**kwargs: Any,
) -> ProfileRecordcreate_harness_profile(
self,
name: str,
value: Mapping[str, Any],
**kwargs: Any,
) -> ProfileRecordget_profile(
self,
profile_id: str,
*,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRecord | Noneresolve_profile(
self,
profile_id: str,
selection: RevisionSelection,
*,
kind: Literal['server', 'harness'],
) -> tuple[ProfileRecord | None, ProfileRevisionRecord | None]Read one profile and its selected immutable revision.
list_profiles(
self,
kind: str,
*,
include_archived: bool = False,
) -> tuple[ProfileRecord, ...]List one profile family in stable name/id order.
list_profile_revisions(
self,
profile_id: str,
*,
kind: Literal['server', 'harness'] | None = None,
) -> tuple[ProfileRevisionRecord, ...]Return every revision ordered by revision number then id.
update_profile(
self,
profile_id: str,
*,
name: str | None = None,
description: str | None = None,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRecordUpdate mutable metadata without changing the immutable revision.
add_revision(
self,
profile_id: str,
value: Mapping[str, Any],
*,
revision_id: str | None = None,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRevisionRecordarchive_profile(
self,
profile_id: str,
*,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRecordrestore_profile(
self,
profile_id: str,
*,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRecordresolve_revision(
self,
profile_id: str,
selection: RevisionSelection | str = 'latest',
*,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRevisionRecordget_revision(
self,
revision_id: RevisionId | str,
*,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRevisionRecord | Noneenqueue_command(
self,
execution_id: ExecutionId | str,
kind: str = 'execution',
payload: Mapping[str, Any] | None = None,
*,
command_id: str | None = None,
session_id: SessionId | str | None = None,
turn_id: TurnId | str | None = None,
) -> Commandenqueue(
self,
execution_id: ExecutionId | str,
kind: str = 'execution',
payload: Mapping[str, Any] | None = None,
*,
command_id: str | None = None,
session_id: SessionId | str | None = None,
turn_id: TurnId | str | None = None,
) -> Commandget_command(
self,
command_id: str,
) -> Command | Noneclaim_next(
self,
owner_id: str,
*,
lease_seconds: float = 30.0,
) -> tuple[Command, Lease] | Noneclaim(
self,
owner_id: str,
*,
lease_seconds: float = 30.0,
) -> tuple[Command, Lease] | Noneheartbeat(
self,
lease: Lease | str,
*,
owner_id: str | None = None,
lease_seconds: float = 30.0,
) -> boolrenew_lease(
self,
lease: Lease | str,
*,
owner_id: str | None = None,
lease_seconds: float = 30.0,
) -> boolmark_interrupted_if_lease_lost(
self,
lease: Lease,
*,
reason: str = 'worker lease lost',
) -> boolAtomically close work whose owner can no longer renew its lease.
mark_managed_recovery_unavailable_if_lease_lost(
self,
lease: Lease,
*,
reason: str = 'managed interaction cannot be resumed safely after worker loss',
) -> boolTerminalize managed input when its owning worker is lost.
release_lease(
self,
lease: Lease | str,
*,
owner_id: str | None = None,
) -> boolcomplete_command(
self,
command_id: str,
*,
owner_id: str,
lease_token: str,
status: str = 'done',
) -> boolMark one claimed command terminal under its current lease.
request_cancel(
self,
execution_id: ExecutionId | str,
reason: str | None = None,
) -> boolcancel(
self,
execution_id: ExecutionId | str,
reason: str | None = None,
) -> boolfinalize_cancelled(
self,
execution_id: ExecutionId | str,
*,
reason: str = 'cancelled',
) -> boolPersist a cancellation terminal event when a worker observed it.
cancellation_requested(
self,
execution_id: ExecutionId | str,
) -> boolmark_stale_interrupted(
self,
*,
now: datetime | None = None,
) -> tuple[ExecutionId, ...]delete_execution(
self,
execution_id: ExecutionId | str,
) -> Nonedelete(
self,
execution_id: ExecutionId | str,
) -> Noneput_raw_evidence(
self,
event_id: EventId | str,
content: bytes,
*,
media_type: str,
) -> EvidenceCaptureRedact, bound, and durably associate evidence with one event.
read_raw_evidence(
self,
reference: EvidenceRef,
*,
max_bytes: int = 1048576,
) -> RawEvidenceclone_execution(
self,
execution_id: ExecutionId | str,
*,
use_latest: bool = False,
) -> ExecutionIdclone(
self,
execution_id: ExecutionId | str,
*,
use_latest: bool = False,
) -> ExecutionIdresolved_bindings(
self,
execution_id: ExecutionId | str,
) -> Mapping[str, Any]Return the immutable, submission-time profile binding snapshot.
close(
self,
) -> NoneSQLiteManagedInputStore
m3.storage.SQLiteManagedInputStore(
database: str | Path,
*,
busy_timeout_ms: int = 5000,
clock: Callable[[], datetime] = m3.storage.managed_input._now,
) -> NoneSQLite-backed managed-input storage sharing a database path safely.
get_round(
self,
execution_id: str,
round_id: str,
) -> ManagedInputRecord | Nonelist_rounds(
self,
execution_id: str,
) -> tuple[ManagedInputRecord, ...]create_round(
self,
pending: PendingElicitationRound,
*,
round_index: int,
round_limit: int,
owner_id: str,
lease_seconds: float,
lease_token: str | None = None,
harness_session_id: str | None = None,
native_resume_token: str | None = None,
delivery_idempotency_key: str | None = None,
session_id: str | None = None,
turn_id: str | None = None,
operation_parameters: Mapping[str, object] | None = None,
) -> ManagedInputRecordclaim_round(
self,
execution_id: str,
round_id: str,
*,
owner_id: str,
lease_seconds: float,
expected_lease_token: str | None = None,
delivery_state: Literal['not_started', 'not_delivered', 'delivered'] | None = None,
idempotent_delivery: bool = False,
) -> ManagedInputRecordrenew(
self,
execution_id: str,
round_id: str,
*,
owner_id: str,
lease_token: str,
lease_seconds: float,
) -> ManagedInputRecordRenew an active local wait without changing its compare-and-set token.
submit_responses(
self,
execution_id: str,
round_id: str,
responses: Mapping[str, ElicitationResponse],
*,
owner_id: str,
lease_token: str,
response_idempotency_key: str,
) -> ManagedInputRecordstart_delivery(
self,
execution_id: str,
round_id: str,
*,
owner_id: str,
lease_token: str,
) -> ManagedInputRecordmark_delivered(
self,
execution_id: str,
round_id: str,
*,
owner_id: str,
lease_token: str,
) -> ManagedInputRecordresolve(
self,
execution_id: str,
round_id: str,
*,
owner_id: str,
lease_token: str,
) -> ManagedInputRecordfail(
self,
execution_id: str,
round_id: str,
*,
owner_id: str,
lease_token: str,
code: str,
message: str,
) -> ManagedInputRecordfail_recovery(
self,
execution_id: str,
round_id: str,
*,
expected_lease_token: str,
message: str,
) -> ManagedInputRecordredacted_responses(
self,
execution_id: str,
round_id: str,
) -> Mapping[str, object] | NoneSQLiteStore
m3.storage.SQLiteStore(
database: str | Path,
*,
blob_root: str | Path | None = None,
config: RedactionConfig | None = None,
capture_config: CaptureOptions | None = None,
payload_blob_threshold: int = 65536,
**kwargs: Any,
) -> NoneSQLite implementation of the public :class:ExecutionStore protocol.
managed_input_store(property): Lazily open the managed-input tables for opted-in executions.
resolve_managed_input_store(
self,
) -> SQLiteManagedInputStoreResolve managed-input storage lazily for opted-in executions.
ensure_project(
self,
project_id: str,
project_name: str,
) -> tuple[str, str]Register a stable project identity and refresh its display name.
get_project(
self,
project_id: str,
) -> tuple[str, str] | Noneensure_suite(
self,
suite_name: str,
project_id: str | None = None,
) -> Suiteget_suite(
self,
suite_id: SuiteId | str,
) -> Suite | Noneget_suite_by_name(
self,
suite_name: str,
project_id: str | None = None,
) -> Suite | Nonelist_suites(
self,
) -> tuple[Suite, ...]create(
self,
snapshot: ExecutionState,
*,
specification: Mapping[str, Any] | None = None,
provenance: Mapping[str, Any] | None = None,
server_bindings: Sequence[Mapping[str, Any]] = (),
harness_binding: Mapping[str, Any] | None = None,
parent_execution_id: ExecutionId | str | None = None,
run_id: RunId | str | None = None,
) -> Nonecreate_execution(
self,
snapshot: ExecutionState,
*,
specification: Mapping[str, Any] | None = None,
provenance: Mapping[str, Any] | None = None,
server_bindings: Sequence[Mapping[str, Any]] = (),
harness_binding: Mapping[str, Any] | None = None,
parent_execution_id: ExecutionId | str | None = None,
run_id: RunId | str | None = None,
) -> Nonesave_test_run(
self,
run_id: str,
value: Mapping[str, object],
) -> Noneget_test_run(
self,
run_id: str,
) -> Mapping[str, object] | Nonelist_test_runs(
self,
) -> tuple[Mapping[str, object], ...]list_test_run_page(
self,
*,
limit: int | None = None,
offset: int = 0,
suite_id: int | None = None,
project_id: str | None = None,
q: str | None = None,
) -> tuple[tuple[Mapping[str, object], ...], int]Return newest-first run manifests with their suites, and the total.
save_test_result(
self,
run_id: str,
attempt_id: str,
value: Mapping[str, object],
) -> Nonelist_test_results(
self,
run_id: str,
) -> tuple[Mapping[str, object], ...]get_snapshot(
self,
execution_id: ExecutionId | str,
) -> ExecutionState | Noneget_execution_spec(
self,
execution_id: ExecutionId | str,
) -> ExecutionSpec | NoneReturn the immutable typed submission spec, if one was saved.
get_spec(
self,
execution_id: ExecutionId | str,
) -> ExecutionSpec | NoneReturn the immutable typed submission spec, if one was saved.
list_executions(
self,
*,
limit: int = 50,
offset: int = 0,
lifecycle: ExecutionStatus | str | None = None,
outcome: ExecutionOutcome | str | None = None,
run_id: str | None = None,
suite_id: int | None = None,
project_id: str | None = None,
) -> ExecutionPageget_report(
self,
execution_id: ExecutionId | str,
*,
after_sequence: int = -1,
event_limit: int | None = None,
artifact_limit: int | None = None,
) -> ExecutionReport | Noneget_trace(
self,
execution_id: ExecutionId | str,
) -> TraceResult | Noneget_trace_view(
self,
execution_id: ExecutionId | str,
) -> TraceView | Nonesave_snapshot(
self,
snapshot: ExecutionState,
) -> Noneupdate_snapshot(
self,
snapshot: ExecutionState,
) -> Noneappend_events(
self,
events: Sequence[Event],
) -> Noneappend(
self,
events: Sequence[Event],
) -> Noneappend_event(
self,
event: Event,
content: bytes,
*,
media_type: str,
) -> EventCommit an event and its raw blob in one SQLite transaction.
iter_events(
self,
execution_id: ExecutionId | str,
*,
after_sequence: int = -1,
) -> Iterator[Event]events(
self,
execution_id: ExecutionId | str,
*,
after_sequence: int = -1,
) -> tuple[Event, ...]create_session(
self,
execution_id: ExecutionId | str,
session_id: SessionId | str,
*,
state: str = 'open',
) -> SessionIdclose_session(
self,
session_id: SessionId | str,
) -> Nonesave_turn(
self,
snapshot: TurnState,
result: TurnResult | Mapping[str, Any] | None = None,
) -> Noneappend_turn(
self,
snapshot: TurnState,
result: TurnResult | Mapping[str, Any] | None = None,
) -> Noneturns(
self,
execution_id: ExecutionId | str,
) -> tuple[tuple[TurnState, TurnResult | None], ...]save_evaluation(
self,
execution_id: ExecutionId | str,
result: Mapping[str, Any] | Any,
*,
evaluation_id: str | None = None,
turn_id: TurnId | str | None = None,
) -> strreserve_judge_request(
self,
run_id: str,
limit: int,
) -> boolAtomically reserve one judge request for a run across workers.
save(
self,
result: EvaluationResult,
) -> Noneget(
self,
evaluation_id: str,
) -> EvaluationResult | Noneall(
self,
) -> tuple[EvaluationResult, ...]evaluations(
self,
execution_id: ExecutionId | str,
*,
turn_id: TurnId | str | None = None,
) -> tuple[EvaluationRecord, ...]aggregate_evaluations(
self,
query: EvaluationQuery,
) -> EvaluationReportCalculate summaries from persisted evaluations and execution traces.
persisted_evaluations(
self,
execution_id: ExecutionId | str,
*,
turn_id: TurnId | str | None = None,
) -> tuple[EvaluationRecord, ...]evaluation_json(
self,
execution_id: ExecutionId | str,
*,
turn_id: TurnId | str | None = None,
) -> tuple[Mapping[str, Any], ...]transaction(
self,
execution_id: ExecutionId | str,
) -> ExecutionTransactionallocate(
self,
execution_id: ExecutionId | str,
*,
count: int = 1,
) -> tuple[int, ...]allocate_sequence(
self,
execution_id,
)allocate_sequences(
self,
execution_id: ExecutionId | str,
*,
count: int = 1,
) -> tuple[int, ...]release(
self,
execution_id: ExecutionId | str,
sequences: Sequence[int],
) -> Nonerelease_sequences(
self,
execution_id: ExecutionId | str,
sequences: Sequence[int],
) -> Nonesubscribe(
self,
execution_id: ExecutionId | str,
callback: EventCallback,
) -> Callable[[], None]save_acp_probe(
self,
result: ACPProbeResult,
) -> ACPProbeResultPersist one redacted ACP probe result and return its safe copy.
get_acp_probe(
self,
probe_id: str,
) -> ACPProbeResult | Nonelist_acp_probes(
self,
dimension: ACPProbeDimension | None = None,
*,
include_inflight: bool = True,
) -> tuple[ACPProbeResult, ...]latest_acp_probe(
self,
dimension: ACPProbeDimension,
) -> ACPProbeResult | Nonecreate_profile(
self,
kind: str,
name: str,
value: Mapping[str, Any],
*,
description: str = '',
profile_id: str | None = None,
revision_id: str | None = None,
) -> ProfileRecordcreate_server_profile(
self,
name: str,
value: Mapping[str, Any],
**kwargs: Any,
) -> ProfileRecordcreate_harness_profile(
self,
name: str,
value: Mapping[str, Any],
**kwargs: Any,
) -> ProfileRecordget_profile(
self,
profile_id: str,
*,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRecord | Noneresolve_profile(
self,
profile_id: str,
selection: RevisionSelection,
*,
kind: Literal['server', 'harness'],
) -> tuple[ProfileRecord | None, ProfileRevisionRecord | None]Read one profile and its selected immutable revision.
list_profiles(
self,
kind: str,
*,
include_archived: bool = False,
) -> tuple[ProfileRecord, ...]List one profile family in stable name/id order.
list_profile_revisions(
self,
profile_id: str,
*,
kind: Literal['server', 'harness'] | None = None,
) -> tuple[ProfileRevisionRecord, ...]Return every revision ordered by revision number then id.
update_profile(
self,
profile_id: str,
*,
name: str | None = None,
description: str | None = None,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRecordUpdate mutable metadata without changing the immutable revision.
add_revision(
self,
profile_id: str,
value: Mapping[str, Any],
*,
revision_id: str | None = None,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRevisionRecordarchive_profile(
self,
profile_id: str,
*,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRecordrestore_profile(
self,
profile_id: str,
*,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRecordresolve_revision(
self,
profile_id: str,
selection: RevisionSelection | str = 'latest',
*,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRevisionRecordget_revision(
self,
revision_id: RevisionId | str,
*,
kind: Literal['server', 'harness'] | None = None,
) -> ProfileRevisionRecord | Noneenqueue_command(
self,
execution_id: ExecutionId | str,
kind: str = 'execution',
payload: Mapping[str, Any] | None = None,
*,
command_id: str | None = None,
session_id: SessionId | str | None = None,
turn_id: TurnId | str | None = None,
) -> Commandenqueue(
self,
execution_id: ExecutionId | str,
kind: str = 'execution',
payload: Mapping[str, Any] | None = None,
*,
command_id: str | None = None,
session_id: SessionId | str | None = None,
turn_id: TurnId | str | None = None,
) -> Commandget_command(
self,
command_id: str,
) -> Command | Noneclaim_next(
self,
owner_id: str,
*,
lease_seconds: float = 30.0,
) -> tuple[Command, Lease] | Noneclaim(
self,
owner_id: str,
*,
lease_seconds: float = 30.0,
) -> tuple[Command, Lease] | Noneheartbeat(
self,
lease: Lease | str,
*,
owner_id: str | None = None,
lease_seconds: float = 30.0,
) -> boolrenew_lease(
self,
lease: Lease | str,
*,
owner_id: str | None = None,
lease_seconds: float = 30.0,
) -> boolmark_interrupted_if_lease_lost(
self,
lease: Lease,
*,
reason: str = 'worker lease lost',
) -> boolAtomically close work whose owner can no longer renew its lease.
mark_managed_recovery_unavailable_if_lease_lost(
self,
lease: Lease,
*,
reason: str = 'managed interaction cannot be resumed safely after worker loss',
) -> boolTerminalize managed input when its owning worker is lost.
release_lease(
self,
lease: Lease | str,
*,
owner_id: str | None = None,
) -> boolcomplete_command(
self,
command_id: str,
*,
owner_id: str,
lease_token: str,
status: str = 'done',
) -> boolMark one claimed command terminal under its current lease.
request_cancel(
self,
execution_id: ExecutionId | str,
reason: str | None = None,
) -> boolcancel(
self,
execution_id: ExecutionId | str,
reason: str | None = None,
) -> boolfinalize_cancelled(
self,
execution_id: ExecutionId | str,
*,
reason: str = 'cancelled',
) -> boolPersist a cancellation terminal event when a worker observed it.
cancellation_requested(
self,
execution_id: ExecutionId | str,
) -> boolmark_stale_interrupted(
self,
*,
now: datetime | None = None,
) -> tuple[ExecutionId, ...]delete_execution(
self,
execution_id: ExecutionId | str,
) -> Nonedelete(
self,
execution_id: ExecutionId | str,
) -> Noneput_raw_evidence(
self,
event_id: EventId | str,
content: bytes,
*,
media_type: str,
) -> EvidenceCaptureRedact, bound, and durably associate evidence with one event.
read_raw_evidence(
self,
reference: EvidenceRef,
*,
max_bytes: int = 1048576,
) -> RawEvidenceclone_execution(
self,
execution_id: ExecutionId | str,
*,
use_latest: bool = False,
) -> ExecutionIdclone(
self,
execution_id: ExecutionId | str,
*,
use_latest: bool = False,
) -> ExecutionIdresolved_bindings(
self,
execution_id: ExecutionId | str,
) -> Mapping[str, Any]Return the immutable, submission-time profile binding snapshot.
close(
self,
) -> NoneSQLiteStoreWorker
m3.storage.SQLiteStoreWorker(
store: Any,
runner: Callable[[Any, Any, Any], Any],
*,
worker_id: str | None = None,
lease_seconds: float = 30.0,
cancel_runner: Callable[[Any], Any] | None = None,
) -> NoneEmbedded worker facade for :class:SQLiteExecutionStore.
run_once(
self,
) -> boolstop(
self,
) -> NoneStorageConflict
m3.storage.StorageConflictThe append or snapshot operation conflicts with committed state.
StorageError
m3.storage.StorageErrorBase class for expected ephemeral storage failures.
TemporaryArtifactStore
m3.storage.TemporaryArtifactStore(
root: str | os.PathLike[str] | None = None,
*,
config: RedactionConfig | None = None,
) -> NoneFilesystem-backed temporary artifact store with atomic blob placement.
root(property)
put(
self,
execution_id: ExecutionId | str,
name: str,
content: bytes,
*,
media_type: str | None = None,
) -> ArtifactRefget(
self,
artifact: ArtifactRef | ArtifactId | str,
) -> bytesdelete(
self,
artifact: ArtifactRef | ArtifactId | str,
) -> Nonecleanup(
self,
) -> Noneserialize_durable
m3.storage.serialize_durable(
value: Any,
*,
config: RedactionConfig | None = None,
path: str = '$',
) -> AnyReturn strict JSON-compatible durable data while preserving references.