Skip to content

corral API reference

Auto-generated from package docstrings via mkdocstrings. Every public symbol gets a stable anchor (#<dotted.qualname>) that matches the entries in ai/api-index.json.

For task-oriented recipes see the cookbook. For design rationale see architecture.

Top-level

corral.dataset.Package(spec, tables=dict(), engine=None, source=None, dirty_tracker=None, metadata=dict()) dataclass

Lazy wrapper around a multi-table Frictionless data package.

Composes the spec, the table mapping, the engine, the source locator, and (optionally) the sync-state :class:DirtyTracker. The primary constructors are :meth:from_source (load from a path/URL) and :meth:from_tables (compose from already-built :class:~corral.dataset.Table instances).

Attributes:

Name Type Description
spec DataPackage

The parsed :class:~corral.spec.model.DataPackage.

tables dict[str, Table]

Mapping {name: Table}. Insertion order is preserved (Python 3.7+) so iteration is deterministic.

engine Engine | None

The :class:~corral.engines.base.Engine that produced the table expressions. All materialisation routes back through this same engine.

source str | None

Original source identifier (path / URL). None for in-memory packages built via :meth:from_tables.

dirty_tracker Any

Optional sync-state tracker. When None, the sync-state validation pass is a no-op and :meth:write consults only the per-table dirty flag.

metadata dict[str, Any]

Free-form bag of extras (engine name, write timestamp, etc.); echoed into the JSON validation report.

Examples:

Load the bundled generic sample fixture, validate, then write to a fresh parquet directory::

>>> import tempfile, pathlib
>>> from corral.fixtures import sample
>>> from corral.dataset import Package
>>> from corral.engines.ibis_engine import IbisEngine
>>> pkg = Package.from_source(
...     sample.csv_dir(),
...     engine=IbisEngine(),
...     spec=sample.DATAPACKAGE,
...     tables=["book", "author"],
... )
>>> "book" in pkg
True
>>> report = pkg.validate()
>>> isinstance(report.issues, list)
True

tables = field(default_factory=dict) class-attribute instance-attribute

from_source(source, *, engine=None, spec=None, tables=None, format=None, credentials=None) classmethod

Build a :class:Package from a source path/URL.

Dispatch order (per docs/architecture.md §6.1 + the :mod:corral.io package docstring):

1. Resolve ``engine`` (default: registry default).
2. Resolve ``spec`` — explicit :class:`DataPackage` wins;
   else look for a ``datapackage.json`` alongside ``source``;
   else synthesise a minimal spec from whatever resources
   are discovered.
3. Resolve ``source`` to a :class:`ResourceListing`:

    a. explicit ``format=`` short-circuits sniffing;
    b. URL scheme (``duckdb://``, ``s3://``, ...);
    c. file extension (``.csv``, ``.parquet``, ``.duckdb``,
       ``.csv.zip``);
    d. directory-of-known-formats walk (``csv/``,
       ``parquet/``) — directories have no extension to
       dispatch on, so we walk children and ask each
       matching adapter to scan one. Mirrors
       :func:`corral.validation.structural._scan_directory_of_known_formats`
       so the structural validator and the package loader
       can't drift apart on "what counts as a directory of
       tables".
    e. ``probe`` chain (any adapter that recognises the
       source);
    f. :class:`~corral.io.FormatNotDetected`.

4. For each :class:`~corral.io.base.ResourceRef`, ask
   the owning adapter to ``read`` it lazily through
   ``engine``, then wrap the resulting
   :class:`~corral.engines.base.TableExpr` in a
   :class:`~corral.dataset.Table`. The spec is
   consulted for the matching
   :class:`~corral.spec.model.Schema` so the validation
   layer doesn't have to re-load it.

Parameters:

Name Type Description Default
source str | Path

Path / URL / directory pointing at the data package. Anything :func:corral.io.dispatch knows about, plus the directory-of-files convention.

required
engine Engine | None

Engine to materialise through. Defaults to the registered default (typically :class:~corral.engines.ibis_engine.IbisEngine).

None
spec DataPackage | str | Path | None

Either an in-memory :class:DataPackage, a path / URL to a datapackage.json, or None to auto-discover.

None
tables Iterable[str] | None

Optional subset of resource names to load. Useful for memory-efficient partial loads.

None
format str | None

Optional adapter-name override for the top-level source. Used when the source has no extension to sniff on (e.g. an API endpoint that returns CSV but has no .csv suffix). For s3:// / https:// URLs this becomes the inner-format hint passed through :class:~corral.io.remote.RemoteAdapter.

None
credentials dict[str, str] | None

Optional explicit credentials dict forwarded to :class:~corral.io.remote.RemoteAdapter. Only the remote adapter consumes it; local-fs reads ignore the kwarg. Caller-provided keys win over env / keyring / netrc — see :func:corral.io.credentials.resolve_credentials.

None

Returns:

Type Description
Package

A populated :class:Package with lazy :class:Table

Package

entries.

Raises:

Type Description
FormatNotDetected

If neither extension dispatch nor the directory walk could resolve source.

Examples:

>>> from corral.fixtures import sample
>>> from corral.dataset import Package
>>> from corral.engines.ibis_engine import IbisEngine
>>> pkg = Package.from_source(
...     sample.csv_dir(),
...     engine=IbisEngine(),
...     spec=sample.DATAPACKAGE,
...     tables=["book"],
... )
>>> list(pkg.keys())
['book']

from_tables(tables, *, spec=None, engine=None) classmethod

Build a :class:Package from already-constructed :class:Table instances.

Useful for tests, ad-hoc compositions, and adapter-less in-memory sources. When spec is not provided, a minimal one is synthesised so :meth:validate and the dict-like surface still work.

Parameters:

Name Type Description Default
tables Mapping[str, Table]

Mapping {name: Table}.

required
spec DataPackage | None

Optional explicit :class:DataPackage. Synthesised from tables when None.

None
engine Engine | None

Optional explicit engine. Defaults to the first table’s engine.

None

Returns:

Type Description
Package

A new :class:Package.

Examples:

>>> from corral.engines.ibis_engine import IbisEngine
>>> from corral.dataset import Package, Table
>>> e = IbisEngine()
>>> link = Table(name="link", expr=e.from_records([{"id": 1}]), engine=e)
>>> p = Package.from_tables({"link": link})
>>> "link" in p
True

validate(*, schema=True, structural=True, foreign_keys=True, sync_state=True, strict=False)

Run the selected validation passes and return the combined report.

Order: structural → schema (per table) → foreign_keys → sync_state. Each pass appends issues to the same report. Skip any pass by passing False for its flag.

Sync-state contract: after a clean FK pass, the :class:DirtyTracker (when attached) stamps each non-violating FK via :meth:DirtyTracker.stamp_fk_from_exprs so subsequent edits surface sync.fk_stale instead of silently writing the broken FK. The sync_state pass then consults :meth:DirtyTracker.check to report any stamps that no longer match the live tables. No-op when dirty_tracker is None.

Parameters:

Name Type Description Default
schema bool

Run per-table schema validation (:func:corral.validation.check_schema).

True
structural bool

Run structural validation (:func:corral.validation.check_structural).

True
foreign_keys bool

Run FK validation (:func:corral.validation.check_foreign_keys).

True
sync_state bool

Consult the :class:DirtyTracker (when attached) via its check method. No-op when no tracker is attached.

True
strict bool

Forwarded to validators that accept it (FK unverifiable becomes ERROR instead of WARNING).

False

Returns:

Name Type Description
A ValidationReport
ValidationReport

with one :class:~corral.validation.Issue per finding.

Examples:

>>> from corral.fixtures import sample
>>> from corral.dataset import Package
>>> from corral.engines.ibis_engine import IbisEngine
>>> pkg = Package.from_source(
...     sample.csv_dir(),
...     engine=IbisEngine(),
...     spec=sample.DATAPACKAGE,
...     tables=["book", "author"],
... )
>>> r = pkg.validate()
>>> isinstance(r.issues, list)
True

write(dest, *, format=None, overwrite=False, strict_sync=False)

Persist every table in the package to dest.

Sync-state contract: before writing, consults two signals per table — the per-:class:Table dirty flag (always available) and :meth:DirtyTracker.is_table_dirty (when a tracker is attached). EITHER signal flags the table as stale — the union is the conservative thing to surface. Stale tables trigger :class:OutOfSyncWarning, or :class:OutOfSyncError under strict_sync=True. When no tracker is attached, the per-table flag alone is consulted, so a caller can still flag a mutation via :meth:Table.invalidate without installing the tracker.

Format dispatch goes through :mod:corral.io exactly as on the read path:

  • format explicit → adapter resolved by name;
  • else inferred from dest extension;
  • else default to "parquet" when dest looks like a directory.

For partitioned/multi-table formats (parquet directory, duckdb file), each :class:Table writes to a per-table location under dest. The exact layout matches the read-side scan so a round-trip via :meth:from_source returns the same logical package. A zip (.zip / .csv.zip, or format="zip") holds one flat <table>.csv member per table plus a datapackage.json with inlined schemas, and is written atomically (staged beside dest, then renamed over it).

Parameters:

Name Type Description Default
dest str | Path

Target directory or file.

required
format str | None

Optional explicit format name (an adapter name; "zip" is accepted as an alias for "zipcsv").

None
overwrite bool

When False (the default), refuses to clobber an existing dest. When True, removes the existing target first.

False
strict_sync bool

When True, stale tables raise :class:OutOfSyncError instead of warning.

False

Raises:

Type Description
FileExistsError

When dest exists and overwrite=False.

OutOfSyncError

When strict_sync=True and any table is dirty.

PackageError

When the package has no engine attached (so materialisation can’t run).

FormatNotDetected

When format is not given and the dest extension is unknown — the user must pass format= explicitly to disambiguate.

WriteUnsupportedForSchemeError

When dest is a URL — packages are written to local paths only.

Examples:

Roundtrip the bundled sample fixture through parquet::

>>> import tempfile, pathlib
>>> from corral.fixtures import sample
>>> from corral.dataset import Package
>>> from corral.engines.ibis_engine import IbisEngine
>>> with tempfile.TemporaryDirectory() as tmp:
...     pkg = Package.from_source(
...         sample.csv_dir(),
...         engine=IbisEngine(),
...         spec=sample.DATAPACKAGE,
...         tables=["book"],
...     )
...     out = pathlib.Path(tmp) / "out.pkg"
...     pkg.write(out, format="parquet")
...     out.exists()
True

safe_count(name)

Return self[name].count() or None if absent / transiently uncountable.

Convenience used by previews (notebook _repr_html_, CLI info, HTTP /networks/{id}) where a missing or transiently-broken table should degrade to a friendly “?” rather than crash the surrounding render.

The contract: returns None when the table is absent OR when table.count() raises any :class:Exception (e.g. a transient backend hiccup on an ibis-against-remote-duckdb connection). Hard exceptions in surrounding code still propagate — this only shields the count call itself.

Examples:

>>> from corral.dataset import Package, Table
>>> from corral.engines.ibis_engine import IbisEngine
>>> e = IbisEngine()
>>> pkg = Package.from_tables({"x": Table(name="x", expr=e.from_records([{"a": 1}]), engine=e)})
>>> pkg.safe_count("x")
1
>>> pkg.safe_count("absent_table") is None
True

corral.dataset.Table(name, expr, engine, schema=None, source=None, format=None, dirty=False, metadata=dict()) dataclass

Lazy wrapper around one table in a data package.

The wrapper is mutable (its expr may be replaced — for example by an editing op or a scoped view), but each underlying TableExpr is immutable per the engine’s semantics. Lazy ops (:meth:filter, :meth:select, :meth:head) therefore return a new :class:Table, leaving the original unchanged.

Mutation: this class intentionally does not expose an in-place update shim. For row-level edits with full diff / rollback / audit, use :class:corral.editing.Edit + :class:corral.editing.Session. For low-level mutations that bypass the editing framework, build a new expression directly and call :meth:invalidate on the table to flip its dirty flag so the sync-state tracker sees the change.

Attributes:

Name Type Description
name str

Logical resource name. Matches :attr:corral.spec.model.Resource.name and is used as the dict key inside a :class:~corral.dataset.Package.

expr TableExpr

Engine-native lazy expression. Concrete type depends on the engine (ibis.expr.types.Table for ibis, polars.LazyFrame for polars, pandas.DataFrame for pandas — pandas is eager by design).

engine Engine

The engine that produced expr. Materialisation always routes back through this same engine to preserve the cross-engine dtype contract.

schema Schema | None

Optional resolved Frictionless schema. Carried so the validation layer can run per-field checks without re-loading the spec.

source str | None

Optional source locator (path / URL / sub-locator like "net.duckdb::link"). Set by :meth:Package.from_source; None for tables constructed inline.

format str | None

Optional short format identifier ("csv", "parquet", "duckdb"). Mirrors :attr:corral.io.base.ResourceRef.format.

dirty bool

True when the table has been mutated (or its source hash differs from a sync-state stamp). Read by the sync-state tracker (task 2.6) before any write.

metadata dict[str, Any]

Free-form bag of extras. Carried through writes via the package-level metadata sidecar in Phase 3.

Examples:

Construct directly from in-memory records via an engine::

>>> from corral.engines.ibis_engine import IbisEngine
>>> from corral.dataset import Table
>>> e = IbisEngine()
>>> expr = e.from_records([{"a": 1}, {"a": 2}, {"a": 3}])
>>> t = Table(name="t", expr=expr, engine=e)
>>> t.count()
3
>>> t.dirty
False

filter(predicate)

Return a new :class:Table whose expression is predicate(expr).

The predicate receives the engine-native expression and returns a new engine-native expression (typically a filtered one). The caller writes the predicate in the engine’s own dialect — this method doesn’t translate between engines.

Parameters:

Name Type Description Default
predicate Callable[[TableExpr], TableExpr]

Callable taking the current expr and returning a new expression of the same engine type.

required

Returns:

Type Description
Table

A new :class:Table wrapping the transformed expression;

Table

the original :class:Table is unchanged.

Examples:

>>> from corral.engines.ibis_engine import IbisEngine
>>> from corral.dataset import Table
>>> e = IbisEngine()
>>> expr = e.from_records([{"a": 1}, {"a": 2}, {"a": 3}])
>>> t = Table(name="t", expr=expr, engine=e)
>>> t2 = t.filter(lambda df: df.filter(df["a"] > 1))
>>> t2 is t
False
>>> t.count(), t2.count()
(3, 2)

select(*columns)

Return a new :class:Table projected to the named columns.

Parameters:

Name Type Description Default
*columns str

Column names to keep. Must all exist on the current expr.

()

Returns:

Type Description
Table

A new :class:Table carrying only the projected columns.

Examples:

>>> from corral.engines.ibis_engine import IbisEngine
>>> from corral.dataset import Table
>>> e = IbisEngine()
>>> expr = e.from_records([{"a": 1, "b": 2}])
>>> Table(name="t", expr=expr, engine=e).select("a").columns()
['a']

head(n=5)

Return a new :class:Table containing only the first n rows.

Lazy — the underlying engine builds a head expression, but no materialisation runs until the caller asks for it.

Parameters:

Name Type Description Default
n int

Row count to keep (default 5).

5

Returns:

Type Description
Table

A new :class:Table whose expression is the head.

Examples:

>>> from corral.engines.ibis_engine import IbisEngine
>>> from corral.dataset import Table
>>> e = IbisEngine()
>>> expr = e.from_records([{"a": i} for i in range(10)])
>>> Table(name="t", expr=expr, engine=e).head(3).count()
3

count()

Return the row count, pushing down to the engine where possible.

Delegates to :meth:Engine.count, which routes through each backend’s native cheap-count path (ibis pushes a SELECT COUNT(*) to duckdb; polars selects pl.len() inside the lazy plan; pandas calls len on the already- materialised frame). No table is materialised just to count it.

Examples:

>>> from corral.engines.ibis_engine import IbisEngine
>>> from corral.dataset import Table
>>> e = IbisEngine()
>>> Table(name="t", expr=e.from_records([{"a": 1}]), engine=e).count()
1

to_pandas()

Materialise as a pandas.DataFrame using the engine’s converter.

The returned frame uses the cross-engine nullable numpy-backed dtype family (Int64 / Float64 / string / boolean) regardless of which engine produced the expression — see :meth:corral.engines.base.Engine.to_pandas for the full contract.

Examples:

>>> from corral.engines.ibis_engine import IbisEngine
>>> from corral.dataset import Table
>>> e = IbisEngine()
>>> t = Table(name="t", expr=e.from_records([{"a": 1}]), engine=e)
>>> t.to_pandas()["a"].tolist()
[1]

to_polars()

Materialise as a polars.DataFrame using the engine’s converter.

Raises :class:~corral.engines.errors.EngineNotAvailableError if polars is not installed and the engine cannot satisfy the conversion.

Examples:

>>> import pytest
>>> _ = pytest.importorskip("polars")
>>> from corral.engines.ibis_engine import IbisEngine
>>> from corral.dataset import Table
>>> e = IbisEngine()
>>> t = Table(name="t", expr=e.from_records([{"a": 1}]), engine=e)
>>> t.to_polars().shape[0]
1

collect()

Force eager materialisation; return the engine-native frame.

For ibis this triggers backend execution and returns a pyarrow Table; for polars this collects the LazyFrame to a DataFrame; for pandas this is essentially an identity (pandas is already eager).

Examples:

>>> from corral.engines.ibis_engine import IbisEngine
>>> from corral.dataset import Table
>>> e = IbisEngine()
>>> t = Table(name="t", expr=e.from_records([{"a": 1}]), engine=e)
>>> out = t.collect()
>>> out is not t
True

columns()

Return the column names of the underlying expression.

Delegates to :meth:Engine.columns, which uses the engine’s lazy schema introspection (ibis expr.schema().names, polars expr.collect_schema().names(), pandas list(expr.columns)). No row materialisation runs.

Examples:

>>> from corral.engines.ibis_engine import IbisEngine
>>> from corral.dataset import Table
>>> e = IbisEngine()
>>> Table(name="t", expr=e.from_records([{"a": 1, "b": 2}]), engine=e).columns()
['a', 'b']

Engines

corral.engines.Engine

Bases: Protocol

Execution engine for tabular operations on a Frictionless data package.

Implementations are concrete classes (one per backend: ibis, polars, pandas) that satisfy this protocol structurally. The registry (corral.engines) holds singleton instances and resolves them by name.

The contract has three layers:

  1. Per-format primitives (read_csv / read_parquet / read_duckdb_table / from_records / their write counterparts). Adapters call these directly. An engine that doesn’t natively support a primitive (e.g. :class:~corral.engines.ibis_engine.IbisEngine.write_duckdb_table) raises :class:~corral.engines.errors.EngineNotAvailableError with a clear pointer at the engine that does.
  2. Schema casting (cast_schema). Promoted from a per-engine internal helper so adapters can apply Frictionless schema casts uniformly after a primitive read.
  3. Convenience delegators (scan / write). 3-line wrappers that route through corral.io.dispatch. They exist so callers who don’t care about adapters can still write engine.scan(path).

Attributes:

Name Type Description
name str

Short identifier used as the registry key ("ibis" / "polars" / "pandas").

read_csv(source, schema=None, **kwargs)

Read a single CSV file into an engine-native lazy expression.

Adapters (CsvAdapter, ZipCsvAdapter after member extraction, RemoteAdapter after URL resolution) call this primitive directly. Engines should apply cast_schema at the tail when schema is non-None so behaviour matches the convenience scan(..., schema=...) path.

Parameters:

Name Type Description Default
source SourceRef

Filesystem path / URL / Path.

required
schema Any | None

Optional Frictionless :class:Schema to cast columns with after the read.

None
**kwargs Any

Forwarded verbatim to the underlying library reader. Names follow the library, not a normalized vocabulary — e.g. polars wants separator=, pandas wants sep=, duckdb wants delim=.

{}

read_parquet(source, schema=None, *, hive_partitioning=False, **kwargs)

Read parquet (single file or Hive-partitioned directory).

When hive_partitioning=True the engine must enable partition discovery (so partition columns survive into the result and downstream filters become true partition prunes). For single-file parquet hive_partitioning is ignored.

Parameters:

Name Type Description Default
source SourceRef

Path to a .parquet file or partitioned directory.

required
schema Any | None

Optional Frictionless schema, applied after read.

None
hive_partitioning bool

Enable Hive-style partition discovery.

False
**kwargs Any

Forwarded to the underlying reader.

{}

read_duckdb_table(source, *, table, schema=None, **kwargs)

Read one named table out of a .duckdb file.

Parameters:

Name Type Description Default
source SourceRef

Path / URL to a .duckdb file.

required
table str

The table name to read. Required (no defaulting — the duckdb adapter’s scan() enumerates tables).

required
schema Any | None

Optional Frictionless schema, applied after read.

None
**kwargs Any

Forwarded to the engine-specific reader.

{}

Raises:

Type Description
EngineNotAvailableError

If the engine can’t read duckdb files (none of the stock engines raise — all three support this primitive).

from_records(records, schema=None)

Build a table expression from in-memory records.

Accepts the two shapes a caller might naturally write:

  • List of row dicts: [{"a": 1, "b": 2}, {"a": 3, "b": 4}]
  • Columnar dict: {"a": [1, 3], "b": [2, 4]}

Adapters do not call this primitive (no on-disk format produces in-memory records); Engine.scan calls it directly for the {"data": ...} dict-source contract documented on :meth:scan.

Parameters:

Name Type Description Default
records list[dict[str, Any]] | dict[str, list[Any]]

Either shape above.

required
schema Any | None

Optional Frictionless schema, applied after build.

None

from_arrow(arrow_table)

Wrap a :class:pyarrow.Table as an engine-native lazy expression.

Type-preserving counterpart to :meth:from_records. Use this whenever the source is already an Arrow table — going through records = arrow.to_pylist() then from_records(records) is lossy (binary becomes bytes, decimals become Decimal, timestamps lose tz, all numerics widen via JSON-ish coercion) and is the round-trip the editing layer formerly used. Engines materialise the Arrow buffer directly:

  • Ibis registers it as a duckdb temp table via con.create_table(name, obj=arrow_table, temp=True).
  • Polars wraps it as a :class:polars.LazyFrame via pl.from_arrow(arrow_table).lazy().
  • Pandas materialises with arrow_table.to_pandas().convert_dtypes() so the numpy-backed nullable dtype contract documented on :meth:to_pandas is preserved end-to-end.

Parameters:

Name Type Description Default
arrow_table Table

Already-materialised :class:pyarrow.Table.

required

Returns:

Type Description
TableExpr

An engine-native lazy expression; the contents match

TableExpr

arrow_table 1:1 modulo each engine’s documented dtype

TableExpr

convention.

Examples:

>>> import pyarrow as pa
>>> from corral.engines.ibis_engine import IbisEngine
>>> tbl = pa.table({"a": [1, 2], "b": ["x", "y"]})
>>> e = IbisEngine()
>>> expr = e.from_arrow(tbl)
>>> list(expr.columns), expr.count().to_pyarrow().as_py()
(['a', 'b'], 2)

write_csv(expr, dest, **kwargs)

Write expr to dest as a single CSV file.

Parameters:

Name Type Description Default
expr TableExpr

Engine-native expression to materialize and write.

required
dest SourceRef

Target path.

required
**kwargs Any

Forwarded to the underlying CSV writer (compression, line terminator, …). Names follow the library.

{}

write_parquet(expr, dest, *, partition_by=None, **kwargs)

Write expr to dest as parquet (single file or partitioned).

When partition_by is given the engine writes a Hive-style partitioned dataset under dest. The parquet adapter currently handles partitioned writes through pyarrow regardless of engine; engines may either implement the partitioned-write path natively or document that they only support the single-file case.

Parameters:

Name Type Description Default
expr TableExpr

Engine-native expression.

required
dest SourceRef

Target path or directory.

required
partition_by list[str] | None

Optional Hive partition columns. None writes a single file.

None
**kwargs Any

Forwarded to the underlying writer.

{}

write_duckdb_table(expr, dest, *, table, **kwargs)

Write expr into a .duckdb file as a named table.

Engines that cannot natively write duckdb without raw SQL (notably :class:~corral.engines.ibis_engine.IbisEngine) raise :class:~corral.engines.errors.EngineNotAvailableError and point callers at :class:~corral.engines.ibis_engine.IbisEngine for this primitive. The architecture’s “no raw SQL outside ibis_engine” rule is the constraint that forces this split.

Parameters:

Name Type Description Default
expr TableExpr

Engine-native expression.

required
dest SourceRef

Target duckdb file path.

required
table str

Destination table name (required).

required
**kwargs Any

Forwarded to the underlying writer.

{}

cast_schema(expr, schema)

Cast columns of expr per a Frictionless schema.

Fields whose Frictionless type does not map cleanly to an engine-native type are left untouched (the v0.3 apply_schema_to_df bug was a reminder that silent-but-noisy is the safer default for partial schemas). Fields named in the schema but absent from expr are skipped — that’s a validation concern, not a scan-time concern.

Parameters:

Name Type Description Default
expr TableExpr

Engine-native expression.

required
schema Any

Frictionless :class:~corral.spec.model.Schema.

required

Returns:

Type Description
TableExpr

A new expression with the casts applied, or expr

TableExpr

unchanged if no field type was castable.

scan(source, format=None, schema=None, **kwargs)

Open source as a lazy table expression.

Convenience method — the body is a ~3-line delegation to corral.io.dispatch(source, format=format).read(source, engine=self, schema=schema, **kwargs). Power users and adapters skip this method entirely and call the primitives (:meth:read_csv, etc.) directly.

Parameters:

Name Type Description Default
source SourceRef

A path, URL, Path, or handle dict pointing at a tabular file or table. Format dispatch is delegated to the corral.io FormatAdapter layer. The two handle-dict shapes every engine MUST accept are:

  • {"data": [...]} — inline data. Value is either a list of row dicts ([{"a": 1}, {"a": 2}]) or a columnar dict ({"a": [1, 2]}). Used for in-memory test fixtures and small synthetic frames. The dispatcher can’t sniff a dict source so this case is short-circuited inside :meth:scan and routed directly to :meth:from_records.
  • {"format": "duckdb", "path": "net.duckdb", "table": "link"} — a duckdb table handle. "format" is optional when "path" ends in .duckdb. Also short-circuited inside :meth:scan (the dispatcher rejects dict sources).

Any other dict shape MUST raise :class:~corral.engines.errors.UnsupportedSourceError with a message listing these supported shapes.

required
format str | None

Optional explicit format hint forwarded to corral.io.dispatch(source, format=format). When None the dispatcher uses URL-scheme / extension / probe resolution. Use this when the source is ambiguous (e.g. an extensionless file, an http:// URL that returns parquet bytes).

None
schema Any | None

Optional Frictionless Schema to apply at scan time (column types, missing-value handling). If None, the engine infers from file metadata.

None
**kwargs Any

Adapter-specific options forwarded verbatim to the resolved FormatAdapter.read (delimiter, compression, partition pruning predicate, etc.). Engines must not strip or mutate this mapping.

{}

Returns:

Type Description
TableExpr

An engine-native lazy expression. Concrete type depends on

TableExpr

the engine — see module docstring.

materialize(expr)

Execute expr and return the engine’s native materialized frame.

For ibis this triggers backend execution; for polars this .collect()-s a LazyFrame; for pandas this is an identity (pandas is already eager).

columns(expr)

Return the column names of expr without materialising rows.

Each engine has a cheap lazy schema path (ibis expr.schema().names, polars expr.collect_schema().names(), pandas list(expr.columns)). Implementations route to that path so the wrapper in :class:corral.dataset.Table is one line.

Parameters:

Name Type Description Default
expr TableExpr

Engine-native expression.

required

Returns:

Type Description
list[str]

A new list of column names in declaration order.

count(expr)

Return the row count of expr, pushing the aggregate to the engine.

ibis pushes a SELECT COUNT(*) to duckdb; polars selects pl.len() inside the lazy plan; pandas calls len(expr). Implementations must not full-materialise the table to count it.

Parameters:

Name Type Description Default
expr TableExpr

Engine-native expression.

required

Returns:

Type Description
int

Number of rows in expr.

head(expr, n)

Return a new lazy expression with only the first n rows.

Lazy — the engine builds a head expression but does not execute it. The pandas engine is the only one where this materialises (pandas has no lazy mode).

Parameters:

Name Type Description Default
expr TableExpr

Engine-native expression.

required
n int

Row count to keep.

required

Returns:

Type Description
TableExpr

A new engine-native expression of the same type as expr.

select(expr, columns)

Return a new lazy expression projected to columns.

Parameters:

Name Type Description Default
expr TableExpr

Engine-native expression.

required
columns list[str]

Column names to keep. Must all exist on expr.

required

Returns:

Type Description
TableExpr

A new engine-native expression carrying only columns.

order_by(expr, columns, descending=False)

Return a new lazy expression sorted by columns.

ibis emits ORDER BY (expr.order_by); polars expr.sort; pandas sort_values (stable). descending applies to all columns. With :meth:limit this is the server-side paging path.

Parameters:

Name Type Description Default
expr TableExpr

Engine-native expression.

required
columns list[str]

Column names to sort by, in priority order.

required
descending bool

Sort descending instead of ascending.

False

Returns:

Type Description
TableExpr

A new engine-native expression of the same type as expr.

limit(expr, n, offset=0)

Return a new lazy expression of n rows starting at offset.

ibis emits LIMIT n OFFSET offset; polars expr.slice; pandas positional iloc. Pairs with :meth:order_by for stable paging.

Parameters:

Name Type Description Default
expr TableExpr

Engine-native expression.

required
n int

Maximum rows to keep.

required
offset int

Leading rows to skip.

0

Returns:

Type Description
TableExpr

A new engine-native expression of the same type as expr.

to_pandas(expr)

Materialize expr and return it as a pandas.DataFrame.

Cross-engine convergence point. Implementations may raise :class:~corral.engines.errors.EngineNotAvailableError if pandas is not installed.

Dtype convention (cross-engine contract). The returned DataFrame uses pandas numpy-backed nullable dtypes: Int64 (capital I), Float64, string, boolean. We standardize on this family because:

  • Null semantics are preserved without silently upcasting integer columns with missing values to float64 (the numpy default’s footgun).
  • Numpy-backed (not pyarrow-backed) dtypes keep compatibility with downstream libraries (sklearn, older matplotlib, geopandas pre-1.0) that don’t understand pd.ArrowDtype columns.

Implementations achieve this by post-processing with :meth:pandas.DataFrame.convert_dtypes (the universal path), regardless of which engine produced expr. The Leavenworth fixture’s link.from_node_id column round-trips as Int64 from all three stock engines under this convention — that is the regression locked in by the cross-engine parity test.

to_polars(expr)

Materialize expr and return it as a polars.DataFrame.

Cross-engine convergence point. Implementations may raise EngineNotAvailableError if polars is not installed.

write(expr, dest, fmt, **kwargs)

Write expr to dest in format fmt.

Convenience method — the body is a ~3-line delegation to corral.io.get_adapter(fmt).write(expr, dest, engine=self, **kwargs). Adapters skip this and call the write primitives (:meth:write_csv, :meth:write_parquet, :meth:write_duckdb_table) directly.

Parameters:

Name Type Description Default
expr TableExpr

The expression to materialize and write.

required
dest SourceRef

A path, URL, or handle dict — same contract as scan’s source.

required
fmt str

Format name ("parquet", "csv", "duckdb", etc.). Dispatched through the corral.io FormatAdapter layer.

required
**kwargs Any

Format-specific options (compression, partitioning, etc.). Adapter-defined.

{}

corral.engines.get_engine(name=None)

Return the registered engine for name, or the default if None.

Parameters:

Name Type Description Default
name str | None

The engine name ("ibis" / "duckdb" — the one compute engine). None returns the default.

None

Returns:

Type Description
Engine

The registered engine instance.

Raises:

Type Description
EngineNotAvailableError

If the registry is empty, or the requested name is not registered. The message lists the currently-registered engine names so the caller can correct the typo.

Examples:

Look up the auto-registered default (ibis on duckdb):

>>> from corral.engines import get_engine
>>> default = get_engine()
>>> default.name
'ibis'

Looking up a name that is not registered raises a clear error:

>>> from corral.engines import EngineNotAvailableError
>>> try:
...     get_engine("not-a-real-engine")
... except EngineNotAvailableError as exc:
...     "not-a-real-engine" in str(exc)
True

corral.engines.register_engine(engine, *, default=False)

Register an engine instance under its name attribute.

Re-registering an existing name overwrites the previous registration (no warning). This is intentional — it lets tests and notebooks swap in fakes without ceremony, and lets a downstream package replace the stock implementation.

Parameters:

Name Type Description Default
engine Engine

An instance satisfying the Engine protocol. Must expose a non-empty name string and the engine methods (scan / materialize / to_pandas / to_polars / write).

required
default bool

If True, also make this the default engine returned by get_engine() with no argument.

False

Raises:

Type Description
TypeError

If engine does not satisfy the Engine protocol (missing one or more required methods).

ValueError

If engine.name is empty.

Examples:

Register a minimal fake engine and look it up. The fake uses a unique name to avoid colliding with auto-registered engines:

>>> from corral.engines import (
...     register_engine, get_engine, list_engines
... )
>>> class _DoctestEngine:
...     name = "doctest-register"
...     def read_csv(self, source, schema=None, **kw): return None
...     def read_parquet(self, source, schema=None, **kw): return None
...     def read_duckdb_table(self, source, table, schema=None, **kw): return None
...     def from_records(self, records, schema=None): return None
...     def from_arrow(self, arrow_table): return None
...     def write_csv(self, expr, dest, **kw): return None
...     def write_parquet(self, expr, dest, **kw): return None
...     def write_duckdb_table(self, expr, dest, table, **kw): return None
...     def cast_schema(self, expr, schema): return expr
...     def scan(self, source, schema=None, **kw): return None
...     def write(self, expr, dest, fmt, **kw): return None
...     def materialize(self, expr): return None
...     def to_pandas(self, expr): return None
...     def to_polars(self, expr): return None
...     def columns(self, expr): return []
...     def count(self, expr): return 0
...     def head(self, expr, n): return expr
...     def select(self, expr, columns): return expr
...     def order_by(self, expr, cols, descending=False): return expr
...     def limit(self, expr, n, offset=0): return expr
>>> fake = _DoctestEngine()
>>> try:
...     register_engine(fake)
...     get_engine("doctest-register") is fake
... finally:
...     from corral import engines as _eng
...     _ = _eng._REGISTRY.pop("doctest-register", None)
True

corral.engines.resolve_engine(name)

Return the compute :class:Engine (DuckDB via ibis).

DuckDB is the single compute engine; pandas / polars / pyarrow are I/O formats, not compute backends. None, "ibis" and "duckdb" all return the one engine; any other name raises :class:ValueError.

Parameters:

Name Type Description Default
name str | None

None / "ibis" / "duckdb".

required

Returns:

Type Description
Engine

The registered DuckDB engine.

Raises:

Type Description
ValueError

If name is not one of the accepted values.

Examples:

>>> from corral.engines import resolve_engine
>>> isinstance(resolve_engine(None), Engine)
True

Reports

corral.reports.ValidationReport(spec_version=None, source=None, issues=list(), metadata=dict(), created_at=datetime.now()) dataclass

Result of one or more validation passes over a data package.

Mutable — multiple checks (schema + FK + structural + sync_state + data_quality) build it up over a single run by calling :meth:add_issue or :meth:add. When the run is complete, hand it to a renderer (:func:~corral.validation.render_rich, :func:~corral.validation.render_json, or the HTML renderer in :mod:corral.reports.render).

The report carries the spec version and source identifier so the rendered output is self-describing — a saved JSON or HTML report tells you which spec it was validated against and where the data came from.

Attributes:

Name Type Description
spec_version str | None

The spec version this run was validated against (e.g. "0.97"). Optional — populated by the caller.

source str | None

Path / URL identifier of the package being validated. Used by renderers in the header.

issues list[Issue]

All findings recorded so far. Order is insertion order; renderers re-sort by severity.

metadata dict[str, Any]

Free-form metadata about the run — engine name, scope, timestamps. Echoed in the JSON output.

created_at datetime

When the report was constructed (timezone-naive local time, matching datetime.now()).

Examples:

>>> report = ValidationReport(spec_version="0.97", source="leavenworth.gmns")
>>> issue = report.add(
...     severity=Severity.ERROR,
...     category=Category.SCHEMA,
...     code="schema.required",
...     message="link.from_node_id row 0: value is null",
...     table="link",
... )
>>> report.has_errors
True
>>> report.count(Severity.ERROR)
1

has_errors property

True when the report contains at least one ERROR.

Examples:

>>> r = ValidationReport()
>>> r.has_errors
False

has_warnings property

True when the report contains at least one WARNING.

Examples:

>>> r = ValidationReport()
>>> r.has_warnings
False

is_clean property

True when no errors AND no warnings are present.

INFO issues do not break is_clean — they’re informational, and surfaced for awareness. Data-quality findings (Category.DATA_QUALITY) carry whichever severity the rule chose; WARNING ones flip is_clean to False, INFO ones don’t.

Examples:

>>> r = ValidationReport()
>>> r.is_clean
True
>>> r.add(severity=Severity.INFO, category=Category.STRUCTURAL,
...       code="structural.optional_missing", message="x")
Issue(...)
>>> r.is_clean
True

add_issue(issue)

Append an already-constructed :class:Issue to the report.

Parameters:

Name Type Description Default
issue Issue

The issue to record.

required

Examples:

>>> report = ValidationReport()
>>> report.add_issue(Issue(
...     severity=Severity.WARNING,
...     category=Category.SYNC_STATE,
...     code="sync.fk_stale",
...     message="link FK to node is stale",
... ))
>>> len(report.issues)
1

add(*, severity, category, code, message, **kw)

Construct an :class:Issue, append it, and return it.

This is the convenience builder validators usually call directly — the returned issue is handy when the validator wants to log or annotate the same finding elsewhere.

Parameters:

Name Type Description Default
severity Severity

Severity of the finding.

required
category Category

Category of the finding.

required
code str

Stable dotted identifier (e.g. "schema.required").

required
message str

Human-readable description that names the broken input.

required
**kw Any

Any of table, column, row, fix_hint, extra from :class:Issue.

{}

Returns:

Type Description
Issue

The newly created (and stored) issue.

Examples:

>>> report = ValidationReport()
>>> issue = report.add(
...     severity=Severity.ERROR,
...     category=Category.SCHEMA,
...     code="schema.required",
...     message="link.from_node_id row 0: value is null",
...     table="link",
...     column="from_node_id",
...     row=0,
... )
>>> issue.code
'schema.required'

by_severity(severity)

Return all issues at the given severity, preserving insertion order.

Examples:

>>> r = ValidationReport()
>>> r.add(severity=Severity.ERROR, category=Category.SCHEMA,
...       code="schema.required", message="x")
Issue(severity=<Severity.ERROR: 'error'>, ...)
>>> len(r.by_severity(Severity.ERROR))
1
>>> r.by_severity(Severity.WARNING)
[]

by_category(category)

Return all issues in the given category, preserving insertion order.

Examples:

>>> r = ValidationReport()
>>> r.add(severity=Severity.ERROR, category=Category.FOREIGN_KEY,
...       code="fk.missing_target", message="x")
Issue(...)
>>> [i.code for i in r.by_category(Category.FOREIGN_KEY)]
['fk.missing_target']

by_table(table)

Return all issues attached to table, preserving insertion order.

Cross-cutting issues (Issue.table is None) are excluded — they would match every table query otherwise.

Examples:

>>> r = ValidationReport()
>>> r.add(severity=Severity.ERROR, category=Category.SCHEMA,
...       code="schema.required", message="x", table="link")
Issue(...)
>>> r.add(severity=Severity.ERROR, category=Category.STRUCTURAL,
...       code="structural.missing_table", message="y")  # no table
Issue(...)
>>> [i.table for i in r.by_table("link")]
['link']

count(severity=None)

Count issues, optionally filtered to a single severity.

Parameters:

Name Type Description Default
severity Severity | None

If given, count only issues at this severity. None (the default) counts every issue.

None

Examples:

>>> r = ValidationReport()
>>> r.add(severity=Severity.ERROR, category=Category.SCHEMA,
...       code="schema.required", message="x")
Issue(...)
>>> r.add(severity=Severity.WARNING, category=Category.SYNC_STATE,
...       code="sync.fk_stale", message="y")
Issue(...)
>>> r.count()
2
>>> r.count(Severity.ERROR)
1

to_dict()

Return a JSON-serialisable dict snapshot of the report.

Schema:

.. code-block:: python

{
    "report_version": "1",
    "spec_version": "0.97" | None,
    "source": str | None,
    "created_at": "ISO-8601",
    "metadata": {...},
    "summary": {"error": N, "warning": N, "info": N,
                "data_quality": N, "is_clean": bool},
    "issues": [
        {"severity": "error", "category": "schema",
         "code": "...", "message": "...", "table": ..., ...},
        ...
    ],
}

Enum values flatten to their .value strings; datetime flattens via :meth:datetime.isoformat. The key set is stable; downstream consumers (the HTML renderer, the MCP server, the FastAPI server) rely on it.

Examples:

>>> r = ValidationReport(spec_version="0.97", source="x.gmns")
>>> d = r.to_dict()
>>> d["report_version"]
'1'
>>> d["summary"]["is_clean"]
True

to_json(path=None, *, indent=2)

Return the report as a JSON string. Optionally write to path.

Thin wrapper around :func:json.dumps on :meth:to_dict.

Parameters:

Name Type Description Default
path str | Path | None

Optional file path. When given, the rendered JSON is also written to path (parent dirs created if needed). The string is still returned so chained report.to_json("r.json") works.

None
indent int

json.dumps indent setting. Default 2.

2

Examples:

>>> import json
>>> r = ValidationReport()
>>> data = json.loads(r.to_json())              # in-memory
>>> data["summary"]["is_clean"]
True
>>> import tempfile, pathlib
>>> with tempfile.TemporaryDirectory() as d:    # write-to-file
...     out = pathlib.Path(d) / "r.json"
...     _ = r.to_json(out)
...     out.is_file()
True

to_rich()

Return the rich-console rendering of this report.

Convenience wrapper around :func:corral.reports.render_rich. Kept here so that str(report) round-trips through the rich renderer without the caller importing :mod:corral.reports.render.

Examples:

>>> r = ValidationReport(source="empty.gmns")
>>> "empty.gmns" in r.to_rich()
True

to_html(path=None, *, title=None, include_map=True)

Return the interactive single-file HTML rendering of this report.

Shortcut for :func:corral.reports.render_html. See its docstring for the offline-mode trade-off around the optional Vega-Lite map section.

Parameters:

Name Type Description Default
path str | Path | None

Optional file path. When given, the rendered HTML is also written to path (parent directories created if needed). The string is still returned so chained report.to_html("r.html").splitlines() works.

None
title str | None

Optional override for the <title> and <h1>.

None
include_map bool

If False, skip the map section even when geo-located issues are present.

True

Returns:

Type Description
str

A single self-contained HTML string.

Examples:

>>> r = ValidationReport(source="empty.gmns")
>>> html = r.to_html()                              # in-memory
>>> html.lstrip().startswith("<!DOCTYPE html>")
True
>>> import tempfile, pathlib
>>> with tempfile.TemporaryDirectory() as d:        # write-to-file
...     out = pathlib.Path(d) / "r.html"
...     _ = r.to_html(out)
...     out.is_file()
True

__str__()

Alias for :meth:to_rich — usable from print(report).

corral.reports.Issue(severity, category, code, message, table=None, column=None, row=None, fix_hint=None, extra=dict()) dataclass

A single validation finding.

Frozen so a report can be safely held by callers (or shown in a UI) after the underlying data changes — once an issue is recorded, it snapshots the failure. Hashable for the same reason, so consumers can dedupe via set membership without writing a custom __hash__.

The code field is a stable, dotted, namespaced identifier — the string callers grep tracebacks for and filter reports by. Examples: "schema.required", "schema.enum", "fk.missing_target", "structural.missing_table", "sync.fk_stale", "quality.high_speed_residential". The leading namespace (schema., fk., structural., sync., quality.) mirrors :class:Category and keeps codes greppable per rule family.

The message field MUST name the input that broke — table, column, row, value — not a generic phrase. The v0.3 line of bugs where “FK violation” was the entire user-facing string is what this contract is designed to avoid.

The optional fix_hint is a single short sentence telling the user what to do. Renderers display it on a second line when present.

Attributes:

Name Type Description
severity Severity

How bad it is.

category Category

Which rule family produced it.

code str

Stable dotted identifier (e.g. "schema.required").

message str

Human-readable, names the input that broke.

table str | None

Table name, or None for cross-cutting / structural issues that don’t belong to a single table.

column str | None

Column / field name, or None if not field-specific.

row int | None

Zero-based row index, or None if not row-specific.

fix_hint str | None

Optional one-sentence remediation hint.

extra dict[str, Any]

Adapter-specific extras — geo coordinates, target table for FK violations, etc. Renderers may surface known keys.

Examples:

>>> issue = Issue(
...     severity=Severity.ERROR,
...     category=Category.FOREIGN_KEY,
...     code="fk.missing_target",
...     message="link row 12: from_node_id=99 not found in node.node_id",
...     table="link",
...     column="from_node_id",
...     row=12,
...     fix_hint="Add a node row with node_id=99, or remove the link.",
... )
>>> issue.severity
<Severity.ERROR: 'error'>
>>> issue.row
12

__hash__()

Hash on the stable identity fields, ignoring the mutable extra dict.

frozen=True normally generates __hash__, but the extra dict field is mutable and unhashable. Hashing on the identity fields (everything except extra) means equal issues — same severity, location, code, message, fix hint — collapse in a set(). This matches how a human would dedupe a report: two findings with the same code+message at the same row are the same finding, regardless of which validator’s debug payload they happen to carry.

corral.reports.Severity

Bases: StrEnum

Severity of a single validation finding.

The order matters: renderers display issues grouped ERROR -> WARNING -> INFO. The str base (rather than int) was chosen for JSON-serialisability — a validation report dumped to JSON is the same on Python 3.11 and Python 3.13 with no custom encoder, and severity == "error" works for users who read the JSON without re-importing the enum.

For ordering, use :func:severity_rank.

Severity is orthogonal to :class:Category. A data-quality finding carries Category.DATA_QUALITY AND one of the three severity values — a missing speed limit is typically INFO, while a 70mph residential street is WARNING. The pre-1.0 Severity.DATA_QUALITY value conflated the two dimensions and was removed; quality rules now pick the real severity that matches how urgently the finding wants attention.

Attributes:

Name Type Description
ERROR

Spec or contract violation. The data is wrong; downstream consumers should not trust it.

WARNING

Likely problem that does not break correctness. The OutOfSyncWarning family lives here. Also the default severity for a suspicious data-quality finding.

INFO

Informational only. No action required. Also the default severity for an awareness-only data-quality finding.

Examples:

>>> Severity("error") is Severity.ERROR
True
>>> Severity.WARNING.value
'warning'

corral.reports.Category

Bases: StrEnum

Broad category of a validation finding.

Used to group findings in reports and to filter the JSON / HTML output by rule family. Categories map to the validation modules that produce them — schema checks emit SCHEMA, FK checks emit FOREIGN_KEY, and so on.

Attributes:

Name Type Description
SCHEMA

Field-level constraints — type, required, enum, regex, min/max. Produced by :mod:corral.validation.schema_check.

STRUCTURAL

Package-level structure — missing required table, extra unknown table, missing file on disk. Produced by :mod:corral.validation.structural.

FOREIGN_KEY

Cross-table referential integrity. Produced by :mod:corral.validation.foreign_keys.

SYNC_STATE

A previously validated FK is now stale because one side has been mutated since the last check. Produced by :mod:corral.validation.sync_state.

DATA_QUALITY

A configurable quality rule (e.g. high-speed on residential road). Produced by quality plugins registered under the corral.quality.rules entry point — see :mod:corral.quality.

Examples:

>>> Category("foreign_key") is Category.FOREIGN_KEY
True

Editing

corral.editing.Edit(op, table, payload=dict(), metadata=dict()) dataclass

One atomic mutation against one table — domain-free.

Supported ops + their payload shapes:

  • "add_rows" — {"rows": [{...}, {...}]}
  • "update_rows" — {"predicate": <ibis predicate>, "set": {col: value}}
  • "delete_rows" — {"predicate": <ibis predicate>}
  • "replace_table" — {"expr": <engine TableExpr>}

Attributes:

Name Type Description
op str

One of the four op names above. Unknown ops raise :class:~corral.editing.errors.UnsupportedEditOp at apply time.

table str

Logical table name (matches :attr:corral.dataset.Table.name).

payload dict

Op-specific arguments — see above.

metadata dict[str, Any]

Free-form bag carried into the rollback log.

Examples:

>>> from corral.editing import Edit
>>> e = Edit(op="add_rows", table="link", payload={"rows": [{"link_id": 99}]})
>>> e.op, e.table
('add_rows', 'link')

corral.editing.EditResult(edit, diff, rollback_data, applied_at, session_id=None) dataclass

The outcome of one :class:Edit + how to roll it back.

Mutable so the owning :class:~corral.editing.session.Session can stamp session_id after construction. rollback_data is opaque per-op (callers pass it to :func:corral.editing.rollback rather than interpreting it directly).

Attributes:

Name Type Description
edit / diff

Inputs and outputs of the apply.

rollback_data Any

Op-specific blob the reverse handler consumes.

applied_at datetime

Wall-clock time the edit landed.

session_id str | None

Owning :class:Session’s id (None for standalone applies).

Examples:

>>> from datetime import datetime
>>> from corral.editing import Diff, Edit, EditResult
>>> e = Edit(op="add_rows", table="x", payload={"rows": [{}]})
>>> r = EditResult(
...     edit=e,
...     diff=Diff(edit=e, rows_added=1, rows_removed=0, rows_changed=0),
...     rollback_data=None,
...     applied_at=datetime(2026, 1, 1),
... )
>>> r.diff.rows_added
1

corral.editing.Session(package, *, log_path=None, session_id=None)

Atomic batch of :class:Edit s with chronological rollback log.

Inside the with block, :meth:add_edit applies each :class:Edit immediately and appends its :class:EditResult to the in-memory log; the package’s :class:DirtyTracker (when attached) is notified on every apply. On exception, every applied edit is reversed in LIFO order before the exception propagates (atomicity); on clean exit, the log is persisted to log_path (when set) so :func:corral.editing.rollback can replay later.

Parameters:

Name Type Description Default
package Package

The :class:~corral.dataset.Package to mutate.

required
log_path Path | str | None

Optional sidecar parquet path. None keeps the log in memory only.

None
session_id str | None

Optional explicit session id; default is a uuid4.

None

Examples:

>>> import tempfile, pathlib
>>> from corral.dataset import Package, Table
>>> from corral.editing import Edit, Session
>>> from corral.engines.ibis_engine import IbisEngine
>>> e = IbisEngine()
>>> pkg = Package.from_tables({"x": Table(name="x", expr=e.from_records([{"a": 1}]), engine=e)})
>>> with tempfile.TemporaryDirectory() as tmp:
...     log = pathlib.Path(tmp) / "history.parquet"
...     with Session(pkg, log_path=log) as s:
...         _ = s.add_edit(Edit(op="add_rows", table="x",
...                             payload={"rows": [{"a": 2}, {"a": 3}]}))
...     pkg["x"].count(), log.exists()
(3, True)

Construct an unopened session — call :meth:__enter__ (or use with) before edits.

add_edit(edit)

Apply edit immediately, stamp with session_id, append to the log.

Raises :class:EditingError if called outside the with block.

Examples:

>>> from corral.dataset import Package, Table
>>> from corral.editing import Edit, Session
>>> from corral.engines.ibis_engine import IbisEngine
>>> e = IbisEngine()
>>> pkg = Package.from_tables({"x": Table(name="x", expr=e.from_records([{"a": 1}]), engine=e)})
>>> with Session(pkg) as s:
...     r = s.add_edit(Edit(op="add_rows", table="x", payload={"rows": [{"a": 9}]}))
>>> r.diff.rows_added
1

rollback()

Reverse every edit applied in this session in LIFO order (without raising).

corral.editing.rollback

Replay a persisted rollback log to reverse past edits.

Companion to :class:~corral.editing.session.Session. Reads the parquet sidecar the session wrote, filters by the to selector (session_id / timestamp / all), and reverses the selected edits in LIFO order.

rollback(package, log_path, *, to=None)

Reverse every edit in log_path matching the to selector.

Reads the parquet sidecar written by :class:Session, filters per to, and reverses the selected edits in LIFO order so dependent edits unwind cleanly.

Parameters:

Name Type Description Default
package Package

The :class:~corral.dataset.Package to mutate.

required
log_path Path | str

Path to the parquet log written by a :class:Session.

required
to str | datetime | None

Selector.

  • None (default) — reverse every edit.
  • str — treat as session_id; reverse only that session.
  • :class:datetime — reverse edits applied strictly after this time.
None

Returns:

Name Type Description
The list[EditResult]

Raises:

Type Description
RollbackError

log_path missing or unparseable, or to of the wrong type.

Examples:

>>> import tempfile, pathlib
>>> from corral.dataset import Package, Table
>>> from corral.editing import Edit, Session, rollback
>>> from corral.engines.ibis_engine import IbisEngine
>>> e = IbisEngine()
>>> pkg = Package.from_tables({"x": Table(name="x", expr=e.from_records([{"a": 1}]), engine=e)})
>>> with tempfile.TemporaryDirectory() as tmp:
...     log = pathlib.Path(tmp) / "history.parquet"
...     with Session(pkg, log_path=log) as s:
...         _ = s.add_edit(Edit(op="add_rows", table="x", payload={"rows": [{"a": 2}]}))
...     _ = rollback(pkg, log)
...     pkg["x"].count()
1

Quality (generic framework)

corral.quality.Rule

Bases: Protocol

A data-quality rule.

Domain packages register concrete rules under the corral.quality.rules entry-point group. Each rule emits :class:~corral.reports.Issue records with category=Category.DATA_QUALITY into the :class:~corral.reports.ValidationReport it is handed.

Attributes:

Name Type Description
code str

Stable dotted identifier (e.g. "quality.high_speed_residential"). Used by config lookup and report filters.

description str

One-line plain-English explanation.

severity Severity

Default severity. Callers may override via :attr:RuleConfig.severity_override.

applies_to(package)

Cheap pre-check: should this rule run on this package at all?

run(package, report)

Execute the rule; populate report with issues.

corral.quality.RuleConfig(enabled=True, severity_override=None, thresholds=dict()) dataclass

Per-rule configuration.

Domain rules may define their own Config subclass with strongly typed thresholds; this base is sufficient for the generic framework.

Attributes:

Name Type Description
enabled bool

Set to False to skip this rule entirely.

severity_override Severity | None

Optional severity override (e.g., demote ERROR to WARNING).

thresholds dict[str, Any]

Rule-specific thresholds (e.g., {"speed_limit_mph": 45}). Opaque to the framework — the rule reads its own keys.

corral.quality.run_quality(package, *, config=None, report=None)

Run all registered quality rules on package.

Skips rules where :attr:RuleConfig.enabled is False or :meth:Rule.applies_to returns False. Applies :attr:RuleConfig.severity_override post-hoc by rewriting the severity of every issue the rule emitted in this call.

Parameters:

Name Type Description Default
package Package

The data package to evaluate. Opaque to the framework; rules consume it via the lazy Package / Table API.

required
config dict[str, RuleConfig] | None

Optional per-rule config keyed by :attr:Rule.code. Missing keys use :class:RuleConfig defaults.

None
report ValidationReport | None

Existing report to append issues to. When None, a fresh :class:ValidationReport is constructed.

None

Returns:

Type Description
ValidationReport

The (possibly newly constructed) report, populated with one

ValidationReport

Operations (cost model + gating)

corral.operations.OperationCost(op_name, n_rows, n_tables=1, fmt=None) dataclass

Heuristic estimate of an operation’s wall time.

Attributes:

Name Type Description
op_name str

Logical operation kind, e.g. "read", "validate_fk", "scope_bbox". Used (with fmt) to look up a coefficient in :data:COEFFICIENTS.

n_rows int

Approximate row count of the input. May be a planner estimate — exact counts are not required.

n_tables int

Number of tables involved. Multi-table ops (FK checks, joins) cost roughly linearly in table count for the v1.0 model.

fmt str | None

Source format if known (e.g. "parquet", "csv", "duckdb"). Only consulted for read / write ops.

Examples:

Estimate a 5M-row parquet read::

>>> from corral.operations import OperationCost
>>> cost = OperationCost(op_name="read", n_rows=5_000_000, fmt="parquet")
>>> round(cost.est_seconds(), 2)
2.5

est_seconds()

Return the estimated wall time in seconds.

The estimate is coefficient * (n_rows / 1_000_000) * n_tables where coefficient comes from :data:COEFFICIENTS (see :func:_resolve_coefficient for resolution rules). This is a documented heuristic, not a guarantee — see the module docstring.

corral.operations.gate(op_cost, *, approve=False, estimate_threshold_s=30.0, approval_threshold_s=180.0)

Check whether an op should proceed; emit estimate or block on approval.

Parameters:

Name Type Description Default
op_cost OperationCost

The cost record to inspect.

required
approve bool

If True, bypass the approval block. The CLI sets this from --yes or NETSTEAD_AUTO_APPROVE=1; programmatic callers pass it explicitly.

False
estimate_threshold_s float

Estimates at or above this value are logged.

30.0
approval_threshold_s float

Estimates at or above this value require approve=True.

180.0

Returns:

Type Description
OperationCost

The same op_cost (returned for fluent chaining at call sites).

Raises:

Type Description
ApprovalRequired

If op_cost.est_seconds() >= approval_threshold_s and approve is False.

Examples:

Run a cheap op without prompting::

>>> from corral.operations import OperationCost, gate
>>> _ = gate(OperationCost(op_name="read", n_rows=1_000, fmt="parquet"))

corral.operations.ApprovalRequired(cost, threshold_s)

Bases: Exception

Raised when an op exceeds the approval threshold and approve=False.

The CLI / MCP / notebook surface catches this, prompts the user (or honours --yes / NETSTEAD_AUTO_APPROVE=1), and re-calls the gated function with approve=True.

Record the offending cost + threshold and build a human message.

corral.operations.Batch(package, *, log_path=None, strict=False)

Defer + coalesce :class:Edit s on a :class:Package (architecture §6.5).

:meth:add_edit only queues edits — nothing is applied until __exit__ (or an explicit :meth:flush). On clean exit the queue is coalesced via :func:coalesce, applied atomically through a single :class:~corral.editing.Session, then validated once. On exception, the queue is discarded without ever opening a Session, so package state is guaranteed unchanged.

Parameters:

Name Type Description Default
package Package

The :class:~corral.dataset.Package to mutate.

required
log_path Path | str | None

Optional sidecar path forwarded to the underlying :class:Session (rollback-log file).

None
strict bool

When True, an ERROR-severity issue in the post-commit validation report triggers a Session rollback and re-raises as :class:BatchValidationError.

False

Examples:

>>> from corral.dataset import Package, Table
>>> from corral.editing import Edit
>>> from corral.engines.ibis_engine import IbisEngine
>>> from corral.operations import Batch
>>> e = IbisEngine()
>>> pkg = Package.from_tables(
...     {"t": Table(name="t", expr=e.from_records([{"id": 1}]), engine=e)}
... )
>>> with Batch(pkg) as b:
...     b.add_edit(Edit(op="add_rows", table="t", payload={"rows": [{"id": 2}]}))
...     b.add_edit(Edit(op="add_rows", table="t", payload={"rows": [{"id": 3}]}))
>>> pkg["t"].count()
3

Construct an unopened batch — open via with before queueing.

__enter__()

Open the batch for :meth:add_edit calls.

__exit__(exc_type, exc_val, exc_tb)

Commit on clean exit; discard the queue on exception (no Session opened).

add_edit(edit)

Queue edit for application at commit time.

Raises:

Type Description
RuntimeError

When called outside the with block.

Examples:

>>> from corral.dataset import Package, Table
>>> from corral.editing import Edit
>>> from corral.engines.ibis_engine import IbisEngine
>>> from corral.operations import Batch
>>> e = IbisEngine()
>>> pkg = Package.from_tables(
...     {"t": Table(name="t", expr=e.from_records([{"id": 1}]), engine=e)}
... )
>>> with Batch(pkg) as b:
...     b.add_edit(Edit(op="add_rows", table="t", payload={"rows": [{"id": 9}]}))
>>> pkg["t"].count()
2

flush()

Apply the pending queue immediately (mid-block manual commit).

Returns None when the queue is empty. Otherwise opens a Session, applies the coalesced edits, validates once, clears the queue, and stashes the result on :attr:last_result.

corral.operations.coalesce(edits)

Merge compatible same-table edits, preserving original order otherwise.

See module docstring for the full rule list.

Parameters:

Name Type Description Default
edits list[Edit]

The pending queue, in insertion order.

required

Returns:

Type Description
list[Edit]

A new list with coalesced edits in the original relative order.

Examples:

>>> from corral.editing import Edit
>>> from corral.operations import coalesce
>>> a = Edit(op="add_rows", table="t", payload={"rows": [{"id": 1}]})
>>> b = Edit(op="add_rows", table="t", payload={"rows": [{"id": 2}]})
>>> out = coalesce([a, b])
>>> len(out), out[0].payload["rows"]
(1, [{'id': 1}, {'id': 2}])

API server primitives

corral.api.build_app(settings=None, *, package_loader=None, extra_router_factory=None)

Return a :class:FastAPI wired with the generic corral endpoints.

Parameters:

Name Type Description Default
settings ServerSettings | None

:class:ServerSettings to mount. Defaults to :class:ServerSettings defaults (localhost, no packages).

None
package_loader PackageLoader | None

Callable that loads a package from a source string. Domain extensions pass their own loader so the registry caches the right type (e.g. Network.from_source). Defaults to Package.from_source.

None
extra_router_factory ExtraRouterFactory | None

Optional callable used by domain packages (netstead) to attach extra routers. See :class:ExtraRouterFactory for the contract + an example.

None

Returns:

Name Type Description
FastAPI

A fully wired :class:FastAPI instance. The caller passes it

to FastAPI

corral.api.PackageRegistry(settings, *, loader=None)

Lazy registry mapping public id → :class:Package (or subclass).

Built once at app startup from :class:ServerSettings.packages; each :meth:get materialises the package on first access via the configured loader (default :meth:Package.from_source) and caches the result. Hot-reload + cache invalidation are out of scope for v1 — restart the server to pick up a config change.

Domain extensions inject a loader to cache the right type. For example, :func:netstead.server.build_app passes Network.from_source so the cache holds :class:Network instances and the /networks/{id} handler can read pkg.spec_version directly without a second load.

Index settings by public id and remember which loader to use on first access.

Parameters:

Name Type Description Default
settings ServerSettings

The server config; only the packages list is consumed here.

required
loader PackageLoader | None

Callable that takes a source string and returns a :class:Package (or subclass). Defaults to :meth:Package.from_source.

None

require(pkg_id)

Return the package for pkg_id or raise :class:fastapi.HTTPException(404).

Public surface for endpoint handlers — both the generic corral routes and domain extensions (netstead.server). Replaces the previous private _safe_get helper.

source_for(pkg_id)

Return the configured source string for pkg_id.

Public surface for domain routers that need to re-resolve the original source (e.g. via a domain-specific factory like :meth:Network.from_source). Use this instead of reaching into the private _refs dict.

describe(pkg_id)

Return {id, source, description} without loading the package.

list_ids()

Return all configured public ids (insertion order).

corral.api.ServerSettings

Bases: BaseModel

Top-level server configuration.

Attributes:

Name Type Description
bind str

Interface to bind. Default "127.0.0.1" — operator must explicitly opt into a public bind.

port int

TCP port. Default 8000.

auth AuthSettings

Auth configuration (default bearer-token).

packages list[PackageRef]

List of packages to expose under public ids. Empty list means “no packages mounted”; the app still starts (e.g. for /health smoke tests) but every /packages/{id} request 404s.

is_public_bind()

Return True when :attr:bind exposes the server beyond localhost.

warn_on_unsafe_combinations()

Emit a loud log warning when the combination is risky.

Today: auth=none + non-localhost bind. Future: missing TLS, weak token, etc.

corral.api.AuthSettings

Bases: BaseModel

How the server authenticates incoming requests.

Attributes:

Name Type Description
kind Literal['none', 'bearer']

Either "none" (no auth — only safe on localhost) or "bearer" (require a Bearer <token> header matching :attr:token). Default "bearer".

token str | None

Required when kind == "bearer". Compared against the incoming Authorization header in constant time.

corral.api.ExtraRouterFactory

Bases: Protocol

Callable shape accepted by :func:build_app for domain router extensions.

A factory receives the live :class:PackageRegistry and the auth dependency, and returns a :class:fastapi.APIRouter mounted on top of the generic routes.

Examples:

Minimal extension that adds a /extras/{id} endpoint::

from fastapi import APIRouter, Depends
from corral.api import build_app, AuthDep, PackageRegistry

def my_factory(registry: PackageRegistry, auth_dep: AuthDep) -> APIRouter:
    router = APIRouter(prefix="/extras", tags=["extras"])

    @router.get("/{pkg_id}", dependencies=[Depends(auth_dep)])
    def get_extra(pkg_id: str) -> dict:
        return registry.describe(pkg_id)

    return router

app = build_app(settings, extra_router_factory=my_factory)

__call__(registry, auth_dep)

Return a router to mount alongside the generic routes.

corral.api.AuthDep = Callable[[str | None], None] module-attribute

corral.api.PackageLoader = Callable[[str], Package] module-attribute

MCP primitives

corral.mcp.build_server(name='corral', *, state=None)

Return a configured :class:FastMCP exposing generic corral tools.

The returned server is ready to run(transport="stdio"). Tools:

  • describe_package(source) — package metadata: table list + row counts + engine name.
  • validate_package(source) — full validation; returns the :class:~corral.reports.ValidationReport as a JSON dict via :meth:~corral.reports.ValidationReport.to_dict (canonical wire shape).
  • list_tables(source) — short list of table names (cheap, no row counts).

Parameters:

Name Type Description Default
name str

MCP server display name. Default "corral".

'corral'
state dict[str, Any] | None

Optional shared-state dict for stateful tools (sessions, indexed scope caches, etc.) that compose on this server via :func:netstead.mcp.build_server. Stored on the returned server’s settings under the key "corral_state" so domain extensions can read / write keys cooperatively. None (default) means the server is purely stateless — today’s contract. The kwarg is documented + accepted now so the public signature can grow stateful tools later without a breaking change.

None

Returns:

Name Type Description
A FastMCP
FastMCP

invokes .run(transport="stdio") from the CLI entry point.

Examples:

>>> import pytest
>>> pytest.importorskip("mcp")
<module ...>
>>> server = build_server()
>>> # The server object is ready; ``server.run(...)`` would block
>>> # on stdio. Tests exercise the tools via the lower-level
>>> # registry instead.
>>> isinstance(server.name, str)
True
>>> # Stateful seam: pass a dict that future stateful tools
>>> # (netstead.mcp.edit_session etc.) will read + write into.
>>> shared = {}
>>> server2 = build_server(state=shared)
>>> isinstance(server2.name, str)
True

See also