Skip to content

SQL

The path most callers use: connect once, configure a policy once, hand the agent a tool.

Connecting

gantry.sql.connect

connect(
    provider: str,
    *,
    policy: Policy | None = None,
    **config: object,
) -> SQLConnection

Resolve a provider preset and create a governed SQL connection.

policy attaches a reusable policy to everything this connection runs. It is trusted configuration: per-operation constraints still apply on top, and neither can be reached from a tool argument.

Source code in gantry/sql/api.py
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
def connect(provider: str, *, policy: Policy | None = None, **config: object) -> SQLConnection:
    """Resolve a provider preset and create a governed SQL connection.

    `policy` attaches a reusable policy to everything this connection runs. It
    is trusted configuration: per-operation constraints still apply on top, and
    neither can be reached from a tool argument.
    """

    if policy is not None and not isinstance(policy, Policy):
        raise TypeError("policy must be a gantry.Policy")
    preset = resolve_provider(provider)
    preset.validate_config(config)
    target = SQLTarget(
        provider=preset.name,
        dialect=preset.dialect,
        driver=preset.driver,
        config=dict(config),
        metadata=preset.metadata,
    )
    return SQLConnection(
        target,
        preset.adapter_factory(target),
        resolve_dialect(preset.dialect),
        policy=policy,
    )

gantry.sql.SQLConnection

A provider-neutral, governed SQL connection.

Configuration and the native adapter remain private so an agent tool cannot accidentally expose credentials or a raw database connection.

Source code in gantry/sql/api.py
 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
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
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
class SQLConnection:
    """A provider-neutral, governed SQL connection.

    Configuration and the native adapter remain private so an agent tool cannot
    accidentally expose credentials or a raw database connection.
    """

    def __init__(
        self,
        target: SQLTarget,
        adapter: SQLAdapter,
        dialect: SQLDialect,
        policy: Policy | None = None,
    ) -> None:
        self._target = target
        self._adapter = adapter
        self._dialect = dialect
        self._policy = policy
        self._plane = ControlPlane(store=MemoryExecutionStore())

    @property
    def policy(self) -> Policy | None:
        """The reusable policy this connection admits against, if any.

        Readable so an application can log what it attached. There is no
        setter: policy is configuration, and something that could be swapped at
        runtime is something an agent-reachable code path could swap.
        """
        return self._policy

    @property
    def provider(self) -> str:
        return self._target.provider

    @property
    def dialect(self) -> str:
        return self._target.dialect

    def capabilities(self) -> SQLCapabilities:
        """Return the adapter's declared SQL capabilities without credentials."""

        return self._adapter.capabilities()

    async def describe(self) -> DatabaseSchema:
        """Return normalized catalogs, schemas, tables, and columns."""

        if not self._adapter.capabilities().describe_schema:
            raise NotImplementedError(f"{self.provider} does not support schema discovery")
        return await self._adapter.describe(self._target)

    async def validate(
        self,
        sql: str,
        *,
        policy: SQLPolicy | None = None,
        context: Context | None = None,
    ) -> ValidationResult:
        active_policy = policy or SQLPolicy()
        bridge = self._bridge(active_policy)
        return await bridge.validate(
            artifact=Artifact(sql, "sql"),
            target=self._execution_target(),
            context=context or Context(),
            policy=self._adapter.capabilities().policy_requirements(active_policy),
        )

    async def explain(self, sql: str) -> ExplainResult:
        if not self._adapter.capabilities().explain:
            return ExplainResult(supported=False)
        return await self._adapter.explain(sql, self._target)

    async def submit(
        self,
        sql: str,
        *,
        policy: SQLPolicy | None = None,
        context: Context | None = None,
    ) -> ExecutionHandle:
        active_policy = policy or SQLPolicy()
        bridge = self._install_bridge(active_policy)
        return await self._plane.submit(
            Artifact(sql, "sql"),
            target=self._execution_target(),
            context=context or Context(),
            policy=bridge.adapter.capabilities().policy_requirements(active_policy),
        )

    async def execute(
        self,
        sql: str,
        *,
        policy: SQLPolicy | None = None,
        context: Context | None = None,
        verify: Sequence[Verifier] = (),
        poll_interval_seconds: float = 0.05,
    ) -> Result:
        active_policy = policy or SQLPolicy()
        bridge = self._install_bridge(active_policy)
        return await self._plane.run(
            Artifact(sql, "sql"),
            target=self._execution_target(),
            context=context or Context(),
            policy=bridge.adapter.capabilities().policy_requirements(active_policy),
            verify=verify,
            poll_interval_seconds=poll_interval_seconds,
        )

    def query(
        self,
        *,
        read_only: bool = True,
        schemas: Sequence[str] = (),
        tables: Sequence[str] = (),
        denied_tables: Sequence[str] = (),
        max_rows: int = 1_000,
        timeout: float = 30,
        max_bytes_scanned: int | None = None,
        max_cost_usd: float | None = None,
        allow_multiple_statements: bool = False,
        checks: Sequence[Verifier | MaterializationCheck] = (),
        verify: Sequence[Verifier | MaterializationCheck] | None = None,
    ) -> SQLQuery:
        """Configure a governed query operation.

        `checks` takes the same `gantry.verify` checks a materialization takes,
        evaluated against the rows the query returns rather than a destination,
        plus any `Verifier` that inspects the execution itself.
        """

        from gantry.sql.query import SQLQuery

        configured_policy = SQLPolicy(
            read_only=read_only,
            allowed_schemas=schemas,
            allowed_tables=tables,
            denied_tables=denied_tables,
            max_rows=max_rows,
            timeout_seconds=timeout,
            max_bytes_scanned=max_bytes_scanned,
            max_cost_usd=max_cost_usd,
            allow_multiple_statements=allow_multiple_statements,
        )
        if verify is not None and checks:
            raise TypeError("pass trusted checks with checks= (verify= is a compatibility alias)")
        trusted = checks if verify is None else verify
        return SQLQuery(self, configured_policy, tuple(trusted))

    async def _query(
        self,
        sql: str,
        *,
        policy: SQLPolicy,
        context: Context | None = None,
        trusted_verify: Sequence[Verifier | MaterializationCheck] = (),
        agent_verify: Sequence[MaterializationCheck] = (),
        agent_capabilities: frozenset[str] = frozenset(),
    ) -> Run:
        """Execute SQL for a configured query operation, and record the run.

        Returns the `Run`: the durable record of what was asked, what was
        allowed, what ran and what was decided. A query's rows travel on it as
        `run.rows`, and are dropped on the way to storage — a run record is not
        a place for result sets to accumulate.

        Two kinds of check arrive here. A `Verifier` inspects the execution and
        runs inside the control plane, as it always has. A `MaterializationCheck`
        — the `gantry.verify` library — is evaluated afterwards against the rows
        that came back, described as a table, so the same `row_count` means the
        same thing whether the caller is querying or materializing.
        """
        verifiers = tuple(item for item in trusted_verify if isinstance(item, Verifier))
        trusted_checks = tuple(item for item in trusted_verify if not isinstance(item, Verifier))

        # Recorded before anything reaches the engine. If this raises, nothing
        # is submitted — an engine job without a record of why it was allowed
        # is the one outcome this ordering exists to prevent.
        recorder = RunRecorder(
            kind=OperationKind.QUERY,
            engine="sql",
            provider=self.provider,
            proposal=sql,
            agent_verification=tuple(type(check).__name__ for check in agent_verify),
        )

        # Admission, in the sense the spec means: is this contract even
        # satisfiable? A contradictory one is settled here, before the engine is
        # asked to do work that could never be accepted.
        from gantry.verify import (
            VerificationConflict,
            VerificationInputError,
            VerificationUnsupported,
            detect_conflicts,
            validate_agent_checks,
        )

        try:
            agent_checks = validate_agent_checks(
                tuple(agent_verify), capabilities=agent_capabilities
            )
            detect_conflicts(trusted_checks, agent_checks)
        except VerificationUnsupported as error:
            return recorder.rejected(
                (str(error),), status=RunStatus.VERIFICATION_UNSUPPORTED
            ).with_failure(Failure(FailureKind.VERIFICATION_UNSUPPORTED, False, str(error)))
        except VerificationConflict as error:
            return recorder.rejected(
                (str(error),), status=RunStatus.VERIFICATION_CONFLICT
            ).with_failure(Failure(FailureKind.VERIFICATION_CONFLICT, False, str(error)))
        except VerificationInputError as error:
            return recorder.rejected((str(error),), status=RunStatus.POLICY_REJECTED).with_failure(
                Failure(FailureKind.VALIDATION_ERROR, False, str(error))
            )
        agent_checks_validated = cast("tuple[MaterializationCheck, ...]", tuple(agent_checks))

        # Authority, before anything external happens. What the SQL touches
        # comes from the dialect, not from the proposal's own account of itself.
        request = query_request(sql, dialect=self._dialect, provider=self.provider, policy=policy)

        async def run() -> Run:
            return await self._execute_query(
                sql,
                policy=policy,
                context=context,
                recorder=recorder,
                verifiers=verifiers,
                trusted_checks=trusted_checks,
                agent_checks=agent_checks_validated,
                agent_verify=agent_verify,
            )

        return await gate(recorder, policy=self._policy, request=request, resume=run)

    async def _execute_query(
        self,
        sql: str,
        *,
        policy: SQLPolicy,
        context: Context | None,
        recorder: RunRecorder,
        verifiers: Sequence[Verifier],
        trusted_checks: Sequence[MaterializationCheck],
        agent_checks: Sequence[MaterializationCheck],
        agent_verify: Sequence[MaterializationCheck],
    ) -> Run:
        """Everything after admission: execute, verify, decide, record.

        Split out so the gate can hold it back. A confirmation that parks the
        run parks exactly this, and resuming runs it unchanged — the same
        proposal against the same already-admitted run.
        """
        result = await self.execute(sql, policy=policy, context=context, verify=verifiers)
        recorder.running(result.handle)
        inline = _find_inline(result, policy.max_rows)
        base = _source_existing(result.verification, CheckSource.TRUSTED)
        verification = _merge(
            base,
            (
                *_result_set_checks(trusted_checks, inline, CheckSource.TRUSTED),
                *_result_set_checks(agent_checks, inline, CheckSource.AGENT),
            ),
        )
        status, failure = _decide(result, verification)
        evidence = _query_evidence(
            sql, result, inline, verification, status, agent_verify=agent_verify
        )
        if status is ResultStatus.REJECTED:
            return recorder.rejected(
                _reasons(failure), status=RunStatus.POLICY_REJECTED
            ).with_failure(failure)
        if status is ResultStatus.VERIFICATION_UNSUPPORTED:
            return recorder.rejected(
                (failure.message if failure else "unsupported verification",),
                status=RunStatus.VERIFICATION_UNSUPPORTED,
                verification=verification,
                evidence=evidence,
            ).with_failure(failure)
        if status is ResultStatus.VERIFICATION_CONFLICT:
            return recorder.rejected(
                (failure.message if failure else "verification conflict",),
                status=RunStatus.VERIFICATION_CONFLICT,
                verification=verification,
                evidence=evidence,
            ).with_failure(failure)
        if status not in {ResultStatus.ACCEPTED, ResultStatus.VERIFICATION_FAILED}:
            return recorder.execution_failed().with_failure(failure)
        recorder.verifying()
        return recorder.decided(
            verification=verification,
            evidence=evidence,
            outputs=_output_refs(self.provider, result.outputs),
            result_ref=None
            if inline is None
            else QueryResultRef(rows=len(inline.rows), inline=True, truncated=inline.truncated),
            inline=inline,
        ).with_failure(failure)

    async def status(self, handle: ExecutionHandle) -> Execution:
        return await self._adapter.status(handle)

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

    async def result(self, handle: ExecutionHandle) -> SQLResult:
        """Recover a completed result directly from its provider-native handle."""

        engine_result = await self._adapter.result(handle)
        maximum = handle.metadata.get("max_rows", 1_000)
        max_rows = maximum if isinstance(maximum, int) and maximum > 0 else 1_000
        return _from_engine_result(engine_result, max_rows)

    async def wait(
        self,
        handle: ExecutionHandle,
        *,
        poll_interval_seconds: float = 1.0,
    ) -> SQLResult:
        """Observe a submitted job and recover its result, including after reconnect."""

        if poll_interval_seconds < 0:
            raise ValueError("poll interval must not be negative")
        while True:
            execution = await self.status(handle)
            if execution.state is ExecutionState.SUCCEEDED:
                return await self.result(handle)
            if execution.terminal:
                status = {
                    ExecutionState.CANCELLED: ResultStatus.CANCELLED,
                    ExecutionState.UNKNOWN: ResultStatus.UNKNOWN,
                }.get(execution.state, ResultStatus.FAILED)
                return SQLResult(
                    status,
                    handle=handle,
                    metrics=execution.metrics,
                    failure=execution.failure,
                )
            timeout = handle.metadata.get("timeout_seconds")
            if (
                isinstance(timeout, (int, float))
                and not isinstance(timeout, bool)
                and (datetime.now(UTC) - handle.submitted_at).total_seconds() >= timeout
            ):
                await self.cancel(handle)
                return SQLResult(
                    ResultStatus.FAILED,
                    handle=handle,
                    metrics=execution.metrics,
                    failure=Failure(
                        FailureKind.TIMEOUT,
                        True,
                        f"SQL execution exceeded {timeout} seconds",
                    ),
                )
            await asyncio.sleep(poll_interval_seconds)

    def materialize(
        self,
        *,
        sources: Sequence[str] = (),
        destinations: Sequence[str],
        create_only: bool = True,
        max_bytes_scanned: int | None = None,
        max_cost_usd: float | None = None,
        timeout: float = 300,
        checks: Sequence[MaterializationCheck] = (),
        verify: Sequence[MaterializationCheck] | None = None,
    ) -> SQLMaterializer:
        """Configure a governed, create-only native SQL materialization operation."""

        from gantry.sql.materialization import (
            MaterializationPolicy,
            SQLMaterializer,
        )
        from gantry.verify import MaterializationCheck

        if verify is not None and checks:
            raise TypeError("pass trusted checks with checks= (verify= is a compatibility alias)")
        configured = checks if verify is None else verify
        trusted_checks: list[MaterializationCheck] = []
        for check in configured:
            if not isinstance(check, MaterializationCheck):
                raise TypeError("verify must contain materialization verification checks")
            trusted_checks.append(check)
        policy = MaterializationPolicy(
            sources=sources,
            destinations=destinations,
            create_only=create_only,
            timeout_seconds=timeout,
            max_bytes_scanned=max_bytes_scanned,
            max_cost_usd=max_cost_usd,
        )
        return SQLMaterializer(self, policy, trusted_checks)

    async def _inspect_table(
        self,
        reference: TableRef,
        *,
        include_row_count: bool = False,
    ) -> Table | None:
        from gantry.sql.materialization import MaterializationAdapter

        if not isinstance(self._adapter, MaterializationAdapter):
            raise NotImplementedError(f"{self.provider} does not support table introspection")
        return await self._adapter.inspect_table(
            reference,
            self._target,
            include_row_count=include_row_count,
        )

    def _materialization_adapter(self) -> SQLAdapter:
        """The native adapter, for observations beyond `inspect_table`.

        Kept private: an agent tool must not reach the driver, and this is the
        only reason anything outside the connection needs it.
        """
        return self._adapter

    def _bridge(self, policy: SQLPolicy) -> SQLExecutionAdapter:
        return SQLExecutionAdapter(self._adapter, self._dialect, self._target, policy)

    def _install_bridge(self, policy: SQLPolicy) -> SQLExecutionAdapter:
        bridge = self._bridge(policy)
        self._plane.register_adapter(self.provider, bridge)
        return bridge

    def _execution_target(self) -> ExecutionTarget:
        # Deliberately exclude provider configuration and credentials.
        return ExecutionTarget(self.provider, {"dialect": self.dialect})

