Skip to content

Runs

Every governed operation produces a run: the durable record of who asked for work, what was allowed, what executed, what was checked, and what Gantry decided. It outlives the agent and the process that started it, which is what makes it possible to answer those questions later.

gantry.runs.configure(gantry.runs.SQLiteRunStore(".gantry/runs.db"))

with gantry.actor.context(actor=gantry.actor.actor("agent", "migration-agent")):
    run = await db.query(schemas=["analytics"], verify=[...])(sql)

run.id, run.status, run.rows

# …in another process, holding only the id
print(gantry.runs.get(run_id).render())

Recorded before anything runs

The run is created before the proposal reaches the engine. If that first write fails, nothing is submitted — an engine job that exists without a record of why it was allowed to is the one outcome this ordering prevents.

gantry.Run dataclass

One governed operation, from proposal to decision.

Built by the operation, not by the caller: run_id is allocated and the record persisted before anything external happens, so an engine job cannot exist without a control-plane record of why it was allowed to.

Source code in gantry/runs/model.py
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
@dataclass(frozen=True, slots=True)
class Run:
    """One governed operation, from proposal to decision.

    Built by the operation, not by the caller: `run_id` is allocated and the
    record persisted before anything external happens, so an engine job cannot
    exist without a control-plane record of why it was allowed to.
    """

    id: str
    status: RunStatus
    actor: ActorRef
    operation: OperationRef
    proposal: ProposalRecord | None = None
    inputs: tuple[ResourceRef, ...] = ()
    outputs: tuple[ResourceRef, ...] = ()
    result_ref: QueryResultRef | None = None
    admission: AdmissionRecord | None = None
    confirmation: ConfirmationRecord | None = None
    execution: ExecutionRecord | None = None
    verification: VerificationResult | None = None
    evidence: EvidenceBundle | None = None
    created_at: datetime = field(default_factory=lambda: datetime.now(UTC))
    updated_at: datetime = field(default_factory=lambda: datetime.now(UTC))
    completed_at: datetime | None = None

    # Live-only. Excluded from `as_dict`, and therefore from storage: a run
    # record is not a place for query results to accumulate.
    inline: object | None = field(default=None, compare=False, repr=False)
    """The bounded result, in whatever shape the engine returned it.

    Rows for SQL, documents for MongoDB. Read it through `rows`, `columns` or
    `documents` rather than directly — those narrow it, and they are empty
    rather than wrong when the operation produced nothing.
    """

    failure: Failure | None = field(default=None, compare=False, repr=False)
    handle: ExecutionHandle | None = field(default=None, compare=False, repr=False)
    """The engine handle, for callers that submit and wait separately.

    Live-only, like `inline`. What persists is the engine's own id inside the
    execution record; a handle is a thing to act on now, not a thing to read
    back later.
    """

    @property
    def ok(self) -> bool:
        """Accepted means executed *and* verified. Nothing else is acceptance."""
        return self.status is RunStatus.ACCEPTED

    @property
    def safe_to_retry(self) -> bool:
        """Whether submitting this same proposal again is safe and worth doing.

        `failure.retryable` answers half the question — whether the condition
        may pass. This answers the half that depends on what was being done and
        how far it got, which is the half that can destroy something.

        A query is idempotent: re-running a `SELECT` that timed out costs
        another attempt and nothing else. A write is not. Materializations are
        create-only, and a run that reached the engine may have left its
        destination behind even though it failed — so the retry does not
        succeed, it comes back `DESTINATION_EXISTS`, turning a transient failure
        into one that looks permanent. A batch or streaming job that was
        submitted is worse: running it again does not replace the first one, it
        adds a second.

        So a write is retry-safe only when nothing reached the engine at all.
        `False` here is not a claim that retrying will fail — it is a claim that
        Gantry cannot promise it is harmless, which for a write is the answer
        that matters.
        """
        if self.failure is None or not self.failure.retryable:
            return False
        if self.operation.kind is OperationKind.QUERY:
            return True
        return self.execution is None or self.execution.native_id is None

    @property
    def run_id(self) -> str:
        """Alias for `id`, for callers that hold results from several systems."""
        return self.id

    @property
    def rows(self) -> tuple[tuple[object, ...], ...]:
        """A SQL query's rows, when it returned any and they came back inline."""
        return tuple(getattr(self.inline, "rows", ()) or ())

    @property
    def columns(self) -> tuple[str, ...]:
        """The column names a SQL query returned."""
        return tuple(getattr(self.inline, "columns", ()) or ())

    @property
    def truncated(self) -> bool:
        """Whether a bound clipped the result.

        Recorded rather than implied: a truncated answer changes what every
        other observation about it means.
        """
        if self.result_ref is not None:
            return self.result_ref.truncated
        return bool(getattr(self.inline, "truncated", False))

    @property
    def documents(self) -> tuple[object, ...]:
        """A document query's results, when it returned any."""
        return tuple(getattr(self.inline, "documents", ()) or ())

    @property
    def uri(self) -> str | None:
        """The first durable output, in the engine-owned URI form.

        A resource is recorded the way the engine names it —
        `reporting.rollup` — because that is what someone reads. The URI keeps
        the slash form the rest of the library already returns, so a caller
        that was matching on it still matches.
        """
        ref = next(iter(self.outputs), None)
        if ref is None:
            return None
        if "://" in ref.resource:
            return ref.resource
        return f"{ref.system}://{ref.resource.replace('.', '/')}"

    @property
    def trusted_checks(self) -> tuple[CheckResult, ...]:
        return self._checks("trusted")

    @property
    def agent_checks(self) -> tuple[CheckResult, ...]:
        return self._checks("agent")

    def _checks(self, source: str) -> tuple[CheckResult, ...]:
        """Split checks by provenance, without losing any.

        `CheckResult.source` carries two kinds of answer: who asked for the
        check (`trusted`, `agent`) and where its observation came from
        (`postgres`, `result set`). Only `agent` means agent-proposed, so
        everything else is trusted — a check whose source names a provider is
        still one the application required, and grouping on equality alone
        dropped it from the record entirely.
        """
        checks = () if self.verification is None else self.verification.checks
        if source == "agent":
            return tuple(check for check in checks if str(check.source or "") == "agent")
        return tuple(check for check in checks if str(check.source or "") != "agent")

    def with_failure(self, failure: Failure | None) -> Run:
        """Attach the failure a caller needs to read, without persisting it twice.

        The reason is already recorded — in the admission record, the execution
        record, or the failing check. This is the live object carrying it in the
        shape callers already expect.
        """
        return self if failure is None else replace(self, failure=failure)

    def advanced(self, status: RunStatus, **changes: object) -> Run:
        """The same run, moved on. One run evolves; stages are not separate runs."""
        now = datetime.now(UTC)
        completed = now if status.terminal else self.completed_at
        return replace(
            self,
            status=status,
            updated_at=now,
            completed_at=completed,
            **changes,  # type: ignore[arg-type]
        )

    def as_dict(self) -> dict[str, object]:
        """The stored form. Never includes result rows."""
        return {
            "id": self.id,
            "status": self.status.value,
            "actor": self.actor.as_dict(),
            "operation": self.operation.as_dict(),
            "proposal": None if self.proposal is None else self.proposal.as_dict(),
            "inputs": [ref.as_dict() for ref in self.inputs],
            "outputs": [ref.as_dict() for ref in self.outputs],
            "result_ref": None if self.result_ref is None else self.result_ref.as_dict(),
            "admission": None if self.admission is None else self.admission.as_dict(),
            "confirmation": None if self.confirmation is None else self.confirmation.as_dict(),
            "execution": None if self.execution is None else self.execution.as_dict(),
            "verification": _verification_as_dict(self.verification),
            "evidence": None if self.evidence is None else self.evidence.as_dict(),
            "created_at": _time(self.created_at),
            "updated_at": _time(self.updated_at),
            "completed_at": _time(self.completed_at),
        }

    def to_json(self, *, indent: int | None = None) -> str:
        return json.dumps(self.as_dict(), indent=indent)

    def render(self) -> str:
        return render(self)

