Skip to content

Handles and execution

For work that outlives the call that started it. A handle is a value you can store and come back to from another process — polling it, or cancelling it.

gantry.ExecutionHandle dataclass

A durable reference to one engine job, safe to store and return to.

gantry_id identifies the run to Gantry and native_id identifies it to the engine; keeping both is what makes reconnection possible after the submitting process is gone. All four identifying fields must be non-empty and submitted_at must be timezone-aware, enforced at construction.

Source code in gantry/handle.py
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
@dataclass(frozen=True, slots=True)
class ExecutionHandle:
    """A durable reference to one engine job, safe to store and return to.

    `gantry_id` identifies the run to Gantry and `native_id` identifies it to
    the engine; keeping both is what makes reconnection possible after the
    submitting process is gone. All four identifying fields must be non-empty
    and `submitted_at` must be timezone-aware, enforced at construction.
    """

    gantry_id: str
    engine: str
    target: str
    native_id: str
    submitted_at: datetime = field(default_factory=lambda: datetime.now(UTC))
    metadata: Mapping[str, object] = field(default_factory=dict)

    def __post_init__(self) -> None:
        values = {
            "gantry_id": self.gantry_id,
            "engine": self.engine,
            "target": self.target,
            "native_id": self.native_id,
        }
        for field_name, value in values.items():
            if not value.strip():
                raise ValueError(f"execution handle {field_name} must not be empty")
        if self.submitted_at.tzinfo is None:
            raise ValueError("execution handle submitted_at must be timezone-aware")

gantry.Execution dataclass

Gantry's normalized live view of an engine job.

Source code in gantry/execution.py
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
@dataclass(frozen=True, slots=True)
class Execution:
    """Gantry's normalized live view of an engine job."""

    handle: ExecutionHandle
    state: ExecutionState
    started_at: datetime | None = None
    updated_at: datetime | None = None
    metrics: ExecutionMetrics = field(default_factory=ExecutionMetrics)
    failure: Failure | None = None
    native: Mapping[str, object] = field(default_factory=dict)

    @property
    def terminal(self) -> bool:
        return self.state in {
            ExecutionState.SUCCEEDED,
            ExecutionState.FAILED,
            ExecutionState.CANCELLED,
            ExecutionState.UNKNOWN,
        }

gantry.ExecutionState

Bases: StrEnum

Gantry's normalized job states, mapped from each engine's own.

UNKNOWN is a real state, not an error case: it means Gantry could not establish what happened, which a caller must treat differently from a known failure because the work may still be running.

Source code in gantry/execution.py
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
class ExecutionState(StrEnum):
    """Gantry's normalized job states, mapped from each engine's own.

    `UNKNOWN` is a real state, not an error case: it means Gantry could not
    establish what happened, which a caller must treat differently from a known
    failure because the work may still be running.
    """

    PENDING = "PENDING"
    SUBMITTED = "SUBMITTED"
    RUNNING = "RUNNING"
    SUCCEEDED = "SUCCEEDED"
    FAILED = "FAILED"
    CANCELLED = "CANCELLED"
    UNKNOWN = "UNKNOWN"

gantry.ValidationResult dataclass

What the engine's own planner said about a statement before running it.

This is native validation — EXPLAIN or the engine's validate endpoint — not Gantry's policy check. errors makes a statement inadmissible; warnings do not. Build one with accepted() or rejected().

Source code in gantry/execution.py
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
@dataclass(frozen=True, slots=True)
class ValidationResult:
    """What the engine's own planner said about a statement before running it.

    This is native validation — `EXPLAIN` or the engine's validate endpoint —
    not Gantry's policy check. `errors` makes a statement inadmissible;
    `warnings` do not. Build one with `accepted()` or `rejected()`.
    """

    ok: bool
    errors: tuple[str, ...] = ()
    warnings: tuple[str, ...] = ()
    metadata: Mapping[str, object] = field(default_factory=dict)

    @classmethod
    def accepted(
        cls,
        *,
        warnings: tuple[str, ...] = (),
        metadata: Mapping[str, object] | None = None,
    ) -> ValidationResult:
        return cls(ok=True, warnings=warnings, metadata={} if metadata is None else metadata)

    @classmethod
    def rejected(cls, *errors: str) -> ValidationResult:
        return cls(ok=False, errors=errors)

Running work

The module-level functions act on a process-default control plane. run submits and waits; the others exist for work that outlives the call.

