Skip to content

Batch and stream jobs

Two entry points named for intent rather than for the engine behind them. gantry.batch is for work that ends; gantry.stream is for work that does not, where "did it succeed" is the wrong question and "is it healthy" is the right one.

See Flink backends for what each guarantees, and for the setup a real cluster needs.

Batch

gantry.batch.connect

connect(
    provider: str,
    *,
    endpoint: str,
    config: Mapping[str, object] | None = None,
    policy: Policy | None = None,
    **options: object,
) -> BatchConnection
Source code in gantry/batch/api.py
57
58
59
60
61
62
63
64
65
66
67
68
69
70
def connect(
    provider: str,
    *,
    endpoint: str,
    config: Mapping[str, object] | None = None,
    policy: Policy | None = None,
    **options: object,
) -> BatchConnection:
    normalized = provider.strip().lower()
    if normalized != "flink":
        raise ValueError(f"unsupported batch provider: {provider}")
    return BatchConnection(
        normalized, connect_runtime(endpoint, config=config, policy=policy, **options)
    )

gantry.batch.BatchConnection

A configured batch provider connection.

Source code in gantry/batch/api.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
47
48
49
50
51
52
53
54
class BatchConnection:
    """A configured batch provider connection."""

    def __init__(self, provider: str, runtime: FlinkRuntime) -> None:
        self._provider = provider
        self._runtime = runtime

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

    def capabilities(self) -> BatchCapabilities:
        capabilities = self._runtime.capabilities()
        return BatchCapabilities(
            insert_into=capabilities.write_execution,
            insert_overwrite=capabilities.write_execution,
            durable_job_id=capabilities.reconnect,
            cancel=capabilities.cancellation,
            metrics=capabilities.metrics,
            output_verification=capabilities.result_reference,
        )

    def job(
        self,
        *,
        inputs: Collection[str],
        outputs: Collection[str],
        checks: Sequence[FlinkHealthCheck | MaterializationCheck] = (),
        timeout: float | None = None,
        poll_interval: float = 1.0,
    ) -> FlinkBatchJob:
        return FlinkBatchJob(
            self._runtime,
            inputs=inputs,
            outputs=outputs,
            checks=checks,
            timeout=timeout,
            poll_interval=poll_interval,
        )

Stream

gantry.stream.connect

connect(
    provider: str,
    *,
    endpoint: str,
    config: Mapping[str, object] | None = None,
    policy: Policy | None = None,
    **options: object,
) -> StreamConnection
Source code in gantry/stream/api.py
56
57
58
59
60
61
62
63
64
65
66
67
68
69
def connect(
    provider: str,
    *,
    endpoint: str,
    config: Mapping[str, object] | None = None,
    policy: Policy | None = None,
    **options: object,
) -> StreamConnection:
    normalized = provider.strip().lower()
    if normalized != "flink":
        raise ValueError(f"unsupported stream provider: {provider}")
    return StreamConnection(
        normalized, connect_runtime(endpoint, config=config, policy=policy, **options)
    )

gantry.stream.StreamConnection

A configured streaming provider connection.

Source code in gantry/stream/api.py
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
class StreamConnection:
    """A configured streaming provider connection."""

    def __init__(self, provider: str, runtime: FlinkRuntime) -> None:
        self._provider = provider
        self._runtime = runtime

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

    def capabilities(self) -> StreamCapabilities:
        capabilities = self._runtime.capabilities()
        return StreamCapabilities(
            insert_into=capabilities.write_execution,
            durable_job_id=capabilities.reconnect,
            cancel=capabilities.cancellation,
            health=capabilities.remote_status,
            metrics=capabilities.metrics,
            watermark_metrics=capabilities.metrics,
        )

    def job(
        self,
        *,
        inputs: Collection[str],
        outputs: Collection[str],
        checks: Sequence[FlinkHealthCheck] = (),
        timeout: float | None = None,
        poll_interval: float = 1.0,
    ) -> FlinkStreamJob:
        return FlinkStreamJob(
            self._runtime,
            inputs=inputs,
            outputs=outputs,
            checks=checks,
            timeout=timeout,
            poll_interval=poll_interval,
        )

Verification checks

Declarative verification checks and safe agent-input parsing.

MaterializationCheck

Bases: Protocol

A trusted check evaluated against a table's metadata.

The same checks serve a materialization and a query. A materialization is checked against the destination it created; a query is checked against the shape of the rows it returned, described as a table so one check can do both. Two of them do need a destination and say so with requires_destination, which makes them unsupported on a query rather than quietly true.

Source code in gantry/verify.py
22
23
24
25
26
27
28
29
30
31
32
33
34
@runtime_checkable
class MaterializationCheck(Protocol):
    """A trusted check evaluated against a table's metadata.

    The same checks serve a materialization and a query. A materialization is
    checked against the destination it created; a query is checked against the
    shape of the rows it returned, described as a table so one check can do
    both. Two of them do need a destination and say so with
    `requires_destination`, which makes them unsupported on a query rather than
    quietly true.
    """

    def evaluate(self, table: Table | None) -> CheckResult: ...

VerificationInputError

