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: |
tables |
dict[str, Table]
|
Mapping |
engine |
Engine | None
|
The :class: |
source |
str | None
|
Original source identifier (path / URL). |
dirty_tracker |
Any
|
Optional sync-state tracker. When |
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: |
required |
engine
|
Engine | None
|
Engine to materialise through. Defaults to the
registered default (typically
:class: |
None
|
spec
|
DataPackage | str | Path | None
|
Either an in-memory :class: |
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
|
None
|
credentials
|
dict[str, str] | None
|
Optional explicit credentials dict forwarded
to :class: |
None
|
Returns:
| Type | Description |
|---|---|
Package
|
A populated :class: |
Package
|
entries. |
Raises:
| Type | Description |
|---|---|
FormatNotDetected
|
If neither extension dispatch nor the
directory walk could resolve |
Examples:
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 |
required |
spec
|
DataPackage | None
|
Optional explicit :class: |
None
|
engine
|
Engine | None
|
Optional explicit engine. Defaults to the first table’s engine. |
None
|
Returns:
| Type | Description |
|---|---|
Package
|
A new :class: |
Examples:
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: |
True
|
structural
|
bool
|
Run structural validation
(:func: |
True
|
foreign_keys
|
bool
|
Run FK validation
(:func: |
True
|
sync_state
|
bool
|
Consult the :class: |
True
|
strict
|
bool
|
Forwarded to validators that accept it (FK
|
False
|
Returns:
| Name | Type | Description |
|---|---|---|
A |
ValidationReport
|
|
ValidationReport
|
with one :class: |
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:
formatexplicit → adapter resolved by name;- else inferred from
destextension; - else default to
"parquet"whendestlooks 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;
|
None
|
overwrite
|
bool
|
When |
False
|
strict_sync
|
bool
|
When |
False
|
Raises:
| Type | Description |
|---|---|
FileExistsError
|
When |
OutOfSyncError
|
When |
PackageError
|
When the package has no engine attached (so materialisation can’t run). |
FormatNotDetected
|
When |
WriteUnsupportedForSchemeError
|
When |
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:
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: |
expr |
TableExpr
|
Engine-native lazy expression. Concrete type depends on
the engine ( |
engine |
Engine
|
The engine that produced |
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 |
format |
str | None
|
Optional short format identifier ( |
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 |
required |
Returns:
| Type | Description |
|---|---|
Table
|
A new :class: |
Table
|
the original :class: |
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 |
()
|
Returns:
| Type | Description |
|---|---|
Table
|
A new :class: |
Examples:
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: |
Examples:
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:
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:
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:
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:
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:
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:
- 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.EngineNotAvailableErrorwith a clear pointer at the engine that does. - Schema casting (
cast_schema). Promoted from a per-engine internal helper so adapters can apply Frictionless schema casts uniformly after a primitive read. - Convenience delegators (
scan/write). 3-line wrappers that route throughcorral.io.dispatch. They exist so callers who don’t care about adapters can still writeengine.scan(path).
Attributes:
| Name | Type | Description |
|---|---|---|
name |
str
|
Short identifier used as the registry key
( |
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 / |
required |
schema
|
Any | None
|
Optional Frictionless :class: |
None
|
**kwargs
|
Any
|
Forwarded verbatim to the underlying library reader.
Names follow the library, not a normalized vocabulary —
e.g. polars wants |
{}
|
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 |
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 |
required |
table
|
str
|
The table name to read. Required (no defaulting —
the duckdb adapter’s |
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.LazyFrameviapl.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_pandasis preserved end-to-end.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
arrow_table
|
Table
|
Already-materialised :class: |
required |
Returns:
| Type | Description |
|---|---|
TableExpr
|
An engine-native lazy expression; the contents match |
TableExpr
|
|
TableExpr
|
convention. |
Examples:
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
|
**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: |
required |
Returns:
| Type | Description |
|---|---|
TableExpr
|
A new expression with the casts applied, or |
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,
Any other dict shape MUST raise
:class: |
required |
format
|
str | None
|
Optional explicit format hint forwarded to
|
None
|
schema
|
Any | None
|
Optional Frictionless |
None
|
**kwargs
|
Any
|
Adapter-specific options forwarded verbatim to the
resolved |
{}
|
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 |
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 |
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 |
required |
Returns:
| Type | Description |
|---|---|
TableExpr
|
A new engine-native expression carrying only |
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 |
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 |
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 understandpd.ArrowDtypecolumns.
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
|
required |
fmt
|
str
|
Format name ( |
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 ( |
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):
Looking up a name that is not registered raises a clear error:
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 |
required |
default
|
bool
|
If |
False
|
Raises:
| Type | Description |
|---|---|
TypeError
|
If |
ValueError
|
If |
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
|
|
required |
Returns:
| Type | Description |
|---|---|
Engine
|
The registered DuckDB engine. |
Raises:
| Type | Description |
|---|---|
ValueError
|
If |
Examples:
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. |
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 |
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
¶
has_warnings
property
¶
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:
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:
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. |
required |
message
|
str
|
Human-readable description that names the broken input. |
required |
**kw
|
Any
|
Any of |
{}
|
Returns:
| Type | Description |
|---|---|
Issue
|
The newly created (and stored) issue. |
Examples:
by_severity(severity)
¶
Return all issues at the given severity, preserving insertion order.
Examples:
by_category(category)
¶
Return all issues in the given category, preserving insertion order.
Examples:
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
|
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:
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 |
None
|
indent
|
int
|
|
2
|
Examples:
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:
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 |
None
|
title
|
str | None
|
Optional override for the |
None
|
include_map
|
bool
|
If |
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. |
message |
str
|
Human-readable, names the input that broke. |
table |
str | None
|
Table name, or |
column |
str | None
|
Column / field name, or |
row |
int | None
|
Zero-based row index, or |
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
|
|
INFO |
Informational only. No action required. Also the default severity for an awareness-only data-quality finding. |
Examples:
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: |
|
STRUCTURAL |
Package-level structure — missing required table,
extra unknown table, missing file on disk. Produced by
:mod: |
|
FOREIGN_KEY |
Cross-table referential integrity. Produced by
:mod: |
|
SYNC_STATE |
A previously validated FK is now stale because one
side has been mutated since the last check. Produced by
:mod: |
|
DATA_QUALITY |
A configurable quality rule (e.g. |
Examples:
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: |
table |
str
|
Logical table name (matches
:attr: |
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: |
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: |
required |
log_path
|
Path | str | None
|
Optional sidecar parquet path. |
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: |
required |
log_path
|
Path | str
|
Path to the parquet log written by a :class: |
required |
to
|
str | datetime | None
|
Selector.
|
None
|
Returns:
| Name | Type | Description |
|---|---|---|
The |
list[EditResult]
|
|
Raises:
| Type | Description |
|---|---|
RollbackError
|
|
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.
|
description |
str
|
One-line plain-English explanation. |
severity |
Severity
|
Default severity. Callers may override via
:attr: |
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 |
severity_override |
Severity | None
|
Optional severity override (e.g., demote
|
thresholds |
dict[str, Any]
|
Rule-specific thresholds (e.g.,
|
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 |
required |
config
|
dict[str, RuleConfig] | None
|
Optional per-rule config keyed by :attr: |
None
|
report
|
ValidationReport | None
|
Existing report to append issues to. When |
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. |
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. |
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 |
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
|
180.0
|
Returns:
| Type | Description |
|---|---|
OperationCost
|
The same |
Raises:
| Type | Description |
|---|---|
ApprovalRequired
|
If |
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: |
required |
log_path
|
Path | str | None
|
Optional sidecar path forwarded to the underlying
:class: |
None
|
strict
|
bool
|
When |
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 |
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: |
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.
|
None
|
extra_router_factory
|
ExtraRouterFactory | None
|
Optional callable used by domain packages
(netstead) to attach extra routers. See
:class: |
None
|
Returns:
| Name | Type | Description |
|---|---|---|
FastAPI
|
A fully wired :class: |
|
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 |
required |
loader
|
PackageLoader | None
|
Callable that takes a source string and returns a
:class: |
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 |
port |
int
|
TCP port. Default |
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 |
corral.api.AuthSettings
¶
Bases: BaseModel
How the server authenticates incoming requests.
Attributes:
| Name | Type | Description |
|---|---|---|
kind |
Literal['none', 'bearer']
|
Either |
token |
str | None
|
Required when |
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.ValidationReportas 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'
|
state
|
dict[str, Any] | None
|
Optional shared-state dict for stateful tools (sessions,
indexed scope caches, etc.) that compose on this server
via :func: |
None
|
Returns:
| Name | Type | Description |
|---|---|---|
A |
FastMCP
|
|
FastMCP
|
invokes |
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 ¶
shared/ai/index.md— explains the api-index.json + llms.txt artifacts.- corral cookbook
- Architecture