gantry.submit async

submit(
    artifact: Artifact,
    *,
    target: ExecutionTarget,
    context: Context,
    policy: PolicyRequirements,
) -> ExecutionHandle

Validate, admit and start an artifact, returning a durable handle.

The handle outlives this call, so work can be polled or cancelled from another process. Admission happens here: a policy the adapter cannot enforce is refused before anything runs.

Source code in gantry/runtime.py
461
462
463
464
465
466
467
468
469
470
471
472
473
474
async def submit(
    artifact: Artifact,
    *,
    target: ExecutionTarget,
    context: Context,
    policy: PolicyRequirements,
) -> ExecutionHandle:
    """Validate, admit and start an artifact, returning a durable handle.

    The handle outlives this call, so work can be polled or cancelled from
    another process. Admission happens here: a policy the adapter cannot
    enforce is refused before anything runs.
    """
    return await _default.submit(artifact, target=target, context=context, policy=policy)

gantry.wait async

wait(
    handle: ExecutionHandle,
    *,
    verify: Sequence[Verifier] = (),
    poll_interval_seconds: float = 1.0,
) -> Result

Block until an execution finishes, then verify it.

A terminal engine state is not the answer on its own — the verifiers decide whether the result may be believed, and their outcome is part of the Result.

Source code in gantry/runtime.py
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
async def wait(
    handle: ExecutionHandle,
    *,
    verify: Sequence[Verifier] = (),
    poll_interval_seconds: float = 1.0,
) -> Result:
    """Block until an execution finishes, then verify it.

    A terminal engine state is not the answer on its own — the verifiers decide
    whether the result may be believed, and their outcome is part of the
    `Result`.
    """
    return await _default.wait(
        handle,
        verify=verify,
        poll_interval_seconds=poll_interval_seconds,
    )

gantry.run async

run(
    artifact: Artifact,
    *,
    target: ExecutionTarget,
    context: Context,
    policy: PolicyRequirements,
    verify: Sequence[Verifier] = (),
    poll_interval_seconds: float = 1.0,
) -> Result

Submit an artifact and wait for it, returning the verified result.

The common case. Use submit and wait separately when the work outlives the request that started it.

Source code in gantry/runtime.py
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
async def run(
    artifact: Artifact,
    *,
    target: ExecutionTarget,
    context: Context,
    policy: PolicyRequirements,
    verify: Sequence[Verifier] = (),
    poll_interval_seconds: float = 1.0,
) -> Result:
    """Submit an artifact and wait for it, returning the verified result.

    The common case. Use `submit` and `wait` separately when the work outlives
    the request that started it.
    """
    return await _default.run(
        artifact,
        target=target,
        context=context,
        policy=policy,
        verify=verify,
        poll_interval_seconds=poll_interval_seconds,
    )

gantry.get async

get(handle: ExecutionHandle) -> Execution

Return the current Execution for handle without waiting.

Refreshes state from the adapter registered for the handle's target. Never raises: a missing adapter, an adapter that fails, or a reply for a different handle all come back as an Execution in state UNKNOWN carrying the reason as its failure.

Source code in gantry/runtime.py
477
478
479
480
481
482
483
484
485
async def get(handle: ExecutionHandle) -> Execution:
    """Return the current `Execution` for `handle` without waiting.

    Refreshes state from the adapter registered for the handle's target. Never
    raises: a missing adapter, an adapter that fails, or a reply for a
    different handle all come back as an `Execution` in state `UNKNOWN`
    carrying the reason as its failure.
    """
    return await _default.get(handle)

gantry.cancel async

cancel(
    handle: ExecutionHandle, *, mode: str = "default"
) -> Execution

Ask the adapter owning handle to stop that execution.

mode is passed through to the adapter; engines that distinguish a graceful stop from a hard kill interpret it. Returns the Execution the adapter reports after the request, which may still be running if the engine stops asynchronously. Like get, a missing or failing adapter yields an Execution in state UNKNOWN rather than an exception.

Source code in gantry/runtime.py
507
508
509
510
511
512
513
514
515
516
async def cancel(handle: ExecutionHandle, *, mode: str = "default") -> Execution:
    """Ask the adapter owning `handle` to stop that execution.

    `mode` is passed through to the adapter; engines that distinguish a
    graceful stop from a hard kill interpret it. Returns the `Execution` the
    adapter reports after the request, which may still be running if the
    engine stops asynchronously. Like `get`, a missing or failing adapter
    yields an `Execution` in state `UNKNOWN` rather than an exception.
    """
    return await _default.cancel(handle, mode=mode)