policy property

policy: Policy | None

The reusable policy this connection admits against, if any.

Readable so an application can log what it attached. There is no setter: policy is configuration, and something that could be swapped at runtime is something an agent-reachable code path could swap.

capabilities

capabilities() -> SQLCapabilities

Return the adapter's declared SQL capabilities without credentials.

Source code in gantry/sql/api.py
93
94
95
96
def capabilities(self) -> SQLCapabilities:
    """Return the adapter's declared SQL capabilities without credentials."""

    return self._adapter.capabilities()

describe async

describe() -> DatabaseSchema

Return normalized catalogs, schemas, tables, and columns.

Source code in gantry/sql/api.py
 98
 99
100
101
102
103
async def describe(self) -> DatabaseSchema:
    """Return normalized catalogs, schemas, tables, and columns."""

    if not self._adapter.capabilities().describe_schema:
        raise NotImplementedError(f"{self.provider} does not support schema discovery")
    return await self._adapter.describe(self._target)

query

query(
    *,
    read_only: bool = True,
    schemas: Sequence[str] = (),
    tables: Sequence[str] = (),
    denied_tables: Sequence[str] = (),
    max_rows: int = 1000,
    timeout: float = 30,
    max_bytes_scanned: int | None = None,
    max_cost_usd: float | None = None,
    allow_multiple_statements: bool = False,
    checks: Sequence[Verifier | MaterializationCheck] = (),
    verify: Sequence[Verifier | MaterializationCheck]
    | None = None,
) -> SQLQuery

Configure a governed query operation.

checks takes the same gantry.verify checks a materialization takes, evaluated against the rows the query returns rather than a destination, plus any Verifier that inspects the execution itself.

Source code in gantry/sql/api.py
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
def query(
    self,
    *,
    read_only: bool = True,
    schemas: Sequence[str] = (),
    tables: Sequence[str] = (),
    denied_tables: Sequence[str] = (),
    max_rows: int = 1_000,
    timeout: float = 30,
    max_bytes_scanned: int | None = None,
    max_cost_usd: float | None = None,
    allow_multiple_statements: bool = False,
    checks: Sequence[Verifier | MaterializationCheck] = (),
    verify: Sequence[Verifier | MaterializationCheck] | None = None,
) -> SQLQuery:
    """Configure a governed query operation.

    `checks` takes the same `gantry.verify` checks a materialization takes,
    evaluated against the rows the query returns rather than a destination,
    plus any `Verifier` that inspects the execution itself.
    """

    from gantry.sql.query import SQLQuery

    configured_policy = SQLPolicy(
        read_only=read_only,
        allowed_schemas=schemas,
        allowed_tables=tables,
        denied_tables=denied_tables,
        max_rows=max_rows,
        timeout_seconds=timeout,
        max_bytes_scanned=max_bytes_scanned,
        max_cost_usd=max_cost_usd,
        allow_multiple_statements=allow_multiple_statements,
    )
    if verify is not None and checks:
        raise TypeError("pass trusted checks with checks= (verify= is a compatibility alias)")
    trusted = checks if verify is None else verify
    return SQLQuery(self, configured_policy, tuple(trusted))

result async

result(handle: ExecutionHandle) -> SQLResult

Recover a completed result directly from its provider-native handle.

Source code in gantry/sql/api.py
357
358
359
360
361
362
363
async def result(self, handle: ExecutionHandle) -> SQLResult:
    """Recover a completed result directly from its provider-native handle."""

    engine_result = await self._adapter.result(handle)
    maximum = handle.metadata.get("max_rows", 1_000)
    max_rows = maximum if isinstance(maximum, int) and maximum > 0 else 1_000
    return _from_engine_result(engine_result, max_rows)

wait async

wait(
    handle: ExecutionHandle,
    *,
    poll_interval_seconds: float = 1.0,
) -> SQLResult

Observe a submitted job and recover its result, including after reconnect.

Source code in gantry/sql/api.py
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
async def wait(
    self,
    handle: ExecutionHandle,
    *,
    poll_interval_seconds: float = 1.0,
) -> SQLResult:
    """Observe a submitted job and recover its result, including after reconnect."""

    if poll_interval_seconds < 0:
        raise ValueError("poll interval must not be negative")
    while True:
        execution = await self.status(handle)
        if execution.state is ExecutionState.SUCCEEDED:
            return await self.result(handle)
        if execution.terminal:
            status = {
                ExecutionState.CANCELLED: ResultStatus.CANCELLED,
                ExecutionState.UNKNOWN: ResultStatus.UNKNOWN,
            }.get(execution.state, ResultStatus.FAILED)
            return SQLResult(
                status,
                handle=handle,
                metrics=execution.metrics,
                failure=execution.failure,
            )
        timeout = handle.metadata.get("timeout_seconds")
        if (
            isinstance(timeout, (int, float))
            and not isinstance(timeout, bool)
            and (datetime.now(UTC) - handle.submitted_at).total_seconds() >= timeout
        ):
            await self.cancel(handle)
            return SQLResult(
                ResultStatus.FAILED,
                handle=handle,
                metrics=execution.metrics,
                failure=Failure(
                    FailureKind.TIMEOUT,
                    True,
                    f"SQL execution exceeded {timeout} seconds",
                ),
            )
        await asyncio.sleep(poll_interval_seconds)

materialize