inline class-attribute instance-attribute

inline: object | None = field(
    default=None, compare=False, repr=False
)

The bounded result, in whatever shape the engine returned it.

Rows for SQL, documents for MongoDB. Read it through rows, columns or documents rather than directly — those narrow it, and they are empty rather than wrong when the operation produced nothing.

handle class-attribute instance-attribute

handle: ExecutionHandle | None = field(
    default=None, compare=False, repr=False
)

The engine handle, for callers that submit and wait separately.

Live-only, like inline. What persists is the engine's own id inside the execution record; a handle is a thing to act on now, not a thing to read back later.

ok property

ok: bool

Accepted means executed and verified. Nothing else is acceptance.

safe_to_retry property

safe_to_retry: bool

Whether submitting this same proposal again is safe and worth doing.

failure.retryable answers half the question — whether the condition may pass. This answers the half that depends on what was being done and how far it got, which is the half that can destroy something.

A query is idempotent: re-running a SELECT that timed out costs another attempt and nothing else. A write is not. Materializations are create-only, and a run that reached the engine may have left its destination behind even though it failed — so the retry does not succeed, it comes back DESTINATION_EXISTS, turning a transient failure into one that looks permanent. A batch or streaming job that was submitted is worse: running it again does not replace the first one, it adds a second.

So a write is retry-safe only when nothing reached the engine at all. False here is not a claim that retrying will fail — it is a claim that Gantry cannot promise it is harmless, which for a write is the answer that matters.