gantry.configure

configure(
    *,
    adapters: Mapping[str, ExecutionAdapter],
    store: ExecutionStore | None = None,
) -> None

Replace the default control plane's adapters and execution store.

Source code in gantry/runtime.py
453
454
455
456
457
458
def configure(
    *, adapters: Mapping[str, ExecutionAdapter], store: ExecutionStore | None = None
) -> None:
    """Replace the default control plane's adapters and execution store."""
    global _default
    _default = ControlPlane(adapters, store)

gantry.register_adapter

register_adapter(
    target_kind: str, adapter: ExecutionAdapter
) -> None

Register adapter on the process-default control plane for target_kind.

Raises ValueError if the target kind is empty. A later registration for the same kind replaces the earlier one.

Source code in gantry/runtime.py
444
445
446
447
448
449
450
def register_adapter(target_kind: str, adapter: ExecutionAdapter) -> None:
    """Register `adapter` on the process-default control plane for `target_kind`.

    Raises `ValueError` if the target kind is empty. A later registration for
    the same kind replaces the earlier one.
    """
    _default.register_adapter(target_kind, adapter)

The control plane

Construct one directly to keep adapters and execution history isolated from the process default — two control planes do not see each other's runs.

gantry.ControlPlane

Routes executions to adapters by target kind and records their state.

Holds the adapter registry and the execution store behind the module-level submit/get/wait/cancel/run helpers. Construct one directly to keep adapters and history isolated from the process default.

Source code in gantry/runtime.py
 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
 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