materialize(
    *,
    sources: Sequence[str] = (),
    destinations: Sequence[str],
    create_only: bool = True,
    max_bytes_scanned: int | None = None,
    max_cost_usd: float | None = None,
    timeout: float = 300,
    checks: Sequence[MaterializationCheck] = (),
    verify: Sequence[MaterializationCheck] | None = None,
) -> SQLMaterializer

Configure a governed, create-only native SQL materialization operation.

Source code in gantry/sql/api.py
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
def materialize(
    self,
    *,
    sources: Sequence[str] = (),
    destinations: Sequence[str],
    create_only: bool = True,
    max_bytes_scanned: int | None = None,
    max_cost_usd: float | None = None,
    timeout: float = 300,
    checks: Sequence[MaterializationCheck] = (),
    verify: Sequence[MaterializationCheck] | None = None,
) -> SQLMaterializer:
    """Configure a governed, create-only native SQL materialization operation."""

    from gantry.sql.materialization import (
        MaterializationPolicy,
        SQLMaterializer,
    )
    from gantry.verify import MaterializationCheck

    if verify is not None and checks:
        raise TypeError("pass trusted checks with checks= (verify= is a compatibility alias)")
    configured = checks if verify is None else verify
    trusted_checks: list[MaterializationCheck] = []
    for check in configured:
        if not isinstance(check, MaterializationCheck):
            raise TypeError("verify must contain materialization verification checks")
        trusted_checks.append(check)
    policy = MaterializationPolicy(
        sources=sources,
        destinations=destinations,
        create_only=create_only,
        timeout_seconds=timeout,
        max_bytes_scanned=max_bytes_scanned,
        max_cost_usd=max_cost_usd,
    )
    return SQLMaterializer(self, policy, trusted_checks)

Configuring an operation

gantry.sql.SQLQuery dataclass

A query policy configured once for direct calls or agent tools.

Source code in gantry/sql/query.py
 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
 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
@dataclass(frozen=True, slots=True)
class SQLQuery:
    """A query policy configured once for direct calls or agent tools."""

    _connection: SQLConnection = field(repr=False)
    _policy: SQLPolicy = field(repr=False)
    _verify: tuple[Verifier | MaterializationCheck, ...] = field(default=(), repr=False)

    _agent_capabilities = frozenset({"not_empty", "row_count", "required_columns", "null_rate"})

    async def __call__(
        self,
        sql: str,
        *,
        context: Context | None = None,
        verify: Sequence[MaterializationCheck] = (),
    ) -> Run:
        return await self._connection._query(
            sql,
            policy=self._policy,
            agent_capabilities=self._agent_capabilities,
            context=context,
            trusted_verify=self._verify,
            agent_verify=tuple(verify),
        )

    def tool(
        self,
        *,
        name: str = "query_sql",
        description: str = "Run a governed read-only SQL query.",
    ) -> Tool[Run]:
        """Return the narrow framework-neutral form of this query operation."""

        if not self._policy.read_only:
            raise ValueError(
                "query.tool() requires read_only=True; use db.materialize() for agent writes"
            )
        return Tool(
            name=name,
            description=description,
            input_schema={
                "type": "object",
                "properties": {
                    "sql": {"type": "string"},
                    "verify": verification_schema(self._agent_capabilities),
                },
                "required": ["sql"],
                "additionalProperties": False,
            },
            _handler=self._invoke_tool,
        )

    async def _invoke_tool(self, arguments: Mapping[str, object]) -> Run:
        _require_query_arguments(arguments)
        sql = arguments.get("sql")
        if not isinstance(sql, str):
            raise TypeError("sql must be a string")
        try:
            checks = parse_agent_checks(
                arguments.get("verify"), capabilities=self._agent_capabilities
            )
        except VerificationUnsupported as error:
            return self._refused(sql, error, RunStatus.VERIFICATION_UNSUPPORTED)
        except VerificationInputError as error:
            return self._refused(sql, error, RunStatus.POLICY_REJECTED)
        return await self(sql, verify=checks)  # type: ignore[arg-type]

    def _refused(self, sql: str, error: Exception, status: RunStatus) -> Run:
        """Record a run for a proposal refused before admission.

        A tool response carries a run id, so a malformed proposal has to get one
        too — an agent that is told "rejected" with nothing to refer to cannot
        be asked about it later.
        """
        recorder = RunRecorder(
            kind=OperationKind.QUERY,
            engine="sql",
            provider=self._connection.provider,
            proposal=sql,
        )
        kind = (
            FailureKind.VERIFICATION_UNSUPPORTED
            if status is RunStatus.VERIFICATION_UNSUPPORTED
            else FailureKind.VALIDATION_ERROR
        )
        return recorder.rejected((str(error),), status=status).with_failure(
            Failure(kind, False, str(error))
        )

tool

tool(
    *,
    name: str = "query_sql",
    description: str = "Run a governed read-only SQL query.",
) -> Tool[Run]

Return the narrow framework-neutral form of this query operation.

Source code in gantry/sql/query.py
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
def tool(
    self,
    *,
    name: str = "query_sql",
    description: str = "Run a governed read-only SQL query.",
) -> Tool[Run]:
    """Return the narrow framework-neutral form of this query operation."""

    if not self._policy.read_only:
        raise ValueError(
            "query.tool() requires read_only=True; use db.materialize() for agent writes"
        )
    return Tool(
        name=name,
        description=description,
        input_schema={
            "type": "object",
            "properties": {
                "sql": {"type": "string"},
                "verify": verification_schema(self._agent_capabilities),
            },
            "required": ["sql"],
            "additionalProperties": False,
        },
        _handler=self._invoke_tool,
    )

gantry.sql.SQLPolicy dataclass

The bounds application code puts around agent-written SQL.

Defaults are the safe end: read-only, 1,000 rows, 30 seconds, one statement. An empty allowed_schemas/allowed_tables means no allow-list is applied, while denied_tables always wins. Names are lower-cased and frozen at construction, and a malformed policy raises there rather than at submission — passing a bare string where a collection belongs is a TypeError, not a policy that matches one table.

Declaring a field is not the same as it being enforced: a bound is admitted only when the adapter can apply it.

Source code in gantry/sql/policy.py
 8
 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
@dataclass(frozen=True, slots=True)
class SQLPolicy:
    """The bounds application code puts around agent-written SQL.

    Defaults are the safe end: read-only, 1,000 rows, 30 seconds, one
    statement. An empty `allowed_schemas`/`allowed_tables` means no allow-list
    is applied, while `denied_tables` always wins. Names are lower-cased and
    frozen at construction, and a malformed policy raises there rather than at
    submission — passing a bare string where a collection belongs is a
    `TypeError`, not a policy that matches one table.

    Declaring a field is not the same as it being enforced: a bound is admitted
    only when the adapter can apply it.
    """

    read_only: bool = True
    allowed_schemas: Collection[str] = frozenset()
    allowed_tables: Collection[str] = frozenset()
    denied_tables: Collection[str] = frozenset()
    max_rows: int = 1_000
    timeout_seconds: float = 30
    max_bytes_scanned: int | None = None
    max_cost_usd: float | None = None
    allow_multiple_statements: bool = False

    def __post_init__(self) -> None:
        if not isinstance(self.read_only, bool):
            raise TypeError("read_only must be a boolean")
        if not isinstance(self.allow_multiple_statements, bool):
            raise TypeError("allow_multiple_statements must be a boolean")
        for field_name in ("allowed_schemas", "allowed_tables", "denied_tables"):
            values = getattr(self, field_name)
            if isinstance(values, str):
                raise TypeError(f"{field_name} must be a collection of names, not a string")
            if any(not isinstance(value, str) for value in values):
                raise TypeError(f"{field_name} must contain only strings")
            if any(not value.strip() for value in values):
                raise ValueError(f"{field_name} must not contain empty names")
            object.__setattr__(self, field_name, frozenset(value.lower() for value in values))
        if isinstance(self.max_rows, bool) or not isinstance(self.max_rows, int):
            raise TypeError("max rows must be an integer")
        if self.max_rows <= 0:
            raise ValueError("max rows must be positive")
        if isinstance(self.timeout_seconds, bool) or not isinstance(
            self.timeout_seconds, (int, float)
        ):
            raise TypeError("timeout must be numeric")
        if self.timeout_seconds <= 0:
            raise ValueError("timeout must be positive")
        if self.max_bytes_scanned is not None and (
            isinstance(self.max_bytes_scanned, bool) or not isinstance(self.max_bytes_scanned, int)
        ):
            raise TypeError("max bytes scanned must be an integer")
        if self.max_bytes_scanned is not None and self.max_bytes_scanned < 0:
            raise ValueError("max bytes scanned must not be negative")
        if self.max_cost_usd is not None and (
            isinstance(self.max_cost_usd, bool) or not isinstance(self.max_cost_usd, (int, float))
        ):
            raise TypeError("max cost must be numeric")
        if self.max_cost_usd is not None and self.max_cost_usd < 0:
            raise ValueError("max cost must not be negative")

A policy is a claim about what the adapter will enforce

Declaring a field is not the same as it being enforced. read_only=True is admitted only when the adapter can hold a read-only session in the engine; a row limit is admitted only when the adapter can bound the result. See the capability matrix for which provider enforces what, and what happens when it cannot.

Materialization

Create-only: one CREATE TABLE AS or CREATE VIEW AS, to a schema-qualified destination that does not already exist.

gantry.sql.SQLMaterializer

Credential-free, create-only SQL materializer.