run_id property

run_id: str

Alias for id, for callers that hold results from several systems.

rows property

rows: tuple[tuple[object, ...], ...]

A SQL query's rows, when it returned any and they came back inline.

columns property

columns: tuple[str, ...]

The column names a SQL query returned.

truncated property

truncated: bool

Whether a bound clipped the result.

Recorded rather than implied: a truncated answer changes what every other observation about it means.

documents property

documents: tuple[object, ...]

A document query's results, when it returned any.

uri property

uri: str | None

The first durable output, in the engine-owned URI form.

A resource is recorded the way the engine names it — reporting.rollup — because that is what someone reads. The URI keeps the slash form the rest of the library already returns, so a caller that was matching on it still matches.

with_failure

with_failure(failure: Failure | None) -> Run

Attach the failure a caller needs to read, without persisting it twice.

The reason is already recorded — in the admission record, the execution record, or the failing check. This is the live object carrying it in the shape callers already expect.

Source code in gantry/runs/model.py
367
368
369
370
371
372
373
374
def with_failure(self, failure: Failure | None) -> Run:
    """Attach the failure a caller needs to read, without persisting it twice.

    The reason is already recorded — in the admission record, the execution
    record, or the failing check. This is the live object carrying it in the
    shape callers already expect.
    """
    return self if failure is None else replace(self, failure=failure)

advanced

advanced(status: RunStatus, **changes: object) -> Run

The same run, moved on. One run evolves; stages are not separate runs.

Source code in gantry/runs/model.py
376
377
378
379
380
381
382
383
384
385
386
def advanced(self, status: RunStatus, **changes: object) -> Run:
    """The same run, moved on. One run evolves; stages are not separate runs."""
    now = datetime.now(UTC)
    completed = now if status.terminal else self.completed_at
    return replace(
        self,
        status=status,
        updated_at=now,
        completed_at=completed,
        **changes,  # type: ignore[arg-type]
    )

as_dict

as_dict() -> dict[str, object]

The stored form. Never includes result rows.

Source code in gantry/runs/model.py
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
def as_dict(self) -> dict[str, object]:
    """The stored form. Never includes result rows."""
    return {
        "id": self.id,
        "status": self.status.value,
        "actor": self.actor.as_dict(),
        "operation": self.operation.as_dict(),
        "proposal": None if self.proposal is None else self.proposal.as_dict(),
        "inputs": [ref.as_dict() for ref in self.inputs],
        "outputs": [ref.as_dict() for ref in self.outputs],
        "result_ref": None if self.result_ref is None else self.result_ref.as_dict(),
        "admission": None if self.admission is None else self.admission.as_dict(),
        "confirmation": None if self.confirmation is None else self.confirmation.as_dict(),
        "execution": None if self.execution is None else self.execution.as_dict(),
        "verification": _verification_as_dict(self.verification),
        "evidence": None if self.evidence is None else self.evidence.as_dict(),
        "created_at": _time(self.created_at),
        "updated_at": _time(self.updated_at),
        "completed_at": _time(self.completed_at),
    }

gantry.RunStatus

Bases: StrEnum

Where a run got to, and why it stopped there.

The distinctions that matter are between the ways a run can fail. A proposal Gantry refused, work the engine could not do, a check that could not be evaluated, and a result that was measured and rejected are four different events with four different remedies. Collapsing any of them into "failed" throws away the only part anyone can act on.

Source code in gantry/runs/status.py
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
class RunStatus(StrEnum):
    """Where a run got to, and why it stopped there.

    The distinctions that matter are between the ways a run can fail. A
    proposal Gantry refused, work the engine could not do, a check that
    could not be evaluated, and a result that was measured and rejected are
    four different events with four different remedies. Collapsing any of
    them into "failed" throws away the only part anyone can act on.
    """

    PENDING = "PENDING"
    """Recorded, not yet admitted. A run exists here before anything external does."""

    POLICY_REJECTED = "POLICY_REJECTED"
    """Refused before execution. Nothing ran."""

    AWAITING_CONFIRMATION = "AWAITING_CONFIRMATION"
    """Policy allowed it, and a rule asked that the user be told first.

    Not a refusal and not terminal: the run is parked, nothing external has
    started, and it resumes if the host confirms. A run left here forever is the
    honest record of a question nobody answered.
    """

    CONFIRMATION_DECLINED = "CONFIRMATION_DECLINED"
    """The host said the user declined. Terminal, and nothing ran."""

    VERIFICATION_CONFLICT = "VERIFICATION_CONFLICT"
    """The verification contract contradicted itself, so no outcome could satisfy it."""

    RUNNING = "RUNNING"
    """Submitted to the engine and in flight."""

    EXECUTION_FAILED = "EXECUTION_FAILED"
    """The engine could not carry out work that Gantry had admitted."""

    VERIFYING = "VERIFYING"
    """Execution finished, or a stream reached its required state; checks are running."""

    VERIFICATION_UNSUPPORTED = "VERIFICATION_UNSUPPORTED"
    """A required check could not be evaluated, so acceptance cannot be claimed."""

    REJECTED = "REJECTED"
    """It ran, it was measured, and the result is not acceptable."""

    ACCEPTED = "ACCEPTED"
    """Executed and verified. For a stream, this means it reached a verified healthy
    state — not that it has finished, and not that it will stay healthy."""

    @property
    def awaiting(self) -> bool:
        """Parked, waiting on someone outside Gantry. Not terminal, not running."""
        return self is RunStatus.AWAITING_CONFIRMATION

    @property
    def terminal(self) -> bool:
        return self in _TERMINAL

    @property
    def accepted(self) -> bool:
        return self is RunStatus.ACCEPTED