Bases: ValueError

An agent supplied a check outside the declarative verification DSL.

Source code in gantry/verify.py
37
38
class VerificationInputError(ValueError):
    """An agent supplied a check outside the declarative verification DSL."""

VerificationUnsupported

Bases: VerificationInputError

A valid check is not supported by this governed operation.

Source code in gantry/verify.py
41
42
class VerificationUnsupported(VerificationInputError):  # noqa: N818
    """A valid check is not supported by this governed operation."""

VerificationConflict

Bases: VerificationInputError

The trusted and agent commitments cannot both be satisfied.

Source code in gantry/verify.py
45
46
class VerificationConflict(VerificationInputError):  # noqa: N818
    """The trusted and agent commitments cannot both be satisfied."""

DocumentCount dataclass

Bases: RowCount

MongoDB spelling of a result-size assertion.

Source code in gantry/verify.py
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
@dataclass(frozen=True, slots=True)
class DocumentCount(RowCount):
    """MongoDB spelling of a result-size assertion."""

    requires_document_count: ClassVar[bool] = True

    def evaluate(self, table: object | None) -> CheckResult:
        metadata = getattr(table, "metadata", None)
        actual = metadata.get("document_count") if isinstance(metadata, Mapping) else None
        if not isinstance(actual, int) or isinstance(actual, bool):
            return CheckResult(
                "document_count",
                False,
                {"min": self.minimum, "max": self.maximum},
                actual,
                "document count is unavailable",
                supported=False,
            )
        ok = (self.minimum is None or actual >= self.minimum) and (
            self.maximum is None or actual <= self.maximum
        )
        return CheckResult(
            "document_count",
            ok,
            {"min": self.minimum, "max": self.maximum},
            actual,
            None if ok else f"document count {actual} is outside the accepted range",
        )

RequiredFields dataclass

Bases: RequiredColumns

MongoDB spelling of a required-shape assertion.

Source code in gantry/verify.py
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
@dataclass(frozen=True, slots=True)
class RequiredFields(RequiredColumns):
    """MongoDB spelling of a required-shape assertion."""

    def evaluate(self, table: object | None) -> CheckResult:
        actual = () if table is None else tuple(getattr(table, "fields", ()))
        present = {field.lower() for field in actual}
        missing = tuple(field for field in self.columns if field not in present)
        return CheckResult(
            name="required_fields",
            ok=not missing,
            expected={"fields": list(self.columns)},
            actual={"fields": list(actual)},
            message=None if not missing else f"required fields are missing: {', '.join(missing)}",
        )

NullRate dataclass

The fraction of rows where one column is null, bounded above.

Measured at the destination by the provider, not reported by the statement that wrote it. A provider that cannot measure it makes the check unsupported rather than passing it, because an unmeasured bound would be indistinguishable from a satisfied one.

Source code in gantry/verify.py
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
@dataclass(frozen=True, slots=True)
class NullRate:
    """The fraction of rows where one column is null, bounded above.

    Measured at the destination by the provider, not reported by the statement
    that wrote it. A provider that cannot measure it makes the check
    unsupported rather than passing it, because an unmeasured bound would be
    indistinguishable from a satisfied one.
    """

    requires_row_count: ClassVar[bool] = False

    column: str
    maximum: float

    def __post_init__(self) -> None:
        if not self.column.strip():
            raise ValueError("null rate column must not be empty")
        if isinstance(self.maximum, bool) or not isinstance(self.maximum, (int, float)):
            raise TypeError("null rate maximum must be numeric")
        if not 0 <= self.maximum <= 1:
            raise ValueError("null rate maximum must be a fraction between 0 and 1")

    @property
    def null_rate_columns(self) -> tuple[str, ...]:
        return (self.column,)

    def evaluate(self, table: Table | None) -> CheckResult:
        expected = {"column": self.column, "max": self.maximum}
        if table is None:
            return CheckResult(
                "null_rate",
                False,
                expected,
                None,
                "there is nothing to measure",
            )
        rates = table.metadata.get("null_rates")
        observed = rates.get(self.column) if isinstance(rates, Mapping) else None
        if not isinstance(observed, (int, float)) or isinstance(observed, bool):
            return CheckResult(
                "null_rate",
                False,
                expected,
                None,
                f"null rate for {self.column} is unavailable from this provider",
                supported=False,
            )
        ok = observed <= self.maximum
        return CheckResult(
            "null_rate",
            ok,
            expected,
            {"value": observed},
            None if ok else f"{self.column} is null in {observed:.1%} of rows",
        )

null_rate

null_rate(column: str, *, max: float) -> NullRate

Reject a destination where a column is emptier than it should be.

The failure this catches is a query that runs, produces the right number of rows, and joins wrongly — so the column everyone downstream keys on is null in most of them. Row count says the table is fine. This does not.

Source code in gantry/verify.py
243
244
245
246
247
248
249
250
def null_rate(column: str, *, max: float) -> NullRate:  # noqa: A002
    """Reject a destination where a column is emptier than it should be.

    The failure this catches is a query that runs, produces the right number of
    rows, and joins wrongly — so the column everyone downstream keys on is null
    in most of them. Row count says the table is fine. This does not.
    """
    return NullRate(column=column, maximum=max)