146
147
148
149
150
151
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
188
189
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
217
218
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
class ControlPlane:
    """Routes executions to adapters by target kind and records their state.

    Holds the adapter registry and the execution store behind the
    module-level `submit`/`get`/`wait`/`cancel`/`run` helpers. Construct one
    directly to keep adapters and history isolated from the process default.
    """

    def __init__(
        self,
        adapters: Mapping[str, ExecutionAdapter] | None = None,
        store: ExecutionStore | None = None,
    ) -> None:
        self._adapters = dict(adapters or {})
        self._store = store or MemoryExecutionStore()

    def register_adapter(self, target_kind: str, adapter: ExecutionAdapter) -> None:
        if not target_kind.strip():
            raise ValueError("target kind must not be empty")
        self._adapters[target_kind] = adapter

    def _adapter(self, target_kind: str) -> ExecutionAdapter | None:
        return self._adapters.get(target_kind)

    async def submit(
        self,
        artifact: Artifact,
        *,
        target: ExecutionTarget,
        context: Context,
        policy: PolicyRequirements,
    ) -> ExecutionHandle:
        adapter = self._adapter(target.kind)
        if adapter is None:
            failure = Failure(
                kind=FailureKind.VALIDATION_ERROR,
                retryable=False,
                message=f"no adapter registered for target: {target.kind}",
            )
            admission = AdmissionDecision(
                allowed=False,
                validation=ValidationResult.rejected(failure.message),
                capabilities=_empty_capabilities(),
                reasons=(failure,),
            )
            raise SubmissionError(Result.rejected(admission))

        try:
            capabilities_value: object = adapter.capabilities()
        except Exception as error:
            failure = _exception_failure(
                FailureKind.VALIDATION_ERROR, "adapter capabilities", error
            )
            admission = AdmissionDecision(
                allowed=False,
                validation=ValidationResult.rejected(failure.message),
                capabilities=_empty_capabilities(),
                reasons=(failure,),
            )
            raise SubmissionError(Result.rejected(admission)) from error
        if not isinstance(capabilities_value, AdapterCapabilities):
            failure = Failure(
                kind=FailureKind.VALIDATION_ERROR,
                retryable=False,
                message="adapter returned invalid capabilities",
            )
            admission = AdmissionDecision(
                allowed=False,
                validation=ValidationResult.rejected(failure.message),
                capabilities=_empty_capabilities(),
                reasons=(failure,),
            )
            raise SubmissionError(Result.rejected(admission))
        capabilities = capabilities_value

        try:
            validation_value: object = await adapter.validate(
                artifact=artifact,
                target=target,
                context=context,
                policy=policy,
            )
        except Exception as error:
            validation_value = ValidationResult.rejected(
                _exception_failure(
                    FailureKind.VALIDATION_ERROR, "adapter validation", error
                ).message
            )
        validation = (
            validation_value
            if isinstance(validation_value, ValidationResult)
            else ValidationResult.rejected("adapter returned an invalid validation result")
        )
        admission = admit(validation, capabilities, policy)
        if not admission.allowed:
            raise SubmissionError(Result.rejected(admission))

        try:
            handle_value: object = await adapter.submit(
                artifact=artifact,
                target=target,
                context=context,
            )
        except Exception as error:
            failure = _exception_failure(FailureKind.SUBMISSION_ERROR, "adapter submission", error)
            raise SubmissionError(
                Result.terminal_failure(ResultStatus.FAILED, failure, admission=admission)
            ) from error
        if not isinstance(handle_value, ExecutionHandle):
            failure = Failure(
                kind=FailureKind.SUBMISSION_ERROR,
                retryable=False,
                message="adapter returned an invalid execution handle",
            )
            raise SubmissionError(
                Result.terminal_failure(ResultStatus.FAILED, failure, admission=admission)
            )
        handle = handle_value
        if handle.target != target.kind:
            failure = Failure(
                kind=FailureKind.SUBMISSION_ERROR,
                retryable=False,
                message="execution handle target does not match the submitted target",
            )
            raise SubmissionError(
                Result.terminal_failure(
                    ResultStatus.FAILED,
                    failure,
                    handle=handle,
                    admission=admission,
                )
            )

        record = RunRecord(handle, artifact, target, context, policy, admission)
        try:
            await self._store.put(record)
        except Exception as error:
            with suppress(Exception):
                await adapter.cancel(handle=handle)
            failure = _exception_failure(FailureKind.SUBMISSION_ERROR, "execution store", error)
            raise SubmissionError(
                Result.terminal_failure(
                    ResultStatus.FAILED,
                    failure,
                    handle=handle,
                    admission=admission,
                )
            ) from error
        return handle

    async def get(self, handle: ExecutionHandle) -> Execution:
        adapter = self._adapter(handle.target)
        if adapter is None:
            return _unknown_execution(handle, f"no adapter registered for target: {handle.target}")
        try:
            execution_value: object = await adapter.status(handle=handle)
        except Exception as error:
            return _unknown_execution(
                handle,
                _exception_failure(FailureKind.UNKNOWN, "adapter status", error).message,
            )
        if not isinstance(execution_value, Execution):
            return _unknown_execution(handle, "adapter returned an invalid execution")
        if execution_value.handle != handle:
            return _unknown_execution(handle, "execution handle does not match the requested job")
        return execution_value

    async def cancel(self, handle: ExecutionHandle, *, mode: str = "default") -> Execution:
        adapter = self._adapter(handle.target)
        if adapter is None:
            return _unknown_execution(handle, f"no adapter registered for target: {handle.target}")
        try:
            execution_value: object = await adapter.cancel(handle=handle, mode=mode)
        except Exception as error:
            return _unknown_execution(
                handle,
                _exception_failure(FailureKind.UNKNOWN, "adapter cancellation", error).message,
            )
        if not isinstance(execution_value, Execution) or execution_value.handle != handle:
            return _unknown_execution(handle, "adapter returned an invalid cancelled execution")
        return execution_value

    async def wait(
        self,
        handle: ExecutionHandle,
        *,
        verify: Sequence[Verifier] = (),
        poll_interval_seconds: float = 1.0,
    ) -> Result:
        if poll_interval_seconds < 0:
            raise ValueError("poll interval must not be negative")
        try:
            record = await self._store.get(handle.gantry_id)
        except Exception as error:
            return _unknown(
                _exception_failure(FailureKind.UNKNOWN, "execution store", error).message,
                handle=handle,
            )
        if record is None or record.handle != handle:
            return _unknown("no persisted execution record matches the handle", handle=handle)

        adapter = self._adapter(handle.target)
        if adapter is None:
            return _unknown(f"no adapter registered for target: {handle.target}", handle=handle)

        while True:
            execution = await self.get(handle)
            if execution.state is ExecutionState.SUCCEEDED:
                break
            if execution.state is ExecutionState.CANCELLED:
                failure = execution.failure or Failure(
                    kind=FailureKind.CANCELLED,
                    retryable=False,
                    message="execution was cancelled",
                )
                return Result.terminal_failure(
                    ResultStatus.CANCELLED,
                    failure,
                    handle=handle,
                    execution=execution,
                    admission=record.admission,
                )
            if execution.state is ExecutionState.FAILED:
                failure = execution.failure or Failure(
                    kind=FailureKind.ENGINE_ERROR,
                    retryable=False,
                    message="engine execution failed",
                )
                return Result.terminal_failure(
                    ResultStatus.FAILED,
                    failure,
                    handle=handle,
                    execution=execution,
                    admission=record.admission,
                )
            if execution.state is ExecutionState.UNKNOWN:
                return Result.terminal_failure(
                    ResultStatus.UNKNOWN,
                    execution.failure
                    or Failure(FailureKind.UNKNOWN, False, "execution state is unknown"),
                    handle=handle,
                    execution=execution,
                    admission=record.admission,
                )
            if _timed_out(execution, record.policy, record.admission.capabilities):
                await self.cancel(handle)
                failure = Failure(
                    kind=FailureKind.TIMEOUT,
                    retryable=True,
                    message=f"execution exceeded {record.policy.max_runtime_seconds} seconds",
                )
                return Result.terminal_failure(
                    ResultStatus.FAILED,
                    failure,
                    handle=handle,
                    execution=execution,
                    admission=record.admission,
                )
            await asyncio.sleep(poll_interval_seconds)

        try:
            engine_result_value: object = await adapter.result(handle=handle)
        except Exception as error:
            failure = _exception_failure(FailureKind.ENGINE_ERROR, "adapter result", error)
            return Result.terminal_failure(
                ResultStatus.FAILED,
                failure,
                handle=handle,
                execution=execution,
                admission=record.admission,
            )
        if not isinstance(engine_result_value, ExecutionResult):
            failure = Failure(FailureKind.ENGINE_ERROR, False, "adapter returned an invalid result")
            return Result.terminal_failure(
                ResultStatus.FAILED,
                failure,
                handle=handle,
                execution=execution,
                admission=record.admission,
            )
        engine_result = engine_result_value
        if engine_result.handle != handle:
            failure = Failure(
                FailureKind.ENGINE_ERROR, False, "result handle does not match execution"
            )
            return Result.terminal_failure(
                ResultStatus.FAILED,
                failure,
                handle=handle,
                execution=execution,
                admission=record.admission,
            )
        if not engine_result.ok:
            return Result.terminal_failure(
                ResultStatus.FAILED,
                engine_result.failure
                or Failure(FailureKind.ENGINE_ERROR, False, "engine result failed"),
                handle=handle,
                execution=execution,
                admission=record.admission,
            )

        verification = await _verify_all(
            verify,
            artifact=record.artifact,
            context=record.context,
            execution=execution,
            engine_result=engine_result,
        )
        return Result.from_execution(
            execution=execution,
            engine_result=engine_result,
            verification=verification,
            admission=record.admission,
        )

    async def run(
        self,
        artifact: Artifact,
        *,
        target: ExecutionTarget,
        context: Context,
        policy: PolicyRequirements,
        verify: Sequence[Verifier] = (),
        poll_interval_seconds: float = 1.0,
    ) -> Result:
        try:
            handle = await self.submit(artifact, target=target, context=context, policy=policy)
        except SubmissionError as error:
            return error.result
        return await self.wait(
            handle,
            verify=verify,
            poll_interval_seconds=poll_interval_seconds,
        )