PENDING class-attribute instance-attribute

PENDING = 'PENDING'

Recorded, not yet admitted. A run exists here before anything external does.

POLICY_REJECTED class-attribute instance-attribute

POLICY_REJECTED = 'POLICY_REJECTED'

Refused before execution. Nothing ran.

AWAITING_CONFIRMATION class-attribute instance-attribute

AWAITING_CONFIRMATION = 'AWAITING_CONFIRMATION'

Policy allowed it, and a rule asked that the user be told first.

Not a refusal and not terminal: the run is parked, nothing external has started, and it resumes if the host confirms. A run left here forever is the honest record of a question nobody answered.

CONFIRMATION_DECLINED class-attribute instance-attribute

CONFIRMATION_DECLINED = 'CONFIRMATION_DECLINED'

The host said the user declined. Terminal, and nothing ran.

VERIFICATION_CONFLICT class-attribute instance-attribute

VERIFICATION_CONFLICT = 'VERIFICATION_CONFLICT'

The verification contract contradicted itself, so no outcome could satisfy it.

RUNNING class-attribute instance-attribute

RUNNING = 'RUNNING'

Submitted to the engine and in flight.

EXECUTION_FAILED class-attribute instance-attribute

EXECUTION_FAILED = 'EXECUTION_FAILED'

The engine could not carry out work that Gantry had admitted.

VERIFYING class-attribute instance-attribute

VERIFYING = 'VERIFYING'

Execution finished, or a stream reached its required state; checks are running.

VERIFICATION_UNSUPPORTED class-attribute instance-attribute

VERIFICATION_UNSUPPORTED = 'VERIFICATION_UNSUPPORTED'

A required check could not be evaluated, so acceptance cannot be claimed.

REJECTED class-attribute instance-attribute

REJECTED = 'REJECTED'

It ran, it was measured, and the result is not acceptable.

ACCEPTED class-attribute instance-attribute

ACCEPTED = 'ACCEPTED'

Executed and verified. For a stream, this means it reached a verified healthy state — not that it has finished, and not that it will stay healthy.

awaiting property

awaiting: bool

Parked, waiting on someone outside Gantry. Not terminal, not running.

Identity

gantry.runs.OperationKind

Bases: StrEnum

What was asked for, at the granularity a reader cares about.

Deliberately not provider API shapes: a BigQuery load job and a Postgres CREATE TABLE AS are both materializations, and a record that says so is readable across engines.

Source code in gantry/runs/model.py
52
53
54
55
56
57
58
59
60
61
62
63
class OperationKind(StrEnum):
    """What was asked for, at the granularity a reader cares about.

    Deliberately not provider API shapes: a BigQuery load job and a Postgres
    `CREATE TABLE AS` are both materializations, and a record that says so is
    readable across engines.
    """

    QUERY = "query"
    MATERIALIZE = "materialize"
    BATCH = "batch"
    STREAM = "stream"

gantry.runs.OperationRef dataclass

Source code in gantry/runs/model.py
73
74
75
76
77
78
79
80
@dataclass(frozen=True, slots=True)
class OperationRef:
    kind: OperationKind
    engine: str
    provider: str | None = None

    def as_dict(self) -> dict[str, object]:
        return {"kind": self.kind.value, "engine": self.engine, "provider": self.provider}

gantry.ActorRef dataclass

The thing that initiated a run.

metadata is for the few facts that help identify a caller later — which framework, which deployment. It is deliberately not somewhere to put a conversation: a run record is not a transcript store, and anything written here outlives the process and lands in a durable file.

