Metadata-Version: 2.5
Name: deltacat-client
Version: 0.1.32
Summary: Lightweight REST client and thin MCP scaffolding for the DeltaCAT API server.
Project-URL: Homepage, https://github.com/ray-project/deltacat
Author: Ray Team
License-Expression: Apache-2.0
Classifier: Development Status :: 4 - Beta
Classifier: Intended Audience :: Developers
Classifier: Operating System :: OS Independent
Classifier: Programming Language :: Python :: 3 :: Only
Classifier: Programming Language :: Python :: 3.10
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Requires-Python: <3.13,>=3.10
Requires-Dist: attrs>=22.2.0
Requires-Dist: deltacat-io-core<0.2,>=0.1.32
Requires-Dist: httpx<0.29.0,>=0.23.0
Requires-Dist: python-dateutil>=2.8.0
Requires-Dist: requests>=2.31
Provides-Extra: all
Requires-Dist: deltacat-io-core[all]<0.2,>=0.1.32; extra == 'all'
Provides-Extra: daft
Requires-Dist: deltacat-io-core[daft]<0.2,>=0.1.32; extra == 'daft'
Provides-Extra: dataloader
Requires-Dist: deltacat-io-core[dataloader]<0.2,>=0.1.32; extra == 'dataloader'
Requires-Dist: ray[default]==2.48.0; extra == 'dataloader'
Provides-Extra: io
Requires-Dist: deltacat-io-core[io]<0.2,>=0.1.32; extra == 'io'
Provides-Extra: lance
Requires-Dist: deltacat-io-core[lance]<0.2,>=0.1.32; extra == 'lance'
Provides-Extra: lerobot
Requires-Dist: deltacat-io-core[lerobot]<0.2,>=0.1.32; extra == 'lerobot'
Provides-Extra: mcp
Requires-Dist: mcp<2,>=1.10; extra == 'mcp'
Provides-Extra: packds
Requires-Dist: deltacat-io-core[packds]<0.2,>=0.1.32; extra == 'packds'
Provides-Extra: pandas
Requires-Dist: deltacat-io-core[pandas]<0.2,>=0.1.32; extra == 'pandas'
Provides-Extra: polars
Requires-Dist: deltacat-io-core[polars]<0.2,>=0.1.32; extra == 'polars'
Description-Content-Type: text/markdown

<p align="center">
  <img src="https://github.com/ray-project/deltacat/raw/2.0/media/deltacat-logo-alpha-750.png" alt="deltacat logo" style="width:55%; height:auto; text-align: center;">
</p>