Source code in gantry/sql/materialization.py
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
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
class SQLMaterializer:
    """Credential-free, create-only SQL materializer."""

    name = "materialize_sql"
    description = "Create a governed derived table from approved SQL sources."
    _base_agent_capabilities = frozenset(
        {
            "destination_exists",
            "not_empty",
            "row_count",
            "required_columns",
            "null_rate",
        }
    )

    @property
    def _agent_capabilities(self) -> frozenset[str]:
        capabilities = self._base_agent_capabilities
        if not isinstance(self._connection._materialization_adapter(), ColumnStatistics):
            capabilities = capabilities - {"null_rate"}
        return capabilities

    def __init__(
        self,
        connection: SQLConnection,
        policy: MaterializationPolicy,
        verify: Sequence[MaterializationCheck],
    ) -> None:
        if any(not isinstance(check, MaterializationCheck) for check in verify):
            raise TypeError("verify must contain materialization verification checks")
        self._connection = connection
        self._policy = policy
        self._trusted_checks = tuple(verify)
        self._plans: dict[str, MaterializationPlan] = {}
        self._proposals: dict[str, str] = {}
        self._agent_checks: dict[str, tuple[MaterializationCheck, ...]] = {}

    @property
    def input_schema(self) -> dict[str, object]:
        return {
            "type": "object",
            "properties": {
                "sql": {"type": "string"},
                "verify": verification_schema(self._agent_capabilities),
            },
            "required": ["sql"],
            "additionalProperties": False,
        }

    def tool(
        self,
        *,
        name: str = "materialize_sql",
        description: str = "Create a derived table from approved SQL sources.",
    ) -> Tool[Run]:
        """Return the narrow framework-neutral form of this operation."""

        return Tool(
            name=name,
            description=description,
            input_schema=self.input_schema,
            _handler=self._invoke_tool,
        )

    @property
    def capabilities(self) -> MaterializationCapabilities:
        capabilities = self._connection.capabilities()
        return MaterializationCapabilities(
            create_table_as=capabilities.create_table_as,
            create_view_as=capabilities.create_view_as,
            durable_jobs=capabilities.reconnect,
            cancel=capabilities.cancellation,
            estimate_bytes_scanned=capabilities.bytes_scanned,
            destination_introspection=capabilities.destination_introspection,
            result_reference=capabilities.materialization_reference,
        )

    def inspect(self, proposal: str | MaterializationProposal) -> MaterializationPlan:
        value = (
            proposal
            if isinstance(proposal, MaterializationProposal)
            else MaterializationProposal(proposal)
        )
        return parse_materialization(value.sql)

    async def submit(
        self,
        proposal: str | MaterializationProposal,
        *,
        verify: Sequence[MaterializationCheck] = (),
    ) -> ExecutionHandle:
        value = (
            proposal
            if isinstance(proposal, MaterializationProposal)
            else MaterializationProposal(proposal)
        )
        agent_checks = validate_agent_checks(verify, capabilities=self._agent_capabilities)
        detect_conflicts(self._trusted_checks, agent_checks)
        plan = self.inspect(value)
        await self._admit(plan)
        context = Context(metadata={_PLAN_METADATA: plan})
        try:
            handle = await self._connection.submit(
                value.sql,
                policy=self._sql_policy(),
                context=context,
            )
        except SubmissionError as error:
            failure = error.result.failure or Failure(
                FailureKind.SUBMISSION_ERROR,
                False,
                "materialization submission failed",
            )
            raise MaterializationError(
                _normalize_submission_failure(failure),
                error.result.status,
            ) from error
        self._plans[handle.gantry_id] = plan
        # Identifies the proposal without storing it. Evidence should say which
        # statement was run without becoming a place agent-written SQL
        # accumulates, and a digest is enough to match a run to its proposal.
        proposal_hash = sha256(value.sql.encode()).hexdigest()
        handle = replace(
            handle,
            metadata={
                **handle.metadata,
                "gantry.proposal_hash": proposal_hash,
                "gantry.agent_verification": [check_config(check) for check in agent_checks],
            },
        )
        self._proposals[handle.gantry_id] = proposal_hash
        self._agent_checks[handle.gantry_id] = agent_checks  # type: ignore[assignment]
        return handle

    async def wait(
        self,
        handle: ExecutionHandle,
        *,
        poll_interval_seconds: float = 1.0,
    ) -> MaterializationResult:
        plan = self._plans.get(handle.gantry_id) or _plan_from_handle(handle)
        if plan is None:
            return _failed(
                ResultStatus.UNKNOWN,
                Failure(
                    FailureKind.UNKNOWN,
                    False,
                    "execution handle does not identify a materialization destination",
                ),
            )
        sql_result = await self._connection.wait(
            handle,
            poll_interval_seconds=poll_interval_seconds,
        )
        execution = await self._connection.status(handle)
        if not sql_result.ok:
            verification = VerificationResult.passed()
            result = MaterializationResult(
                status=sql_result.status,
                execution=execution,
                failure=sql_result.failure,
                evidence=_evidence(
                    plan,
                    execution,
                    None,
                    verification,
                    sql_result.status,
                    self._proposal_hash(handle),
                    self._agent_checks_for(handle),
                ),
            )
            return _recorded(result)

        output = _materialized_output(sql_result.outputs, plan.destination)
        if output is None:
            result = MaterializationResult(
                status=ResultStatus.FAILED,
                execution=execution,
                failure=Failure(
                    FailureKind.ENGINE_ERROR,
                    False,
                    "adapter did not return a reference to the materialized destination",
                ),
                evidence=_evidence(
                    plan,
                    execution,
                    None,
                    VerificationResult.passed(),
                    ResultStatus.FAILED,
                    self._proposal_hash(handle),
                    self._agent_checks_for(handle),
                ),
            )
            return _recorded(result)

        agent_checks = self._agent_checks_for(handle)
        verification = await self._verification(plan.destination, agent_checks)
        if not verification.ok:
            unsupported = verification.unsupported_checks
            reasons = unsupported or verification.failed_checks
            message = next(
                (check.message for check in reasons if check.message),
                "materialization verification failed",
            )
            status = ResultStatus.VERIFICATION_FAILED
            result = MaterializationResult(
                status=status,
                output=output,
                execution=execution,
                verification=verification,
                failure=Failure(
                    FailureKind.UNSUPPORTED_VERIFICATION
                    if unsupported
                    else FailureKind.VERIFICATION_FAILED,
                    False,
                    message,
                ),
                evidence=_evidence(
                    plan,
                    execution,
                    output,
                    verification,
                    status,
                    self._proposal_hash(handle),
                    agent_checks,
                ),
            )
            return _recorded(result)
        result = MaterializationResult(
            status=ResultStatus.ACCEPTED,
            output=output,
            execution=execution,
            verification=verification,
            evidence=_evidence(
                plan,
                execution,
                output,
                verification,
                ResultStatus.ACCEPTED,
                self._proposal_hash(handle),
                agent_checks,
            ),
        )
        return _recorded(result)

    async def status(self, handle: ExecutionHandle) -> Execution:
        return await self._connection.status(handle)

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

    async def __call__(
        self,
        proposal: str | MaterializationProposal,
        *,
        verify: Sequence[MaterializationCheck] = (),
    ) -> Run:
        """Materialize, and record the run.

        The run is created before admission, so a proposal that never reaches
        the engine still leaves a record of having been refused — which is the
        case anyone asks about later.
        """
        body = proposal.sql if isinstance(proposal, MaterializationProposal) else str(proposal)
        recorder = RunRecorder(
            kind=OperationKind.MATERIALIZE,
            engine="sql",
            provider=self._connection.provider,
            proposal=body,
            agent_verification=tuple(type(check).__name__ for check in verify),
        )

        # Authority first, and from the parsed plan rather than the proposal's
        # own description of itself. A refusal here means nothing was submitted.
        try:
            plan = self.inspect(proposal)
        except MaterializationError as error:
            return recorder.rejected(
                (error.failure.message,), status=RunStatus.POLICY_REJECTED
            ).with_failure(error.failure)
        except (TypeError, ValueError) as error:
            return _refuse(recorder, error, RunStatus.POLICY_REJECTED)

        request = self._policy_request(plan)

        async def run() -> Run:
            return await self._execute_materialization(proposal, verify=verify, recorder=recorder)

        return await gate(recorder, policy=self._connection.policy, request=request, resume=run)

    async def _execute_materialization(
        self,
        proposal: str | MaterializationProposal,
        *,
        verify: Sequence[MaterializationCheck],
        recorder: RunRecorder,
    ) -> Run:
        """Everything after admission: submit, wait, verify, record.

        Held back by the gate when a rule asks for confirmation, and run
        unchanged when the host confirms — the same proposal, the same run.
        """
        try:
            handle = await self.submit(proposal, verify=verify)
        except VerificationUnsupported as error:
            return _refuse(recorder, error, RunStatus.VERIFICATION_UNSUPPORTED)
        except VerificationConflict as error:
            return _refuse(recorder, error, RunStatus.VERIFICATION_CONFLICT)
        except VerificationInputError as error:
            return _refuse(recorder, error, RunStatus.POLICY_REJECTED)
        except MaterializationError as error:
            return recorder.rejected(
                (error.failure.message,), status=RunStatus.POLICY_REJECTED
            ).with_failure(error.failure)

        plan = self._plans.get(handle.gantry_id)
        recorder.running(handle)
        result = await self.wait(handle, poll_interval_seconds=0.05)
        if result.status is ResultStatus.VERIFICATION_UNSUPPORTED:
            return recorder.rejected(
                (result.failure.message if result.failure else "unsupported verification",),
                status=RunStatus.VERIFICATION_UNSUPPORTED,
            ).with_failure(result.failure)
        if result.status is ResultStatus.VERIFICATION_CONFLICT:
            return recorder.rejected(
                (result.failure.message if result.failure else "verification conflict",),
                status=RunStatus.VERIFICATION_CONFLICT,
            ).with_failure(result.failure)
        if result.status is not ResultStatus.ACCEPTED and result.verification is None:
            return recorder.execution_failed().with_failure(result.failure)
        recorder.verifying()
        return recorder.decided(
            verification=result.verification,
            evidence=result.evidence,
            outputs=()
            if result.output is None
            else (
                ResourceRef(
                    system=self._connection.provider,
                    resource=plan.destination.qualified_name
                    if plan is not None
                    else result.output.uri,
                ),
            ),
        ).with_failure(result.failure)

    async def _invoke_tool(self, arguments: Mapping[str, object]) -> Run:
        unexpected = set(arguments) - {"sql", "verify"}
        if unexpected:
            names = ", ".join(sorted(unexpected))
            raise ValueError(f"unexpected materialization tool arguments: {names}")
        sql = arguments.get("sql")
        if not isinstance(sql, str):
            raise TypeError("sql must be a string")
        try:
            checks = parse_agent_checks(
                arguments.get("verify"), capabilities=self._agent_capabilities
            )
        except (VerificationUnsupported, VerificationInputError) as error:
            return _refuse(
                RunRecorder(
                    kind=OperationKind.MATERIALIZE,
                    engine="sql",
                    provider=self._connection.provider,
                    proposal=sql,
                ),
                error,
                RunStatus.VERIFICATION_UNSUPPORTED
                if isinstance(error, VerificationUnsupported)
                else RunStatus.POLICY_REJECTED,
            )
        return await self(sql, verify=checks)  # type: ignore[arg-type]

    async def _admit(self, plan: MaterializationPlan) -> None:
        capabilities = self._connection.capabilities()
        if (
            plan.operation is MaterializationOperation.CREATE_TABLE_AS
            and not capabilities.create_table_as
        ):
            _reject(FailureKind.OPERATION_NOT_ALLOWED, "adapter does not support CREATE TABLE AS")
        if (
            plan.operation is MaterializationOperation.CREATE_VIEW_AS
            and not capabilities.create_view_as
        ):
            _reject(FailureKind.OPERATION_NOT_ALLOWED, "adapter does not support CREATE VIEW AS")
        if not capabilities.write_execution:
            _reject(
                FailureKind.UNSUPPORTED_POLICY_REQUIREMENT,
                "adapter cannot enforce governed write execution",
            )
        if not capabilities.materialization_reference:
            _reject(
                FailureKind.UNSUPPORTED_POLICY_REQUIREMENT,
                "adapter cannot return a materialized output reference",
            )
        if not capabilities.destination_introspection:
            _reject(
                FailureKind.UNSUPPORTED_POLICY_REQUIREMENT,
                "adapter cannot inspect materialization destinations",
            )
        if self._policy.max_bytes_scanned is not None and not capabilities.bytes_scanned:
            _reject(
                FailureKind.UNSUPPORTED_POLICY_REQUIREMENT,
                "adapter cannot enforce maximum bytes scanned",
            )
        if self._policy.max_cost_usd is not None and not capabilities.cost_limit:
            _reject(
                FailureKind.UNSUPPORTED_POLICY_REQUIREMENT,
                "adapter cannot enforce maximum cost",
            )
        disallowed = [
            source for source in plan.sources if not _matches(source, self._policy.sources)
        ]
        if disallowed:
            names = ", ".join(source.qualified_name for source in disallowed)
            _reject(FailureKind.SOURCE_NOT_ALLOWED, f"source is not allowed: {names}")
        if not _matches(plan.destination, self._policy.destinations):
            _reject(
                FailureKind.DESTINATION_NOT_ALLOWED,
                f"destination is not allowed: {plan.destination.qualified_name}",
            )
        try:
            existing = await self._connection._inspect_table(plan.destination)
        except NotImplementedError as error:
            _reject(
                FailureKind.UNSUPPORTED_POLICY_REQUIREMENT,
                f"adapter cannot inspect materialization destinations: {error}",
            )
        except Exception as error:
            _reject(
                FailureKind.VALIDATION_ERROR,
                f"destination inspection raised {type(error).__name__}: {error}",
            )
        if existing is not None:
            _reject(
                FailureKind.DESTINATION_EXISTS,
                f"destination already exists: {plan.destination.qualified_name}",
            )

    async def _verification(
        self,
        destination: TableRef,
        agent_checks: Sequence[MaterializationCheck] = (),
    ) -> VerificationResult:
        effective = (*self._trusted_checks, *agent_checks)
        if not effective:
            return VerificationResult.passed()
        try:
            table = await self._connection._inspect_table(
                destination,
                include_row_count=any(
                    getattr(verifier, "requires_row_count", False) for verifier in effective
                ),
            )
        except Exception as error:
            failure = VerificationResult.failed(
                f"destination inspection raised {type(error).__name__}: {error}",
                name="destination_introspection",
            )
            check = sourced_result(
                failure.checks[0],
                CheckSource.TRUSTED,
                observation_source="destination",
            )
            return replace(failure, checks=(check,))
        table = await self._with_null_rates(table, destination, effective)
        checks: list[CheckResult] = []
        for index, verifier in enumerate(effective):
            source = CheckSource.TRUSTED if index < len(self._trusted_checks) else CheckSource.AGENT
            try:
                checks.append(
                    sourced_result(
                        verifier.evaluate(table), source, observation_source="destination"
                    )
                )
            except Exception as error:
                checks.append(
                    CheckResult(
                        name=type(verifier).__name__,
                        ok=False,
                        message=f"verification raised {type(error).__name__}: {error}",
                        source=source,
                    )
                )
        return VerificationResult(
            ok=all(check.ok for check in checks),
            checks=tuple(checks),
        )

    def _proposal_hash(self, handle: ExecutionHandle) -> str | None:
        stored = handle.metadata.get("gantry.proposal_hash")
        return self._proposals.get(handle.gantry_id) or (
            stored if isinstance(stored, str) else None
        )

    def _agent_checks_for(self, handle: ExecutionHandle) -> tuple[MaterializationCheck, ...]:
        local = self._agent_checks.get(handle.gantry_id)
        if local is not None:
            return local
        try:
            checks = parse_agent_checks(
                handle.metadata.get("gantry.agent_verification"),
                capabilities=self._agent_capabilities,
            )
        except VerificationInputError:
            return ()
        return checks  # type: ignore[return-value]

    async def _with_null_rates(
        self,
        table: Table | None,
        destination: TableRef,
        checks: Sequence[MaterializationCheck],
    ) -> Table | None:
        """Attach null rates to the destination metadata, if any check needs them.

        Gathered only when asked for: it is a scan per column, and a check that
        nobody configured should cost nothing.
        """
        wanted: list[str] = []
        for verifier in checks:
            wanted.extend(getattr(verifier, "null_rate_columns", ()))
        if table is None or not wanted:
            return table
        adapter = self._connection._materialization_adapter()
        if not isinstance(adapter, ColumnStatistics):
            # Leaves the metadata absent, which makes the check unsupported.
            return table
        try:
            rates = await adapter.column_null_rates(
                destination, self._connection._target, tuple(dict.fromkeys(wanted))
            )
        except Exception:
            return table
        return replace(table, metadata={**table.metadata, "null_rates": dict(rates)})

    def _policy_request(
        self, proposal: str | MaterializationProposal | MaterializationPlan
    ) -> PolicyRequest:
        """The normalized request for one materialization proposal.

        Parsing is what determines the destination, so a proposal that will not
        parse has no determinable effect — reported as unresolved rather than
        guessed at.
        """
        plan = proposal if isinstance(proposal, MaterializationPlan) else self.inspect(proposal)
        return materialize_request(
            provider=self._connection.provider,
            sources=[source.qualified_name for source in plan.sources],
            destination=plan.destination.qualified_name,
            policy=self._sql_policy(),
        )

    def _sql_policy(self) -> SQLPolicy:
        return SQLPolicy(
            read_only=False,
            max_rows=1,
            timeout_seconds=self._policy.timeout_seconds,
            max_bytes_scanned=self._policy.max_bytes_scanned,
            max_cost_usd=self._policy.max_cost_usd,
        )