Source code in gantry/actor.py
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
@dataclass(frozen=True, slots=True)
class ActorRef:
    """The thing that initiated a run.

    `metadata` is for the few facts that help identify a caller later — which
    framework, which deployment. It is deliberately not somewhere to put a
    conversation: a run record is not a transcript store, and anything written
    here outlives the process and lands in a durable file.
    """

    type: ActorType = ActorType.UNKNOWN
    id: str | None = None
    session_id: str | None = None
    metadata: Mapping[str, object] = field(default_factory=dict)

    def __post_init__(self) -> None:
        if self.id is not None and not self.id.strip():
            raise ValueError("actor id must not be empty when given")
        oversized = [key for key, value in self.metadata.items() if len(str(value)) > 1024]
        if oversized:
            raise ValueError(
                f"actor metadata values must stay small; too long: {', '.join(sorted(oversized))}"
            )

    @property
    def label(self) -> str:
        """`agent:migration-agent`, or just the type when nothing identified it."""
        return self.type.value if self.id is None else f"{self.type.value}:{self.id}"

    def as_dict(self) -> dict[str, object]:
        payload: dict[str, object] = {"type": self.type.value, "id": self.id}
        if self.session_id is not None:
            payload["session_id"] = self.session_id
        if self.metadata:
            payload["metadata"] = {str(key): value for key, value in self.metadata.items()}
        return payload

label property

label: str

agent:migration-agent, or just the type when nothing identified it.

gantry.ActorType

Bases: StrEnum

What kind of thing initiated a run.

UNKNOWN is the honest default rather than a failure: a library that guessed would put a wrong name in a durable record, and an unattributed run is more useful than a misattributed one.

Source code in gantry/actor.py
19
20
21
22
23
24
25
26
27
28
29
30
class ActorType(StrEnum):
    """What kind of thing initiated a run.

    `UNKNOWN` is the honest default rather than a failure: a library that
    guessed would put a wrong name in a durable record, and an unattributed run
    is more useful than a misattributed one.
    """

    AGENT = "agent"
    USER = "user"
    SERVICE = "service"
    UNKNOWN = "unknown"

gantry.actor.actor

actor(
    type: str | ActorType = AGENT,
    id: str | None = None,
    *,
    session_id: str | None = None,
    metadata: Mapping[str, object] | None = None,
) -> ActorRef

Name the caller. gantry.actor("agent", "research-agent").

Source code in gantry/actor.py
77
78
79
80
81
82
83
84
85
86
87
88
89
90
def actor(
    type: str | ActorType = ActorType.AGENT,  # noqa: A002
    id: str | None = None,  # noqa: A002
    *,
    session_id: str | None = None,
    metadata: Mapping[str, object] | None = None,
) -> ActorRef:
    """Name the caller. `gantry.actor("agent", "research-agent")`."""
    return ActorRef(
        type=ActorType(type) if not isinstance(type, ActorType) else type,
        id=id,
        session_id=session_id,
        metadata=dict(metadata or {}),
    )

gantry.actor.context

context(
    *,
    actor: ActorRef | None = None,
    environment: str | None = None,
) -> Iterator[ActorRef]

Establish the trusted context every run inside the block is judged in.

Both facts come from the host application and neither is reachable from a proposal: an agent that could name its own actor could name someone else's, and one that could name its own environment could call production staging.

Context variables rather than globals: concurrent requests in one process each keep their own caller and environment, which a module-level assignment would not.

Source code in gantry/actor.py
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
@contextmanager
def context(*, actor: ActorRef | None = None, environment: str | None = None) -> Iterator[ActorRef]:
    """Establish the trusted context every run inside the block is judged in.

    Both facts come from the host application and neither is reachable from a
    proposal: an agent that could name its own actor could name someone else's,
    and one that could name its own environment could call production staging.

    Context variables rather than globals: concurrent requests in one process
    each keep their own caller and environment, which a module-level assignment
    would not.
    """
    if actor is None and environment is None:
        raise ValueError("a trusted context must set an actor, an environment, or both")
    if environment is not None and not environment.strip():
        raise ValueError("environment must not be empty when given")
    actor_token = None if actor is None else _current.set(actor)
    environment_token = (
        None if environment is None else _environment.set(environment.strip().lower())
    )
    try:
        yield actor or _current.get()
    finally:
        if environment_token is not None:
            _environment.reset(environment_token)
        if actor_token is not None:
            _current.reset(actor_token)

What a run records

gantry.runs.ProposalRecord dataclass

What the caller asked Gantry to run.

The hash is always kept; the body depends on configuration. Agent-written SQL can carry values from the data it is filtering — a customer id in a WHERE clause is the ordinary case — so hash mode exists for callers who would rather a durable file did not accumulate them.