gantry.ExecutionStore

Bases: Protocol

Persistence for RunRecords, keyed by gantry_id.

Implement this to let handles outlive the process that created them. The default is MemoryExecutionStore, which does not.

Source code in gantry/store.py
35
36
37
38
39
40
41
42
43
44
class ExecutionStore(Protocol):
    """Persistence for `RunRecord`s, keyed by `gantry_id`.

    Implement this to let handles outlive the process that created them. The
    default is `MemoryExecutionStore`, which does not.
    """

    async def put(self, record: RunRecord) -> None: ...

    async def get(self, gantry_id: str) -> RunRecord | None: ...

gantry.RunRecord dataclass

Everything needed to resume governing a run in another process.

Persisted at submission. It keeps the artifact, policy and admission decision beside the handle, because reconnecting to a job is not enough: deciding whether to accept its result requires knowing what was promised when it was admitted.

Source code in gantry/store.py
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
@dataclass(frozen=True, slots=True)
class RunRecord:
    """Everything needed to resume governing a run in another process.

    Persisted at submission. It keeps the artifact, policy and admission
    decision beside the handle, because reconnecting to a job is not enough:
    deciding whether to accept its result requires knowing what was promised
    when it was admitted.
    """

    handle: ExecutionHandle
    artifact: Artifact
    target: ExecutionTarget
    context: Context
    policy: PolicyRequirements
    admission: AdmissionDecision