tool

tool(
    *,
    name: str = "materialize_sql",
    description: str = "Create a derived table from approved SQL sources.",
) -> Tool[Run]

Return the narrow framework-neutral form of this operation.

Source code in gantry/sql/materialization.py
370
371
372
373
374
375
376
377
378
379
380
381
382
383
def tool(
    self,
    *,
    name: str = "materialize_sql",
    description: str = "Create a derived table from approved SQL sources.",
) -> Tool[Run]:
    """Return the narrow framework-neutral form of this operation."""

    return Tool(
        name=name,
        description=description,
        input_schema=self.input_schema,
        _handler=self._invoke_tool,
    )

gantry.sql.MaterializationPolicy dataclass

Where a materialization may read from and write to.

destinations must be non-empty — there is no default place to write — and patterns are fnmatch-style, lower-cased at construction. create_only is fixed at True in v0 and passing anything else raises, so replacement cannot be enabled by a config change before the semantics exist.

Source code in gantry/sql/materialization.py
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
@dataclass(frozen=True, slots=True)
class MaterializationPolicy:
    """Where a materialization may read from and write to.

    `destinations` must be non-empty — there is no default place to write — and
    patterns are `fnmatch`-style, lower-cased at construction. `create_only` is
    fixed at `True` in v0 and passing anything else raises, so replacement
    cannot be enabled by a config change before the semantics exist.
    """

    sources: Collection[str]
    destinations: Collection[str]
    create_only: bool = True
    timeout_seconds: float = 300
    max_bytes_scanned: int | None = None
    max_cost_usd: float | None = None

    def __post_init__(self) -> None:
        for field_name in ("sources", "destinations"):
            values = getattr(self, field_name)
            if isinstance(values, str):
                raise TypeError(f"{field_name} must be a collection of patterns, not a string")
            if any(not isinstance(value, str) for value in values):
                raise TypeError(f"{field_name} must contain only strings")
            if any(not value.strip() for value in values):
                raise ValueError(f"{field_name} must not contain empty patterns")
            object.__setattr__(self, field_name, tuple(value.lower() for value in values))
        if not self.destinations:
            raise ValueError("at least one materialization destination must be allowed")
        if self.create_only is not True:
            raise ValueError("SQL materialization v0 supports create_only=True only")
        if isinstance(self.timeout_seconds, bool) or not isinstance(
            self.timeout_seconds, (int, float)
        ):
            raise TypeError("timeout must be numeric")
        if self.timeout_seconds <= 0:
            raise ValueError("timeout must be positive")
        if self.max_bytes_scanned is not None and (
            isinstance(self.max_bytes_scanned, bool) or not isinstance(self.max_bytes_scanned, int)
        ):
            raise TypeError("max bytes scanned must be an integer")
        if self.max_bytes_scanned is not None and self.max_bytes_scanned < 0:
            raise ValueError("max bytes scanned must not be negative")
        if self.max_cost_usd is not None and (
            isinstance(self.max_cost_usd, bool) or not isinstance(self.max_cost_usd, (int, float))
        ):
            raise TypeError("max cost must be numeric")
        if self.max_cost_usd is not None and self.max_cost_usd < 0:
            raise ValueError("max cost must not be negative")

gantry.sql.MaterializationResult dataclass

The outcome of a materialization: the destination, and whether to trust it.

ok is ACCEPTED only, which requires verification to have passed — a CREATE TABLE AS that ran and produced the wrong table is not accepted. uri names the destination when one was created.

Source code in gantry/sql/materialization.py
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
@dataclass(frozen=True, slots=True)
class MaterializationResult:
    """The outcome of a materialization: the destination, and whether to trust it.

    `ok` is `ACCEPTED` only, which requires verification to have passed — a
    `CREATE TABLE AS` that ran and produced the wrong table is not accepted.
    `uri` names the destination when one was created.
    """

    status: ResultStatus
    output: OutputRef | None = None
    execution: Execution | None = None
    verification: VerificationResult | None = None
    failure: Failure | None = None
    evidence: EvidenceBundle | None = None
    """What Gantry observed while deciding, serializable and outliving the run."""

    @property
    def ok(self) -> bool:
        return self.status is ResultStatus.ACCEPTED

    @property
    def handle(self) -> ExecutionHandle | None:
        return None if self.execution is None else self.execution.handle

    @property
    def uri(self) -> str | None:
        """Return the materialized destination URI when one is available."""

        return None if self.output is None else self.output.uri

evidence class-attribute instance-attribute

evidence: EvidenceBundle | None = None