Source code in gantry/runs/model.py
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
@dataclass(frozen=True, slots=True)
class ProposalRecord:
    """What the caller asked Gantry to run.

    The hash is always kept; the body depends on configuration. Agent-written
    SQL can carry values from the data it is filtering — a customer id in a
    `WHERE` clause is the ordinary case — so `hash` mode exists for callers who
    would rather a durable file did not accumulate them.
    """

    hash: str
    kind: str = "sql"
    body: str | None = None
    storage: ProposalStorage = ProposalStorage.FULL
    agent_verification: tuple[str, ...] = ()

    @classmethod
    def of(
        cls,
        body: str,
        *,
        kind: str = "sql",
        storage: ProposalStorage = ProposalStorage.FULL,
        agent_verification: Sequence[str] = (),
    ) -> ProposalRecord:
        return cls(
            hash=f"sha256:{sha256(body.encode()).hexdigest()}",
            kind=kind,
            body=body if storage is ProposalStorage.FULL else None,
            storage=storage,
            agent_verification=tuple(agent_verification),
        )

    def as_dict(self) -> dict[str, object]:
        payload: dict[str, object] = {
            "hash": self.hash,
            "kind": self.kind,
            "storage": self.storage.value,
        }
        if self.body is not None:
            payload["body"] = self.body
        if self.agent_verification:
            payload["agent_verification"] = list(self.agent_verification)
        return payload

gantry.runs.ProposalStorage

Bases: StrEnum

How much of the proposal a run keeps.

Source code in gantry/runs/model.py
66
67
68
69
70
class ProposalStorage(StrEnum):
    """How much of the proposal a run keeps."""

    FULL = "full"
    HASH = "hash"

gantry.runs.ResourceRef dataclass

A table, collection or topic, named by the system that holds it.

Source code in gantry/runs/model.py
129
130
131
132
133
134
135
136
137
@dataclass(frozen=True, slots=True)
class ResourceRef:
    """A table, collection or topic, named by the system that holds it."""

    system: str
    resource: str

    def as_dict(self) -> dict[str, object]:
        return {"system": self.system, "resource": self.resource}

gantry.runs.QueryResultRef dataclass

A bounded query's output, as a reference rather than its rows.

Source code in gantry/runs/model.py
140
141
142
143
144
145
146
147
148
149
@dataclass(frozen=True, slots=True)
class QueryResultRef:
    """A bounded query's output, as a reference rather than its rows."""

    rows: int
    inline: bool = True
    truncated: bool = False

    def as_dict(self) -> dict[str, object]:
        return {"rows": self.rows, "inline": self.inline, "truncated": self.truncated}

gantry.runs.AdmissionRecord dataclass

Gantry's authority decision, and what it rested on.

policy and policy_hash name the exact configuration that decided, so a run stays explainable after the policy changes: v0 never re-evaluates work that was already admitted, and a record that pointed at "the policy" rather than at one version of it would quietly start lying the next time someone edited a rule.

request is the normalized question — actor, operation, engine, inputs, outputs, environment — as provider inspection derived it, not as the proposal described itself. codes are the machine-readable reasons; reasons the same thing in a sentence.

Source code in gantry/runs/model.py
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
@dataclass(frozen=True, slots=True)
class AdmissionRecord:
    """Gantry's authority decision, and what it rested on.

    `policy` and `policy_hash` name the exact configuration that decided, so a
    run stays explainable after the policy changes: v0 never re-evaluates work
    that was already admitted, and a record that pointed at "the policy" rather
    than at one version of it would quietly start lying the next time someone
    edited a rule.

    `request` is the normalized question — actor, operation, engine, inputs,
    outputs, environment — as provider inspection derived it, not as the
    proposal described itself. `codes` are the machine-readable reasons;
    `reasons` the same thing in a sentence.
    """

    allowed: bool
    reasons: tuple[str, ...] = ()
    policy: str | None = None
    policy_hash: str | None = None
    matched_rules: tuple[str, ...] = ()
    codes: tuple[str, ...] = ()
    request: Mapping[str, object] | None = None
    decided_at: datetime = field(default_factory=lambda: datetime.now(UTC))

    def as_dict(self) -> dict[str, object]:
        return {
            "allowed": self.allowed,
            "reasons": list(self.reasons),
            "policy": self.policy,
            "policy_hash": self.policy_hash,
            "matched_rules": list(self.matched_rules),
            "codes": list(self.codes),
            "request": None if self.request is None else dict(self.request),
            "decided_at": self.decided_at.isoformat(),
        }

gantry.runs.ExecutionRecord dataclass

What the engine did, as the engine reported it.