`deltacat-client` is the primary Python client package for
[DeltaCAT](https://github.com/ray-project/deltacat). It lets you read and
write tables, run jobs, and build data pipelines against a DeltaCAT API
server without installing the full storage, compute, or server runtime
stack.

The client talks to a DeltaCAT server over HTTP. Metadata operations
(schema validation, transaction management, compaction) stay server-side,
while large data reads and writes go directly between the client and cloud
storage using short-lived credentials the server vends on demand.


## Overview

The client is organized around a root `Client` object with resource-oriented subclients:

| Subclient | Purpose |
|-----------|---------|
| `client.catalog` | Namespaces, tables, read/write, transactions |
| `client.jobs` | Job submission, claiming, lifecycle, progress |
| `client.publications` | Incremental producers into the DeltaCAT Lakehouse (pipeline root nodes) |
| `client.subscriptions` | Incremental consumers from the DeltaCAT Lakehouse (pipeline leaf nodes) |
| `client.transforms` | Incrementally transform data within the DeltaCAT Lakehouse (pipeline intermediate nodes) |
| `client.pipelines` | Wire publications, transforms, and subscriptions into a connected DAG |
| `client.training` | Create converters, reprocess selected episodes, observe runs, and explicitly promote staged repairs |

To copy files between storage locations while preserving their paths, start with
[incremental exact-path copies](#incremental-exact-path-copies) using `create_placer(...)`.


## Installation

```bash
pip install deltacat-client
```

Optional extras for local data materialization:

```bash
pip install "deltacat-client[all]"    # Full client-side read/write stack
pip install "deltacat-client[pandas]" # Pandas DataFrame support
pip install "deltacat-client[polars]" # Polars DataFrame support
pip install "deltacat-client[daft]"   # Daft DataFrame support
pip install "deltacat-client[lance]"  # Lance dataset support
pip install "deltacat-client[mcp]"    # Typed async MCP HTTP client
```

## Incremental exact-path copies

`create_placer(...)` builds a managed **crawler → differ → relay** pipeline.
Its `ExactPathRelaySubscriber` copies new and changed files byte-for-byte to
the same relative paths beneath a destination root. Unchanged files are skipped,
and destination copies are retained when source files are deleted.

Use an existing source discovery-root binding and a destination DataRoot with
`usage=["relay_target"]`, configured by your catalog administrator. This example
copies a source bucket into a replica prefix and checks for changes nightly:

```python
from deltacat_client import (
    Client,
    PlacementDestination,
    PlacementSource,
    create_placer,
)

client = Client("https://deltacat-api.example.com", bearer_token="your-api-token")
replica = create_placer(
    client,
    sources=[
        PlacementSource(
            name="source",
            uri="s3://source-bucket/",
            binding="source-discovery",
            region="us-east-1",
            reflect_deletions=False,  # Retain copies if source files are deleted.
        ),
    ],
    destinations={
        "source": [
            PlacementDestination(
                name="replica",
                binding="replica-s3",
                uri="s3://replica-bucket/source-bucket/",
            ),
        ],
    },
    schedule_cron="0 0 * * *",
    schedule_timezone="America/Los_Angeles",
    name="nightly-replica",
    run_id="v1",
    include_version_changes=True,  # Also relay timestamp-only changes.
    trigger_now=True,  # Start the initial copy, then check nightly.
)
```

For example, `s3://source-bucket/cameras/front.mp4` is copied to
`s3://replica-bucket/source-bucket/cameras/front.mp4`. The pipeline adds no
`__external__` or catalog/table directories to the destination paths.
Creation returns pipeline IDs; copying continues on the managed workers.
Use `dry_run=True` to preview the IDs without creating or starting a pipeline.
Each source's `reflect_deletions` policy defaults to `False`. `True` is currently
rejected because destination deletion is not supported.

The examples set `include_version_changes=True` to also copy objects whose
crawler-observed modification stamp changes while their content hash and size
stay equal, such as metadata-only updates. The default is `False`, which keeps
content-based change detection. This option compares existing crawl metadata;
it adds no object-body checksum scan. Changes invisible to the crawler's
observed stamp, hash, and size require an explicit source manifest or a fresh
inventory containing the changed metadata.

See [multiple sources and destination fanout](#copying-source-trees-fanout)
below and the [relay documentation](../deltacat/compute/relay/README.md#incremental-exact-path-placement) for
`ExactPathRelaySubscriber` and transfer behavior.

## Registering prewritten files

For files copied from an older table definition, `session.entry(...)` and
`session.entry_for(...)` accept `schema_id` and `sort_scheme_id`. Supply the IDs
that describe the actual payload, not the current table defaults. An explicit
`sort_scheme_id=None` or an omitted sort ID records unknown order. Data-bearing
writes record the scheme actually applied by the writer; a destination's current
scheme cannot certify a prewritten file. Registration preserves historical
schema/sort metadata and does not sort or rewrite a copied file. It does not
retain referenced external media.

Retain the original `session.session_id` before commit and recover it with
`client.catalog.get_write_session(...)` after an interruption. An explicit
`write_key` can recover an OPEN session, but staging with that key after commit
creates a new session for a later write. Retry the original session with the
same operation ID and entries; do not treat a recycled key as proof of completion.

For a destination captured earlier, pass `expected_table_id` and an explicit
`table_version` to `client.catalog.stage(...)`. The existing table and version
must match; `create_if_missing` cannot create a target for this conditional
session. The server retains the condition through replay and checks it inside
commit validation, refreshing mutable names so rename/recreate cannot redirect
the write. Omission keeps ordinary staging behavior. Use a fully upgraded server
and recovery-worker fleet: older servers do not understand the stored condition.

## Training converters

`client.training.create_converter(...)` provides a source-neutral namespace for
converter creation. Omitted source/output options retain the existing PackDS
pipeline defaults. Existing keyword options (loader configuration, dispatch,
schedules, placement, plugins, and explicit node/table IDs) remain available;
use `source_format` in place of `loader_format`.

Start with [creating a converter](#creating-converters),
[repairing selected episodes](#reprocessing-selected-episodes-preview), or
[observing conversion and required copies](#observing-conversion-and-required-copies-preview).
Creation success means graph installation, not completed conversion. The same
converter can be reused for later repairs; no per-episode pipeline is required.

The submitting client does not run the source loader. Installing
`deltacat-client` (even `[all]`) does not install robotics SDKs on a worker.
For local or custom **Dexmate worker** execution, install `deltacat[compute,dexmate]`
in that worker's environment; a combined local server/worker also needs the
`server` extra. This supplies the pinned HDF5 and OpenCV dependencies alongside
the shared codecs, without requiring the larger XDOF robotics stack. The existing
managed-worker image and per-job HDF5 overlay remain unchanged. See
[source runtime installation](../deltacat/compute/training/README.md#source-runtime-installation)
for plugin and source-specific prerequisites.

Repairs do not universally require unreleased-version promotion. Choose the
appropriate existing catalog write mode for the intended change and output
format; use isolated version promotion when the application needs that cutover
boundary. See [Choosing a repair write strategy](../deltacat/compute/training/README.md#choosing-a-repair-write-strategy)
for ADD/APPEND, keyed MERGE, partition/table replacement, and format limitations.
The preview `training.reprocess` follows this contract, with the currently
implemented capabilities described below; recognizing a mode does not imply
that every output writer supports it.

### Reprocessing selected episodes (preview)

```python
from deltacat_client import EpisodeSelection

operation = client.training.reprocess(
    converter.pipeline_id,
    episodes=EpisodeSelection.query(
        namespace="iris",
        table="episodes_to_repair",
        filter_predicate={"eq": ["needs_repair", True]},
    ),
    request_id="repair-2026-09-11-1",  # Retain and reuse with identical arguments.
    conversion_mode="full",
    write_mode="replace_partition",
)
result = client.training.wait_reprocess(operation, timeout=300)
assert result.phase == "succeeded"
```

`episodes` also accepts an ID sequence or `EpisodeSelection.manifest` with a
declared data root, path, format, size and SHA-256. `plan_reprocess` builds a
request locally, and `submit_reprocess` submits that same request on retry.
Admission resolves the selection/facts once and retains them in the ordinary
conversion job. `get_reprocess` and `wait_reprocess` include required outputs,
the rich episodes companion (inline or its reconciliation child), and required
post-conversion copies. A completed primary conversion alone is not success.
A client timeout does not cancel either job. `failed` is a terminal status,
not an exception from the wait method; inspect the ordinary job for its error.
Failure does not promise automatic rollback of earlier successful batches.

By default, `source_mode="recorded"` replays the converter's recorded source
facts. To acquire current placed files after source discovery/placement has
completed, pass `source_mode="refresh"`. Use `conversion_mode="full"` for new or
changed media, or `conversion_mode="reuse_media"` with `write_mode="replace_partition"`
for changed non-media features with unchanged media/frame alignment. A matching
worker captures the selected inputs and output plan once on the original job;
retry does not acquire newer files, and the standing pipeline watermark is not
advanced. This supports primary PackDS, native exports, LeRobot fanout and
LeRobot-only outputs with capable file-set loaders. It does not infer source/frame
alignment or retain deleted source objects. Mutable-source drift refuses rather
than silently changing the repair; submit a new request after correcting inputs.
Reuse keeps captured predecessor media and refuses incompatible alignment; it
does not silently fall back to full conversion. See the
[source-capture contract](../deltacat/compute/training/README.md#public-reprocess-admission).

Current execution supports a retained converter with PackDS and/or LeRobot outputs,
`full` conversion, `in_place` or `unreleased` publication and explicit `add`, `append`, or
`replace_partition`. Inline selections accept up to 4,096 episode IDs; use a
query or manifest for larger selections. Admission bounds selected IDs to
1,000,000 rows /64 MiB of native ID values and retained facts to 128 MiB,
with bounded per-record decoding. These are independent limits, not a promise
that every source/output combination can execute a million-episode repair:
source acquisition, media-reuse dependencies, output sessions and publication
retain their own enforced bounds. FULL conversion also accepts `replace`: it
preserves unaffected rows and
output sessions in private complete-table images and atomically publishes the
declared tables (the primary steps/episodes pair, LeRobot inventory tables, or
both) after any required primary reconciliation
and output generation finish. This entails full-image population, not only selected-partition
I/O. Explicit source-filter outcomes can remove selected episodes, including all
episodes; schemas and existing media are retained, and the companion must finish
before any replacement becomes visible. Any intervening source-table write
refuses the cutover, preserving that writer's result. Combined-output replacement
requires workers advertising `training-reprocess-output-replace-v1` in addition
to the existing private-replacement capability. LeRobot-only full-table replacement
also requires `training-reprocess-output-only-replace-v1`. It reuses the same native
publication machinery without primary tables or a companion job. Row MERGE remains
unsupported; [explicit staged-repair promotion](#explicit-staged-repair-promotion)
is a separate optional action after validation. ADD never
silently upgrades to replacement. Primary APPEND writes ordered catalog deltas
for complete new episode partitions, with an atomic absence condition; it cannot
append frames to an existing episode through this repair API. A missing table is
created without an AUTO fallback to ADD. LeRobot `append` keeps published episode
and task indices, then adds new episodes
after existing ones. ADD instead rebuilds canonical source ordering. Both refuse
existing episode identities; neither is a keyed upsert. Row MERGE is not
currently supported by this execution path (the PackDS layout rejects merge keys).
Source history uses the existing latest-delivery,
modification-time and path ordering, not append order. Exact duplicates collapse;
equal-version rows with conflicting contents and unsupported modes/outputs refuse
before job submission. Annotation converters instead reuse the standing join's
committed-history reader: annotation-only changes share a raw source version, so
the latest SUCCESS-visible Delta timestamp selects the retained annotated input.
Ambiguous timestamps or incomparable streams refuse; retry never reloads current
annotations. The reader's bounded physical scope may refuse oversized facts
files even for small selections; see the [training runtime bounds](../deltacat/compute/training/README.md#public-reprocess-admission).
Existing output roots
and converter options are retained; new media is isolated under the output root.
Workers must advertise `training-reprocess-v1`; older workers cannot claim these
jobs. Primary-plus-LeRobot generation repairs additionally require
`training-reprocess-generations-v1`. They retain original output inventories and
partition preconditions in the admitted job, preserving unselected session members
without old B1 staging or survivor video decoding. Output summaries must complete
before the primary result is checkpointed, and status includes the companion.
LeRobot-only repairs use the existing `source_node_process` handler and the same
output-only unit executor as ordinary conversion, with the additional
`training-reprocess-output-only-v1` worker tag. They do not create primary tables,
run B2, or request a primary reconciliation child. The standing source watermark
is not advanced or reset. Local-dispatch nodes execute on the local worker;
remote nodes still require the ordinary distributed catalog bootstrap.

For LeRobot-only outputs using the standard LeRobot or UMI source loader, set
`conversion_mode="reuse_media"` and `write_mode="replace_partition"` to rebuild
numeric features while retaining video content. Identity, camera membership,
frame count, rate and the complete timestamp/frame-index alignment must match the
original output. The job requires `training-reprocess-media-reuse-v1` workers.
Native LeRobot and UMI sources also support a primary PackDS output, with or without
additional LeRobot outputs, through `training-reprocess-primary-reuse-v1` workers.
This primary path is qualified locally, including changed-feature and repeated
repairs. Recorded predecessor storage revisions survive ordinary job claims and
credential renewal separately from output-write access; cloud/distributed
qualification remains separate. Source-specific loader overrides are not
implicitly reuse-capable. UMI feature-only acquisition skips source MP4
downloads and binds the original
camera-frame mapping in output metadata. Changed mappings or older output lacking
that metadata require full conversion. This mode avoids video decoding/encoding; final
complete-session assembly still copies videos and rebuilds Parquet, so it is
not an end-to-end zero-copy operation. Required placement/retry semantics are
the same as full conversion.

An explicit reuse repair rebuilds the selected features even when source values
are unchanged. This is separate from incremental pipeline no-op detection and
request-ID replay, which still returns the original job. Repeated primary repairs
continue the existing committed media-dependency edges and preserve their depth;
missing or inconsistent predecessor metadata refuses instead of resetting history.

Inventories use one projected read per declared output under a shared snapshot;
assembly rebuilds complete sessions, not only selected files. The complete
predecessor is streamed; it is no longer limited to 512 episodes. Retained
inventories now use original snapshot-bound native file references, with up to
1,048,576 files under the runtime's byte/index bounds. They must remain available
for the job; missing or changed original files refuse without selecting current
table data. Retained selected job input uses the shared 1,000,000-row /128-MiB
bounds; compact legacy inputs retain their original 4,096-row /4-MiB encoding.
Explicit selected receipts also stream, without the old 512-episode
execution-unit limit; the job-input limit still applies. Real local tests cover
two-output repairs within a 513-episode predecessor and public repair selecting
513 of 514 episodes, plus two of 2,049 episodes with a 4,104-file inventory,
with original staging removed. A larger local query-selected repair converts
4,097 of 4,098 generated-media episodes with three projected admission reads;
separate primary-reuse metadata controls capture 4,097 real partitions through
one companion scan. Neither result is a vendor-media or distributed throughput
claim. See the training
README's generation-repair section for physical work and byte/index bounds.

Explicit repairs and their direct copy/companion jobs are retained beyond ordinary
job/event expiry, so status and same-request replay keep using the original inputs.
A different request under the same ID conflicts. This retains job metadata until
deliberate removal, not source/media files; deleting the job/control store removes
the guarantee, and memory mode does not survive process loss. Historical layouts
without retained source receipts, complete runtime pinning and distributed
qualification remain unfinished. Unsupported configurations
refuse rather than silently omitting outputs/copies. History processing is
bounded to one million rows / 128 MiB after the projected read; larger history
acquisition and unversioned same-path source conflicts need additional handling.
This preview is not production-qualified.

### Reading an allocated private repair image

For an allocated private repair image, `catalog.describe_table`, `catalog.plan`
and `catalog.read` accept the same `staged_stream_id`, `expected_table_id` and
explicit `table_version` as its write sessions. Use `metadata_only=True` for
its description. These calls require an upgraded server and writer access;
ordinary calls remain active-image reads. Missing or published private targets
refuse rather than select another stream. These are internal selectors used by
the public complete-image repair machinery, not a replacement for
`training.reprocess` or caller-chosen repair destination overrides.

A saved private `Plan` retains the selection in its request context and JSON
representation. `catalog.read(plan=...)` executes the original files even after
publication or a later replacement; it does not reread the active alias. SDK
and shared IO-core execution reject dropped or conflicting destination fields.
Referenced-media credential refresh retains the same selection and bounded
predicate, without a per-chunk catalog request. Upgrade client, IO-core, server
and recovery workers together before persisting private plans; older clients
that ignore these fields are not supported private-plan consumers.

### Copying pre-conversion files to custom paths (preview)

`pre_conversion_placements` maps exact crawler seed URIs to one or more
`PlacementDestination` values. For example, add this to `plan_converter` or
`create_converter`:

```python
from deltacat_client import PlacementDestination

pre_conversion_placements = {
    "file:///datasets/my_dataset/": [
        PlacementDestination(
            name="archive",
            binding="archive-root",
            uri="s3://my-archive/team-selected-layout/my_dataset/",
        ),
    ],
}
```

The key must match a declared seed (a trailing slash is optional). Copies retain
each source-relative filename; no catalog/table directory hierarchy is added.
Each destination is an ordinary relay subscription with its own watermark,
sharing the existing crawler and file diff. A loader-discovery recipe without
a diff gets one shared diff, not another crawler or one diff per destination.
Omitting this option preserves the original version-1 declaration bytes.

Plans and dry-runs require explicit destination URIs and do not look up storage
roots. Creation can resolve an omitted URI from its named DataRoot. The binding
must authorize relay writes at that exact path; a local root also needs its
normal server-execution write policy. Configuration does not grant credentials
or access. Overlapping source/destination and destination/destination prefixes
are rejected. Source DELETE events do not delete archive copies.

These are independent copies of **pre-conversion files**, not a barrier that
delays conversion until copying finishes or a redirection of conversion input.
Creation status means graph installation, not successful copies. Observe the
source crawl job with `get_run`/`wait_run` to include its required copy branches;
a conversion-leaf job does not include independent copies from an earlier crawl.
See the run observer below for required versus nonblocking completion. Remote
expiry/locality and distributed Alpha qualification remain open.

### Copying converted outputs to custom paths (preview)

`post_conversion_placements` uses the same destination declarations, keyed by
the `output_root_name` of a declared LeRobot or PackDS inventory output:

```python
post_conversion_placements = {
    "training-output": [
        PlacementDestination("archive", "archive-root", "s3://my-archive/training/"),
        PlacementDestination("local-copy", "local-root", "file:///datasets/training/"),
    ],
}
```

Pass this to `training.plan_converter` or `training.create_converter` alongside
the output declaration. When output formats are omitted, a single mapping key
selects the primary inventory root, retaining the source's existing conversion
settings. Put all desired destinations under that key. The same opt-in works
for an explicitly declared primary output without an inventory root, including
beside LeRobot outputs. Selecting only a declared LeRobot root does not also
copy the primary, and an explicitly rooted primary is never retargeted. Multiple
otherwise undeclared roots refuse as ambiguous. Empty/no placement mappings
leave the original definition unchanged.

Copies require their original files to remain available through completion and
the intended retry window; retained job metadata is not a media backup. A
missing original fails the copy instead of copying a newer generation. See
[placement input lifetime](../deltacat/compute/training/README.md#placement-input-lifetime)
for retention-policy boundaries and same-job retry after restoring originals.

Each committed output inventory has its own subscription
per destination; no output-directory crawler or extra encoder
is created. Relative session and `.generations/<id>` paths are preserved, including
the Parquet, metadata and video layout. The copied generation directory is the
dataset, not the outer destination root; no current-generation alias is implied.
Outputs are not pruned from archives when the source inventory advances.
This matches ordinary per-table dispatch: one output cannot consume another
output's pending copy under a shared watermark. Graph size is bounded and scales
with declared outputs times destinations, not episodes or inventory files. New
post-conversion event jobs retain the original publication inputs in their job
context and do not consume or advance the subscription watermark. A later
conversion cannot silently substitute its files for an older pending copy.

Named source roots resolve from the executing job's catalog before inventory reads.
This is routing metadata only: relay jobs still enforce scoped source/destination
access. Missing roots and overlaps with any declared file-output root refuse before
copying. Copies use the same generation publisher as conversion and repair;
the presence of post-placement is not required for a single LeRobot output to
have a committed inventory. Placement itself does not reinterpret already
persisted converter definitions or legacy output files.

Pre/post copies may be combined. Destination prefixes must not overlap across
branches. Custom-dispatch jobs require `training-placement-v1`, advertised by
updated subscriber/conversion workers; LOCAL qualification does not prove managed
Ray/OSMO or distributed Alpha routing. Each destination advances only after its
ordinary relay child jobs finish. Use the run observer below for an exact
source-job run; `wait_reprocess` also includes the repair's required post-conversion
copies. Native copies relocate Lance payload/media references and return their
completed inventory locations in the copy job's `training_output_copy.inventories`;
they are not verbatim index copies. Native-copy workers also require
`training-native-copy-v1`. Omitted copy parallelism uses the relay's host-aware
single-worker sizing (12 per host core by default), subject to the existing
native-copy bound of 4,096; configured relay caps and smaller explicit requests
remain effective. Real local tests cover delayed original copies after a newer
conversion, selected FULL/REUSE_MEDIA/REPLACE repair copies, native destination
retry, and empty inventories after complete retirement. Required copies keep
the original run pending until they finish; a later run's copies remain separate.
These tests use generated media and do not qualify post-GC file retention,
remote credential expiry/locality, distributed execution or throughput.

### Observing conversion and required copies (preview)

```python
started = client.transforms.run(converter.conversion_id)
run = client.training.get_run(converter.pipeline_id, started.job_id)
finished = client.training.wait_run(run, timeout=300)
if finished.phase == "failed":
    failed_jobs = [job for job in finished.jobs if job.state == "failed"]
```

Creation status means the graph was installed, not that data was converted or
copied. `get_run` follows exact downstream jobs retained by ordinary completion
processing. It reports `running`, `placing`, `succeeded`, `failed`, or `incomplete`.
Required copies default to `required=True` on `PlacementDestination`. Set
`required=False` for an independent, nonblocking archive: pending/failed/absent
optional copies stay visible but do not delay required work. `job_id=None` with
`state="not_admitted"` means no exact submission was retained; it is never
reported as a completed copy. Missing or changed required evidence is incomplete.

The root matters: a conversion job covers its descendants, not an earlier
independent raw-file copy. Pass the source crawl job ID for that complete run;
the creation operation ID itself is not a source-job ID. An older run remains
bound to its recorded children even when watermarks advance. Waiting checks the
root incarnation and retained definition; it never starts/retries jobs, and a
timeout leaves server work running. `succeeded` means all required jobs in this
run completed, not all historical runs or every optional archive. This is job
completion evidence, not a substitute for validating media/content correctness.

One request reads the pipeline plus at most 256 exact, strong job lookups, memoized
across shared children. It never scans job history, catalog data or episode rows,
and writes nothing. Cost is O(recorded jobs + graph edges), independent of episode
cardinality. Missing/expired root metadata refuses; missing downstream metadata
is incomplete. A cleared watermark or terminal conversion alone is not success.

Internal generation-publication operations expose their original `prepared`
native inventory entries as an immutable IO-core value. `wait_publication_unit`
checks that these inputs remain unchanged; they are not storage credentials, and
only a committed operation establishes publication. They survive publication-job
and write-session expiry without looking up a newer inventory. Older compact
results can return `prepared=None`. New post-conversion jobs bind these original
inputs into their ordinary child context and completion, including selected
repairs. Missing inputs refuse instead of selecting today's files. This does not
preserve files against GC. Standing conversions keep their ordinary root/child
retention window; explicit repairs and their direct dependencies are retained.
Reprocess admission does not require the worker's plugin on the API server.
The first worker checks its loader and records the conversion implementation
fingerprint in the original job checkpoint; retries with different code refuse.
This does not claim complete transitive dependency pinning or qualify the
distributed plugin-upload/Ray-task environment.

Primary companion reconciliation may complete inline or through an existing
task-only, outbox, or checkpoint-recovery child. Status checks the retained
source lifecycle and target in either case. An inline completed task does not
require a synthetic child; a deferred child must actually complete before the
run can succeed. `get_reprocess` uses this same interpretation. Neither status
path re-reads companion-table data to make that decision.

A no-change tick reports success only when its source session explicitly selected
no new input; a zero processed count is insufficient. This outcome is preserved
through remote completion, still persists the source watermark, and creates no
new downstream copies. It does not certify an earlier run's unfinished copies.

To retry a failed destination on this same run, correct its input/access problem,
then use the observed snapshot:

```python
failed = client.training.get_run(converter.pipeline_id, started.job_id)
copy = next(job for job in failed.jobs if job.placement and job.state == "failed")
retried = client.training.retry_placement(failed, copy.job_id)
finished = client.training.wait_run(retried, timeout=300)
```

This uses the existing job requeue, keeping the same child ID/incarnation, input,
retry counters and original publication inputs. It does not rerun conversion
or successful destinations. Pipeline owners/managers/admins can request it; the
worker still needs ordinary storage permission. Keep the same `failed` snapshot
when repeating an ambiguous request: it can observe an already retried attempt,
but cannot rearm a later failed claim. Inspect a fresh snapshot before choosing
to retry that new failure. Optional failed copies may be retried even after
required work succeeds; `wait_run` does not wait for those optional copies.

Each explicit retry uses at most two bounded run reads and one strong target-row
read, plus the existing requeue/index/event writes. It scans no catalog data.
No-input completion of a recorded copy is incomplete, not evidence that its old
generation was copied. New post-conversion jobs read their saved, hash-checked
inventories through the original storage bindings rather than selecting through
today's watermark. Delayed copies and failed-destination retries after a newer
conversion are qualified locally. Explicit repair jobs survive automatic cleanup;
ordinary conversion jobs still expire. The API does not reconstruct deleted
root/child jobs or garbage-collected files; missing original inputs refuse.

Original-inventory reads stream through one bound input file and replay projected
native batches locally; they do not rescan the catalog for each batch. Decoder
chunk size does not increase relay job count. This supports larger copy indexes
without enlarging job JSON, but does not remove the repair-input limits above or
extend file retention. See the relay README for exact reader resource bounds.

For a repair with post-placement, `get_reprocess` reports `placing` until required
copies finish, or `incomplete` when their completion is missing. Inspect and retry
a failed destination through the same run API:

```python
repair_run = client.training.get_run(repair.pipeline_id, repair.operation_id)
failed_copy = next(job for job in repair_run.jobs if job.placement and job.state == "failed")
client.training.retry_placement(repair_run, failed_copy.job_id)
finished = client.training.wait_reprocess(repair)
```

Optional pending/failed copies do not delay the repair's success, and may still
be retried. New repairs retain the converter definition hash and refuse to
interpret a changed definition as the original run. Legacy post-placement
definitions using only a watermark need recreation with original-publication
inputs; they are not silently upgraded during repair admission.

### Creating converters

```python
from deltacat_client import Client, ConversionOutput

client = Client("http://localhost:8080", catalog="training")
converter = client.training.create_converter(
    dataset="my_dataset",
    source_format="xdof",
    seed_prefixes=["file:///datasets/my_dataset/"],
    staged_plugin_working_dir="file:///plugins/xdof/",
    output_formats=[ConversionOutput(format="packds", version="6")],
    trigger_now=False,
    dry_run=True,
    allow_duplicate=True,  # Skip server duplicate lookup for this offline preview.
)
```

The supported source adapters are XDOF, Dexmate, Gear YAM raw v1, Mecka, DROID,
ABC LeRobot, UMI aligned Zarr, and a bounded LeRobot v2 file set. DROID/ABC retain
their canonical-source and manual-schedule requirements. PackDS format 6 is
supported by every adapter; UMI, the generic LeRobot file-set input, DROID/ABC
LeRobot prebuilt facts, XDOF, Dexmate, Gear YAM raw v1 and Mecka sources additionally
support simultaneous LeRobot 2.1 output, or LeRobot-only conversion without
PackDS tables or payloads. Unknown versions,
unsupported combinations, and an explicit empty output list refuse instead of
selecting a different writer. `ConversionOutput` preserves omitted options so
source-specific defaults (such as embodiment) still apply. Output versions are
on-disk format versions, not converter builds or API versions.

Explicit `run_id`, `pipeline_id`, `crawler_id`, `diff_id`, `placement_id`, and
`conversion_id` select orchestration identities; `dataset` remains the source
label stored in output rows, not the run ID. `crawl_table`, `diff_table`,
`ext_table`, and `facts_table` override incremental-state table names. File-source
recipes with optional outputs retain descriptors in the separately declared
`<facts_table>_training_source` sink; their default PackDS-only recipe and the
UMI/file-set recipe use the requested facts table directly. Gear YAM's optional
annotation root also accepts `annotation_crawler_id`, `annotation_diff_id`,
`annotation_placement_id` and the corresponding annotation table overrides.
These are catalog identities, not filesystem destinations: choose arbitrary
authorized output paths through the output and placement options below.

The advanced `xdof_identity_activation_receipt` option retains its existing
production-only XDOF contract: an exact S3 URI, version ID and SHA-256 plus
worker-side receipt, runtime and table-incarnation authentication. Planning a
definition does not authenticate or activate that receipt. The existing
`n6_finalization_authority` option remains an exact, append-only 20-episode
canary contract; it is not a general repair mode and cannot be combined with
additional training outputs. The unified namespace does not widen either
option's authority or remove its runtime checks.

XDOF and Gear YAM optional outputs keep the ordinary crawl/diff/place/folded-B0
pipeline and require their source embodiment (`YAM` by default) on LeRobot
declarations. They share conversion, incremental preservation and publication
with the other source adapters; there is no second file-output executor.
Their workers require `training-file-output-sessions-v1`. Incomplete Gear YAM
episodes produce a distinct `no_convertible_episodes` job result, advance the
source watermark without publishing empty outputs, and become eligible when a
later crawl supplies the missing files. With `enable_yam_annotation_join=True`,
Gear YAM additionally supports delayed annotation arrivals and annotation-only
corrections both with default PackDS output and with explicit outputs. Optional
`annotation_seed_prefixes`,
`annotation_crawl_binding`, and `annotation_crawl_region` select the annotation
source; omitted values retain the canonical S3 defaults. A local source should
explicitly pass its `file://` seeds, authorized local binding and region `None`.
These jobs additionally require `training-annotation-outputs-v1` workers.
New default annotation converters retain raw input and annotations in the same
source snapshot used by multi-output jobs, without creating LeRobot outputs.
Existing stored legacy jobs are not upgraded in place. `node_dispatch_mode="local"`
runs B1 inline for both default and explicit outputs; cluster modes retain fanout.
If annotations arrive before raw facts, `inputs_not_ready` consumes no input
deltas and `get_run` reports `incomplete`; a later source tick can process them.
Sessions exceeding the current complete-session repair bounds remain unsupported.
Local generated-media qualification is not distributed Alpha proof.

Mecka optional outputs default to the same ordinary crawl/diff/place pipeline. Set
`embodiment` and the LeRobot output's `embodiment_tag` consistently (for example
`"mecka"`); omitting outputs preserves the established primary-only recipe.
Task captions and shard-local vendor IDs stay intact in the source descriptor,
separate from the dataset-level output session and globally scoped episode ID.
Local `file://` placement is supported, including escaped filenames. The legacy
development-only `mecka_loader_discovers=True` route also supports optional outputs
without adding diff/placement nodes. It captures discovered descriptors and their
source frame rate in the existing durable source snapshot; retry does not repeat
discovery or adopt changed source metadata. Remote discovery and subsequent source
reads require the catalog's scoped source credentials; local file paths are not
S3 credential scopes. Real GR00T/Sharpa retargeting and local multi-output readback
are tested with generated pose/video inputs; this is not vendor-fleet or Alpha
qualification.

For Mecka outputs created with recorded training alignment, selected
`reprocess(conversion_mode="reuse_media", write_mode="replace_partition", ...)`
rebuilds numeric features and annotations without reading or encoding source
video. Exact source-frame selection, static trim, camera membership and timing
must remain unchanged. Older outputs lacking that proof require full conversion;
the repair does not infer their alignment from current source data. The existing
Mecka evolution planner remains separate, and compatible client/server/worker
versions are still required. This does not change the default creation recipe.

Dexmate optional outputs also retain the ordinary crawl/diff/place pipeline.
Use consistent `embodiment`/`embodiment_tag` values and pass source behavior
through `loader_kwargs` (for example `apply_fk`, `apply_remove_static`, `fps`,
and camera sizing). A local source must explicitly set the loader's
`source_root` and `expected_binding_name` to its local discovery root/binding,
in addition to the pipeline's crawl settings. Acquired HDF5 features and video
decoding are shared across outputs; optional outputs do not rerun kinematics.
Source captions remain separate from output grouping. Missing or variable-width
optional producer fields retain their original primary behavior but cannot be
exported as complete fixed-width LeRobot features. Generated local HDF5/video
tests cover primary-only, combined, and output-only jobs; actual vendor inputs,
large source catalogs and distributed execution require separate qualification.

Explicit DROID and ABC LeRobot multi-output declarations use the same durable source snapshot
and publication runtime. Supply the original shared loader metadata once in
`loader_kwargs`, `crawl_schedule=None`, and the canonical embodiment
(`OXE_DROID_DREAM` for DROID, `YAM` for ABC) on the output declarations.
ABC's additional outputs preserve its six annotation modalities, human action
source and final-step marker using the same acquired source features and camera
decode pass as the primary output. Authorized `file://` seeds are
also accepted in this explicit mode. Omitted remote-source options retain the
canonical PDX root and binding. For an authorized physical copy, declare all four
source options together (also supported for default PackDS-only output):

```python
copy_root = "s3://my-training-input/isolated/droid"
source_options = dict(
    seed_prefixes=[copy_root + "/"],
    loader_kwargs={"source_root_uri": copy_root},  # plus shared loader metadata
    crawl_binding="my_copied_source",  # existing, explicitly authorized binding
    crawl_region="us-east-2",  # explicitly None for an endpoint-scoped binding
)
# Pass **source_options to training.plan_converter/create_converter, together
# with the canonical dataset, source_format, staged plugin and manual schedule.
```

Crawler and loader must name the same single concrete prefix. Partial redirects
refuse before provider access; even a binding/region equal to an old default is
an explicit choice, not a request to infer PDX authority. The binding still
supplies exact revision, scope, endpoint and credentials through SourceReadAccess.
A copied URI is not proof of dataset equivalence: source generations/hashes,
episode identity and source/frame correspondence remain the loader's contract.
For DROID native originals, use the plugin's separate `copy_uri` field while
preserving original `source_uri` provenance. Copying inputs does not permit a
new writer to overlap another pipeline's episode IDs; use isolated sinks for
qualification or an explicitly coordinated shared-sink repair.

Prebuilt input is left unchanged and captured into a separate facts
sink. Native source batching now supports sessions larger than a single Python
decode batch; 513-episode jobs are tested with and without primary output. This
does not qualify either entire canonical dataset or remove the separate
complete-session repair bounds. See the
[training runtime guide](../deltacat/compute/training/README.md) for qualification
and remaining sharding work.

`client.training.plan_converter(...)` takes the same options and returns an
immutable `ConverterDefinition` with `to_dict()`, `canonical_json`, and `sha256`.
It compiles the exact crawler/transform/subscription declarations without calling the server,
checking live duplicates, creating namespaces, uploading plugins, or starting
work. Source plugins must already be staged (the built-in UMI and LeRobot file-set
loaders need no plugin). Unlike the legacy console-only dry run, this plan uses live routing
validation and rejects `dry_run=True`. Server permission, collision, and runtime
admission are still required. Plans are bounded to 256 nodes and 1 MiB of JSON;
their hash identifies a declaration, not a successful conversion or authority.

`create_converter` uses service-owned creation transport for every source adapter.
For explicit asynchronous control, submit a compiled declaration and
retain it along with the returned operation ID:

```python
definition = client.training.plan_converter(
    dataset="my_umi_dataset",
    source_format="umi_aligned_zarr",
    seed_prefixes=["file:///datasets/my_umi_dataset/"],
    trigger_now=False,
)
operation = client.training.submit_converter(definition)
operation = client.training.wait_operation(operation)
assert operation.phase == "created", operation.message
```

This path uses `POST /v1/training/converters` and
`GET /v1/training/operations/{operation_id}`. The server retains the resolved
declaration, original principal, and exact initial job identities in its JobStore.
The existing executor owns installation and retry. After a transport error,
resubmit the **same definition**, not a newly compiled random pipeline ID.
Changed input under the same pipeline ID refuses. Status is owner/admin scoped;
it does not return credential envelopes or retained node configs.

`created` means the complete graph is ready and requested initial jobs have been
admitted. It does **not** mean conversion or placement has completed; inspect
`initial_job_ids` and pipeline progress separately. A wait timeout does not cancel
server work. Persistent restart requires a persistent JobStore/SystemTables
backend; memory mode is process-local. If an operation's job retention expires,
creation currently refuses to recreate an existing pipeline rather than repeating
its initial jobs. Long-lived operation receipt/retention remains an acceptance
item for the final unified API.

The unified convenience creator retains the source-specific result dataclasses.
It waits for graph installation and initial job admission; local dispatch also waits for
the original initial crawl jobs and returns their actual `triggered_nodes`.
Remote dispatch does not wait for conversion. Downstream completion repair may
be deferred even when a local crawl completes; an empty `triggered_nodes` list
does not mean no downstream work is pending. Plugin staging remains a client
prelude, scoped to the selected catalog. Namespace and graph mutations belong
to the service. PackDS's duplicate-source preflight and explicit `force` /
`allow_duplicate` overrides are preserved (the preflight is not an atomic
cross-pipeline uniqueness constraint).

If a convenience creator raises `ConverterCreationError`, its `definition`
preserves generated IDs and staged plugin references; `operation` is the last
acknowledged operation, if any. Inspect `__cause__` and replay the same definition
with `submit_converter`. A timeout or missing acknowledgement is not cancellation.
Do not log the definition indiscriminately: it can contain private configuration.
These APIs require a server with the training creation endpoints; there is no
silent fallback to partial client-driven graph creation on an older server.

For UMI dual output, add
`ConversionOutput(format="lerobot", version="2.1", options={"mode": "pretraining",
"output_root_name": "training-output"})` alongside PackDS. The existing
`PackDSOutput` and `LeRobotV2Output` convenience classes remain supported and
retain their historical UMI defaults.

PackDS can also retain a generation inventory of its original native payloads
and companion rows: set `output_root_name` on `PackDSOutput`, or in the generic
PackDS output's `options`. This names a registered index root, independently of
`steps_data_root` and `object_store_root`; omission preserves the default pipeline.
It uses the same publication and repair jobs as the other outputs, without a
second conversion. Its original index references original payload/media files,
which must remain readable until optional post copies finish. Add that declared
root to `post_conversion_placements` for independent native destination datasets;
the existing copy jobs relocate media/Lance references and retain fresh inventories.
Native-only/mixed local two-destination readback and binding-change retry are
qualified; the remaining native repair/lifetime matrix and Alpha are not.

For UMI or generic LeRobot-only conversion, omit the PackDS declaration and
specify only the LeRobot output. Up to 16 distinct LeRobot destinations can share one decode;
each can select a registered output root and its own `layout_prefix_template`.
Their source eligibility mode and embodiment must match. A single destination
keeps the existing manifest-table name; multiple destinations get distinct,
configuration-derived inventory tables. A PackDS leg can also accompany multiple
LeRobot destinations (at most 16 total outputs). Newly created pipelines use
immutable generation inventories for every LeRobot destination, including one
LeRobot output alongside PackDS. Read the committed manifest inventory to locate
the current `.generations/<id>` files; do not infer the active generation from a
directory listing. This makes selected repair independent of output count.
Primary-only defaults are unchanged. Previously created single-LeRobot-plus-PackDS
pipelines retain their legacy finalizer/layout; upgrading the client does not
convert their existing files into a generation inventory, and selected repair
must not infer a predecessor inventory from those files.

Output-only jobs retain their original source snapshot, complete session units,
implementation identity and destination preconditions before conversion. A retry
reuses admitted publication jobs, and a later crawl replaces only affected
session inventories while retaining old generation files. Native source sessions
use <=512-row/4-MiB descriptor batches inside a <=1,048,576-row/256-MiB native
session, with an independent 64-MiB ID index. The shared streamed assembler
retains one complete session/publication rather than inventing a job per batch.
Selected repair and source-specific discovery have separate bounds.
Local HTTP/media tests cover this path; distributed Alpha qualification and
standing-conversion recovery after job-retention expiry remain open.

Primary jobs retain the same output plans and recover committed B2 receipts
without re-encoding the primary data. Their source result may delegate rich
companion-table materialization to an explicit recovery child; wait for that child
before treating all tables as repaired. LOCAL partial-B2 continuation still
refuses before B1 and requires the existing distributed bounded-wave coordinator.
This local qualification does not prove distributed partial recovery.

The new namespace defaults to Beta, a midnight UTC crawl, and an immediate first
crawl. When migrating historical UMI declarations, explicitly preserve
`env="prod"`, `schedule_timezone="America/Los_Angeles"`, and `trigger_now=False`
when those are the intended settings. Catalog selection follows the root client
unless explicitly overridden.

For a generic LeRobot v2 input, use `source_format="lerobot_v2"` with one exact
dataset-root seed and explicit output declarations. Each output must specify
its robot identity (`embodiment` for PackDS, `embodiment_tag` for LeRobot), rather
than inheriting UMI labels. An omitted output rate preserves the source rate.
This shares the UMI crawl/diff/place/conversion machinery; it does not change the
separate canonical DROID/ABC recipes. The initial file-set input requires
`info.json`, `modality.json`, `tasks.jsonl` and `episodes.jsonl` under `meta/`.
Native file-set discovery uses <=256-MiB inventory/index/facts buffers and
<=1,048,576 JSONL records; individual metadata objects remain <=16 MiB. A shared
task vocabulary is limited to 512 entries. Only <=512-row/4-MiB metadata batches
are decoded, and one episode's file facts remain <=32,768 files/4 MiB. See the
training README for the distinct source/session and repair bounds. Distributed
qualification remains open.

This surface shares the reprocessing and optional pre/post fanout machinery
described above. Remaining source-option, lifetime, scale and distributed
qualification is tracked separately; local examples are not a production release
claim. The former general module-level creation APIs have been removed; all
source recipes are reached through the training namespace.

Internal conversion coordinators can use `training.submit_publication_unit` and
`training.get_publication_unit` to hand an already prepared generation commit to
the server. See the [generation publication contract](../deltacat/compute/training/README.md#generation-publication-internal).
These bounded metadata units are not a second user-facing repair workflow:
source selection, encoding and parent completion still belong to the converter
or repair job. Successful completion has a bounded immutable record in the
existing JobStore, so ordinary job/session expiry does not trigger reconversion.
Persistent recovery requires SystemTables/DynamoDB; process-local mode still
loses its state on restart. This does not extend media retention or reconstruct
the parent's original source selection.


## Getting Started

Before using the client, you need a running DeltaCAT API server. See
[Server Setup](docs/server-setup.md) for instructions.

DeltaCAT lets you manage **Tables** across one or more **Catalogs**. A
**Table** is a named collection of data files. A **Catalog** is a named
data lake that contains tables. For the full data model, see the
[DeltaCAT README](../README.md#getting-started).

### Quick Start

```python
from deltacat_client import Client
import pyarrow as pa

# Connect to a DeltaCAT server
client = Client("http://localhost:8080")

# Write data to a table. The table is created automatically; by default
# the client uses mode="auto" and lets the server infer the file format.
data = pa.table({
    "id": [1, 2, 3],
    "name": ["Cheshire", "Dinah", "Felix"],
    "age": [3, 7, 5],
})
client.catalog.write(data, table="cool_cats")

# Read the data back
df = client.catalog.read(table="cool_cats", read_as="pandas")
print(df)
```

### Core Concepts

Expand the sections below to see examples of core client operations.

Connection recovery is separate from request replay. After a server 5xx or
transport failure, the synchronous client discards its failed connection pool
so an explicit recovery request uses a fresh connection. The original failure
still reaches the caller; this does not make an unkeyed mutation automatically
retryable. Successful requests and ordinary 4xx refusals retain their pool.
Existing bounded read and keyed-write retry policies are unchanged.

<details>

<summary><b>Authentication</b></summary>

When connecting to a production server with auth enabled, provide a bearer token:

```python
from deltacat_client import Client

client = Client(
    "https://deltacat-api.example.com",
    bearer_token="your-api-token",
)
```

Operators running catalog/controller operations against one immutable server
artifact can opt in to an exact response fence. When enabled, each response
on the shared REST request path must carry the exact
`X-DeltaCAT-Server-Build-Id` value, including redirects, error responses, and
retries; missing or different values fail without retrying the identity error.

```python
client = Client(
    "https://deltacat-api.example.com",
    bearer_token="your-api-token",
    expected_server_build_id="2.1.0-immutable-build-id",
)
```

The option defaults to `None`, preserving the existing client behavior.

Admins can also onboard a new user, grant their initial access, and receive a
one-time token to share over a secure out-of-band channel:

```python
from deltacat_client import Client

admin = Client(
    "https://deltacat-api.example.com",
    bearer_token="admin-api-token",
)

created = admin.auth.create_user(
    user_id="newadmin@example.com",
    display_name="New Admin",
    email="newadmin@example.com",
    initial_role="ADMIN",
    resource_type="catalog",
    resource_name="*",
    issue_token_label="bootstrap",
    idempotency_key="create-newadmin-001",
)

bootstrap_token = created.token.token
print("Share this token once via a secure secret channel:", bootstrap_token)
```

The server validates your token and maps it to a user identity for permission checks (read, write, admin). The client automatically includes the token on every request. For data file access, the server vends short-lived STS credentials that the client uses to read from and write to cloud storage directly.

See [Configuration](docs/configuration.md) for the detailed auth model, bootstrap tokens, and role-based access control.

</details>

<details>

<summary><b>Reading Data</b></summary>

```python
# Read a table as a PyArrow table (the default format)
arrow_table = client.catalog.read(namespace="robotics", table="episodes")

# Read as Pandas, Polars, or Daft
df = client.catalog.read(namespace="robotics", table="episodes", read_as="pandas")
polars_result = client.catalog.read(namespace="robotics", table="episodes", read_as="polars")
daft_df = client.catalog.read(namespace="robotics", table="episodes", read_as="daft")

# Filter rows and limit results
df = client.catalog.read(
    namespace="robotics",
    table="episodes",
    read_as="pandas",
    filter_predicate={"eq": ["task", "pick_screwdriver"]},
    limit=5000,
)

# Time travel: read the table as it existed at a prior point in time.
# The as_of value is a nanosecond-precision Unix epoch timestamp.
df = client.catalog.read(
    namespace="robotics",
    table="episodes",
    read_as="pandas",
    as_of=1712697600000000000,
)
```

See [Reading and Writing](docs/reading-and-writing.md) for scalable reads,
lazy vs. eager materialization, and format-specific behavior.
For PackDS-backed training tables, see [Training Data](docs/training-data.md).

</details>

<details>

<summary><b>Computed Fields</b></summary>

Computed fields derive values while a PyArrow read is materialized. They may be
supplied for one request or persisted through schema evolution so later reads
resolve them automatically. Physical inputs can stay hidden from the result:

```python
import pyarrow.compute as pc

from deltacat_client import Client, ComputedField

client = Client("http://localhost:8080")
action = pc.field("action:action")

steps = client.catalog.read(
    namespace="robotics",
    table="training_steps",
    read_as="pyarrow",
    columns=["episode_id", "step_index"],
    computed_fields=[
        ComputedField(
            name="action_42",
            expression=pc.list_element(action, 42),
            source_fields=("action:action",),
        ),
    ],
)

# episode_id, step_index, and the extracted float32 action_42 are returned;
# the fixed_size_list action:action input is not.
print(steps.schema)
```

See [Computed Fields](docs/computed-fields.md) for the PackDS range-extraction
schema-evolution example, expression composition, projection and plan-reuse
semantics, validation, performance, and current limits.

</details>

<details>

<summary><b>Describing a table, with or without size counters</b></summary>

`describe_table` reads only the table, table-version and stream metadata by
default, whatever the table's size. The partition-derived counters
(`total_records`, `total_bytes`, `total_files`, `delta_count`,
`partition_count`) and `event_time_watermark` are `None`, and
`stats_omitted["reason"]` is `"not_requested"`:

```python
description = client.catalog.describe_table(
    namespace="iris",
    table="packds_steps",
    table_version="1",
)
schema = description.schema
```

Pass `include_stats=True` to fold the counters from the partition-stats
system table. This still enumerates no partitions, deltas or manifests and is
bounded; `description.stats_source` is `"partition_stats"`. The stats rows are
maintained by a post-commit background task, so these counters are eventually
consistent with commits (sub-second lag in practice), are not pinned to the
request's transaction, and can undercount when a partition has no stats row
(only the no-rows case is detected). When the stats cannot be trusted the
counters stay `None` and `stats_omitted["reason"]` says why
(`"partition_stats_unavailable"`, `"partition_stats_dirty"`,
`"too_many_partitions"`). Code that fences on the counters must treat an
omission as unknown rather than as zero:

```python
description = client.catalog.describe_table(
    namespace="iris", table="packds_steps", include_stats=True
)
if description.stats_omitted is not None:
    raise RuntimeError(f"no trusted counters: {description.stats_omitted}")
frontier = (description.stream_id, description.delta_count, description.total_records)
```

Pass `exact_stats=True` for exact, snapshot-pinned counters
(`stats_source == "exact"`). That reads one metafile per partition and per
delta, is unbounded, and is the only mode that derives `event_time_watermark`;
use it for before/after fences on small sealed tables or for batch callers
prepared to wait. A failed walk reports `stats_omitted["reason"] ==
"exact_stats_failed"`. `metadata_only=True` remains accepted as an explicit
spelling of the default.

</details>

<details>

<summary><b>Inspecting partition statistics (administrator diagnostics)</b></summary>

`query_partition_stats` requires catalog `ADMIN` and is a best-effort
diagnostic surface for investigating hot, dirty, or growing partitions. It is
not a reader-facing partition-discovery API and must not be used to construct
training-read membership.

```python
# Administrative diagnostics only.
token = None
while True:
    res = client.catalog.query_partition_stats(
        namespace="xdof",
        table="my_steps",
        limit=100,
        after_token=token,
    )
    for row in res.rows:
        print(row)
    token = res.next_token
    if not token:
        break
```

PackDS `_steps` data reads must instead use `plan`/`read` with a bounded
`episode_id` partition predicate derived from the application's authenticated
membership or source frontier. Do not infer a complete read set from this
best-effort statistics page.

</details>

<details>

<summary><b>Fast metadata-only planning (counts, sizes, existence)</b></summary>

An ordinary unfiltered current-snapshot
`plan(..., include_files=False)` request without a manifest-identity or
partition-frontier attestation returns a **metadata-only** plan: the server
skips manifest hydration and file enumeration entirely, answering from catalog
partition statistics instead of resolving per-file URIs. It is the right call
whenever you do **not** intend to read row data — existence checks, row/byte
counts, content-type discovery, and (for PackDS) episode counts. The fast path
proves one complete table-scoped statistics page and one complete bounded
partition page, requires the runtime-obligation scope to be empty, and matches
the same exact `(stream_id, partition_id)` cohort. Every accepted row must be a
full baseline built in one catalog snapshot no later than the reader's MVCC
snapshot; delta-updated, future-snapshot, legacy, or unversioned rows are
refused. Runtime baseline repair examines at most 1,000 visible deltas per
partition; a larger or indeterminate history remains refused. It never turns
an oversized or indeterminate metadata request into an unbounded
partition/file enumeration. Explicit attestation modes use separate contracts
and may resolve the complete snapshot; only a supplied
partition-frontier operation makes that acquisition checkpointed.

```python
# Metadata-only full-table stats, no per-file resolution. PackDS requires the
# explicit full-scan override even when include_files=False.
plan = client.catalog.plan(
    namespace="robotics",
    table="steps",
    include_files=False,
    allow_full_packds_scan=True,
)

has_data   = plan.total_records > 0          # existence check
row_count  = plan.total_records              # summed from partition stats
byte_size  = plan.total_bytes
fmt        = plan.content_type
episodes   = plan.episode_count              # PackDS partition (episode) count
assert plan.total_files == 0 and not plan.scan_tasks   # no files were enumerated
```

A metadata-only plan carries **no `scan_tasks` and no vended credentials**, so it
cannot be used to read data. When you actually need to materialize rows, request
files — `include_files=True` resolves the manifest into `scan_tasks` (file URIs,
per-file stats, and short-lived credentials) so a downstream `read(...)` /
`pl.scan_parquet(...)` has something to scan:

```python
# Data read: resolve files (scope PackDS reads to a partition / episode_id).
plan  = client.catalog.plan(
    namespace="robotics", table="steps",
    filter_predicate={"eq": ["episode_id", "ep_1"]},
    include_files=True,
)
steps = client.catalog.read(namespace="robotics", table="steps", plan=plan, read_as="polars")
```

A file-enumerating request can legitimately resolve no files, for example after
removing every selected episode. The SDK retains `include_files=True` in the
plan's serializable request context, so `read` returns a schema-aligned empty
result instead of treating it as metadata-only. Explicit metadata-only requests
still cannot be executed as data reads. Legacy plans with snapshot/layout
metadata but no recorded file-enumeration intent retain that refusal.

Long-running training brokers should use the bounded preparation facade. It
deduplicates physical planning while preserving the exact logical episode order
and multiplicity, waits for one immutable snapshot, renews server authority,
and returns an executable authorized plan:

```python
import uuid

prepared = client.catalog.prepare_packds_episode_plan(
    namespace="robotics",
    table="steps",
    operation_id=f"prepare-{uuid.uuid4().hex}",
    authority_operation_id=f"authority-{uuid.uuid4().hex}",
    episode_ids=["ep_2", "ep_1", "ep_2"],
    as_of=1712697600000000000,
    columns=["episode_id", "step_index", "state", "action"],
    timeout_seconds=300,
    run_id="n2-training-20260712-001",
)
authorized_plan = prepared.composition.plan
assert prepared.composition.requested_episode_ids == ("ep_2", "ep_1", "ep_2")
assert prepared.composition.composed_episode_ids == ("ep_2", "ep_1")
rows = client.catalog.read(
    plan=authorized_plan,
    read_as="pyarrow",
    plan_authority_verification_keys=prepared.plan_authority_verification_keys,
)
```

Trusted brokers can prepare overlapping consumers with one structural union
while retaining exact authority boundaries:

```python
prepared_by_rank = client.catalog.prepare_packds_episode_plans(
    namespace="robotics",
    table="steps",
    operation_id=f"prepare-bundle-{uuid.uuid4().hex}",
    authority_operation_id=f"authority-bundle-{uuid.uuid4().hex}",
    episode_id_groups={
        "rank-0000": ["ep_2", "ep_1", "ep_2"],
        "rank-0001": ["ep_1", "ep_3"],
    },
    as_of=1712697600000000000,
    columns=["episode_id", "step_index", "state", "action"],
    timeout_seconds=300,
)
rank_zero = prepared_by_rank["rank-0000"]
assert rank_zero.composition.requested_episode_ids == (
    "ep_2",
    "ep_1",
    "ep_2",
)
```

Each named group receives an independently signed exact authority with a common
bundle expiration. The helper never slices or reuses a covering union authority;
the entire grouped refresh fails closed if any child cannot be vended, signed,
persisted, replayed, or reauthorized.

The caller supplies stable operation IDs so transport retries remain
idempotent. `run_id` is an optional, non-secret provenance identifier sent on
every preparation, poll, authority, and verification-key request. Keep it
stable across retries for coherent telemetry. It does not alter operation,
cache, preparation, or authority identity. Launchers can validate their final
Gear ID before submission with
`deltacat_client.validate_packds_fragment_run_id(...)`. Pass a
`threading.Event` as `cancellation_event` for bounded
cooperative cancellation. The low-level fragment facade raises
`UnsupportedFeatureError` explicitly for old servers; the bounded high-level
facade raises `PackDSEpisodePlanPreparationUnsupported` without retrying it, so
callers can opt into a narrow legacy fallback without treating transient HTTP
failures as unsupported servers.
Preparation telemetry contains fixed-cardinality timing and retry counts and
never embeds episode, fragment, credential, or token values.

The lower-level split/compose primitives remain available for trusted brokers
that manage fragment persistence and authority refresh themselves:

```python
from deltacat_client import (
    compose_packds_episode_plan,
    split_packds_episode_plan,
)

plan = client.catalog.plan(
    namespace="robotics",
    table="steps",
    as_of=1712697600000000000,
    filter_predicate={"in": ["episode_id", ["ep_1", "ep_2"]]},
    include_files=True,
)
fragments = split_packds_episode_plan(plan)
cache.update({episode_id: fragment.to_dict() for episode_id, fragment in fragments.items()})
combined = compose_packds_episode_plan([fragments["ep_2"], fragments["ep_1"]])
assert combined.credentials is None
```

Fragment composition is intentionally limited to complete, non-MOR,
nontransactional PackDS episode plans with an explicit `as_of`. The hashes detect
cache corruption; they are not server signatures. The credential-free composed
plan is deliberately non-executable. A trusted broker must obtain a fresh
Ed25519-signed server authority covering the exact fragment fingerprints,
snapshot, roots, bindings, and credentials, then use
`compose_authorized_packds_episode_plan(...)` with a trusted server verification
key before executing it.

The metadata-only fast path fails closed on stale, dirty, incomplete, malformed,
or over-ceiling partition statistics. It raises a bounded-planning error instead
of returning a stale count or silently falling back to unknown-cardinality
partition/delta enumeration. Allow bounded runtime maintenance to rebuild the
statistics, or narrow the request before retrying.

</details>

<details>

<summary><b>Reading PackDS video blobs</b></summary>

A PackDS `_steps` table stores each camera — and each resolution — as its own
`video_blob:<key>` column. Frames are packed into 16-frame H.264 groups whose
*leader* row carries the byte-range reference (the other 15 rows are empty); a
companion `video_blob_frame_index:<key>` column records each row's position
(0..15) within its group. Read blobs through the client: the server vends
short-lived, scoped credentials on demand, so you never need direct S3 (or
S3-Express) access.

`step_index` (a real `_steps` column) is the **canonical selection coordinate** —
group-leader recovery is value-based on it, so selection stays correct under
filtering/reordering. Select by `step_index=` (canonical) or `where=` (a deltacat
filter_predicate dict applied client-side).

```python
# 1. Scope to one episode partition (see "Listing / paginating partitions" for how
#    to discover a partition_id). read(read_as="pyarrow") attaches the scoped
#    storage-authority lease to the returned table.
plan = client.catalog.plan(namespace="xdof", table="my_steps", partition_filter=[partition_id])
steps = client.catalog.read(namespace="xdof", table="my_steps", plan=plan, read_as="pyarrow")

# 2. Discover which camera/resolution blobs the table carries.
#    Bare key = native resolution; the "_320_240" suffix = the downscaled copy.
keys = client.catalog.packds_video_blob_keys(steps)
# ['video_blob:left_camera-images-rgb', 'video_blob:left_camera-images-rgb_320_240', ...]

# 3. Read a SPECIFIC blob for a step by step_index (canonical) or a predicate.
#    The 16-frame group leader is auto-resolved, so any step works. Returns the
#    raw H.264 group slice bytes (no `av` needed):
native = client.catalog.read_packds_video_blob(
    steps, video_blob_key="left_camera-images-rgb", step_index=42, plan=plan)
by_anno = client.catalog.read_packds_video_blob(
    steps, video_blob_key="left_camera-images-rgb",
    where={"eq": ["annotation", "grasp"]}, plan=plan)

# 4. Batch-resolve several cameras across several steps in ONE credential vend
#    -> dict keyed by (row, column):
blobs = client.catalog.resolve_packds_video_blobs(
    steps, video_blob_keys=keys, step_indices=[10, 26, 42], plan=plan)

# 5. The pyarrow accessor reuses ONE plan + fs cache across calls and recovers
#    the lease from the table automatically:
vb = client.catalog.video_blobs(steps)
data   = vb.blob("left_camera-images-rgb", step_index=42)             # bytes
frames = vb.frames("left_camera-images-rgb", step_index=42)           # list[np.ndarray]
many   = vb.resolve(keys, step_indices=[10, 26])                      # dict

# Credential-free, process-local counters for attribution/benchmarking. Bytes
# count successful returned range payloads; wait covers the exact authority-
# created filesystem read. Reset only at a quiescent run boundary.
metrics = client.catalog.packds_video_read_metrics(reset=False)
```

#### High-level episode readers

These do their own episode-scoped read, so the group leader is always present —
you do not pre-read the steps table:

```python
# A single frame by ABSOLUTE frame number (frame is offset via min(step_index);
# fail-closed if outside the episode's step range). Decodes -> np.ndarray:
frame = client.catalog.read_packds_video(
    namespace="xdof", table="my_steps", episode_id=ep_id,
    video_blob_key="top_camera-images-rgb", frame=128, decode=True)

# Or by canonical step_index / predicate (returns the 16-frame group bytes):
group = client.catalog.read_packds_video(
    namespace="xdof", table="my_steps", episode_id=ep_id,
    video_blob_key="top_camera-images-rgb", step_index=128)

# Walk the WHOLE episode's video IN ORDER (group leaders only, sorted by
# step_index) to reassemble it — yields each 16-frame group:
for group_bytes in client.catalog.iter_packds_episode_video(
        namespace="xdof", table="my_steps", episode_id=ep_id,
        video_blob_key="top_camera-images-rgb"):
    ...  # concatenate / write out in order
```

#### Polars: resolve a column lazily

With polars installed, a `packds` namespace resolves the group leader for **every
current row** lazily (only the rows you `.collect()` are resolved). It force-keeps
the `video_blob_frame_index:<key>` companion and fails closed if it was projected
away:

```python
df = client.catalog.read(namespace="xdof", table="my_steps", plan=plan, read_as="polars")
resolved = (
    df.lazy()
      .packds.with_resolved_video("top_camera-images-rgb", out_column="rgb")
      .collect()
)  # adds an `rgb` column of raw group-slice bytes
```

Pass `plan=` (the scoped read plan) to reuse the credentials it already vended; or
omit it and pass `namespace=`/`table=` and the client vends a plan for you. The same
`where=` filter_predicate dict has identical semantics server-side (in `plan(...)`,
prunes files) and client-side (against a materialized table, scans rows). Decoding
needs `pip install "deltacat-client[packds]"` (pulls `av`); the raw-bytes path does not.

</details>

<a id="copying-source-trees-fanout"></a>

<details>

<summary><b>Incrementally copying source trees with create_placer</b></summary>

`create_placer(...)` declares one or more autonomous crawl → diff → relay
pipelines, as shown in the [quick example](#incremental-exact-path-copies).
Each source has a stable name and an independent crawler/diff state;
its destination list may contain one target or a fanout. Each
`ExactPathRelaySubscriber` consumes the shared diff and copies only `ADD` and
`UPDATE` objects, byte-for-byte, to the same relative key below each destination
root. An unchanged crawl
submits no relay jobs. Source deletes intentionally do not delete destination
objects.

```python
from deltacat_client import (
    Client,
    PlacementDestination,
    PlacementSource,
    create_placer,
)

client = Client("https://deltacat-api.example.com", bearer_token="...")
result = create_placer(
    client,
    sources=[
        PlacementSource(
            name="india",
            uri="s3://gear-umi-india/",
            binding="aws-account-discovery",
            region="us-east-2",
        ),
        PlacementSource(
            name="philippines",
            uri="s3://gear-umi-philippines/",
            binding="aws-account-discovery",
            region="us-east-2",
        ),
    ],
    destinations={
        "india": [
            PlacementDestination(name="primary", binding="pdx-india"),
            PlacementDestination(name="audit", binding="archive"),
        ],
        "philippines": [
            PlacementDestination(name="primary", binding="pdx-philippines")
        ],
    },
    schedule_cron="0 6 * * *",
    schedule_timezone="UTC",
    name="umi-vendor",
    run_id="prod-nightly-v1",
    include_version_changes=True,
)
```

An S3 source requires a discovery-root binding. Each destination binding must
be a DataRoot authorized for `relay_target` (legacy `managed_write` and
`placement_target` roots remain compatible); omitting the destination `uri`
uses that DataRoot's root. For split control/worker networks, pass
`worker_server_url=` so subscriber workers call an internal endpoint while the
client submits through an external endpoint. `dry_run=True` returns the
deterministic node IDs without changing the server.

The crawl tables, immutable diff deltas, per-subscription watermarks, and
durable child relay jobs provide the audit trail. This is an exact-object-copy
API, not DeltaCAT placement: it does not introduce `__external__` paths,
convert content, or make the destination a managed table root.

</details>

<details>

<summary><b>Converting a dataset to PackDS</b></summary>

Declare an autonomous crawl → diff → placement → conversion pipeline with one
thin-client call. By default each source converges into a shared corpus in the
auto-created `iris` namespace:

- **`iris.packds_steps`** — the unified steps table, identity-partitioned by `episode_id`;
  conversions contribute episode-scoped ADD/REPLACE partitions here.
- **`iris.packds_episodes`** — the rich application episodes table (`source_format`, `embodiment`,
  length, duration, video-codec metadata, …).
- **`iris.<dataset>_crawl` / `_diff` / `_placement` / `_facts`** — per-`--dataset`
  incremental state: one set per source dataset, reused across that source's runs, so
  re-running only converts what changed.

`client.training.create_converter` supports these source adapters:
XDOF (`source_format="xdof"` or omitted), Dexmate (`"dexmate"`), YAM in-house
(`"gear_yam_raw_v1"`), Mecka (`"mecka"`), canonical DROID (`"droid"`), canonical
ABC LeRobot (`"abc_lerobot"`), UMI (`"umi_aligned_zarr"`), and a generic LeRobot v2
file set (`"lerobot_v2"`). DROID and ABC require
their canonical dataset/source and a manual crawl. Unknown selectors fail closed.
`plugin_source` stages the implementation expected by the selected strategy; it
does not dynamically register an arbitrary new `source_format`.

```python
from deltacat_client import Client

client = Client(server_url="https://deltacat-api.example.com", bearer_token="...")
result = client.training.create_converter(
    dataset="my_xdof_source",
    seed_prefixes=["s3://my-source-bucket/episodes/"],
    plugin_source="./xdof_packds_loader.py",
    source_format="xdof",
    crawl_schedule=None,  # validate manual runs before enabling a schedule
)
print(result.pipeline_id, result.crawler_id)
```

For these file sources, the method returns a `PackDSConverterResult` with node IDs
(`crawler_id` / `diff_id` / `placement_id` / `conversion_id` / `pipeline_id`).
`crawl_schedule` is persisted on the crawler as a cron trigger (see Scheduled
Processing below); the first run still fires immediately by default
(`trigger_now=True`). Pass `dry_run=True` to build the declaration without server
mutation. In prod, pass `conversion_cluster_affinity="xlarge"` for multi-TB or
million+ episode backfills so they can only be claimed by the xlarge conversion
cluster. The convert-only operator CLI (`python -m deltacat.compute.training.submit_x_to_packds`,
for an already-crawled facts table) remains available for that narrower case.

For the complete strategy table, data model, B0/B1/B2 stages, annotation
contract, scheduling policy, and retry behavior, see
[PackDS converter pipelines](../deltacat/docs/packds/CONVERTER_PIPELINES.md).

**Multiple converters, one unified corpus (multi-writer).** Several training converter
pipelines may feed the SAME unified `iris.<out_table>_steps` / `_episodes` tables — this is
the intended end state for consolidating many sources into one queryable corpus. Each
converter flags its shared `_steps` / `_episodes` sinks as co-writable, and the server only
relaxes its single-writer guard when EVERY writer of a table is a mutually-flagged packds
converter (a packds converter never silently collides with an unrelated general transform).
The per-dataset `_facts` table stays single-writer.

**Duplicate rejection.** Before graph mutation, the PackDS adapter reads existing
server declarations and rejects creating a
TRUE duplicate — a pipeline that feeds the SAME `_steps` table from the SAME crawl source
(same seed-prefix set + crawl binding) — and raises `PackDSDuplicatePipelineError` naming the
offending pipeline / conversion / crawler ids. DISTINCT sources feeding the same unified table
is allowed (that is the multi-writer case above). Pass `force=True` (a.k.a.
`allow_duplicate=True`) to create it anyway; under `dry_run=True` the finding is printed
instead of raised.

**Redrive is disabled by default.** Because a redrive does a WHOLE-TABLE rollback per declared
sink, redriving one converter would reverse EVERY co-writer's contributions to the shared
`_steps` / `_episodes` tables. The conversion transform is therefore created with redrive
disabled: the server rejects `client.transforms.redrive(conversion_id)` (409
`transform_redrive_disabled`) and `client.pipelines.redrive(pipeline_id)` (409
`pipeline_redrive_disabled`). To replace episodes, drive the change through the SOURCE
instead — see [PackDS Episode Rewrites](#) below: update the crawled source so the diff emits
the affected `episode_id`s as REPLACE (per-episode `steps_write_mode=REPLACE`), or
delete/deprecate the offending episode-ID partitions. (Proper shared-table redrive is tracked
separately.)

**Run an existing converter on demand.** Re-crawl + reconvert the latest source files
whenever you want by running the converter's root crawler publication. Use the
`crawler_id` from the `PackDSConverterResult` (the cascade fans out from this single
trigger — diff → place → convert drain automatically):

```python
run_result = client.publications.run(result.crawler_id)
print(run_result.triggered_nodes)  # the downstream nodes the crawl completion fanned out to
```

**Update the configured crawl schedule.** Change the cron cadence of an existing
converter's crawler with a new `schedule` trigger; the server persists the new
`schedule_cron` and the schedule scanner picks up the next fire on the new cadence. A
manual `publications.run(...)` still works alongside the schedule:

```python
from deltacat_client import PublicationTrigger

# Re-crawl every 6 hours instead of nightly (pass cron=None / a different trigger to clear).
client.publications.update(
    result.crawler_id,
    trigger=PublicationTrigger.schedule(cron="0 */6 * * *", timezone="UTC"),
)
```

</details>

<details>

<summary><b>Writing Data</b></summary>

The client supports writing PyArrow tables, Pandas DataFrames, Polars DataFrames, Daft DataFrames, NumPy arrays, Ray Datasets, and local files (Parquet, CSV, TSV, PSV, Feather, JSON, ORC, AVRO, Lance).

```python
import pyarrow as pa

# Write data to a table. The defaults are mode="auto" (creates the
# table on first write, appends thereafter) and the client default Parquet format.
data = pa.table({"episode_id": [1, 2], "score": [0.95, 0.87]})
client.catalog.write(data, namespace="robotics", table="predictions")

# Append more data — same call, same defaults.
data2 = pa.table({"episode_id": [3, 4], "score": [0.91, 0.89]})
client.catalog.write(data2, namespace="robotics", table="predictions")

# Write a Pandas DataFrame
import pandas as pd
df = pd.DataFrame({"episode_id": [5], "score": [0.93]})
client.catalog.write(df, namespace="robotics", table="predictions")

# Override the defaults when you need a specific mode or format. Here we
# create a Lance-formatted table with explicit schema, partitioning, and
# sort order.
from deltacat_io_core.schemes import (
    PartitionKey,
    PartitionScheme,
    PartitionTransform,
    SortKey,
    SortOrder,
    SortScheme,
)

client.catalog.create_table(
    namespace="robotics",
    table="scored_episodes",
    schema=pa.schema([
        pa.field("episode_id", pa.int64()),
        pa.field("score", pa.float64()),
        pa.field("episode_day", pa.string()),
    ]),
    partition_scheme=PartitionScheme.of([
        PartitionKey(key=["episode_day"], transform=PartitionTransform.IDENTITY),
    ]),
    sort_scheme=SortScheme.of([
        SortKey(key=["episode_id"], sort_order=SortOrder.DESCENDING),
    ]),
    auto_create_namespace=True,
)

data = pa.table({
    "episode_id": [1, 2],
    "score": [0.95, 0.87],
    "episode_day": ["2025-01-15", "2025-01-15"],
})
client.catalog.write(data, namespace="robotics", table="scored_episodes")

# Evolve a table's schema after creation
from deltacat_io_core.schema_updates import add_field

client.catalog.alter_table(
    namespace="robotics",
    table="predictions",
    schema_updates=[add_field(pa.field("confidence", pa.float64()))],
)
```

See [Reading and Writing](docs/reading-and-writing.md) for write modes,
staged writes, path-based writes with `create_if_missing`, and the full
list of supported formats and schema inputs.
For PackDS training data writes, see [Training Data](docs/training-data.md).

</details>

<details>

<summary><b>Live Feature Enrichment</b></summary>

Declare one or more merge keys on the table schema and subsequent writes
update existing records in place.

```python
import pyarrow as pa
import deltacat as dc
from deltacat_client import Client

client = Client("http://localhost:8080")

# Fields tagged with is_merge_key=True identify the key columns used to
# reconcile updates.
schema = dc.Schema.of([
    dc.Field.of(pa.field("user_id", pa.int64()), is_merge_key=True),
    dc.Field.of(pa.field("name", pa.string())),
    dc.Field.of(pa.field("age", pa.int32())),
    dc.Field.of(pa.field("job", pa.string())),
])

# First write creates a Lance table with the merge-key schema.
initial = pa.table({
    "user_id": pa.array([1, 2, 3], type=pa.int64()),
    "name": ["Jim", "Dinah", "Bob"],
    "age":  pa.array([30, 28, 45], type=pa.int32()),
    "job":  ["Teacher", "Painter", "Sailor"],
})
client.catalog.write(
    initial,
    namespace="demo",
    table="users",
    format="lance",
    create_if_missing={"schema": schema, "auto_create_namespace": True},
)

# Upsert: user_id 1 + 3 are updated; 4 and 5 are inserted.
upsert = pa.table({
    "user_id": pa.array([1, 3, 4, 5], type=pa.int64()),
    "name": ["Cheshire", "Felix", "Tom", "Simpkin"],
    "age":  pa.array([3, 2, 5, 12], type=pa.int32()),
    "job":  ["Tour Guide", "Drifter", "Housekeeper", "Mouser"],
})
client.catalog.write(
    upsert,
    namespace="demo",
    table="users",
    mode="merge",
    format="lance",
)

print(client.catalog.read(namespace="demo", table="users", read_as="pyarrow"))

# Delete: only the merge-key columns need to be supplied.
client.catalog.write(
    pa.table({"user_id": pa.array([3, 5], type=pa.int64())}),
    namespace="demo",
    table="users",
    mode="delete",
    format="lance",
)

# user_id 3 and 5 are gone; 1, 2, 4 remain.
print(client.catalog.read(namespace="demo", table="users", read_as="pyarrow"))
```

PackDS v5 tables do not support DeltaCAT merge keys. For canonical
PackDS training data, use `episode_id` identity partitioning and
`REPLACE_PARTITION`. If you need row-level `MERGE`/`DELETE` semantics,
use a normal DeltaCAT table instead of `layout="packds"`.
See [Training Data](docs/training-data.md).

#### PackDS Episode Rewrites

PackDS v5 tables are identity-partitioned by `episode_id` and reject
merge keys. Live feature enrichment runs episode-at-a-time: select an episode
from application inventory or the converter's rich `_episodes` table, then
rewrite it with
`mode="replace_partition"`. The replaced partition is committed atomically;
sibling episodes are never touched.

```python
import pyarrow as pa
from deltacat_client import Client

client = Client("http://localhost:8080")

# 1. Initial PackDS write across multiple episodes. DeltaCAT creates the
#    table with episode_id identity partitioning and no merge keys.
initial = pa.table({
    "episode_id": [
        "ep_0001", "ep_0001", "ep_0001",
        "ep_0002", "ep_0002",
        "ep_0003", "ep_0003", "ep_0003", "ep_0003",
    ],
    "step_index": [0, 1, 2, 0, 1, 0, 1, 2, 3],
    "task": [
        "pick", "pick", "pick",
        "place", "place",
        "stack", "stack", "stack", "stack",
    ],
    "annotation": [None] * 9,
})
client.catalog.write(
    initial,
    namespace="robotics",
    table="training_steps",
    format="lance",
    create_if_missing={
        "schema_def": [
            {"name": "episode_id", "type": "string"},
            {"name": "step_index", "type": "int64"},
            {"name": "task", "type": "string"},
            {"name": "annotation", "type": "string"},
        ],
        "layout": "packds",
        "default_content_type": "lance",
        "auto_create_namespace": True,
    },
)

# 2. Plans for a terminal _steps name advertise the converter's rich sibling.
#    This metadata does not create or query that table.
plan = client.catalog.plan(
    namespace="robotics",
    table="training_steps",
    filter_predicate={"eq": ["episode_id", "ep_0001"]},
    include_files=False,
)
assert plan.episodes_table == "robotics.training_episodes"
target_id = "ep_0001"
target_len = 3

# 3. Enrich the target episode end-to-end via REPLACE_PARTITION. No
#    merge keys, no merge-on-read. The replacement is per-partition;
#    sibling episodes are unaffected. Note that ``layout="packds"`` is
#    only required at create time; subsequent writes pick the layout up
#    from the table's stored ``dataset_layout`` property.
enriched = pa.table({
    "episode_id": [target_id] * target_len,
    "step_index": list(range(target_len)),
    "task": ["pick"] * target_len,
    "annotation": [f"auto-labeled-step-{i}" for i in range(target_len)],
})
client.catalog.write(
    enriched,
    namespace="robotics",
    table="training_steps",
    mode="replace_partition",
    format="lance",
)

# 4. Scoped read of the enriched episode. PackDS reads must be scoped
#    to a partition filter or a bounded set of episode_ids (equality / IN,
#    capped by ``DELTACAT_PACKDS_EPISODE_ID_FILTER_MAX_FANOUT``, default
#    1024); unscoped full-table PackDS reads are rejected by default and
#    require ``allow_full_packds_scan=True`` (or
#    ``DELTACAT_ALLOW_FULL_PACKDS_SCAN=true`` server-side) to opt in.
print(client.catalog.read(
    namespace="robotics",
    table="training_steps",
    read_as="pyarrow",
    filter_predicate={"eq": ["episode_id", target_id]},
    columns=["episode_id", "step_index", "task", "annotation"],
))
```

Larger enrichment loops typically batch over application-owned episode IDs so
each `replace_partition` write touches exactly one episode while the
worker pool processes many in parallel. Because each episode is its own
DeltaCAT partition, concurrent replace-partition writes to disjoint
episodes do not contend.

See [Training Data](docs/training-data.md) for guidance on choosing between
PackDS episode rewrites and a normal merge-key table.

</details>

<details>

<summary><b>Transactions</b></summary>

Transactions provide atomic multi-step operations with automatic heartbeating and rollback on failure.

```python
import pyarrow as pa

with client.transaction(commit_message="Backfill predictions") as tx:
    client.catalog.write(
        pa.table({"episode_id": [10, 11], "score": [0.88, 0.92]}),
        namespace="robotics",
        table="predictions",
        mode="add",
    )

    # Reads within the transaction see uncommitted writes
    df = client.catalog.read(
        namespace="robotics",
        table="predictions",
        read_as="pandas",
    )
    print(f"Rows visible in transaction: {len(df)}")
# Transaction commits automatically on exit; aborts on exception
```

See [Transactions](docs/transactions.md) for time-travel reads, manual commit/abort, and transaction rules.

</details>

<details>

<summary><b>Jobs</b></summary>

DeltaCAT uses a durable job system for background work (compaction, data relay, subscription processing). The client can submit, monitor, and execute jobs.

```python
# List all jobs
jobs = client.jobs.list()
for job in jobs:
    print(f"{job.job_id}: {job.state}")

# Submit a compaction job
result = client.jobs.submit_compaction(table="predictions")
print(f"Submitted: {result.job_id}")

# Wait for it to complete
status = client.jobs.wait(result.job_id, timeout_seconds=120)
print(f"Final state: {status.state}")
```

Admins can submit a placement metadata hash baseline through
`client.jobs.submit_placement_hash_baseline(...)`, supplying the crawl table
and version, placement table, and source binding. It defaults to
`verify_only=True, apply=False`. An apply request must explicitly set both
flags and provide the existing server confirmation; the client never supplies
confirmation automatically. An active job for the same subject is deduplicated.
Submission is not completion: monitor the returned job through `client.jobs`.

To inspect an existing host-local frontier operation, use
`client.catalog.get_frontier_progress(table=..., frontier_operation_id=...,
frontier_authority_id=...)` on a client connected directly to its owning server.
This read-only call never starts a scan or retries against a different owner.
Unknown totals remain unknown; progress is not proof that a frontier can be
adopted or that a repair has completed.

Workers claim and execute jobs. Any process can act as a worker:

```python
# Claim the next pending job matching our worker tags
job = client.jobs.claim(worker_tags=["subscriber"])
if job:
    print(f"Claimed {job.job_id} (type: {job.context.get('job_type')})")

    # Report fine-grained progress and declare the deadline for the next chunk
    client.jobs.emit_event(
        job,
        event_name="batch_started",
        completed=1,
        expected=4,
        metadata={"batch": 1},
        heartbeat_timeout_seconds=300,
    )

    # Do the work...

    # Mark complete. Source-consuming jobs must include the advanced
    # watermark so the server can persist progress correctly.
    client.jobs.complete(
        job,
        records_processed=1000,
        watermark={
            "partition_watermarks": {"analytics.events": 42},
            "known_partitions": ["analytics.events"],
        },
    )
```

See [Jobs and Workers](docs/jobs-and-workers.md) for job types, worker
routing, heartbeat rules, retry semantics, and dispatch modes.

### Managed Ray Example

For subscription, transform, and publication jobs, `dispatch_mode="ray"`
submits the work onto DeltaCAT-managed Ray clusters instead of the shared
custom-worker claim pool. For large prebuilt environments, prefer
`docker.image` as the primary runtime artifact and use `payload` only for
small code/config overlays. The server still vends the DeltaCAT runtime from
its promoted runtime manifest so the remote cluster gets compatible
`deltacat`, `deltacat_client`, and `deltacat_io_core` packages.

```python
from deltacat_client import Client, DispatchMode, SubscriberType
from deltacat_io_core.triggers import SubscriptionTrigger

client = Client("http://localhost:8080", bearer_token="token")

payload = client.jobs.stage_ray_payload(
    py_modules=["./shared_logic"],
)

ray_dispatch_yaml = """
cluster_shutdown_policy: terminate
docker:
  image: "registry.example.com/ml/runtime@sha256:0123456789abcdef"
payload:
  py_modules:
    - {py_module}
head_node_type: ray.head.default
available_node_types:
  ray.head.default:
    node_config:
      InstanceType: m7i.2xlarge
  ray.worker.default:
    min_workers: 1
    max_workers: 4
    node_config:
      InstanceType: m7i.2xlarge
""".format(
    py_module=payload.payload["py_modules"][0],
)

client.subscriptions.create(
    subscriber_id="episode_processor",
    source_tables=[{"namespace": "robotics", "table": "raw_episodes"}],
    subscriber_type=SubscriberType.CUSTOM,
    dispatch_mode=DispatchMode.RAY,
    dispatch_config=ray_dispatch_yaml,
    trigger=SubscriptionTrigger.schedule(interval_seconds=300),
)
```

The scheduler now auto-triggers the subscription every five minutes, DeltaCAT
auto-provisions the Ray cluster when a run starts, and
`cluster_shutdown_policy: terminate` tears the cluster back down after each run
by default. If the deployment advertises a compatible image profile, DeltaCAT
will execute inside the requested `docker.image` while keeping the promoted
DeltaCAT runtime bundle as the compatibility layer.

For reference, the Ray launcher YAML above is:

```yaml
cluster_shutdown_policy: terminate
docker:
  image: registry.example.com/ml/runtime@sha256:0123456789abcdef
payload:
  py_modules:
    - s3://.../shared_logic.tar.gz
head_node_type: ray.head.default
available_node_types:
  ray.head.default:
    node_config:
      InstanceType: m7i.2xlarge
  ray.worker.default:
    min_workers: 1
    max_workers: 4
    node_config:
      InstanceType: m7i.2xlarge
```

Use `payload` for small overlays only:

- experiment modules
- config bundles
- supplemental wheels

Do not treat `payload` as the primary transport for a large application image.

After launch, owner/admin callers can inspect cluster metadata and then use
the managed-Ray access helpers against a specific job:

```python
# `access` describes the currently discovered cluster shape for the job:
# region, cluster name, head instance id/private IP, worker ids, and the
# launch template metadata DeltaCAT used for the cluster.
access = client.jobs.get_managed_ray_access("job-id")
print(access.head_instance_id, access.head_private_ip)

# By default the SSM session starts on the Ray head node.
# The AWS CLI / Session Manager plugin will open an interactive session in
# your terminal, and the helper asks SSM to launch `bash -l` there so you land
# directly in a login shell on the instance.
client.jobs.start_managed_ray_ssm_session("job-id")

# For non-interactive inspection across the whole cluster, use send-command.
client.jobs.send_managed_ray_ssm_command("job-id", commands=["hostname"], target="all")
```

</details>

<details>

<summary><b>Publications</b></summary>

Publications are incremental producers that write new data into the
DeltaCAT Lakehouse. They sit at the root of a pipeline DAG and can be
triggered manually or fired by an upstream event.

```python
# Create a publication that writes to a sink table
client.publications.create(
    publication_id="episode_publisher",
    name="Episode Publisher",
    sink_tables=[{"namespace": "robotics", "table": "clean_episodes"}],
    dispatch_mode="local",
)

# Run the publication
result = client.publications.run("episode_publisher")
print(f"Published: {result}")
```

See [Pipelines](docs/pipelines.md) for publication configuration and DAG construction.

</details>

<details>

<summary><b>Subscriptions</b></summary>

Subscriptions are incremental consumers of DeltaCAT tables. They sit at
the leaves of a pipeline DAG, tracking a per-partition watermark so each
run picks up only new data.

```python
# Create a subscription that watches for new data in "raw_episodes"
client.subscriptions.create(
    subscriber_id="episode_processor",
    source_tables=[{"namespace": "robotics", "table": "raw_episodes"}],
    subscriber_type="custom",
    dispatch_mode="custom",
)

# Trigger processing (dispatches a job to a subscriber worker)
client.subscriptions.trigger("episode_processor")

# Check watermark state
wm = client.subscriptions.get_watermark("episode_processor")
print(f"Watermark: {wm.watermark}")

# Pause / resume / delete
client.subscriptions.pause("episode_processor")
client.subscriptions.resume("episode_processor")
client.subscriptions.delete("episode_processor")
```

See [Pipelines](docs/pipelines.md) for subscription modes (delta vs. version), triggers, and redrive.

</details>

<details>

<summary><b>Transforms</b></summary>

Transforms are the intermediate nodes of a pipeline DAG. Each transform
reads from one or more source tables, applies processing logic, and
writes to one or more sink tables.

```python
# Create a transform: raw_episodes -> clean_episodes
client.transforms.create(
    transform_id="episode_cleaner",
    name="Episode Cleaner",
    source_tables=[{"namespace": "robotics", "table": "raw_episodes"}],
    sink_tables=[{"namespace": "robotics", "table": "clean_episodes"}],
    dispatch_mode="custom",
)

# Trigger transform processing
client.subscriptions.trigger("episode_cleaner")

# Pause / resume
client.transforms.pause("episode_cleaner")
client.transforms.resume("episode_cleaner")
```

See [Pipelines](docs/pipelines.md) for transform configuration, redrive, and rollback.

</details>

<details>

<summary><b>Pipelines</b></summary>

Pipelines wire publications, transforms, and subscriptions into a
connected DAG. When an upstream node completes, downstream nodes are
triggered automatically.

```python
# First, create connected pipeline nodes
client.publications.create(
    publication_id="ingest_pub",
    name="Raw Ingest Publisher",
    sink_tables=[{"namespace": "robotics", "table": "raw_data"}],
    dispatch_mode="local",
)
client.transforms.create(
    transform_id="clean",
    name="Data Cleaner",
    source_tables=[{"namespace": "robotics", "table": "raw_data"}],
    sink_tables=[{"namespace": "robotics", "table": "clean_data"}],
    dispatch_mode="custom",
)
client.subscriptions.create(
    subscriber_id="consume_clean",
    source_tables=[{"namespace": "robotics", "table": "clean_data"}],
    subscriber_type="custom",
    dispatch_mode="custom",
)

# Preview the connected pipeline from an interior seed node
preview = client.pipelines.discover(seed_node_ids=["clean"])
print(preview.execution_order)

# Option A: persist exactly the previewed node_ids
client.pipelines.create(
    pipeline_id="etl_pipeline_pinned",
    name="ETL Pipeline (Pinned)",
    node_ids=preview.node_ids,
)

# Option B: create directly from seed node ids
client.pipelines.create(
    pipeline_id="etl_pipeline_seeded",
    name="ETL Pipeline (Seeded)",
    seed_node_ids=["clean"],
)

# Check pipeline status
status = client.pipelines.status("etl_pipeline_seeded")

# Pause / resume all nodes at once
client.pipelines.pause("etl_pipeline_seeded")
client.pipelines.resume("etl_pipeline_seeded")
```

See [Pipelines](docs/pipelines.md) for DAG construction, discovery
semantics, redrive, rollback, and stored-order validation.

</details>

<details>

<summary><b>Data Placement and Replication</b></summary>

DeltaCAT catalogs can span multiple storage backends (S3, SwiftStack, Lustre). Data placement lets you replicate tables across roots so readers access data from the closest location.

```python
# List available data roots for the catalog. The keys of this dict
# are the valid values to pass as `roots=` and `root=` below.
available_roots = client.catalog.list_data_roots("default")
for root_name, root_info in available_roots.items():
    print(root_name, root_info)

# Example output (dict[str, DataRootInfoSummary]):
# aws_s3_iad       {'root': 's3://bucket-iad/',        'storage_type': 's3',
#                   'region': 'us-east-1', 'endpoint_url': None}
# aws_s3_pdx       {'root': 's3://bucket-pdx/',        'storage_type': 's3',
#                   'region': 'us-west-2', 'endpoint_url': None}
# swiftstack_pdx   {'root': 's3://swiftstack-bucket/', 'storage_type': 's3',
#                   'region': None,        'endpoint_url': 'https://pdx.swiftstack.example.com'}

# Pick a target root from what the catalog actually advertises.
# In production you'd choose based on region, storage class, or policy.
if not available_roots:
    print("This catalog uses a single default root; no placement root to select.")
else:
    target_root = next(iter(available_roots))

    # Place a table on an additional storage root for replication
    client.catalog.place(
        namespace="robotics",
        table="episodes",
        roots=[target_root],         # Replicate to this root
        backfill=True,               # Copy existing data too
    )

    # Check replication status
    status = client.catalog.replication_status(namespace="robotics", table="episodes")
    print(f"Roots: {status}")

    # Read from the closest root (server resolves automatically)
    df = client.catalog.read(
        namespace="robotics",
        table="episodes",
        read_as="pandas",
        root=target_root,            # Prefer this root for file paths
    )

    # Remove a replication target
    client.catalog.unplace(namespace="robotics", table="episodes", root=target_root)
```

New writes are automatically replicated to all placed roots via a background subscriber. Reads with `root=` get file paths rewritten through the preferred root when data is available there.

See [Configuration](docs/configuration.md) for data root setup and multi-root catalog configuration.

</details>

<details>

<summary><b>Discovery Root Administration</b></summary>

Catalog admins can publish read-only discovery roots for system crawler, GC,
relay, replication, and placement jobs. Discovery roots are separate from
managed data roots, so they can describe authority roots such as `s3://`
without becoming valid table write targets.

```python
client.catalog.create_storage_binding(
    "default",
    name="swiftstack_pdx_root",
    storage_type="s3",
    storage_class="discovery",
    uri_scheme="s3",
    endpoint_url="https://pdx.swiftstack.example.com",
    region="us-west-2",
    credential_source={"kind": "aws_profile", "profile": "pdx"},
)

client.catalog.create_discovery_root(
    "default",
    name="swiftstack_pdx",
    binding_name="swiftstack_pdx_root",
    root_uri="s3://",
    authorization={"principals": ["service:system-crawler"]},
)

roots = client.catalog.list_discovery_roots(
    "default",
    include_disabled=True,
    include_history=True,
)
details = client.catalog.describe_discovery_root(
    "default",
    "swiftstack_pdx",
    include_history=True,
)

# Updates create a new immutable discovery-root version.
client.catalog.update_discovery_root(
    "default",
    "swiftstack_pdx",
    binding_name="swiftstack_pdx_root",
    root_uri="s3://",
    description="PDX authority root",
    authorization={"principals": ["service:system-crawler"]},
)

# Deletes are logical disables and remain visible with include_disabled=True.
client.catalog.delete_discovery_root("default", "swiftstack_pdx")
```

Discovery-root detail responses include the effective binding revision,
endpoint, region, authorization summary, lifecycle state, sanitized
credential-source metadata, and effective storage snapshot.

</details>

<details>

<summary><b>Scheduled Processing</b></summary>

Subscriptions and transforms can run on a schedule instead of being triggered manually. DeltaCAT supports interval-based and cron-based scheduling.

```python
# Process new data every 5 minutes
client.subscriptions.create(
    subscriber_id="metrics_ingester",
    source_tables=[{"namespace": "telemetry", "table": "raw_metrics"}],
    subscriber_type="custom",
    dispatch_mode="custom",
    trigger={"mode": "schedule", "schedule": {"interval_seconds": 300}},
)

# Process at 2am UTC daily using a cron expression
client.subscriptions.create(
    subscriber_id="nightly_aggregator",
    source_tables=[{"namespace": "telemetry", "table": "raw_metrics"}],
    subscriber_type="custom",
    dispatch_mode="custom",
    trigger={"mode": "schedule", "schedule": {"cron": "0 2 * * *", "timezone": "UTC"}},
)

# Event-driven: only runs when triggered manually or by an upstream pipeline node
client.subscriptions.create(
    subscriber_id="on_demand_processor",
    source_tables=[{"namespace": "telemetry", "table": "raw_metrics"}],
    subscriber_type="custom",
    dispatch_mode="custom",
    trigger={"mode": "event"},
)
```

See [Configuration](docs/configuration.md) for trigger options and scheduling details.

</details>


## Agentic Access (MCP)

DeltaCAT ships with a built-in [Model Context
Protocol](https://modelcontextprotocol.io) server so AI agents (for
example, Claude Code) can browse catalogs, inspect schemas, plan reads,
and stage writes through natural language instead of hand-written Python.

Most hand-written code should just use the REST `Client(...)` shown
above. Reach for MCP when you want to:

- let an agent explore and operate on your catalog conversationally
- embed DeltaCAT as a tool in an agentic application
- use a typed async Python wrapper over the MCP HTTP surface
  (`deltacat-client[mcp]`) from code that is already agentic in shape

See the [MCP Server guide](../deltacat/docs/mcp/README.md) for the full
tool reference, the typed async client, and recipes for agent-driven
catalog workflows.


## Server Setup

The DeltaCAT client connects to a DeltaCAT API server. For setup instructions, see:

- **[Server Setup Guide](docs/server-setup.md)** -- start a local or production server
- **[REST API Reference](../deltacat/docs/server/README.md)** -- full REST endpoint documentation
- **[MCP Server Guide](../deltacat/docs/mcp/README.md)** -- agentic access via MCP


## Additional Resources

| Guide | Description |
|-------|-------------|
| [Reading and Writing](docs/reading-and-writing.md) | Read plans, write modes, staged writes, supported formats |
| [Computed Fields](docs/computed-fields.md) | Persisted and request-time PyArrow expressions, schema evolution, projection rules, and validation |
| [Training Data](docs/training-data.md) | PackDS v5 tables, rich episodes metadata, scoped training reads |
| [Transactions](docs/transactions.md) | Transaction lifecycle, time travel, rules and limitations |
| [Jobs and Workers](docs/jobs-and-workers.md) | Job types, claiming, heartbeat, worker routing, authentication |
| [Pipelines](docs/pipelines.md) | Publications, transforms, subscriptions, DAGs, redrive |
| [Configuration](docs/configuration.md) | Auth, dispatch modes, triggers, data placement |
| [Maintainer Workflow](docs/maintainer-workflow.md) | Relationship between REST and MCP, generated bindings, facade updates, and validation guards |
| [Architecture](docs/architecture.md) | Package boundary, generated client, development notes |
| [Client Compatibility Runbook](../deltacat/docs/scratch/design/deploy/CLIENT_COMPATIBILITY_RUNBOOK_2026-04-12.md) | Promotion-time packaged client/server compatibility validation |

For the core DeltaCAT data model, storage architecture, and catalog APIs, see the [DeltaCAT documentation](../README.md#overview).

## Explicit staged-repair promotion

`client.training.reprocess(..., publication_mode="unreleased", write_mode=...)`
keeps ACTIVE readers on their current versions. Wait for conversion, companions
and required placements, then inspect `operation.destinations` by explicit version.
Use `client.training.promote_reprocess(operation.operation_id)` only when ready
to cut over; `wait_reprocess_promotion` / `get_reprocess_promotion` observe the
ordinary server-owned publication job. Repeating the action reuses its original
transaction; it never repeats conversion. Source or destination changes after the
captured snapshots refuse promotion, rather than discard new data. Use a new
`request_id` for a new repair: an ID reused with different retained
inputs cannot claim the historical promotion. ADD/APPEND,
keyed MERGE where supported, and partition/table replacement remain distinct repair
choices; isolated promotion is optional, not a universal repair requirement.