What Gantry observed while deciding, serializable and outliving the run.

uri property

uri: str | None

Return the materialized destination URI when one is available.

gantry.sql.MaterializationError

Bases: Exception

Raised when a proposal is refused, carrying why and at what stage.

failure is the normalized reason and status the outcome it maps to, defaulting to REJECTED — the common case, where nothing ran. It is an exception rather than a returned value because a refused proposal has no result to describe.

Source code in gantry/sql/materialization.py
295
296
297
298
299
300
301
302
303
304
305
306
307
class MaterializationError(Exception):
    """Raised when a proposal is refused, carrying why and at what stage.

    `failure` is the normalized reason and `status` the outcome it maps to,
    defaulting to `REJECTED` — the common case, where nothing ran. It is an
    exception rather than a returned value because a refused proposal has no
    result to describe.
    """

    def __init__(self, failure: Failure, status: ResultStatus = ResultStatus.REJECTED) -> None:
        self.failure = failure
        self.status = status
        super().__init__(failure.message)

Parsing a proposal

What turns proposed SQL into something checkable. Call these directly to inspect what a statement would do before submitting it.

gantry.sql.parse_materialization

parse_materialization(sql: str) -> MaterializationPlan

Parse caller-proposed SQL into a create-only MaterializationPlan.

Accepts exactly one CREATE TABLE ... AS or CREATE VIEW ... AS whose destination names a schema, and returns the operation, the schema-qualified destination, and the source tables the body reads. Raises MaterializationError for anything else: several statements, a non-create operation, OR REPLACE or IF NOT EXISTS, an unqualified destination, a body that is not a SELECT/WITH query, or a body carrying a second effect such as DELETE or DROP.

Source code in gantry/sql/materialization.py
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
def parse_materialization(sql: str) -> MaterializationPlan:
    """Parse caller-proposed SQL into a create-only `MaterializationPlan`.

    Accepts exactly one `CREATE TABLE ... AS` or `CREATE VIEW ... AS` whose
    destination names a schema, and returns the operation, the schema-qualified
    destination, and the source tables the body reads. Raises
    `MaterializationError` for anything else: several statements, a non-create
    operation, `OR REPLACE` or `IF NOT EXISTS`, an unqualified destination, a
    body that is not a `SELECT`/`WITH` query, or a body carrying a second
    effect such as `DELETE` or `DROP`.
    """
    proposal = MaterializationProposal(sql)
    statements = _split_statements(proposal.sql)
    if len(statements) != 1:
        _reject(FailureKind.OPERATION_NOT_ALLOWED, "materialization must contain one SQL statement")
    tokens = _tokenize(statements[0])
    if not tokens or tokens[0].upper != "CREATE":
        _reject(
            FailureKind.OPERATION_NOT_ALLOWED,
            "only CREATE TABLE AS or CREATE VIEW AS is allowed",
        )

    index = 1
    replace = False
    if _words_at(tokens, index, "OR", "REPLACE"):
        replace = True
        index += 2
    if index >= len(tokens) or tokens[index].upper not in {"TABLE", "VIEW"}:
        _reject(
            FailureKind.OPERATION_NOT_ALLOWED,
            "only CREATE TABLE AS or CREATE VIEW AS is allowed",
        )
    object_type = tokens[index].upper
    index += 1
    if _words_at(tokens, index, "IF", "NOT", "EXISTS"):
        _reject(FailureKind.OPERATION_NOT_ALLOWED, "IF NOT EXISTS is not allowed")
    destination, index = _parse_ref(tokens, index)
    if destination.schema is None:
        _reject(
            FailureKind.DESTINATION_NOT_ALLOWED,
            "materialization destination must include a schema or dataset",
        )
    if replace:
        _reject(FailureKind.OPERATION_NOT_ALLOWED, "replacement materialization is not allowed")
    if index >= len(tokens) or tokens[index].upper != "AS":
        _reject(
            FailureKind.OPERATION_NOT_ALLOWED,
            "materialization must use CREATE TABLE AS or CREATE VIEW AS",
        )
    body_start = index + 1
    if body_start >= len(tokens) or tokens[body_start].upper not in {"SELECT", "WITH"}:
        _reject(FailureKind.OPERATION_NOT_ALLOWED, "materialization body must be a SELECT query")
    for token in tokens[body_start:]:
        if token.kind == "word" and token.upper in _DANGEROUS_BODY_WORDS:
            _reject(
                FailureKind.OPERATION_NOT_ALLOWED,
                f"{token.upper} is not allowed inside a materialization query",
            )
    ctes = _cte_names(tokens, body_start)
    sources = _source_refs(tokens, body_start, ctes)
    operation = (
        MaterializationOperation.CREATE_TABLE_AS
        if object_type == "TABLE"
        else MaterializationOperation.CREATE_VIEW_AS
    )
    return MaterializationPlan(operation, sources, destination, replace=False)

gantry.sql.MaterializationProposal dataclass

Caller-proposed SQL before anything has been checked about it.

Construction only rejects a non-string or blank submission; the naming is the point — what arrives here is a proposal, and it is a plan only after parse_materialization has accepted it.

Source code in gantry/sql/materialization.py
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
@dataclass(frozen=True, slots=True)
class MaterializationProposal:
    """Caller-proposed SQL before anything has been checked about it.

    Construction only rejects a non-string or blank submission; the naming is
    the point — what arrives here is a proposal, and it is a plan only after
    `parse_materialization` has accepted it.
    """

    sql: str

    def __post_init__(self) -> None:
        if not isinstance(self.sql, str):
            raise TypeError("materialization SQL must be a string")
        if not self.sql.strip():
            raise ValueError("materialization SQL must not be empty")

gantry.sql.MaterializationPlan dataclass

What a proposal turned out to be: one create, its destination, sources.

Returned by parse_materialization and checked against MaterializationPolicy. sources is what the body reads, excluding CTE names, so the policy matches real tables rather than local aliases.

Source code in gantry/sql/materialization.py
139
140
141
142
143
144
145
146
147
148
149
150
151
@dataclass(frozen=True, slots=True)
class MaterializationPlan:
    """What a proposal turned out to be: one create, its destination, sources.

    Returned by `parse_materialization` and checked against
    `MaterializationPolicy`. `sources` is what the body reads, excluding CTE
    names, so the policy matches real tables rather than local aliases.
    """

    operation: MaterializationOperation
    sources: tuple[TableRef, ...]
    destination: TableRef
    replace: bool = False

gantry.sql.MaterializationOperation

Bases: StrEnum

The two operations v0 materialization allows.

Both create; neither replaces. A destination that already exists is a refusal rather than an overwrite.

Source code in gantry/sql/materialization.py
84
85
86
87
88
89
90
91
92
class MaterializationOperation(StrEnum):
    """The two operations v0 materialization allows.

    Both create; neither replaces. A destination that already exists is a
    refusal rather than an overwrite.
    """

    CREATE_TABLE_AS = "CREATE_TABLE_AS"
    CREATE_VIEW_AS = "CREATE_VIEW_AS"

gantry.sql.TableRef dataclass

A table name parsed out of a materialization, with its parts kept.

qualified_name joins whichever of catalog, schema and name are present. Materialization requires a schema on the destination, so that an agent cannot write to whatever the session's search path happens to be.

Source code in gantry/sql/materialization.py
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
@dataclass(frozen=True, slots=True)
class TableRef:
    """A table name parsed out of a materialization, with its parts kept.

    `qualified_name` joins whichever of catalog, schema and name are present.
    Materialization requires a schema on the destination, so that an agent
    cannot write to whatever the session's search path happens to be.
    """

    name: str
    schema: str | None = None
    catalog: str | None = None

    def __post_init__(self) -> None:
        for field_name in ("name", "schema", "catalog"):
            value = getattr(self, field_name)
            if value is not None and not isinstance(value, str):
                raise TypeError(f"table {field_name} must be a string")
            if value is not None and not value.strip():
                raise ValueError(f"table {field_name} must not be empty")

    @property
    def qualified_name(self) -> str:
        return ".".join(part for part in (self.catalog, self.schema, self.name) if part)

Schema and classification

gantry.sql.DatabaseSchema dataclass

What the connection can see, as the schema handed to an agent.

Produced by SQLConnection.describe(). This is the right shape to put in a prompt: only what the credential can reach, so the agent is not invited to reference a table it cannot read.

Source code in gantry/sql/schema.py
39
40
41
42
43
44
45
46
47
48
49
50
51
@dataclass(frozen=True, slots=True)
class DatabaseSchema:
    """What the connection can see, as the schema handed to an agent.

    Produced by `SQLConnection.describe()`. This is the right shape to put in a
    prompt: only what the credential can reach, so the agent is not invited to
    reference a table it cannot read.
    """

    catalogs: tuple[str, ...] = ()
    schemas: tuple[str, ...] = ()
    tables: tuple[Table, ...] = ()
    metadata: Mapping[str, object] = field(default_factory=dict)

gantry.sql.Table dataclass

One table or view, with the columns an agent may write SQL against.

kind distinguishes a table from a view. columns is empty when the adapter listed the table without describing it.

Source code in gantry/sql/schema.py
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
@dataclass(frozen=True, slots=True)
class Table:
    """One table or view, with the columns an agent may write SQL against.

    `kind` distinguishes a table from a view. `columns` is empty when the
    adapter listed the table without describing it.
    """

    name: str
    schema: str | None = None
    catalog: str | None = None
    columns: tuple[Column, ...] = ()
    primary_key: tuple[str, ...] = ()
    kind: str = "table"
    metadata: Mapping[str, object] = field(default_factory=dict)

gantry.sql.Column dataclass

One column as the engine reports it.

type is the engine's own type name, not a normalized one: an agent writing SQL needs the name the engine will accept.

Source code in gantry/sql/schema.py
 8
 9
10
11
12
13
14
15
16
17
18
19
@dataclass(frozen=True, slots=True)
class Column:
    """One column as the engine reports it.

    `type` is the engine's own type name, not a normalized one: an agent
    writing SQL needs the name the engine will accept.
    """

    name: str
    type: str
    nullable: bool
    metadata: Mapping[str, object] = field(default_factory=dict)

gantry.sql.SQLClassification dataclass

The structure a policy check runs against, never a transpilation.

read_only is the field most decisions turn on, and it is deliberately conservative: a statement must be recognizably read-only to be treated as such. tables is everything referenced, write_targets only what is written, so a policy can allow reading a table it forbids writing.

A SELECT is not automatically a read. SELECT ... FOR UPDATE takes row locks and SELECT nextval(...) advances a sequence, both of which PostgreSQL refuses in a read-only transaction. When operation is SELECT and read_only is false, read_only_reason says which of those it was, so a refusal can explain itself instead of reporting that a SELECT is not allowed by a read-only policy.