Source code in gantry/runs/model.py
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
@dataclass(frozen=True, slots=True)
class ExecutionRecord:
    """What the engine did, as the engine reported it."""

    status: str = "PENDING"
    native_id: str | None = None
    submitted_at: datetime | None = None
    started_at: datetime | None = None
    finished_at: datetime | None = None
    metrics: Mapping[str, object] = field(default_factory=dict)

    @property
    def duration_ms(self) -> int | None:
        if self.started_at is None or self.finished_at is None:
            return None
        return int((self.finished_at - self.started_at).total_seconds() * 1000)

    def as_dict(self) -> dict[str, object]:
        return {
            "status": self.status,
            "native_id": self.native_id,
            "submitted_at": _time(self.submitted_at),
            "started_at": _time(self.started_at),
            "finished_at": _time(self.finished_at),
            "duration_ms": self.duration_ms,
            "metrics": {key: _plain(value) for key, value in self.metrics.items()},
        }

Storage

gantry.runs.RunStore

Bases: Protocol

Persistence for runs, keyed by Gantry run id.

Source code in gantry/runs/store.py
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
@runtime_checkable
class RunStore(Protocol):
    """Persistence for runs, keyed by Gantry run id."""

    def create(self, run: Run) -> None: ...

    def update(self, run: Run) -> None: ...

    def get(self, run_id: str) -> Run | None: ...

    def compare_and_set(self, run_id: str, expected: RunStatus, updated: Run) -> bool:
        """Move a run on only if it is still where the caller last saw it.

        The one operation `update` cannot express. Two hosts may confirm the
        same run at the same moment, and exactly one of them may be the reason
        work gets submitted — so the check and the write have to be one step.
        Returns whether this caller was the one that made the transition.
        """
        ...

compare_and_set

compare_and_set(
    run_id: str, expected: RunStatus, updated: Run
) -> bool

Move a run on only if it is still where the caller last saw it.

The one operation update cannot express. Two hosts may confirm the same run at the same moment, and exactly one of them may be the reason work gets submitted — so the check and the write have to be one step. Returns whether this caller was the one that made the transition.

Source code in gantry/runs/store.py
31
32
33
34
35
36
37
38
39
def compare_and_set(self, run_id: str, expected: RunStatus, updated: Run) -> bool:
    """Move a run on only if it is still where the caller last saw it.

    The one operation `update` cannot express. Two hosts may confirm the
    same run at the same moment, and exactly one of them may be the reason
    work gets submitted — so the check and the write have to be one step.
    Returns whether this caller was the one that made the transition.
    """
    ...

gantry.runs.SQLiteRunStore

Run records in a SQLite file.

One connection behind a lock rather than one per call: SQLite objects are not safe to share between threads by default, and a run is written a handful of times per operation, so contention is not the thing to optimise.

Source code in gantry/runs/sqlite.py
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
class SQLiteRunStore:
    """Run records in a SQLite file.

    One connection behind a lock rather than one per call: SQLite objects are
    not safe to share between threads by default, and a run is written a
    handful of times per operation, so contention is not the thing to optimise.
    """

    def __init__(self, path: str | Path = ".gantry/runs.db") -> None:
        self._path = Path(path)
        self._path.parent.mkdir(parents=True, exist_ok=True)
        self._lock = threading.Lock()
        self._connection = sqlite3.connect(str(self._path), check_same_thread=False)
        self._connection.executescript(_SCHEMA)
        self._connection.commit()

    def create(self, run: Run) -> None:
        self._write(run, insert=True)

    def update(self, run: Run) -> None:
        self._write(run, insert=False)

    def compare_and_set(self, run_id: str, expected: RunStatus, updated: Run) -> bool:
        """One statement, so the check and the write cannot be separated.

        `WHERE status = ?` is the whole mechanism: two hosts confirming the same
        run both run this, SQLite serialises them, and the second one updates no
        rows. Only the caller that changed a row may start work.
        """
        with self._lock:
            cursor = self._connection.execute(
                """
                UPDATE runs SET status = ?, updated_at = ?, document = ?
                WHERE id = ? AND status = ?
                """,
                (
                    updated.status.value,
                    updated.updated_at.isoformat(),
                    updated.to_json(),
                    run_id,
                    expected.value,
                ),
            )
            self._connection.commit()
            return cursor.rowcount == 1

    def get(self, run_id: str) -> Run | None:
        with self._lock:
            row = self._connection.execute(
                "SELECT document FROM runs WHERE id = ?", (run_id,)
            ).fetchone()
        return None if row is None else run_from_dict(json.loads(row[0]))

    def recent(self, *, limit: int = 50, status: str | None = None) -> Sequence[Run]:
        query = "SELECT document FROM runs"
        parameters: tuple[object, ...] = ()
        if status is not None:
            query += " WHERE status = ?"
            parameters = (status,)
        query += " ORDER BY created_at DESC, id DESC LIMIT ?"
        with self._lock:
            rows = self._connection.execute(query, (*parameters, int(limit))).fetchall()
        return tuple(run_from_dict(json.loads(row[0])) for row in rows)

    def close(self) -> None:
        with self._lock:
            self._connection.close()

    def _write(self, run: Run, *, insert: bool) -> None:
        document = run.to_json()
        with self._lock:
            self._connection.execute(
                """
                INSERT INTO runs (id, status, kind, engine, actor, created_at, updated_at, document)
                VALUES (?, ?, ?, ?, ?, ?, ?, ?)
                ON CONFLICT(id) DO UPDATE SET
                    status = excluded.status,
                    updated_at = excluded.updated_at,
                    document = excluded.document
                """,
                (
                    run.id,
                    run.status.value,
                    run.operation.kind.value,
                    run.operation.engine,
                    run.actor.label,
                    run.created_at.isoformat(),
                    run.updated_at.isoformat(),
                    document,
                ),
            )
            self._connection.commit()