gantry.MemoryExecutionStore

Process-local store for development and tests.

Source code in gantry/store.py
47
48
49
50
51
52
53
54
55
56
57
class MemoryExecutionStore:
    """Process-local store for development and tests."""

    def __init__(self) -> None:
        self._records: dict[str, RunRecord] = {}

    async def put(self, record: RunRecord) -> None:
        self._records[record.handle.gantry_id] = record

    async def get(self, gantry_id: str) -> RunRecord | None:
        return self._records.get(gantry_id)

Adapters and targets

gantry.ExecutionAdapter

Bases: Protocol

The contract an engine backend implements to be governed by Gantry.

Six methods: declare what you can enforce (capabilities), ask the engine to check a statement (validate), start it (submit), report on it (status), collect it (result), and stop it (cancel). Gantry stays out of the data path — result returns references to where the engine wrote, not the rows.

Source code in gantry/adapter.py
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
class ExecutionAdapter(Protocol):
    """The contract an engine backend implements to be governed by Gantry.

    Six methods: declare what you can enforce (`capabilities`), ask the engine
    to check a statement (`validate`), start it (`submit`), report on it
    (`status`), collect it (`result`), and stop it (`cancel`). Gantry stays out
    of the data path — `result` returns references to where the engine wrote,
    not the rows.
    """

    def capabilities(self) -> AdapterCapabilities: ...

    async def validate(
        self,
        *,
        artifact: Artifact,
        target: ExecutionTarget,
        context: Context,
        policy: PolicyRequirements,
    ) -> ValidationResult: ...

    async def submit(
        self,
        *,
        artifact: Artifact,
        target: ExecutionTarget,
        context: Context,
    ) -> ExecutionHandle: ...

    async def status(self, *, handle: ExecutionHandle) -> Execution: ...

    async def result(self, *, handle: ExecutionHandle) -> ExecutionResult: ...

    async def cancel(self, *, handle: ExecutionHandle, mode: str = "default") -> Execution: ...

gantry.ExecutionTarget dataclass

The engine-specific environment in which an artifact should run.

Source code in gantry/target.py
10
11
12
13
14
15
16
17
18
19
@dataclass(frozen=True, slots=True)
class ExecutionTarget:
    """The engine-specific environment in which an artifact should run."""

    kind: str
    config: Mapping[str, object] = field(default_factory=dict)

    def __post_init__(self) -> None:
        if not self.kind.strip():
            raise ValueError("execution target kind must not be empty")

gantry.AdapterCapabilities dataclass

What an adapter can actually enforce, declared rather than assumed.

Admission compares this against PolicyRequirements: a requirement with no matching capability is refused, because a bound the engine never applies is worse than no bound. Every field defaults to false, so a new adapter is trusted with nothing until it says otherwise.

Source code in gantry/capabilities.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
@dataclass(frozen=True, slots=True)
class AdapterCapabilities:
    """What an adapter can actually enforce, declared rather than assumed.

    Admission compares this against `PolicyRequirements`: a requirement with no
    matching capability is refused, because a bound the engine never applies is
    worse than no bound. Every field defaults to false, so a new adapter is
    trusted with nothing until it says otherwise.
    """

    reconnect: bool = False
    cancellation: bool = False
    runtime_limit: bool = False
    cost_estimation: bool = False
    cost_limit: bool = False
    read_only_execution: bool = False
    write_execution: bool = False
    scoped_credentials: bool = False
    network_isolation: bool = False
    filesystem_isolation: bool = False
    ephemeral_environment: bool = False
    remote_status: bool = False
    metrics: bool = False
    result_reference: bool = False

gantry.Artifact dataclass

An opaque payload produced by an agent or another upstream system.

Source code in gantry/artifact.py
10
11
12
13
14
15
16
17
18
19
20
21
22
23
@dataclass(frozen=True, slots=True)
class Artifact:
    """An opaque payload produced by an agent or another upstream system."""

    payload: object
    kind: str
    metadata: Mapping[str, object] = field(default_factory=dict)
    dependencies: tuple[str, ...] = ()
    declared_inputs: tuple[str, ...] = ()
    declared_outputs: tuple[str, ...] = ()

    def __post_init__(self) -> None:
        if not self.kind.strip():
            raise ValueError("artifact kind must not be empty")