Source code in gantry/sql/classification.py
58
59
60
61
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 SQLClassification:
    """The structure a policy check runs against, never a transpilation.

    `read_only` is the field most decisions turn on, and it is deliberately
    conservative: a statement must be recognizably read-only to be treated as
    such. `tables` is everything referenced, `write_targets` only what is
    written, so a policy can allow reading a table it forbids writing.

    A `SELECT` is not automatically a read. `SELECT ... FOR UPDATE` takes row
    locks and `SELECT nextval(...)` advances a sequence, both of which
    PostgreSQL refuses in a read-only transaction. When `operation` is `SELECT`
    and `read_only` is false, `read_only_reason` says which of those it was, so
    a refusal can explain itself instead of reporting that a SELECT is not
    allowed by a read-only policy.
    """

    operation: SQLOperation
    read_only: bool
    tables: tuple[SQLObjectRef, ...] = ()
    write_targets: tuple[SQLObjectRef, ...] = ()
    functions: tuple[str, ...] = ()
    statement_count: int = 1
    read_only_reason: str | None = None

gantry.sql.SQLOperation

Bases: StrEnum

What a statement does, at the granularity policy decisions need.

MULTI_STATEMENT is its own operation rather than a list, because a batch is refused as a batch unless the policy allows several statements. UNKNOWN is deny-by-default: an unclassifiable statement is not a SELECT.

Source code in gantry/sql/classification.py
 8
 9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
class SQLOperation(StrEnum):
    """What a statement does, at the granularity policy decisions need.

    `MULTI_STATEMENT` is its own operation rather than a list, because a batch
    is refused as a batch unless the policy allows several statements.
    `UNKNOWN` is deny-by-default: an unclassifiable statement is not a
    `SELECT`.
    """

    SELECT = "SELECT"
    INSERT = "INSERT"
    UPDATE = "UPDATE"
    DELETE = "DELETE"
    MERGE = "MERGE"
    DDL = "DDL"
    MULTI_STATEMENT = "MULTI_STATEMENT"
    UNKNOWN = "UNKNOWN"

gantry.sql.SQLObjectRef dataclass

A possibly-qualified reference to a table as it appeared in the SQL.

schema and catalog are None when the statement did not qualify the name; qualified_name joins whichever parts are present. Policy matching compares against these names, so an unqualified reference is matched as written rather than silently resolved.

Source code in gantry/sql/classification.py
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
@dataclass(frozen=True, slots=True)
class SQLObjectRef:
    """A possibly-qualified reference to a table as it appeared in the SQL.

    `schema` and `catalog` are `None` when the statement did not qualify the
    name; `qualified_name` joins whichever parts are present. Policy matching
    compares against these names, so an unqualified reference is matched as
    written rather than silently resolved.
    """

    name: str
    schema: str | None = None
    catalog: str | None = None

    @property
    def qualified_name(self) -> str:
        return ".".join(part for part in (self.catalog, self.schema, self.name) if part)

gantry.sql.ParsedSQL dataclass

The statements a dialect found in one submitted string.

Splitting is the step that makes "one statement" checkable. Anything that yields more than one statement is a multi-statement submission, whatever it looked like to the caller.

Source code in gantry/sql/classification.py
46
47
48
49
50
51
52
53
54
55
@dataclass(frozen=True, slots=True)
class ParsedSQL:
    """The statements a dialect found in one submitted string.

    Splitting is the step that makes "one statement" checkable. Anything that
    yields more than one statement is a multi-statement submission, whatever it
    looked like to the caller.
    """

    statements: tuple[str, ...]

gantry.sql.ExplainResult dataclass

What the engine's planner estimated, before anything ran.

supported is false when the engine cannot explain the statement, which is distinct from an estimate of zero: a byte-scanned bound cannot be enforced on an estimate that does not exist. Every estimate is optional for the same reason.

Source code in gantry/sql/explain.py
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
@dataclass(frozen=True, slots=True)
class ExplainResult:
    """What the engine's planner estimated, before anything ran.

    `supported` is false when the engine cannot explain the statement, which is
    distinct from an estimate of zero: a byte-scanned bound cannot be enforced
    on an estimate that does not exist. Every estimate is optional for the same
    reason.
    """

    supported: bool
    estimated_rows: int | None = None
    estimated_bytes: int | None = None
    estimated_cost: float | None = None
    referenced_objects: tuple[SQLObjectRef, ...] = ()
    summary: Mapping[str, object] = field(default_factory=dict)
    native: Mapping[str, object] = field(default_factory=dict)

Extending

gantry.sql.register_provider

register_provider(
    name: str,
    *,
    dialect: str,
    driver: str,
    adapter_factory: AdapterFactory,
    validate_config: ConfigValidator | None = None,
    metadata: Mapping[str, object] | None = None,
    replace: bool = False,
) -> None

Register a provider that gantry.sql.connect can open by name.

adapter_factory is called with the resolved SQLTarget each time a connection is opened, so one provider can serve many targets. validate_config runs before the factory and should raise on configuration the adapter cannot honour. Raises ValueError if the name is empty, or already registered and replace is false.

Source code in gantry/sql/registry.py
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
def register_provider(
    name: str,
    *,
    dialect: str,
    driver: str,
    adapter_factory: AdapterFactory,
    validate_config: ConfigValidator | None = None,
    metadata: Mapping[str, object] | None = None,
    replace: bool = False,
) -> None:
    """Register a provider that `gantry.sql.connect` can open by name.

    `adapter_factory` is called with the resolved `SQLTarget` each time a
    connection is opened, so one provider can serve many targets.
    `validate_config` runs before the factory and should raise on
    configuration the adapter cannot honour. Raises `ValueError` if the name
    is empty, or already registered and `replace` is false.
    """
    key = name.strip().lower()
    if not key:
        raise ValueError("provider name must not be empty")
    if key in _providers and not replace:
        raise ValueError(f"SQL provider is already registered: {name}")
    _providers[key] = Provider(
        name=key,
        dialect=dialect.strip().lower(),
        driver=driver,
        adapter_factory=adapter_factory,
        validate_config=validate_config or _accept_config,
        metadata={} if metadata is None else metadata,
    )

gantry.sql.providers

Built-in SQL provider presets.

gantry.sql.SQLAdapter

Bases: Protocol

The contract a SQL backend implements to be governed by Gantry.