compare_and_set

compare_and_set(
    run_id: str, expected: RunStatus, updated: Run
) -> bool

One statement, so the check and the write cannot be separated.

WHERE status = ? is the whole mechanism: two hosts confirming the same run both run this, SQLite serialises them, and the second one updates no rows. Only the caller that changed a row may start work.

Source code in gantry/runs/sqlite.py
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
def compare_and_set(self, run_id: str, expected: RunStatus, updated: Run) -> bool:
    """One statement, so the check and the write cannot be separated.

    `WHERE status = ?` is the whole mechanism: two hosts confirming the same
    run both run this, SQLite serialises them, and the second one updates no
    rows. Only the caller that changed a row may start work.
    """
    with self._lock:
        cursor = self._connection.execute(
            """
            UPDATE runs SET status = ?, updated_at = ?, document = ?
            WHERE id = ? AND status = ?
            """,
            (
                updated.status.value,
                updated.updated_at.isoformat(),
                updated.to_json(),
                run_id,
                expected.value,
            ),
        )
        self._connection.commit()
        return cursor.rowcount == 1

gantry.runs.MemoryRunStore

Process-local storage.

The default, because writing a file into someone's working directory because they imported a library is not a default worth having. It satisfies every part of the contract except the one that matters most — surviving the process — so anything that needs durability configures SQLite.

Source code in gantry/runs/store.py
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
class MemoryRunStore:
    """Process-local storage.

    The default, because writing a file into someone's working directory
    because they imported a library is not a default worth having. It satisfies
    every part of the contract except the one that matters most — surviving the
    process — so anything that needs durability configures SQLite.
    """

    def __init__(self) -> None:
        self._runs: dict[str, Run] = {}
        self._lock = threading.Lock()

    def create(self, run: Run) -> None:
        self._runs[run.id] = run

    def update(self, run: Run) -> None:
        self._runs[run.id] = run

    def get(self, run_id: str) -> Run | None:
        return self._runs.get(run_id)

    def compare_and_set(self, run_id: str, expected: RunStatus, updated: Run) -> bool:
        with self._lock:
            current = self._runs.get(run_id)
            if current is None or current.status is not expected:
                return False
            self._runs[run_id] = updated
            return True

    def recent(self, *, limit: int = 50, status: str | None = None) -> Sequence[Run]:
        runs = sorted(self._runs.values(), key=lambda run: run.created_at, reverse=True)
        if status is not None:
            runs = [run for run in runs if run.status.value == status]
        return tuple(runs[:limit])

gantry.runs.RunPersistenceError

Bases: RuntimeError

The control plane could not record a run.

Raised where the spec says to fail closed: if the initial record cannot be written, no external work is submitted. An engine job that exists without a record of why it was allowed to is the one outcome this design cannot tolerate.

Source code in gantry/runs/store.py
79
80
81
82
83
84
85
86
class RunPersistenceError(RuntimeError):
    """The control plane could not record a run.

    Raised where the spec says to fail closed: if the initial record cannot be
    written, no external work is submitted. An engine job that exists without a
    record of why it was allowed to is the one outcome this design cannot
    tolerate.
    """

gantry.runs.get

get(run_id: str) -> Run | None

The run with this id, or None.

Source code in gantry/runs/store.py
152
153
154
def get(run_id: str) -> Run | None:
    """The run with this id, or `None`."""
    return _default.get(run_id)

gantry.runs.recent

recent(
    *, limit: int = 50, status: str | None = None
) -> Sequence[Run]

Recent runs, newest first. Not required by v0; useful in a terminal.

Source code in gantry/runs/store.py
157
158
159
160
def recent(*, limit: int = 50, status: str | None = None) -> Sequence[Run]:
    """Recent runs, newest first. Not required by v0; useful in a terminal."""
    lister = getattr(_default, "recent", None)
    return () if lister is None else lister(limit=limit, status=status)

gantry.runs.configure

configure(store: RunStore) -> None

Replace the process-default store.

Source code in gantry/runs/store.py
92
93
94
95
def configure(store: RunStore) -> None:
    """Replace the process-default store."""
    global _default
    _default = store