gantry.Context dataclass

Resources and metadata made available for one run.

Source code in gantry/context.py
10
11
12
13
14
15
@dataclass(frozen=True, slots=True)
class Context:
    """Resources and metadata made available for one run."""

    resources: Mapping[str, object] = field(default_factory=dict)
    metadata: Mapping[str, object] = field(default_factory=dict)

gantry.Tool dataclass

A named JSON-schema operation that an agent framework can invoke.

Source code in gantry/tool.py
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
@dataclass(frozen=True, slots=True)
class Tool[T]:
    """A named JSON-schema operation that an agent framework can invoke."""

    name: str
    description: str
    input_schema: Mapping[str, object]
    _handler: Callable[[Mapping[str, object]], Awaitable[T]] = field(repr=False)

    async def invoke(
        self,
        arguments: Mapping[str, object] | None = None,
        /,
        **keyword_arguments: object,
    ) -> T:
        """Invoke the trusted handler with model-supplied arguments."""

        if arguments is not None and keyword_arguments:
            raise TypeError("pass tool arguments as a mapping or keywords, not both")
        values = dict(arguments) if arguments is not None else keyword_arguments
        return await self._handler(values)

invoke async

invoke(
    arguments: Mapping[str, object] | None = None,
    /,
    **keyword_arguments: object,
) -> T

Invoke the trusted handler with model-supplied arguments.

Source code in gantry/tool.py
19
20
21
22
23
24
25
26
27
28
29
30
async def invoke(
    self,
    arguments: Mapping[str, object] | None = None,
    /,
    **keyword_arguments: object,
) -> T:
    """Invoke the trusted handler with model-supplied arguments."""

    if arguments is not None and keyword_arguments:
        raise TypeError("pass tool arguments as a mapping or keywords, not both")
    values = dict(arguments) if arguments is not None else keyword_arguments
    return await self._handler(values)

Admission

The gate between proposing and executing: a policy the adapter cannot enforce is refused rather than warned about.

gantry.PolicyRequirements dataclass

What the application demands of an execution, independent of engine.

Written by application code, never by the agent. The require_* fields name a capability the adapter must have for the work to be admitted at all. Construction rejects a self-contradictory policy (read_only with allow_writes) and non-positive limits, so an impossible policy fails where it is written rather than at admission.

Source code in gantry/policy/requirements.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
@dataclass(frozen=True, slots=True)
class PolicyRequirements:
    """What the application demands of an execution, independent of engine.

    Written by application code, never by the agent. The `require_*` fields
    name a capability the adapter must have for the work to be admitted at all.
    Construction rejects a self-contradictory policy (`read_only` with
    `allow_writes`) and non-positive limits, so an impossible policy fails where
    it is written rather than at admission.
    """

    read_only: bool = False
    allow_writes: bool = False
    max_runtime_seconds: float | None = None
    max_cost_usd: float | None = None
    require_cancel: bool = False
    require_reconnect: bool = False
    require_scoped_credentials: bool = False
    require_network_isolation: bool = False
    require_filesystem_isolation: bool = False
    require_ephemeral_environment: bool = False
    require_metrics: bool = False
    require_result_reference: bool = False

    def __post_init__(self) -> None:
        if self.read_only and self.allow_writes:
            raise ValueError("policy cannot require read-only execution and allow writes")
        if self.max_runtime_seconds is not None and self.max_runtime_seconds <= 0:
            raise ValueError("max runtime must be positive")
        if self.max_cost_usd is not None and self.max_cost_usd < 0:
            raise ValueError("max cost must not be negative")

gantry.admit

admit(
    validation: ValidationResult,
    capabilities: AdapterCapabilities,
    policy: PolicyRequirements,
) -> AdmissionDecision

Decide whether an artifact may run, given what the adapter can enforce.

The decision is the gate between proposing and executing. A policy the adapter cannot enforce is a refusal rather than a warning: a bound nobody applies is worse than no bound, because callers act as though it held.

Source code in gantry/admission.py
 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
 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