parse_agent_checks

parse_agent_checks(
    values: object, *, capabilities: Collection[str]
) -> tuple[object, ...]

Parse untrusted JSON into allowlisted check objects, failing closed.

Source code in gantry/verify.py
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
def parse_agent_checks(
    values: object,
    *,
    capabilities: Collection[str],
) -> tuple[object, ...]:
    """Parse untrusted JSON into allowlisted check objects, failing closed."""
    if values is None:
        return ()
    if not isinstance(values, Sequence) or isinstance(values, str | bytes):
        raise VerificationInputError("verify must be an array of declarative checks")
    checks: list[object] = []
    for index, value in enumerate(values):
        if not isinstance(value, Mapping):
            raise VerificationInputError(f"verify[{index}] must be an object")
        check_type = value.get("type")
        if not isinstance(check_type, str) or check_type not in _CONSTRUCTORS:
            raise VerificationInputError(f"verify[{index}] has an unknown check type")
        if check_type not in capabilities:
            raise VerificationUnsupported(
                f"verification check {check_type!r} is unsupported by this operation"
            )
        unexpected = set(value) - _ALLOWED_KEYS[check_type]
        if unexpected:
            names = ", ".join(sorted(str(key) for key in unexpected))
            raise VerificationInputError(f"verify[{index}] has unexpected fields: {names}")
        try:
            checks.append(_CONSTRUCTORS[check_type](value))
        except (TypeError, ValueError) as error:
            raise VerificationInputError(f"invalid {check_type} check: {error}") from error
    return tuple(checks)

validate_agent_checks

validate_agent_checks(
    values: Sequence[object],
    *,
    capabilities: Collection[str],
) -> tuple[object, ...]

Validate direct-call check objects without allowing callbacks.

Source code in gantry/verify.py
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
def validate_agent_checks(
    values: Sequence[object], *, capabilities: Collection[str]
) -> tuple[object, ...]:
    """Validate direct-call check objects without allowing callbacks."""
    supported = set(capabilities)
    checks = tuple(values)
    for check in checks:
        name = check_type(check)
        if name is None:
            raise VerificationInputError(
                "agent verify must contain gantry.verify declarative checks"
            )
        if name not in supported:
            raise VerificationUnsupported(
                f"verification check {name!r} is unsupported by this operation"
            )
    return checks

check_config

check_config(check: object) -> dict[str, object]

Return the bounded public commitment, never provenance or executable state.

Source code in gantry/verify.py
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
def check_config(check: object) -> dict[str, object]:
    """Return the bounded public commitment, never provenance or executable state."""
    kind = check_type(check)
    if kind is None:
        return {"type": type(check).__name__}
    payload: dict[str, object] = {"type": kind}
    if kind in {"row_count", "document_count"}:
        payload.update(
            {
                "min": getattr(check, "minimum", None),
                "max": getattr(check, "maximum", None),
            }
        )
    elif kind in {"required_columns", "required_fields"}:
        values = getattr(check, "columns", getattr(check, "fields", ()))
        payload["fields" if kind == "required_fields" else "columns"] = list(values)
    elif isinstance(check, NullRate):
        payload.update({"column": check.column, "max": check.maximum})
    elif kind == "restart_count":
        payload["max"] = getattr(check, "maximum")  # noqa: B009
    elif kind == "watermark_lag":
        payload["max_seconds"] = getattr(check, "seconds")  # noqa: B009
    return {key: value for key, value in payload.items() if value is not None}

verification_schema

verification_schema(
    capabilities: Collection[str],
) -> dict[str, object]

JSON Schema for only the checks supported by one operation.

Source code in gantry/verify.py
484
485
486
487
488
489
def verification_schema(capabilities: Collection[str]) -> dict[str, object]:
    """JSON Schema for only the checks supported by one operation."""
    variants = [_check_schema(name) for name in sorted(capabilities)]
    if not variants:
        return {"type": "array", "maxItems": 0}
    return {"type": "array", "items": {"oneOf": variants}}

detect_conflicts

detect_conflicts(
    trusted: Sequence[object], agent: Sequence[object]
) -> None

Reject statically impossible count commitments before execution.

Source code in gantry/verify.py
492
493
494
495
496
497
498
499
500
501
502
503
504
505
def detect_conflicts(trusted: Sequence[object], agent: Sequence[object]) -> None:
    """Reject statically impossible count commitments before execution."""
    trusted_bounds = _count_bounds(trusted)
    agent_bounds = _count_bounds(agent)
    for family in ("rows", "documents"):
        trusted_min, trusted_max = trusted_bounds[family]
        agent_min, agent_max = agent_bounds[family]
        minimum = max(value for value in (trusted_min, agent_min, 0) if value is not None)
        maxima = [value for value in (trusted_max, agent_max) if value is not None]
        maximum = min(maxima) if maxima else None
        if maximum is not None and minimum > maximum:
            raise VerificationConflict(
                f"verification conflict: {family} minimum {minimum} exceeds maximum {maximum}"
            )