Declare what you can enforce (capabilities), expose the schema (describe), ask the engine to check and estimate a statement (validate, explain), then run and track it (submit, status, result, cancel). A capability you do not declare is a bound Gantry will refuse to promise, so under-declaring is safe and over-declaring is not.

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

    Declare what you can enforce (`capabilities`), expose the schema
    (`describe`), ask the engine to check and estimate a statement (`validate`,
    `explain`), then run and track it (`submit`, `status`, `result`, `cancel`).
    A capability you do not declare is a bound Gantry will refuse to promise,
    so under-declaring is safe and over-declaring is not.
    """

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

    async def describe(self, target: SQLTarget) -> DatabaseSchema: ...

    async def validate(
        self,
        sql: str,
        target: SQLTarget,
        context: Context,
        policy: SQLPolicy,
    ) -> ValidationResult: ...

    async def explain(self, sql: str, target: SQLTarget) -> ExplainResult: ...

    async def submit(self, sql: str, target: SQLTarget, 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.sql.SQLCapabilities dataclass

What one SQL adapter can enforce, declared per provider.

The source of truth behind the published capability matrix: core_capabilities() projects these onto the engine-neutral AdapterCapabilities that admission checks, and policy_requirements() derives what a given SQLPolicy demands. Note write_execution defaults to true while everything else defaults to false — a new adapter is assumed able to write and assumed unable to bound.

Source code in gantry/sql/capabilities.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
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
@dataclass(frozen=True, slots=True)
class SQLCapabilities:
    """What one SQL adapter can enforce, declared per provider.

    The source of truth behind the published
    [capability matrix](../api/capabilities.md): `core_capabilities()` projects
    these onto the engine-neutral `AdapterCapabilities` that admission checks,
    and `policy_requirements()` derives what a given `SQLPolicy` demands. Note
    `write_execution` defaults to true while everything else defaults to false —
    a new adapter is assumed able to write and assumed unable to bound.
    """

    describe_schema: bool = False
    explain: bool = False
    dry_run: bool = False
    async_jobs: bool = False
    reconnect: bool = False
    cancellation: bool = False
    read_only_session: bool = False
    write_execution: bool = True
    statement_timeout: bool = False
    row_limit: bool = False
    cost_estimate: bool = False
    cost_limit: bool = False
    bytes_scanned: bool = False
    query_metrics: bool = False
    result_reference: bool = False
    create_table_as: bool = False
    create_view_as: bool = False
    destination_introspection: bool = False
    materialization_reference: bool = False

    def core_capabilities(self) -> AdapterCapabilities:
        return AdapterCapabilities(
            reconnect=self.reconnect,
            cancellation=self.cancellation,
            runtime_limit=self.statement_timeout,
            cost_estimation=self.cost_estimate,
            cost_limit=self.cost_limit,
            read_only_execution=self.read_only_session,
            write_execution=self.write_execution,
            remote_status=self.async_jobs,
            metrics=self.query_metrics,
            result_reference=self.result_reference,
        )

    def policy_requirements(self, policy: SQLPolicy) -> PolicyRequirements:
        return PolicyRequirements(
            read_only=policy.read_only,
            allow_writes=not policy.read_only,
            max_runtime_seconds=policy.timeout_seconds,
            max_cost_usd=policy.max_cost_usd,
            require_reconnect=self.async_jobs,
            require_metrics=policy.max_bytes_scanned is not None,
            require_result_reference=False,
        )

gantry.sql.MaterializationAdapter

Bases: Protocol

Additional adapter operation required only for materialization.

inspect_table is how create-only is enforced and how the destination is verified afterwards: returning None for a table that does not exist is what distinguishes "not created" from "created empty". An adapter without this cannot materialize.

Source code in gantry/sql/materialization.py
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
@runtime_checkable
class MaterializationAdapter(Protocol):
    """Additional adapter operation required only for materialization.

    `inspect_table` is how create-only is enforced and how the destination is
    verified afterwards: returning `None` for a table that does not exist is
    what distinguishes "not created" from "created empty". An adapter without
    this cannot materialize.
    """

    """Additional adapter operation required only for materialization."""

    async def inspect_table(
        self,
        reference: TableRef,
        target: SQLTarget,
        *,
        include_row_count: bool = False,
    ) -> Table | None: ...

gantry.sql.MaterializationCapabilities dataclass

What an adapter can enforce for materialization specifically.

Separate from SQLCapabilities because materialization needs guarantees a query does not: destination_introspection is what makes create-only checkable, and without it the destination cannot be verified after the job runs. All default to false.

Source code in gantry/sql/materialization.py
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
@dataclass(frozen=True, slots=True)
class MaterializationCapabilities:
    """What an adapter can enforce for materialization specifically.

    Separate from `SQLCapabilities` because materialization needs guarantees a
    query does not: `destination_introspection` is what makes create-only
    checkable, and without it the destination cannot be verified after the job
    runs. All default to false.
    """

    create_table_as: bool = False
    create_view_as: bool = False
    durable_jobs: bool = False
    cancel: bool = False
    estimate_bytes_scanned: bool = False
    destination_introspection: bool = False
    result_reference: bool = False

gantry.sql.SQLTarget dataclass

A resolved connection target: the provider, its dialect, and its config.

Built by connect from a registered provider. It holds credentials, so it stays in application code — what reaches the engine adapter is this, and what reaches the agent is only the tool schema.

Source code in gantry/sql/target.py
 8
 9
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 SQLTarget:
    """A resolved connection target: the provider, its dialect, and its config.

    Built by `connect` from a registered provider. It holds credentials, so it
    stays in application code — what reaches the engine adapter is this, and
    what reaches the agent is only the tool schema.
    """

    provider: str
    dialect: str
    driver: str
    config: Mapping[str, object]
    metadata: Mapping[str, object] = field(default_factory=dict)

    def __post_init__(self) -> None:
        for name, value in (
            ("provider", self.provider),
            ("dialect", self.dialect),
            ("driver", self.driver),
        ):
            if not value.strip():
                raise ValueError(f"SQL target {name} must not be empty")

gantry.sql.register

register(
    name: str,
    *,
    adapter: SQLAdapter,
    dialect: str,
    driver: str | None = None,
    replace: bool = False,
) -> None

Register a single already-built adapter as a provider named name.

A shorthand over register_provider for tests and embedded adapters: every target resolved through this name shares the one adapter instance. driver defaults to name.

Source code in gantry/sql/registry.py
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
def register(
    name: str,
    *,
    adapter: SQLAdapter,
    dialect: str,
    driver: str | None = None,
    replace: bool = False,
) -> None:
    """Register a single already-built `adapter` as a provider named `name`.

    A shorthand over `register_provider` for tests and embedded adapters: every
    target resolved through this name shares the one adapter instance. `driver`
    defaults to `name`.
    """
    register_provider(
        name,
        dialect=dialect,
        driver=driver or name,
        adapter_factory=lambda target: adapter,
        replace=replace,
    )

Dialects

A dialect decides how a submission is split and classified. It never rewrites SQL — what the caller wrote is what the engine receives.

gantry.sql.SQLDialect

Bases: Protocol

How one SQL dialect is split and classified for policy checks.

Three methods, none of which rewrite SQL: parse splits a submission into statements, classify says what the single statement does, and referenced_objects lists the tables it names. Register an implementation with register_dialect; ConservativeDialect is the deny-by-default fallback.

Source code in gantry/sql/dialect.py
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
class SQLDialect(Protocol):
    """How one SQL dialect is split and classified for policy checks.

    Three methods, none of which rewrite SQL: `parse` splits a submission into
    statements, `classify` says what the single statement does, and
    `referenced_objects` lists the tables it names. Register an implementation
    with `register_dialect`; `ConservativeDialect` is the deny-by-default
    fallback.
    """

    def parse(self, sql: str) -> ParsedSQL: ...

    def classify(self, sql: str) -> SQLClassification: ...

    def referenced_objects(self, sql: str) -> tuple[SQLObjectRef, ...]: ...

gantry.sql.ConservativeDialect

A small deny-by-default classifier shared until a dialect overrides it.

The two options are lexical rules that genuinely differ between engines, and getting them wrong changes where a statement ends. backslash_escapes is off by default, matching PostgreSQL with standard_conforming_strings, where a backslash is only special inside an E'' string. dollar_quoting is on by default because PostgreSQL has it and MySQL does not — there, $$ is a syntax error and $ is an ordinary identifier character.

Source code in gantry/sql/dialect.py
 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
class ConservativeDialect:
    """A small deny-by-default classifier shared until a dialect overrides it.

    The two options are lexical rules that genuinely differ between engines, and
    getting them wrong changes where a statement ends. `backslash_escapes` is
    off by default, matching PostgreSQL with `standard_conforming_strings`,
    where a backslash is only special inside an `E''` string. `dollar_quoting`
    is on by default because PostgreSQL has it and MySQL does not — there, `$$`
    is a syntax error and `$` is an ordinary identifier character.
    """

    def __init__(self, *, backslash_escapes: bool = False, dollar_quoting: bool = True) -> None:
        self._backslash_escapes = backslash_escapes
        self._dollar_quoting = dollar_quoting

    def parse(self, sql: str) -> ParsedSQL:
        statements = tuple(part.strip() for part in self._split(sql) if part.strip())
        return ParsedSQL(statements)

    def _split(self, sql: str) -> tuple[str, ...]:
        return _split_statements(
            sql,
            backslash_escapes=self._backslash_escapes,
            dollar_quoting=self._dollar_quoting,
        )

    def _blank(self, sql: str) -> str:
        return _blank_comments(
            sql,
            backslash_escapes=self._backslash_escapes,
            dollar_quoting=self._dollar_quoting,
        )

    def classify(self, sql: str) -> SQLClassification:
        parsed = self.parse(sql)
        if not parsed.statements:
            return SQLClassification(SQLOperation.UNKNOWN, False, statement_count=0)
        if len(parsed.statements) > 1:
            return SQLClassification(
                SQLOperation.MULTI_STATEMENT,
                False,
                tables=self.referenced_objects(sql),
                statement_count=len(parsed.statements),
            )

        statement = _LEADING_COMMENTS.sub("", parsed.statements[0])
        # Keyword scanning runs on the statement with comments blanked, so a
        # comment neither hides a keyword nor contributes words of its own.
        scanned = self._blank(statement)
        words = [match.group(0).upper() for match in _WORD.finditer(scanned)]
        first = words[0] if words else ""
        if first == "WITH":
            dangerous = {*_DDL, "INSERT", "UPDATE", "DELETE", "MERGE"}
            first = next((word for word in words[1:] if word in dangerous), "")
            if not first and "SELECT" in words:
                first = "SELECT"
        operation = (
            SQLOperation.DDL if first in _DDL else _OPERATIONS.get(first, SQLOperation.UNKNOWN)
        )
        if operation is SQLOperation.SELECT and "INTO" in words:
            # `SELECT ... INTO new_table` creates a table. PostgreSQL refuses it
            # in a read-only transaction for exactly that reason, and MySQL's
            # `SELECT ... INTO OUTFILE` writes a file. It is a create, not a read.
            operation = SQLOperation.DDL
        tables = self.referenced_objects(statement)
        write_targets = (
            tables[:1]
            if operation
            in {
                SQLOperation.INSERT,
                SQLOperation.UPDATE,
                SQLOperation.DELETE,
                SQLOperation.MERGE,
                SQLOperation.DDL,
            }
            else ()
        )
        functions = tuple(dict.fromkeys(match.group(1) for match in _FUNCTION.finditer(statement)))
        reason = _why_not_read_only(words, functions) if operation is SQLOperation.SELECT else None
        return SQLClassification(
            operation=operation,
            read_only=operation is SQLOperation.SELECT and reason is None,
            read_only_reason=reason,
            tables=tables,
            write_targets=write_targets,
            functions=functions,
            statement_count=1,
        )

    def referenced_objects(self, sql: str) -> tuple[SQLObjectRef, ...]:
        references: list[SQLObjectRef] = []
        seen: set[str] = set()
        # `/* FROM secret */` names no table. Scanning the comment would add a
        # reference the engine never resolves, and let the text of a comment
        # decide whether an allow-list matches.
        for match in _OBJECT.finditer(self._blank(sql)):
            qualified = match.group(1)
            if qualified.lower() in seen:
                continue
            seen.add(qualified.lower())
            parts = [_unquote(part.strip()) for part in qualified.split(".")]
            if len(parts) == 3:
                reference = SQLObjectRef(name=parts[2], schema=parts[1], catalog=parts[0])
            elif len(parts) == 2:
                reference = SQLObjectRef(name=parts[1], schema=parts[0])
            else:
                reference = SQLObjectRef(name=parts[0])
            references.append(reference)
        return tuple(references)

gantry.sql.MySQLDialect

Bases: ConservativeDialect

ConservativeDialect with MySQL's lexical rules rather than PostgreSQL's.

Two differences, both of which move where a statement ends. MySQL escapes backslashes inside ordinary strings — NO_BACKSLASH_ESCAPES is not in the default sql_mode — so 'a\'; SELECT 2' is one string and one statement. Reading it with PostgreSQL's rule splits it in two and refuses it as a batch, which is safe but wrong. And MySQL has no dollar-quoting: $$ is a syntax error there, while $ is an ordinary identifier character, so my$tab$le is one name.

Checked against MySQL 8.4 with the default sql_mode.

Source code in gantry/sql/dialect.py
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
class MySQLDialect(ConservativeDialect):
    """`ConservativeDialect` with MySQL's lexical rules rather than PostgreSQL's.

    Two differences, both of which move where a statement ends. MySQL escapes
    backslashes inside ordinary strings — `NO_BACKSLASH_ESCAPES` is not in the
    default `sql_mode` — so `'a\\'; SELECT 2'` is one string and one statement.
    Reading it with PostgreSQL's rule splits it in two and refuses it as a
    batch, which is safe but wrong. And MySQL has no dollar-quoting: `$$` is a
    syntax error there, while `$` is an ordinary identifier character, so
    `my$tab$le` is one name.

    Checked against MySQL 8.4 with the default `sql_mode`.
    """

    def __init__(self) -> None:
        super().__init__(backslash_escapes=True, dollar_quoting=False)

gantry.sql.register_dialect

register_dialect(
    name: str, dialect: SQLDialect, *, replace: bool = False
) -> None

Register a SQL dialect under a name, for classification and splitting.

Source code in gantry/sql/registry.py
29
30
31
32
33
34
35
36
def register_dialect(name: str, dialect: SQLDialect, *, replace: bool = False) -> None:
    """Register a SQL dialect under a name, for classification and splitting."""
    key = name.strip().lower()
    if not key:
        raise ValueError("dialect name must not be empty")
    if key in _dialects and not replace:
        raise ValueError(f"SQL dialect is already registered: {name}")
    _dialects[key] = dialect