def admit(
    validation: ValidationResult,
    capabilities: AdapterCapabilities,
    policy: PolicyRequirements,
) -> AdmissionDecision:
    """Decide whether an artifact may run, given what the adapter can enforce.

    The decision is the gate between proposing and executing. A policy the
    adapter cannot enforce is a refusal rather than a warning: a bound nobody
    applies is worse than no bound, because callers act as though it held.
    """
    reasons = [
        Failure(
            kind=FailureKind.VALIDATION_ERROR,
            retryable=False,
            message=error,
        )
        for error in validation.errors
    ]
    requirements = (
        (policy.read_only, capabilities.read_only_execution, "read_only_execution"),
        (policy.allow_writes, capabilities.write_execution, "write_execution"),
        (policy.max_cost_usd is not None, capabilities.cost_limit, "cost_limit"),
        (policy.require_cancel, capabilities.cancellation, "cancellation"),
        (policy.require_reconnect, capabilities.reconnect, "reconnect"),
        (
            policy.require_scoped_credentials,
            capabilities.scoped_credentials,
            "scoped_credentials",
        ),
        (policy.require_network_isolation, capabilities.network_isolation, "network_isolation"),
        (
            policy.require_filesystem_isolation,
            capabilities.filesystem_isolation,
            "filesystem_isolation",
        ),
        (
            policy.require_ephemeral_environment,
            capabilities.ephemeral_environment,
            "ephemeral_environment",
        ),
        (policy.require_metrics, capabilities.metrics, "metrics"),
        (policy.require_result_reference, capabilities.result_reference, "result_reference"),
    )
    reasons.extend(
        _unsupported(name)
        for required, supported, name in requirements
        if required and not supported
    )

    if policy.max_runtime_seconds is not None:
        gantry_can_enforce = capabilities.reconnect and capabilities.cancellation
        if not capabilities.runtime_limit and not gantry_can_enforce:
            reasons.append(_unsupported("runtime_limit"))

    if not validation.ok and not validation.errors:
        reasons.append(
            Failure(
                kind=FailureKind.VALIDATION_ERROR,
                retryable=False,
                message="adapter validation failed without an error",
            )
        )
    return AdmissionDecision(
        allowed=validation.ok and not reasons,
        validation=validation,
        capabilities=capabilities,
        reasons=tuple(reasons),
    )

gantry.AdmissionDecision dataclass

The outcome of admit: whether this artifact may be submitted.

reasons carries every refusal, not just the first, so a caller sees all of what would have to change. capabilities is kept because later stages need to know what was promised — wait uses it to decide whether Gantry must enforce a runtime limit the engine cannot.

Source code in gantry/admission.py
14
15
16
17
18
19
20
21
22
23
24
25
26
27
@dataclass(frozen=True, slots=True)
class AdmissionDecision:
    """The outcome of `admit`: whether this artifact may be submitted.

    `reasons` carries every refusal, not just the first, so a caller sees all
    of what would have to change. `capabilities` is kept because later stages
    need to know what was promised — `wait` uses it to decide whether Gantry
    must enforce a runtime limit the engine cannot.
    """

    allowed: bool
    validation: ValidationResult
    capabilities: AdapterCapabilities
    reasons: tuple[Failure, ...] = ()

Engine results

What the engine reported, before verification decides whether to accept it. Most callers read Result instead; this is the adapter-facing envelope.

gantry.ExecutionResult dataclass

An engine completion envelope; success still requires verification.

Source code in gantry/execution.py
 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
@dataclass(frozen=True, slots=True)
class ExecutionResult:
    """An engine completion envelope; success still requires verification."""

    ok: bool
    handle: ExecutionHandle
    outputs: tuple[OutputRef, ...] = ()
    metrics: ExecutionMetrics = field(default_factory=ExecutionMetrics)
    failure: Failure | None = None
    native: Mapping[str, object] = field(default_factory=dict)

    @classmethod
    def succeeded(
        cls,
        handle: ExecutionHandle,
        *,
        outputs: tuple[OutputRef, ...] = (),
        metrics: ExecutionMetrics | None = None,
        native: Mapping[str, object] | None = None,
    ) -> ExecutionResult:
        return cls(
            ok=True,
            handle=handle,
            outputs=outputs,
            metrics=ExecutionMetrics() if metrics is None else metrics,
            native={} if native is None else native,
        )

    @classmethod
    def failed(cls, handle: ExecutionHandle, failure: Failure) -> ExecutionResult:
        return cls(ok=False, handle=handle, failure=failure)