Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Clinker Engine Internals

This book is for engineers working on Clinker — or operators who need to reason about why the engine behaves the way it does under load. It documents the execution model, memory arbitration, correlation-key retraction, operator strategies, storage durability, and CXL compilation in implementation detail.

If you only need to author and run YAML pipelines, you want the Clinker User Guide instead (the separate book under docs/). That book deliberately stays at the level of “what do I type and what happens”; this one explains the machinery beneath it.

What’s in scope here

  • The execution model — which stages stream, which block, and how the memory arbitrator decides who pauses and who spills.
  • Correlation-key retraction — the shadow-column lineage, per-source rollback narrowing, and the retraction protocol that lets a relaxed aggregate drop only the failing records.
  • Operator internals — Combine join-strategy selection, Merge back-pressure, streaming Sink writes, and the schema-drift sidecar.
  • Storage — the staging cache, crash durability, and the locking protocol.
  • CXL compilation — the compiler phases and the type-unification algorithm.

How to read it alongside the User Guide

Most chapters here have a user-facing counterpart in the User Guide. The pattern is consistent: the User Guide tells you the knob and the observable outcome (e.g. “declare sort_order and the aggregate streams”); this book explains the mechanism that makes the outcome true (the streaming-ingest path, the group-table accumulation, the spill trigger). When a chapter has a user-facing sibling, it says so at the top.

Architectural ground rules

Everything here is downstream of three permanent commitments — finite inputs, finite jobs, single process — covered in Overview & Pillars. A mechanism that appears baroque (the single-process memory arbitrator, the in-process Rayon parallelism) usually makes sense only once you hold those three constraints fixed.

Overview & Pillars

Clinker is a bounded-memory batch DAG executor. A pipeline run is a finite job over finite input: Source nodes read until EOF, the DAG drains, the process exits with a status code. It pairs a custom expression language (CXL) with YAML pipeline orchestration.

Within a run, stateless operators (Transform, Route, most Combine probe-side work, Sink) evaluate records one at a time without per-record state accumulation. The DAG executor materializes intermediate buffers between non-fused stages, with retained records accounted across all live stages and spill available for materialized buffers; fused Source → Transform → Sink paths skip materialization entirely. Blocking operators (Aggregate, sort, grace-hash Combine) accumulate state inside the configured RSS budget (default 512 MB) and spill to disk when soft/hard thresholds trip rather than OOM the process.

The three pillars

Every design decision cascades from three commitments. They are permanent — an architectural proposal that violates any of them is rejected at design review, not implementation review.

  1. Finite inputs only. Files (CSV / JSON / XML / fixed-width / EDIFACT / X12 / HL7 v2 / SWIFT MT) and finite-cursor network sources (paginated REST with hard page/record caps) — both reach EOF after exhausting their cursor. Unbounded sources (Kafka, Kinesis, SSE, webhooks, tail -f) are out of scope permanently.

  2. Finite jobs. No daemon mode, no service surface, no infinite event loop. clinker run invokes, drains, exits.

  3. Single process forever. One invocation = one OS process. Parallelism happens inside the process via std::thread and Rayon — no worker-process pools, no multi-machine sharding, no network shuffle, no cluster manager. Scale by adding cores / RAM / disk to one host. If a host genuinely can’t fit the work, partition the input by file or key and run multiple clinker invocations from a shell script.

These pillars are why the memory arbitrator is a single in-process component rather than a distributed scheduler, why there is no network shuffle in Combine, and why spill-to-local-disk is the universal pressure-relief valve.

Crate dependency layers (top → bottom)

Applications:    clinker (CLI) | cxl-cli (CXL tool)
                      |
Edge services:   clinker-channel | clinker-net | clinker-schema | clinker-lineage
                      |
Execution:       clinker-exec (runtime operators, memory, spill, metrics)
                      |
Planning:        clinker-plan (YAML, validation, CXL binding, compiled DAG)
                      |
Language / IO:   cxl | clinker-format
                      |
Foundation:      clinker-record | clinker-core-types

Support:         clinker-scenarios | clinker-bench-support | clinker-benchmarks

The support crates are siblings, not part of the default runtime path. Some edge crates depend on more than one lower layer; the repository’s AI crate map records the detailed dependency edges and their evidence.

The node taxonomy

Pipelines use a single flat nodes: list; each entry’s type: discriminator selects a variant of one homogeneous DAG:

  • Source — finite input endpoint with an inline, generated, or external schema.
  • Transform — record-level CXL projection / filter / lookup (1×1).
  • Aggregate — grouped or windowed reduction.
  • Route — predicate-based fan-out.
  • Merge — streamwise concatenation of inputs.
  • Combine — N-ary record combining with mixed predicates (equi + range + arbitrary CXL); distinct from Merge and Transform+lookup.
  • Reshape — per-group mutate-and-synthesize.
  • Cull — per-group rule evaluation with retained and removed output ports.
  • Envelope — document-level consolidation or expansion at explicit DAG boundaries.
  • Sink — terminal writer.
  • Composition — call-site node referencing a .comp.yaml reusable sub-pipeline, lowered at compile time.

The plan itself is a petgraph DAG (ExecutionPlanDag) of topologically-sorted nodes, each carrying a parallelism strategy and NodeProperties (ordering / partitioning provenance). CXL is typechecked at compile time into a TypedProgram, and schema is propagated across the DAG at plan time.

Planner/runtime handoff

clinker-plan is the execution-admission layer: canonical YAML parsing, topology and path validation, schema binding, CXL typechecking, composition binding, and lowering produce a CompiledPlan. Public executor entry points accept &CompiledPlan, but the current implementation then calls plan.config() and recompiles before dispatch. The stored plan is therefore a typed public boundary today, but its stored DAG and other compiled artifacts are not yet the artifacts the runtime dispatches directly.

The locked D-01 through D-11 contract corrects that mismatch in Phase 5 / PERF-01: the supplied plan must remain authoritative and reusable for sequential in-process runs, while only an enumerated run envelope may refresh. Persistent cache identity, semantic comparison, integrity checks, and source-map refresh are part of the same downstream contract; none of that work is implemented by this chapter. See Stored-plan execution and cache identity and Streaming vs. Blocking Stages.

Terminal destination vocabulary

PipelineNode::Sink, SinkConfig, and YAML type: sink are the current terminal-writer surface; planning lowers them to PlanNode::Sink and runtime execution delegates to executor/sink_dispatch.rs. The retired type: output spelling is rejected with the paste-ready correction type: sink. Output ports, produced artifacts and paths, serialization formats, stdout and machine output, writer results, and OpenLineage output datasets remain distinct and valid output vocabulary. See Sink Nodes, Sink Internals, and Terminal destination vocabulary.

Key engine decisions

  • Memory-aware aggregation. Hash aggregation with disk spill; streaming aggregation when sort order permits; RSS tracking with soft/hard limits. The mechanism is documented in Memory Arbitration & Scheduling.
  • Compile-time CXL typechecking. Type inference produces a TypedProgram; see Compiler Phases & Type Unification.
  • Diagnostics. All user-facing errors use miette for span-annotated reports. Spanned<PipelineNode> covers the YAML side, cxl::Span covers the expression side, and they compose into one report.
  • Pure Rust policy. No crate in the graph invokes a C compiler. The two places that could — TLS and content hashing — are held to Rust implementations: rustls with the graviola provider rather than ring or aws-lc-rs, and blake3’s pure feature, which keeps the std::arch SIMD paths and gives up only the variants its build script would compile through cc. That second trade is not free: pure hashes at roughly 85% of the assembly build’s throughput at the same instruction set, and roughly 57% on a CPU whose AVX-512 kernels it declines to use. The policy accepts that cost rather than a build-time C dependency. A CI job builds the workspace and all its targets with every C-compiler environment variable pointed at a failing program, and then requires a crate that does compile C to fail — so the guarantee is checked, and the check is checked, rather than either being asserted. deny.toml bans cmake alone, and deliberately does not attempt this: a name blocklist cannot separate a build script that runs cc from one that only declares it.

The boundaries available to engine extensions are described in Extension Seams.

Extension Seams

Clinker has explicit places where new behavior joins the engine, but it does not have a general-purpose plug-in system. Some extensions implement a typed streaming contract; others require coordinated changes across the authoring, planning, and runtime layers.

The distinction matters. A transport can join the common ingest path by implementing RecordSource. A new YAML node cannot be dropped into a registry: it changes the language of pipeline topology and must be understood by the planner, executor, diagnostics, and plan consumers.

The central boundary

pipeline YAML + CXL
        |
        v
clinker-plan
  parse -> validate -> bind/typecheck -> lower/enrich
        |
        | CompiledPlan / ExecutionPlanDag
        v
clinker-exec
  ingest -> dispatch -> arbitrate memory -> write/report

The planner owns author input and proves as much as it can before any record flows. It is the sole authority that admits a pipeline for execution. The executor owns runtime effects and consumes compiled artifacts. Raw YAML does not belong in operator code, and byte or network mechanics do not belong in plan lowering. clinker-schema is an advisory discovery and warning tool: its bounded scans and heuristic field extraction do not authorize execution or override a planner rejection (D-17).

File formats

FormatReader turns a byte source into records; FormatWriter turns records into bytes. Both are Send but not Sync because one worker owns each stream. Their less obvious hooks are part of the seam too: schema discovery, multi-file source identity, document preparation and envelope events, non-finalizing byte flushes, document framing, and byte counts.

A format becomes user-selectable only after its typed YAML options and central reader/writer construction arms are wired. It must also state which schema, multi-record, envelope, splitting, and document-cardinality features it can represent. This is deliberate compile-time wiring, not dynamic discovery.

Allocation-aware CSV construction

CSV execution uses CsvReader::from_reader_admitted and MultiRecordReader::new_csv_admitted with finite allocation resources. The explicit legacy reader constructors remain available for non-executing inspection and existing library callers; they do not establish runtime admission. Document preparation returns OwnedMap through both FormatReader and RecordSource, preserving child ownership across adapters. Legacy readers can move-wrap an existing map without claiming that it was admitted.

Direct CSV output construction requires WriterResources: CsvEncoder::into_boxed_writer returns a FormatWriterHandle. Wrappers that allocate another writer use FormatWriterHandle::try_new(value, &scope) with an AllocationScope; all writer hooks and byte counts forward through the handle. AsMut<dyn FormatWriter> borrows the interface without exposing its allocation owner. FormatWriterHandle::from_legacy(Box<dyn FormatWriter>) is an explicit boundary for unchanged implementations, not a CSV fallback.

Split construction uses WriterFactory::try_new(closure, &scope) and calls create(destination, schema) to obtain a FormatWriterHandle for each file. The factory admits the concrete closure layout before type erasure; captured allocations need separate owners. WriterFactory::from_legacy leaves an unchanged codec factory outside CSV/JSON/XML ungoverned. CountedFormatWriter and SplittingWriter retain writer handles, so wrapping and rotating a CSV writer preserve its backing charge. See memory ownership.

Finite native writer construction

JSON and XML expose JsonEncoder/XmlEncoder, admitted shared JsonEncoderConfig/XmlEncoderConfig, and borrowed options. Construct an encoder with finite WriterResources, then wrap it in PreparedWriter<W, E> for a borrowed or owned destination. For type erasure, into_boxed_writer admits the concrete allocation and returns FormatWriterHandle. Pass explicit finite MemoryOnlyResources for direct use or executor WriterResources at runtime; there is no unlimited constructor and no raw JsonWriter/XmlWriter fallback.

FormatEncoder::prepare must leave committed state unchanged. Return pending cache/framing/count state and commit it only after complete stage delivery. Every wrapper forwards document hooks, flush_bytes, finalization, byte counts and failures. Never add buffering outside the prepared writer that can retry bytes from Drop or obscure accepted destination counts. WriterFactory retains the admitted config owner across split rotations.

Schema cache identity is SharedStorageIdentity<Schema>, never a raw address or an owned schema clone. The old cache and its prospective replacement remain charged together until the operation commits or aborts. Preserve the legacy weak-backing reservation and final deallocation order. See native ownership.

Prepared storage extensions

ResourceAuthority::create_stage returns the sealed OperationStage owner. Implement StageStorage and call StorageStage::create(scope, storage) rather than implementing or boxing an operation-stage trait. The adapter admits the concrete backend and its progress buffer before writes; finish consumes the writable stage and moves the same owner into PreparedBytes without another backend allocation.

StageStorage supplies read/write, seal, complete and inline failure evidence. seal rewinds and establishes complete immutable bytes; complete releases readback storage before encoder state commits. resource_failed must not allocate or change the observed error. Recover resource evidence through failure/resource_error at I/O boundaries, and never seal a failed prefix. Providers receive no destination handle. The sealed wrapper keeps the backend and its lease inseparable through delivery and deallocation; see prepared storage for partial delivery and cleanup debt.

Fixed-width repeating groups

Fixed-width repeating groups demonstrate the coordinated form of this seam. The strict Column shape carries fields, occurs, and an optional physical count_field through patching, overlay provenance, canonical identity, and normal planner admission before either format constructor runs. Layout resolution then proves a finite maximum width with checked arithmetic. The reader retains one declared-length row, while the writer validates and encodes one declared-length record before committing bytes. The group remains one logical array-of-records column; a derived count cell is layout machinery and never enters the logical schema or CXL namespace. Delimiter-packed scalar cells remain a distinct split_values representation and cannot bypass the positional-group checks. The executable byte and rejection matrix is in crates/clinker-format/tests/fixed_width_repeating_groups.rs, with planner admission coverage in crates/clinker-plan/src/plan/tests/multi_value_validation.rs.

Non-file transports

RecordSource is the transport-neutral ingest contract. File readers reach it through an adapter; paginated REST implements it directly and enters execution as SourceInput::Records. Once records cross this boundary, common ingest owns schema coercion, provenance, document signals, watermarks, backpressure, and handoff to the DAG.

Every source is finite. A cursor must end after a bounded result set; daemon polling and unbounded streams do not fit this seam.

Pipeline nodes

A pipeline node crosses three representations:

PipelineNode (author shape) -> PlanNode (compiled shape) -> dispatch arm

Adding or changing one therefore requires an end-to-end review:

  • strict span-aware YAML parsing, topology and configuration validation;
  • schema propagation, CXL typing, lowering, and composition-body behavior;
  • ordering, partitioning, correlation-key, streaming, scheduling, and cardinality properties;
  • runtime ports, record/control-event flow, DLQ, metrics, cancellation, memory, spill, and cleanup;
  • explain output, lineage, examples, documentation, and boundary tests.

The compiled DAG also contains synthetic nodes inserted by planning. Those are runtime machinery, not automatically valid YAML node types.

Clinker Expression Language (CXL)

The Clinker Expression Language’s extension seam is its ordered compiler pipeline:

parse -> resolve -> typecheck -> analyze/extract -> evaluate

Each phase consumes stronger input than the previous phase. New syntax or an AST form must keep node identifiers and recursive visitors coherent and must be handled by every later phase that can receive it. CXL remains below planning and execution: it knows records and expression semantics, not pipeline YAML or operator scheduling. It is a per-record ETL expression language, not SQL, and it does not decide whether a pipeline is executable.

Composition example corpus

Committed composition fragments are an executable authoring contract, not parser fixtures. The corpus test follows this boundary:

recursive inventory -> exact case set -> production loader -> generated pipeline -> clinker -> exact bytes

The inventory discovers every .comp.yaml recursively without a manual allowlist and compares normalized, repository-relative paths with the case manifest as exact sets. Duplicate manifest keys, path escapes, missing or empty directories, missing cases, and extra cases are separate failures. Both inventory paths and failure ordering are deterministic.

For each case, the harness materializes only the selected fragment in an isolated workspace, calls the production composition scanner and compiler, and then invokes the built clinker executable. A case succeeds only when the process status, run counters, and output bytes match its committed expectation. clinker run --explain remains a compilation check and cannot replace the runtime byte comparison.

This seam adds no production retention or execution work of its own. Its fixed, finite fixtures exercise the existing compiled-plan, runtime, telemetry, and lineage paths; their existing memory and observability contracts remain the authoritative ones.

Diagnostics and runtime resources

Diagnostic codes are registered centrally, then emitted with source spans and, where useful, typed payloads. The long-form clinker explain --code pages are a second coordinated surface rather than an automatic result of registration.

Runtime resources are run-scoped rather than global. Memory consumers register with one MemoryArbitrator; writers, window indexes, progress callbacks, spill guards, and similar state are passed through explicit handles or registries. A new stateful operator participates in shared memory, pause/spill, disk quota, shutdown, and cleanup rules instead of creating an independent resource model.

Composition resource slots

_compose.resources_schema declares typed slots; the bounded workspace catalog currently admits the file kind. A composition call binds a slot to one logical catalog identity through strict scalar resources: values. Binding rejects missing and undeclared slots, unknown identities, kind or capability mismatches, inline descriptors, credential selectors, and the ordinary call-site outputs and alias fields. Each winning logical binding retains its complete attempted-versus-winning overlay provenance.

An authored body Source names its slot explicitly with resource: <slot> in the Source config. The binder never infers a slot from a Source name or path. It rejects a body Source that combines resource: with path, glob, regex, or paths, and it rejects an authored body Source with no resource link. Input ports remain separate synthetic Source roots seeded by the caller; they do not receive activation entries. Top-level direct file Sources keep their existing matcher surface and reject resource: because no top-level catalog-binding surface exists.

Each bound call compiles its body Sources into scope-qualified CompiledSourceInstance values. The retained requirement contains the slot, logical binding and provenance, kind, finite capabilities, opener family, run-local lifetime, and stable logical dataset identity. It deliberately drops the catalog’s physical path and cannot contain a credential choice, secret, live handle, I/O state, or thread state. Separate calls to the same composition therefore share immutable logical descriptor semantics but have distinct Source identities at activation.

For a data run, the CLI resolves each credential-free file requirement back to the admitted workspace catalog and captures its validated path inside an opaque, single-use factory. The executor receives only the sealed activation bundle, takes each complete group lease before opening any member, opens every member before publishing a bounded Source channel, and retains the sessions until that composition scope ends. A partial open, reader failure, shutdown, or ordinary completion closes sessions, releases leases, and unregisters the channels’ memory consumers. Credential-bearing groups still fail preflight: there is no credential-profile selection surface yet. See the composition resources and call-site surface contract.

Approved transitional exceptions

Four narrow dependency and parser exceptions are recorded by D-20 through D-23. Each approves only the boundary named here:

DecisionCurrent permitted boundaryForbidden expansionOwner
D-20clinker-format -> cxl may import only logical type and document path/index vocabulary: cxl::typecheck::Type and cxl::analyzer::doc_paths::{DocPath, DocIndex}.No parser, resolver, evaluator, planner, or other analyzer dependency.Phase 1 contract; neutral extraction is deferred.
D-21clinker-exec -> clinker-bench-support may remain optional behind bench-alloc, outside default and release graphs.No default-runtime edge or trusted allocation claim until forwarding, allocator identity, plausible measurements, and distortion are qualified.Phase 5 / PERF-07.
D-22Direct serde_saphyr::from_str* calls belong only in clinker-plan::yaml and parser-specific tests.No production or cross-module test bypass of clinker_plan::yaml.Phase 2 repair; Phase 6 / EVID-03 qualification.
D-23A manifest dependency remains only after source, build, generated-code, feature, test, and supported-API use is proven.No speculative dependency or async/runtime coupling.CONT-05 bounded cleanup.

D-20 is implemented as a transitional exception. D-21 and D-23 are only partly implemented, and D-22’s known executor-test bypass has not yet been repaired. The full current evidence and compatibility rules live in the production contract register.

Declarative and read-only edges

Channels alter only declared overlay surfaces and feed the effective result back through normal compilation. Composition bodies are sealed except for declared ports, config, scoped variables, and, after AUTH-01 lands, typed resource slots. Unknown resource kinds, slots, or catalog names must fail closed; channel reachability does not create a second admission authority.

Plan consumers should remain read-only edges. clinker-lineage, for example, walks a CompiledPlan without depending on or invoking runtime operators; live run facts are supplied by the CLI boundary.

For exact change paths, source anchors, focused tests, and unresolved-boundary routing, see the implementer seam map.

Streaming vs. Blocking Stages

User-facing view: the User Guide’s “Streaming vs. Blocking Stages” page.

This page is the engine-internals reference for the runtime classifier that decides whether a node hands its output downstream in bounded batches or accumulates its whole input before emitting. The streaming/blocking split is the mechanism behind Clinker’s bounded-memory guarantee, and the same classifier annotates --explain output and drives the dispatcher at runtime, so the model here is exactly what the executor does — not a simplification of it. Read it alongside Memory Arbitration & Scheduling, which covers how each in-flight batch and materialized slot charges the budget.

Every node in a pipeline plan is one of two kinds at runtime:

  • Streaming stages hand their output downstream in bounded batches over a back-pressured channel, never crossing an inter-stage buffer that charges the memory budget. The two fused streaming paths additionally hold at most one batch of in-flight events at a time, so their inter-stage memory does not grow with input size. The other streaming stages still build their own result before handing it off — streaming spares them the second copy into a charged buffer and overlaps the writer with downstream work, but their own working set is as large as a blocking stage’s would be.
  • Blocking stages must see their whole input before they can produce any output. They accumulate state inside the memory budget and spill to disk when the soft threshold trips, rather than holding everything in RAM.

Source decoding and schema admission

Format decoding and pipeline typing are separate boundaries. A single-record CSV decoder produces text cells; CoercingReader in clinker-exec::pipeline::schema_coerce applies the declared schema before records enter source buffering or downstream operators. Numeric and date columns are therefore typed at source admission, not lazily on their first CXL use. Invalid declared values follow the source’s row-error policy. Positional readers that already parse typed fields carry that proof through the same boundary, which validates it without parsing the text twice.

Materialized input invariant

Every planned materialized edge has an occupied node-buffer slot when its consumer runs. A producer that emitted no rows still admits an explicit empty slot; absence is not another spelling of an empty input. Once a dispatcher has excluded its certified streaming/fused path and any explicit alternate slot, a missing materialized slot is therefore an executor invariant failure. The run returns PipelineError::Internal and stops instead of manufacturing an empty collection and allowing plausible but incomplete output to commit.

The checked retrieval is stage-level work: it performs the same map lookup and moves the same buffer as the successful path. It adds no record-rate allocation, clone, or per-record bookkeeping. Optional lookup remains limited to body seeding, own-slot-versus-predecessor selection, and cleanup paths where absence has defined control-flow meaning.

Aggregate, Reshape, Cull, and planner-synthesized Sort use one address rule: first check their own (consumer, None) slot, which is where a Route branch or Cull port publishes its selected records, then require the incoming (producer, producer_port) slot. A present empty own slot is still authoritative; only the absence of both valid addresses is an invariant failure. Transform and Sink also recognize successor-local slots through their existing specialized input paths. Merge and Combine remain predecessor-slot readers because they select among multiple incoming edges.

Shared-buffer scans and composition scope

Every published materialized slot carries a remaining-reader count keyed by its exact (producer, producer_port) NodeBufferKey. The producer declares that count when it publishes the slot; readers never rediscover ownership from node kind, node index, or dispatch order. The common one-reader path removes the slot and its registration directly with O(1) bookkeeping. With several readers, each earlier reader opens a sequential scan over shared immutable backing while the original stays live for the final reader. This applies uniformly to materialized Transform, Aggregate, Sort, Reshape, Route, Cull, Envelope, Composition, Merge, Combine, and Sink inputs, including successor-local Route/Cull slots and both Sink event paths.

Memory, Spilled, and Mixed all support repeatable scans. A memory scan clones one event at a time; a spill scan opens one chunk at a time, preserving record and punctuation order with O(1) file descriptors per active scan. A consumer that collects the scan into a resident vector first registers the estimated materialized bytes and holds that reservation through its complete synchronous operation; unwinding an error releases it automatically. The reader count changes only after the scan is acquired successfully. The final reader drains the authoritative slot and its ordinary node-buffer registration. Successful scope completion rejects any residual slot, registration, or positive reader count as an internal invariant failure.

An adopted MergeSpilled run set cannot be scanned repeatedly because its merger consumes and unlinks the runs. Its first shared acquisition therefore folds it exactly once into an ordinary spill file. Replacement disk bytes are charged before the input runs are released; an overlap beyond the spill quota returns E320 and cleans up every file and registration from the failed fold.

A composition input uses the same scan/materialization boundary, then transfers the live consumer id and byte handle into the body-local node-buffer registry. Body Source canonicalization briefly needs the seeded events and its prospective output together, so it extends that same reservation before allocating the output and reduces it to the canonicalized footprint as soon as the seed allocation drops. Slot admission then atomically replaces the transient consumer wrapper with the ordinary spill-aware wrapper under the same id. Body execution swaps the parent node_buffers, node-buffer registrations, reader ledger, source-record table, body references, and window state as one scope. Body dispatch and output harvest are captured before a single cleanup path unregisters body residue and restores every parent map, so a successful body, a dispatch error, and a harvest error have the same ownership lifecycle. The transfer keeps one continuous registration and charges both allocations only for their real overlap: there is no unregistered interval and no second consumer charge for the same bytes.

This distinction is what makes Clinker a bounded-memory executor: the budget covers the combined live state of operators, source queues, materialized boundaries, and in-flight batches. A largest-stage estimate alone is not an upper bound, especially when several consumers retain the same upstream boundary. A streaming stage’s output is never separately buffered between dispatch arms, so it is never charged twice: the arbitrator counts each in-flight batch once when the producer flushes it and discharges that charge as the consumer drains it. If RSS still crosses the soft threshold while a single-consumer streaming stage holds batches in flight, the engine spills those batches’ records to disk one batch at a time — the streaming handoff is the per-batch counterpart of a blocking stage’s full-stage spill, not an exemption from spilling.

Plan admission and runtime entry

The current public entry is typed as a compiled-plan boundary, but its complete call path re-enters planning:

PipelineConfig::compile -> CompiledPlan
                             |
PipelineExecutor::run_plan_with_readers_writers(&CompiledPlan)
                             |
                        plan.config()
                             |
          run_with_readers_writers[_in_context]
                             |
                    PipelineConfig::compile
                             |
        newly validated plan.dag() -> runtime dispatch

The _in_context entry supplies a CompileContext so file-size estimates are resolved against the intended workspace instead of the process working directory. The run body also derives memory admission from the embedded config before compiling again. PipelineRunParams carries the current run-scoped envelope: execution and batch IDs, pipeline/static/source/record variable overlays, a shutdown token, and spill root, quota, and compression policy. These parameters affect a run without serving as a second authoring topology.

This is the observed implementation, not the locked destination. D-01 through D-11 require Phase 5 / PERF-01 to execute the supplied plan’s stored DAG, composition bodies, bound schemas, compiled expression artifacts, statistics, and semantic identity directly. Only an explicitly enumerated runtime envelope may refresh; semantic planning changes require replanning. A CompiledPlan is borrowed and survives a call today, but the current recompilation means sequential reuse still repeats immutable planning work. Persistent cache identity, safe misses, atomic replacement, and fresh source-map/provenance handling are also Phase 5 work. See the canonical stored-plan execution and cache identity contract.

Frozen execution wiring

The canonical compile completes every structural graph rewrite before it freezes the two artifacts that runtime dispatch consumes:

  • CompiledConsumerRegistry is keyed by stable producer identity and optional producer port. Each deterministic consumer entry records the consumer’s stable identity and input port, whether it reads the shared producer slot or a pre-forked slot, and whether it crosses a physical writer boundary.
  • ExecutionOrderContract retains source orders, edge guarantees, consumer requirements, terminal guarantees, and topology-derived physical writer boundaries. No later planner pass may change topology after this contract is frozen.

The registry is the single delivery authority for fan-out. A consumer of a shared producer port obtains one complete sequential scan: all but the final observed consumer receive an independent re-readable cursor, and the final consumer removes the authoritative slot and its one MemoryArbitrator registration. Completion is tracked by compiled identities, not graph indexes or dispatch order. A missing entry, duplicate read, missing shared slot, or incomplete final scan is an internal invariant failure instead of a silent short delivery.

For a Source with sort_order, its CompiledSourceOrder binds stable source identity, field indexes and types, direction and null policy, on_unsorted, the sortable event shape, and PerPhysicalFile scope. Runtime constructs the physical-file barrier from that compiled value. Successful rows and declared type failures share one attempt stream; document punctuation is staged with them. The barrier emits the complete attempted/rejected population before the covered attempts. The executor applies that population exactly once and checks its threshold before routing any covered attempt into the DLQ, downstream state or counters, or writer effects.

Each PhysicalWriterBoundary binds one Sink to the producer port, writer mode, partition identity, ordering guarantee, and runtime disposition selected from the finalized topology. The modes are RecordsOnly, PerSourceFile, Envelope, DocumentDlq, CorrelationDeferred, and Streaming; partition identity distinguishes a single writer, split sequence, source file, document, or correlation group. OrderedWriterBoundary verifies that the active dispatch arm matches the compiled mode and that its disposition satisfies the compiled guarantee. Deferred boundaries use the shared stable authored-key sorter and bounded-fan-in spill merge; source-row identities and population indexes remain payload and never become comparison keys. A streaming arm rejects a complete-population disposition rather than pretending to implement it.

All of this remains synchronous, single-process, finite-batch execution. Bounded channels provide back-pressure, MemoryArbitrator remains the one run-scoped memory authority, and source repair, terminal ordering, and shared fan-out reuse the existing spill and merge machinery. Runtime never recounts Sinks or invents a weaker order when a frozen contract and the selected path disagree; it returns PipelineError::Internal.

Ordering evidence and test oracles

Ordering is an explicit plan/runtime contract, not a side effect of scheduling. The frozen ExecutionOrderContract records where stable arrival is preserved, where order is destroyed, and which physical-file sources require verification. Runtime strategies must satisfy that contract; an unordered compiled promise is a conservative lower bound and may be served by a stronger exact runtime path.

Boundary or strategyRuntime behaviorScope and supported oracle
Source without sort_orderPreserves the reader’s arrival sequence but claims no sorted keys.No authored sortedness promise; use multiset and aggregate assertions unless a later operator establishes order.
Source with sort_orderStrict declared-type coercion precedes adjacent-key verification. A memory-arbitrated barrier releases each physical file only after it is verified; default on_unsorted: warn stably repairs and emits one W307, while error releases no prefix.PerPhysicalFile, never global across files. For one fixed path, resident and spill repair must be byte-identical and stable for equal authored keys.
Order-preserving unary paths and Merge concatRetain predecessor arrival; concat drains inputs in declaration order.Exact sequence is valid for the same upstream paths. Matching per-input sorts are not promoted to global order.
Seeded Merge interleaveEstablishes a reproducible stable-arrival schedule for the seed.Exact sequence is valid for the same seed and paths.
Unseeded interleave and every current Combine strategyCross-input or matched-row arrival is incidental and the plan marks it unordered.Compare decoded multisets, aggregate values, counters, and identities; never exact incidental row order.
Terminal Sink sort_orderUses the shared stable resident/spill kernel and orders all records reaching that terminal by exactly the authored fields.Exact authored-key order. Equal-key order is stable within a path, but cross-strategy exact bytes require a total authored business key.

Neither source repair nor terminal sorting adds SourceRowId, physical-file, or canonical-row tie fields. Equal authored keys compare equal. Stability is maintained by arrival/run position within the selected path, so tests must not turn an unpromised upstream tie order into a hidden compatibility contract.

The source barrier stages the whole sortable event shape, including successful and rejected attempts plus file and single-frame punctuation, and charges resident rows, release state, spill I/O buffers, merge cursors, and queued verified rows to the run-scoped memory authority. It first releases one complete population delta, then the attempts covered by that identity. It preserves the exact SourceRowId, source/file provenance, and original document-context Arc while repairing. This is a safe pre-effect barrier; it does not replay a source after downstream writers, DLQ, counters, metrics, lineage, or document state have mutated.

Which stages stream

A stage streams when its output is handed straight to a single downstream consumer instead of crossing a charged inter-stage buffer. The downstream consumer is a Sink writer, an Aggregate’s ingest, or a hash build-probe Combine’s probe (driver) side — see Streaming into an Aggregate and Streaming into a Combine probe below.

Two stages stream and bound their own footprint to one batch, because they pull records off a live upstream channel and forward each batch without ever building a full result:

  • Source → Transform → Sink fused chains. A non-windowed Transform whose only upstream is a single Source and whose only downstream is a single Sink consumes that Source’s records directly and hands each batch to the Sink’s writer thread over a back-pressured channel; neither the Transform nor the Sink materializes the whole record set. A bounded preview (--dry-run -n) runs the Transform unfused so its Sources drain in a fixed order, and the Transform hands its materialized rows to the same Sink writer thread in batches instead of buffering them. A Transform that fans out to multiple consumers, feeds another operator, or roots a window keeps the buffered (materialized) path.
  • Merge in interleave mode fed entirely by Sources. The merge reads each Source’s live stream and forwards records as they arrive.

These stages stream their output to a single downstream consumer too — sparing the second copy and overlapping the consumer — but each still builds its full result first, so its own working set is not bounded to one batch:

  • Single-branch Route. A Route with exactly one branch feeding one Sink streams that branch’s records to the writer thread. A multi-branch Route forks records across several successor buffers and stays materialized.
  • Merge in concat mode, or interleave fed by non-Source inputs, feeding one Sink. The merge drains its predecessors’ buffers in order (concat) or round-robin (interleave) into the merged result, then streams it.
  • streaming-strategy Aggregate feeding one Sink. When the planner certifies the aggregate’s input is pre-sorted on the group key, it finalizes the group rows and streams them rather than buffering them for a downstream arm.
  • Combine probe side (hash build-probe strategy) feeding one Sink. The build relation stays fully materialized in the hash table; the matched probe output streams to the writer.
  • Block-band Combine output (IEJoin and HashPartitionIEJoin) feeding one Sink. The join blocks on its inputs, sorting both sides first, but its bounded, payload-sorted output drain streams to the consumer instead of admitting a node buffer. The sort-merge and grace-hash joins keep their output materialized.

Each of these requires the producer to feed exactly one downstream consumer and to root no window; a producer that roots a window keeps the materialized path because the window arena needs the producer’s full output to build.

  • Every Sink writes records to its configured writer and never buffers a whole stage, except under dlq_granularity: document, where it holds each open document’s records until the document’s verdict. A Sink with an authored sort_order, a split, or a per-source-file path does not certify its producer as a streaming edge (certify_streaming_edge), so that producer keeps a materialized slot.

Document-boundary punctuations (DocumentOpen / DocumentClose, the signals behind the $doc.* context) flow inline with records through streaming stages, preserving their order: a document’s close always trails the document’s last record, even when the document’s records span several batches.

Streaming into an Aggregate

The streaming consumer above is usually a Sink. It can also be an Aggregate’s ingest: when an eligible producer (a fused Source → Transform, a single-branch Route, a non-fused Merge, or a streaming-strategy Aggregate) feeds exactly one downstream Aggregate, the producer streams record-at-a-time into the aggregate’s add_record over a back-pressured channel rather than the aggregate pre-draining the producer’s whole output from a charged buffer. The producer reports buffer: streaming and --explain shows no node_buffer edge between it and the aggregate.

This streams the aggregate’s ingest half only — the producer no longer needs a charged inter-stage slot, and a slow aggregate (one that is spilling, say) paces the producer through the bounded channel. The aggregate’s finalize half stays blocking by nature: a group_by value depends on every member, so the group table accumulates the whole input and emits only once its producer’s input has ended explicitly: the producer’s dispatch returned Ok and the hop’s End arrived. A channel that closes without End (the producer failed, or the run was cancelled) is an incomplete input, and the aggregate finalizes nothing. Spill stays driven by RSS pressure, never by channel depth, exactly as on the materialized path.

Two aggregate shapes keep the materialized ingest, because their finalize is not a single forward pass: a time-windowed aggregate runs a multi-pass per-window algorithm over the whole input, and a relaxed correlation-key aggregate retains its group state for the correlation-commit phase. Both show buffer: materialized on the edge into them.

Streaming into a Combine probe

A producer can also stream into a hash build-probe Combine’s probe (driver) side. When an eligible producer (a fused Source → Transform, a single-branch Route, a non-fused Merge, a streaming-strategy Aggregate, or another hash build-probe Combine) is the Combine’s driver input, the producer streams record-at-a-time into the probe kernel over a back-pressured channel rather than the Combine pre-draining the driver’s whole output from a charged buffer. The driver producer reports buffer: streaming and --explain shows no node_buffer edge between it and the Combine. Only the HashBuildProbe strategy qualifies — the range, sort-merge, and grace-hash kernels re-sort or re-scan the driver and stay materialized.

This streams the Combine’s probe half only. The build side stays fully materialized: the engine builds the complete hash table on the main thread before the driver producer streams its first record, so the probe never matches against an incomplete index. The probe consumer runs on its own thread, so a slow driver paces the probe through the bounded channel and a slow probe (a large fan-out) back-pressures the driver. The build relation’s footprint is the hash table, exactly as on the materialized path; the streaming handoff spares only the driver’s inter-stage slot. Per-source dead-letter rewind, memory accounting, and output are byte-identical to the materialized path.

Which stages block

A stage blocks when its result depends on records it has not seen yet:

  • sort — the full input must be present before the first sorted record is known.
  • Hash Aggregate — a group’s final value depends on every member, so the group table accumulates the whole input. (A streaming-strategy Aggregate over a pre-sorted input is the exception: the planner certifies it can emit a group as soon as the sort key advances.)
  • Combine build side — the build relation is fully indexed before any probe record is matched. The probe side streams against the built index, but the build side materializes.
  • IEJoin / sort-merge Combine — both inputs are sorted before the band/merge step runs, and both are block-spilled so the input axis stays inside the budget, but by different mechanisms. The IEJoin — pure-range and equi+range alike, which share the one block-band path — external-sorts each side to disk on (equality-hash, range-key, …) and slices the sorted stream into min/max-tagged, single-equality-hash blocks, pruning block-pairs on the equality hash and the range bounds before the kernel runs (equality is an added prune axis; each surviving pair re-verifies the canonical equality key, since hashes collide). Its output axis is spill-bounded too — matched rows accumulate in a payload-ordered sort buffer that spills on its own byte threshold (charged through the join’s consumer handle) and drains incrementally, streamed straight to a downstream Sink or folded into a spillable node-buffer, so both axes are bounded with no global-pressure abort. The sort-merge Combine external-sorts each side into runs and merges matching runs; it has no min/max block tags or pruning.
  • CorrelationCommit — a correlation group is held until its commit decision (flush or dead-letter) is known.

A blocking stage keeps its full-stage accumulation inside pipeline.memory.limit and spills to disk past the soft threshold; it does not stream batches.

Seeing the classification

clinker run <pipeline>.yaml --explain annotates every node with its class in the Physical Properties section:

sink.report:
  buffer: streaming

aggregation.dept_totals:
  buffer: materialized

buffer: streaming marks a stage whose output is consumed without an inter-stage buffer — it charges the budget per in-flight batch and, on a single-consumer edge, spills those batches to disk under pressure; buffer: materialized marks a stage whose output crosses a node_buffers slot that charges the memory budget as one full-stage slot and spills the whole stage. Both classes are spill-eligible; they differ in granularity, not in whether they can spill. The explain annotation is derived from the same classifier the executor uses at runtime, so what --explain reports is exactly what the dispatcher does. See Memory Arbitration & Scheduling for the arbitration model that rides alongside the buffer class.

Tuning the batch size

The number of events handed downstream per batch is set by pipeline.batch_size (default 2048), with an optional per-transform override on a Transform’s config.batch_size. For a fused streaming stage — the only kind whose footprint is one batch — smaller batches lower its in-flight footprint at the cost of more per-batch bookkeeping; larger batches do the reverse. For the other streaming stages the batch size sets only the in-flight slice handed across the channel; the producer’s own result is built in full regardless, so batch_size does not cap their footprint. The batch size changes only the memory profile of streaming handoffs — never their output, and never the behavior of blocking stages.

Memory Arbitration & Scheduling

User-facing view: the User Guide’s “Memory Tuning” page.

Interactive companion: the memory system explainer walks through the budget, ledger, arbitration policies, backpressure, spill and scheduler, with a simulator that runs the decision rules on this page. Its script ports select_victim, reconcile_backpressure, spill_reclaimable and next_runnable; update it when those change.

This page is the engine-internals reference for how Clinker tracks, attributes, and reclaims memory at runtime, and how it orders simultaneously-runnable nodes to keep the resident working set bounded. It covers the MemoryConsumer wrapper registry, pull-mode byte attribution, the per-operator arbitration parameters the active policy reads, the bounded-memory contract for materialized stages, the predicted_* values that feed both --explain and the scheduler, and the four ranking rules the scheduler applies (with its fallback to topological order). The user-facing knobs — the memory: block, the --memory-limit flag, the backpressure-policy selection, sizing guidance, and monitoring — live in the User Guide and are intentionally not repeated here. For how each stage’s buffer class (streaming vs materialized) is decided, see Streaming vs. Blocking Stages.

How it works

Exact allocation admission for prepared output

clinker_format::preparation prepares one complete output operation before delivery. CSV, JSON, XML, fixed-width and SWIFT writer construction in the CLI and executor uses this path; EDIFACT, X12 and HL7 retain their existing allocation behavior. MemoryOnlyResources requires an explicit nonzero memory budget. ExecutorResources shares the run’s MemoryArbitrator and takes an explicit optional telemetry producer; disabling telemetry changes no resource limit. Neither provider offers an unlimited memory path.

The allocation vocabulary lives in clinker_record::owned_storage: AllocationAuthority, AllocationResources, AllocationScope and the non-cloneable AllocationLease. Format preparation separately supplies the temporary-storage capability. Both use the same executor admission ledger; record storage does not depend on the format or executor crates.

AllocationLease reserves the complete requested Layout before allocation. ReservedBuffer and ReservedVec retain that grant until the allocation is freed. Growth reserves the replacement while the old block remains charged; allocation failure leaves the old contents and charge intact. Moving a grant or transferring it within one authority moves ownership without a release and reacquire gap. Splitting or merging grants preserves the total; cross-authority transfers are refused. Requested layouts include stage metadata, chunk inventories, and retained progress space, rather than just encoded lengths.

OperationStage retains its metadata lease outside the boxed storage backend. The backend and its allocation are destroyed before that lease is released, including on error and unwinding. Sealing moves this same owner into PreparedBytes; it neither reallocates the backend nor detaches its grant.

FormatWriterHandle and WriterFactory apply the same ordering to the concrete writer and closure backings. They admit the actual concrete Layout before fallible boxing and keep the lease outside the box. Internal buffers and captured values retain their own owners; the outer layout cannot account for their heap allocations. The handles expose neither a detachable box nor its lease. Tests observe the charge at allocator deallocation, including unwinding, rather than treating payload destruction or a final zero balance as proof.

The executor admission ledger is the run’s one charged total; it serializes reservations, consumer charges and limit changes. Governed allocations charge it through MemoryArbitrator::reserve, which checks and charges under the ledger’s one lock and returns a Grant or a Shortfall. Registering a consumer (register_consumer for run-scoped state, register_node_consumer for a node’s state, each with the consumer’s ConsumerHandle and a label) binds the handle to the same ledger: from then until the consumer unregisters, every charge through the handle is a ledger charge under that lock, and unregistering releases what the handle still holds. A handle charges for one consumer at a time: registering a second consumer through a handle still bound to the first is an internal error that registers nothing, and the caller returns it. ConsumerHandle::try_grow and try_resize check a growth against the limit with every other charge, so a handle growth and a governed allocation, or two of either, can never together pass the limit; the older set_bytes / add_bytes charges are applied unchecked. A Shortfall charges nothing and carries a snapshot taken under the same lock: the limit, the charged total, each labelled holder’s current bytes under its node name and an author-vocabulary surface, and the bytes no labelled holder owns, so the holders and that remainder add up to the charged total. A Source’s consumer unregisters once the Source has finished reading, and the drain arm hands its rows on only after that, so the rows it read can still be charged in its name once its consumer is gone. unregister_consumer reads can_back_pressure once (the predicate that makes a listed holder a Source), and for a Source whose grants are still live the ledger keeps its entry, unlabelled and marked as a finished Source’s, until the last of those bytes drops; the snapshot reports them as retired_source, always inside the remainder. Nothing is re-charged. An E310 shows them as memory not held by any one node and never counts them as state that cannot spill, the same rule as for a Source still listed; bytes granted in no consumer’s name, or in the name of another consumer that has unregistered, still count. A consumer that unregisters while it is the walk requester stops being the requester, so the rest of its arm’s walk allocations are charged to no consumer and never join a finished Source’s rows. A reserve made for a consumer is attributed to it for the life of the grant, and the ledger keeps each consumer’s high-water mark over its handle bytes plus its attributed bytes; attribution never changes what is admitted. Every release advances a release epoch. writer_resource_usage() derives current memory, peak memory, disk and descriptor usage from this ledger; its memory figures are what governed allocation grants hold, not the consumer handle charges beside them. set_limit refuses a limit below the charged total and leaves the previous limit unchanged; the disk setter likewise refuses a quota below the sum of outstanding writer disk and legacy spill bytes.

On the walk thread a Shortfall is not yet a refusal. A reserve, a Grant::try_grow or a ConsumerHandle::try_grow / try_resize made on the walk that does not fit runs a reclaim pass without holding the ledger lock: the registered consumers that cannot be paused and hold bytes are taken in the run’s policy order (ties to the older consumer), the requesting consumer last, each consumer’s figures read once when the pass begins (the shipped policies order them with one sort), and each whose state the walk owns is spilled there and then (see “How a pass reaches state” below). A consumer whose state the walk does not own is skipped and never asked to act; one the running dispatch arm holds (a slot out of its scope, a cell its owner is borrowing) frees nothing this pass and has its own spill request raised, which its owner answers at its next boundary or push. A pass aims to bring the ledger, with the request charged, down to the resume watermark, not just to fit the request. The pass’s progress is the sum of its victims’ own releases, recorded by the ledger while each victim spills on the walk (the slot’s charge and the governed allocations its records drop), so a concurrent release by another thread never counts as a victim’s progress; it only marks the pass as having seen a release. The request retries after each pass. It is refused only when a pass freed nothing with no release during it and a final pass then freed nothing too; the refusal’s snapshot is the one taken with the retry after that final pass. Every other thread’s request is checked once and never spills. Governed allocations the walk makes while a dispatch arm runs are charged to that node’s first registered consumer, and release against it however the arm has moved on.

How a pass reaches state

A pass reaches an elected consumer’s state through the walk reclaim set, in one of two ways:

  • Node-buffer slots live in the set’s frames: one frame per dispatch scope, the top level’s and one per composition body running inside it. A pass searches the running scope’s frame first, then each calling scope’s, so a body that falls short can spill a resident slot its callers still hold, though the body itself never reads, replaces or removes a caller’s slot.

  • Walk-owned state lives in a cell of its own (Rc<RefCell<_>>) that is registered under the consumer it charges through register_walk_owned (crates/clinker-exec/src/pipeline/memory/walk.rs). Registered today, each with what its figure leaves out while its owner holds that part:

    • the run’s document dead-letter state (document_dlq.rs): its held rows, resident tails only; a pass that finds the state mid-step raises its spill request instead;
    • an Output’s per-document buckets (document_dlq.rs, one registration per bucket): resident records; a pass that finds the buckets borrowed raises the bucket’s spill request;
    • the rows parked for a deferred consumer (parked_generations.rs, one registration per edge): resident segments, less any an open replay cursor shares;
    • a Cull’s and a Reshape’s group buffers (cull_dispatch.rs, reshape_dispatch.rs): resident groups, less groups already on disk or taken out for routing;
    • a grace-hash Combine’s partition table (grace_hash/mod.rs): the partitions its build is still filling, and nothing once the probe holds them;
    • a hash Aggregate’s group tables on the strict per-document and time-windowed arms (aggregate_dispatch.rs, one registration per table): the groups each table holds resident, and nothing from the moment a finalize takes the table, or ever for a table with no spill directory.

    These seven are every production caller of register_walk_owned. A pass does not reach the following state today:

    • the sort-merge and IEJoin kernels’ state, which spills on thresholds of its own. Their consumers report 0 reclaimable, so no pass elects them, and a refused request’s E310 lists them as cannot spill. The per-operator table and the consumer inventory still class them as spillable at priorities 25 and 20: that is their class once they register, not what a pass can do today;
    • an authored Sort’s buffer, which registers no consumer of its own and spills on a threshold of its own, sized from the limit (sort_dispatch.rs, operator_memory_limit);
    • an inline hash join’s table, which never spills (grace hash is the spillable join strategy): its consumer reports 0 reclaimable;
    • a streaming-ingest Aggregate’s tables, which its worker thread owns (below);
    • a relaxed-key Aggregate’s table, which does not register and ranks by 0 while it ingests and while the commit keeps it, so no pass and no soft-threshold poll elects it. The in-place finalize that ends its ingest, and the commit’s retract and finalize after it, read only resident groups: the in-place finalize fails on a table with spilled groups, so its spill would end the run or lose the Aggregate’s groups rather than free memory (#1288).

    The registry is run-scoped, outside every frame, and holds only a Weak to each cell, so an owner dropped on any exit is never reached. Several cells may register under one consumer; one cell may serve several consumers and spills only what the elected one charges. The owner keeps the registration beside its state, so both drop together. It borrows its cell only for one operation of its own, never across a call that can charge another consumer, a channel wait or a call into another dispatch arm. Its spill never reserves.

State owned by a thread other than the walk (a Source reader, a streaming writer or worker, such as a streaming-ingest Aggregate’s tables) is not reached by a pass: its consumer is skipped, and the pass raises no spill request for it. A streaming-ingest Aggregate’s tables rank by their charge, so a pass can elect them and find them NotOwned; they spill on their own thresholds, and on a spill request the soft-threshold poll raises, which the worker reads as it adds its next row (#1247).

Each victim a pass asks ends in one of three outcomes:

  • Spilled. The walk owns the state and wrote resident state of it to disk, now, on the walk. For walk-owned state, at least one free cell reported that it wrote (OwnedSpillResult::Wrote); a cell’s spill never reports a write it did not make.
  • Busy. The walk owns the state but its owner holds it right now (a slot out of its scope, a slot whose rows a live reader’s cursor or view still shares, a cell its owner is mutating or is itself the requester, or a free cell that holds state for the consumer but had nothing it could write: every part already on disk, taken out for use, or shared with a reader). The pass frees nothing from it and raises the consumer’s spill request, which the owner answers at its next push, yield or batch boundary. A shared slot’s spill writes nothing and leaves its figure and charge as they were; the E310 lists it as in use, never at its floor. The walk’s spill-request sweep at the next node dispatch clears that request and writes nothing while the reader still shares the rows; the next pass that elects the slot raises it again.
  • NotOwned. The walk holds no spillable state for the consumer: a slot its compiled classification keeps in memory, state another thread owns, or a registered owner that is gone or no longer holds that consumer. It is skipped and never asked to act. The round keeps it as evidence that no spill the walk could make would free that consumer’s bytes: when the last pass to elect it found it NotOwned, a refused request’s E310 lists it as cannot spill and counts its bytes as state that cannot spill. A pass that finds it NotOwned does not name it as asked; if an earlier pass of the same round asked it, the reclaim line still names it from that pass. A report with no round has no such evidence, so a consumer another thread owns still lists as in use there.

Spillable state that no pass can reach is a false E310: a request that does not fit is refused while megabytes it could have freed stay resident. So every walk-owned spillable state must register through register_walk_owned, and nothing walk-owned and spillable may be NotOwned; the spillable state listed above as not reached by a pass is the open exception (the inline hash join’s table is listed there too, but it never spills, so no pass could free it). Registering changes none of a consumer’s charge, priority, spill triggers or admission; it only makes the state reachable. A registered group buffer (Cull, Reshape) also records on its consumer’s handle what spilling its resident groups frees now, which is the figure it ranks by. A grace-hash partition table records the bytes of the partitions its build is still filling; a pass that elects it spills every one of them. Once the build finishes the probe holds every in-memory partition, so the figure is 0 and no pass elects the consumer, though the partitions stay charged until they drop. A hash Aggregate’s table records its charge as its figure while it can spill, and 0 after a spill wrote its groups, from the moment a finalize takes it (a walk arm takes it out of its cell first), and always when it has no spill directory or is a relaxed-key Aggregate’s table; a table being finalized or kept for the commit stays charged until it drops, and a refused request’s E310 lists it as cannot spill.

The execution report samples the arbitrator’s spill totals and the ledger’s charged peak after dispatch has finished and every Source worker has joined. Ordered Sources can still release staged spill charges while unwinding cancellation; sampling at dispatch close would report those already-released bytes as live. The total and per-stage spill fields include committed charges minus releases, not every byte ever written to temporary storage.

Two further report figures answer per-node questions those totals cannot. per_stage_spill_bytes_written adds every spill charge a stage records and is never lowered by a release, so a sort whose runs were merged and unlinked before the run ended still shows the bytes it wrote while its on-disk entry is back at zero. per_node_peak_charged_bytes gives, for each node whose state is registered under the node’s name (register_node_consumer, whose label names the node), the highest charge any one of its consumers reached. The ledger keeps each consumer’s mark over its handle’s bytes plus the governed allocations made in its name, and every charge to either raises it, so the figure is exact per consumer rather than sampled, and another node’s state never raises it. A node with several consumers reports the largest single consumer’s mark. Run-scoped state that no node owns (writer output staging, the credential registry, the document dead-letter state) registers through register_consumer and has no entry. The report’s run-wide charged peak is the ledger’s own: the most bytes charged at one instant, consumer handles and governed allocations together, raised by every charge rather than taken only when a streaming batch is charged.

WriterResourceConsumer’s ConsumerHandle charges nothing: its staging chunks are governed allocations the ledger already holds, so a handle charge would count every staged byte twice. It stays registered for its inventory row. It is never backpressureable: parking the synchronous writer would prevent its own release progress. Spill requests are consumed at chunk boundaries, and cancellation is checked before consulting pause state. Grants and cleanup debt issued through the writer’s own admission keep the admission owner registered after the provider handle drops; the final owner unregisters it on success, error or cancellation while the run remains open. A lease issued through a Source’s attributed view is the exception: it releases through its reservation state and holds no release authority, so it does not keep the writer consumer registered. The accounting stays correct because that consumer’s handle charges nothing. Closing the run closes admission and unregisters the consumer even if an allocation escapes the run. Such an allocation retains only the synchronized release state and a weak arbitrator reference, so its eventual drop still settles the ledger without retaining the run or its telemetry producer. Cleanup debt remains visible until the resource is actually released; closing a run does not manufacture a zero balance. Prepared storage describes the separate disk and descriptor ownership.

These are requested-allocation bounds, not whole-process RSS bounds. Allocator metadata/rounding, thread stacks and native I/O internals remain outside them. Named fixed startup allowances are the standalone Arc/mutex control block and the executor authority, admission, consumer and storage control blocks. The environment-derived current-directory lookup is a temporary startup allowance; retained authored paths, path-construction envelopes and descriptor inventories are admitted separately. Each Source channel’s slot array (SourceIngestChannel::DEFAULT_CAPACITY slots, allocated once when the channel is built) is a fixed allowance, constant in input size: a queued attempt’s fixed shell lives in a slot, so only the heap the attempt holds outside the ledger is charged to its Source. CSV’s raw parser buffers and the full intermediate JSON tree used for JSON-encoded cells remain explicit parser allowances. Unchanged readers and legacy operators retain their existing owners; a later deep copy or spill reload is a distinct allocation, not an extension of the original grant.

Native JSON/XML configuration and schema caches

JsonEncoderConfig and XmlEncoderConfig share immutable admitted configuration. Names, envelope policies and schema-derived plans use ReservedText and ReservedVec; the shared configuration backing remains charged through its final alias’s actual deallocation. Construction admits the concrete writer and factory closure layouts before boxing. There are no raw JsonWriter or XmlWriter constructors: direct callers use the finite prepared APIs described in extension seams.

SharedStorageIdentity<Schema> gives each plan cache an opaque identity without retaining the schema payload. Governed storage uses its existing allocation ID; legacy storage retains a weak backing identity with a separately admitted, conservative backing reservation. It cannot upgrade to a strong owner. The weak backing is dropped before its reservation, so even the final weak alias retains accounting until physical deallocation. Equal-content schemas with different storage identities do not share a cached plan accidentally.

A schema change prepares the replacement while the old committed plan remains charged. Failed preparation drops the pending plan; only complete delivery commits it. Encoding borrows the record tree. XML scalar formatting uses a fixed 128-byte stack scratch and escaping uses bounded chunks; neither codec retains a rendered record between operations. Encoded bytes are owned by the shared operation stage and may spill under the same finite authority.

Strict UTF-8 input adds a four-byte probe/incomplete-scalar buffer per open. It changes no input-sized allocation ownership. Existing JSON parser/scanner and XML parser/event/record allocations retain their prior allowance; optional envelope indexes retain their existing cap. Prepared output does not establish whole-reader admission or a whole-process RSS ceiling.

Physical-text configuration, truncation tallies and trailers

FixedWidthEncoderConfig retains admitted derived layout policies, field and group names, and selected envelope names. Layout validation borrows the caller’s column/type trees; the encoder does not clone recursive schemas. Each FixedWidthEncoder owns a truncation: warn tally per warn field and a reused per-record staging array, both allocated once, when the encoder is built, and sized by the layout: a count, a longest length and eight record numbers per field. Preparing a record writes its hits into the staging array; commit folds them into the tally. Neither step allocates or requests budget, so a warn truncation cannot fail a record, and no value text is retained. A record that fails or is never delivered leaves only staged hits, cleared by the next preparation.

SwiftEncoderConfig shares admitted service literals or document-section names; literal precedence avoids retaining an unused section name. Column indices are resolved from the supplied schema without retaining its payload. The first successful body operation commits the document-derived trailer to SwiftEncoder; later record contexts cannot replace it. Literal trailers stay with the shared configuration. Pending trailer growth reserves the old and replacement backing simultaneously, and commit only moves ownership.

Both encoders borrow record strings and document fields. Scalar formatting uses fixed 1 KiB stack scratch, and fixed-width padding uses a fixed chunk. Large strings are not copied into a rendered cell before truncation. Encoded bytes belong to the shared operation stage; truncation tallies and retained trailers remain memory charges even when staged bytes spill. Factory and writer boxes keep the admitted concrete-layout owners described above. There are no raw fixed-width or SWIFT writer constructors: direct callers provide finite WriterResources to the encoder and PreparedWriter.

Strict fixed-width decoding validates each selected physical byte range; ignored gaps and tails retain their existing behavior. SWIFT validates raw block UTF-8 and preserves continuation separators. Its failed initialization releases partial body fields, service text and pending envelope events and remains terminal. These reader corrections do not admit existing line buffers, parser allocations, retained SWIFT fields or legacy document maps. Reader materialization and the remaining EDI writers are outside this writer guarantee.

CSV decoding and document ownership

Runtime CSV constructors select admitted decoding for both single-schema and multi-record input. DecodeWorkspace borrows valid UTF-8 and uses admitted scratch for Latin-1 expansion. Final text, keys, arrays, maps, positional values and schemas use allocation-owned storage. Single-schema repeated-cell parsing reuses the existing split grammar; final nested JSON values are constructed fallibly from the parser’s intermediate tree. Compiled multi-record CSV still rejects repeated input declarations; charset support does not widen that surface.

Multi-record capture retains admitted policy metadata, pending rows, section values and trailer state. FormatReader::prepare_document and the transport adapter return OwnedMap, so section ownership survives the reader boundary. Ingest moves sections into owned envelope values and a shared DocumentContext; source coercion preserves the original decoded record on failure and admits changed values before replacing it. A surviving record, document, string or key alias keeps the corresponding allocation charged. Unchanged format readers wrap their existing maps as legacy storage; the common carrier alone does not admit those allocations.

These boundaries do not establish constant-memory Source or Combine execution: those paths can still materialize whole inputs. Their remaining residency work is tracked in #1183. CSV’s admitted allocations and the existing ownership-relative estimates below must remain distinct from a claim about all readers or whole-process RSS.

Allocation-owned record storage

Governed text, positional values, ordered maps and map keys attach their lease to the allocation’s owner. A record or queue entry is not the lifetime boundary: a detached string or key can remain live after its original container drops. Shared text clones retain one allocation and one charge. A consuming container iterator retains the container charge until its backing is destroyed; yielded children keep their own independent owners.

Complete requested layouts are admitted before allocation, including text bytes, vector capacity, map entries and hash-table backing, and their owner holders. Map bounds follow the pinned container implementation and are checked against actual allocator requests. Growth admits old and replacement backing simultaneously. Budget refusal or allocation failure preserves the original container and returns an unconsumed insertion value.

Shared storage uses a sealed final-owner protocol: the last owner frees the shared allocation before destroying its payload, which in turn releases its lease after its children. Unique holders follow the same destruction order. The public APIs cannot extract an ungoverned backing allocation or grow a governed container without admission.

Physical heap estimates and admission contributions answer different questions. Physical estimates include governed storage for pressure decisions. Runtime admission uses unaccounted_heap_size with the executing run’s live allocation resources. It excludes a backing allocation only when its lease belongs to that same ledger, then classifies each child independently. A custom library source can supply governed storage from a different provider; that foreign allocation contributes its size to the Source’s charge for as long as it is queued in the Source’s channel, and the operator that keeps it downstream charges it under its own rule. The comparison reuses the existing authority identity and never transfers or releases a grant.

The separate legacy-only traversal describes storage representation and does not establish which run owns a charge. Neither traversal subtracts a global managed-byte total from a sampled consumer estimate. An ordinary deep copy or spill reload creates new legacy storage; it cannot reuse the original allocation’s charge. These primitives alone do not establish complete reader admission or replace existing parser-buffer allowances.

Fixed-row node-buffer and streaming estimates still count the record/identity pair and logical value slots; they do not add nested heap to that heuristic. Actual rows omit their value slots only when the vector itself belongs to the executing ledger. Each row is classified independently, including mixed-width batches. Producers and consumers use the same capability and calculate the cost before moving the row. A successful send transfers the producer reservation; the consumer may already have discharged its share before that send returns. Error drains subtract discarded rows individually instead of resetting the shared counter.

A consuming memory scan moves its original storage. A shared scan needs new value slots while retaining the original backing, and disk rows need their full reload forecast. Sort pressure thresholds therefore keep physical size distinct from resident attribution. During a streaming spill, the original run remains charged through serialization and destruction; each decoded row then receives its own charge before publication. Original shared leaves can outlive either representation under their intrinsic allocation grants.

Range-join output checks its initial spill-merge frontier before returning any row to a consumer. The physical footprint includes the open readers, their decoder workspace, merge entries, file inventory and retained metadata. If that frontier alone exceeds the run’s hard limit, execution returns a structured arena memory-budget error and releases its files and charges. The check applies to streaming, materialized, delayed and shared output drains, independently of the sink format. It does not claim pre-allocation admission or a bound on later decoder-table growth; resident attribution still uses the ownership-relative queries above.

Prepared-output telemetry

An executor provider supplied with the existing TelemetryProducer emits closed WriterAdmission, WriterStage, WriterSpill and WriterCleanup spans. Scopes are fixed literals: no filename, field value, owner ID or authored node becomes a metric dimension. Each scope has started/completed/failed/interrupted counters. Admission counts memory-grant attempts; a completed admission grants a layout and does not claim the allocator succeeded. Stage observation spans creation through successful readback and storage release, including metadata, progress-buffer and stage-box refusal. A resource failure is reported with its original kind; cancellation is interrupted, not failed. Cleanup continues under cancellation so resources can be released.

WriterStageDropped counts a stage abandoned without a terminal storage result; it does not infer destination success or failure. The existing SinkRecords, SinkErrors, SinkBytes and Sink lifecycle metrics own destination-level outcomes, so this primitive does not duplicate them. WriterSpillBytes counts bytes actually written to temporary storage, including partial writes. Cleanup attempt counters include retries; they do not claim remaining debt has been released. Read current debt and live bytes from the storage and ledger APIs.

WriterSpillCompleted establishes that an operation spilled successfully. The final execution report’s post-Source-join spill snapshot can be zero after those files have been released; it is not a cumulative spill-event counter. Source parsing failures remain data failures, while cancellation is interrupted. Reader value fidelity and failed-row selection are covered by the Source’s declared-column lineage mapping; writer representation and temporary staging add no dataset or syntax-token edges.

Each span is emitted once after its outcome, with both timestamps closed. The producer’s fixed counters coalesce; span admission may shed load. Full, sampled or contended telemetry never changes preparation, cancellation, delivery or cleanup, and never grows the arena. Producer clones share the existing telemetry arena/counter owner; their inline stage handles and observation state are included in the stage metadata grant. The provider’s one retained producer handle is part of its fixed control block. Format-only providers have no executor telemetry dependency.

Lineage is unchanged: these primitives alter no column values, row routing or dataset boundary. Raw temporary bytes are internal storage, not a new dataset.

Existing consumer attribution

Clinker tracks memory in two layers. The run’s one ledger of charged bytes decides admission and refusal; RSS (resident set size), sampled at chunk boundaries, is a second reading that the soft-threshold poll and the hard-limit backstops also trip on. Alongside RSS, every memory-touching operator except an authored Sort’s buffer, which registers no consumer of its own (Source ingest channels, Aggregate hash maps, grace-hash partitions, sort-merge accumulators, IEJoin arrays, inline-Combine hash tables, the Reshape per-group input buffer, node_buffers slots and their transient scan materializations, and window-runtime arenas) registers a MemoryConsumer wrapper with the pipeline-scoped arbitrator. Each operator owns its live byte counter and updates it on every admit / spill transition; that counter is the consumer handle’s charge on the run’s one ledger. Victims are ranked by reclaimable_bytes(), each consumer’s estimate of what a spill would free now, which the arbitrator reads per consumer at every policy poll and reclaim pass; a pass counts what each spill actually releases toward its target. A Source ingest channel’s handle charges exactly the heap its queued attempts hold outside the run’s ledger (foreign-provider or legacy storage, and a rejection’s box): each attempt carries that charge from just before it is sent until the walk takes it off the channel, or until it is dropped unconsumed. Its records’ admitted bytes are charged when they are allocated, in the Source’s name, so no byte is counted twice. This pull-mode attribution lets the policy distinguish reclaimable bytes (what an operator can give up right now) from currently-held bytes — a grace-hash Combine, for instance, reports only the partitions its build is still filling in memory, and nothing once its probe holds them, and the Reshape and Cull buffers report what spilling the groups still resident would free, each row counted as a node_buffers slot counts it, with rows on disk or taken out for processing counting 0, and a hash Aggregate’s table reports its charge while it can spill and nothing once its groups are on disk or a finalize has taken it, and never for a relaxed-key Aggregate’s table, which the commit retracts from in memory. The sort-merge and IEJoin kernels report nothing reclaimable: no pass can reach them, and each spills on thresholds of its own.

Registrations are scoped to the state they mirror, not to the run: each wrapper is unregistered when the state it attributes drains. A Source’s ingest-channel consumer is released when its stream ends: when the walk takes the reader’s Ended event, which the reader sends only after it has released everything it held (whichever arm consumed it — the Source arm, a fused Merge.interleave, or a fused Transform), or, for a Source the walk never drained to its end, when the walk stops; a Combine branch’s consumer is released when the branch exits — the IEJoin, grace-hash, and sort-merge branches route their clean return and every internal ? early-return through a single unregister, and the inline-hash branch unregisters at its clean exit; and a node_buffers slot’s consumer leaves the registry after its final planned reader. A consumer that collects a sequential scan into a resident vector carries an RAII materialization reservation, and every error return unregisters it. A Combine’s inputs keep theirs only until the rows each one charges have a new owner, are written to disk, or are dropped: the Combine’s kernel owns its inputs’ reservations and ends each one there, never at its own return. Every other operator that collects its input this way holds the reservation for its complete synchronous use: Route, Transform, Sort, Reshape, Aggregate, Cull, Merge and Sink release it when their dispatch returns, and an Envelope releases its body input’s when it returns and its header input’s once the headers are replaced. Composition input seeding transfers that same registration into the body-local node-buffer registry without an unregister/register gap or a second charge. While the body Source canonicalizes its seed, the same byte handle first reserves the prospective output in addition to the still-live seed, then drops back to the output estimate when the seed allocation is gone; admission atomically swaps the wrapper under the existing consumer id. Later stages therefore never see charged bytes from state that has already moved downstream, and the registry the policy polls contains live contributors only.

Window-runtime arenas (the columnar backing store that analytic-window evaluation reads from) are attributed but not independently spillable: an arena is immutable once built and is freed only indirectly, when the operator that consumes its windows drains to disk. Its wrapper reports the arena’s bytes so the arbitrator’s attribution is complete, but ranks last among spill victims so a policy never elects an arena while any consumer that can actually pause or spill remains.

Under dlq_granularity: document the run-scoped document state registers one consumer for the ledgers that record, per rejected document, which rows have been dead-lettered, so a row held by several Sinks is written once. A ledger is a compressed row set, one per Source and document, keyed by absolute row ordinal. It is exact dedup state, so it cannot spill and ranks last: each admission’s worst-case growth is charged through the ledger before the row is recorded, which on the walk first reclaims from every other walk victim; if that falls short the document state flushes its own held rows and retries once, and growth that still does not fit fails the run with an E310 naming the node and its set of rows already dead-lettered. At the end of every rejection pass, and every 65,536 admissions within one, the ledger is compressed and its charge replaced by a bound on the compressed heap; each Sink’s pass ends by settling every ledger it grew, the admissions of its late records included, so no per-admission charge outlives the pass that made it. The consumer is unregistered when the run’s context drops.

The same consumer carries the document state’s held rows (see Document dead-letter state under the bounded-memory contract below). While any held row is resident it reports priority 0, alongside the node_buffers slots, because it spills with one sequential write per document; its try_spill raises its spill request and reports the resident held bytes as what it frees. With no held row resident it reports the last priority and frees nothing, so a state holding only ledgers never shadows a consumer that can spill.

Per-operator arbitration parameters

Each registered consumer carries two parameters the active policy reads: a spill priority (lower is spilled first under Priority) and a back-pressure flag (whether its producer can be paused instead). The defaults are:

Operator classspill_prioritycan_back_pressure
node_buffers slot (inter-stage buffer)0false
rows parked for a deferred (relaxed-key) consumer0false
output staging (writer resources)0false
grace-hash Combine10false
Reshape15false
Cull15false
sort buffer / IEJoin build20false
sort-merge Combine25false
hash Aggregate30false
inline-hash Combine30false
Source ingestN/Atrue
streaming AggregateN/Afalse
credential registrylastfalse
transient scan materializationlastfalse
window arenalastfalse
document dead-letter state0 while it holds resident rows, else lastfalse

A consumer whose state cannot spill is listed as charged-only in crates/clinker-exec/tests/memory_consumer_inventory.rs with the approval that allows it, and a consumer only part of whose charge spills is listed in the partly-charged-only class with its approval. The document dead-letter state is in that class: its held rows spill, while its emitted-row ledgers and each failed document’s verdict slot are charged and never spilled.

Lower priority is spilled first. Among the spillable consumers, node_buffers slots (priority 0) are the cheapest victim class — spilling an inter-stage buffer to disk costs one LZ4 + postcard round-trip and frees the most reclaimable bytes per call. Output staging (writer resources) also sits at 0, but victims rank by an estimate of what a spill would free now (reclaimable_bytes), not by what they have charged, and a consumer whose reclaimable bytes are 0 is never elected: output staging (a fixed floor per writer), the inline-hash build side, the credential registry, the transient scan materialization, the window arena, a Source’s queued-event charge and the sort-merge and IEJoin kernels are counted toward the limit but never chosen as victims. A node_buffers slot ranks by its resident rows’ slot cost and the payload no Source has charged. Text a Source read is charged to that Source and left out of the figure, even when the slot holds its last copy and a spill would free it; the Reshape and Cull buffers, Output’s per-document buckets and a parked edge rank the same way, so a pass can spill one holder, find the target not yet met, and spill the next. That costs extra spill I/O, never a refusal: a pass keeps electing until what it measured covers the target. The blocking operators climb from there: a grace-hash Combine (10) is preferred over Reshape and Cull (both 15), which are preferred over a sort buffer (20), which is preferred over a sort-merge Combine (25), which is preferred over a hash Aggregate or inline-hash Combine (30). The sort buffer row (registered today only by the IEJoin kernel) and the sort-merge row give the order those kernels take once a pass can reach them; until then both report 0 reclaimable and are never elected, and the inline-hash row is never elected either. Reshape sits between grace-hash and sort because its spill round-trip re-runs synthesis on reload — costlier to evict than grace partitions, cheaper than an external-sort merge — and it spills the raw per-group input records rather than post-processed output. Cull shares Reshape’s priority for a similar reason: its grouped record buffer is costlier to evict than grace partitions, because reload re-splits the group, but cheaper than an external-sort merge.

A Source and a streaming Aggregate show spill_priority=N/A because neither operator holds spillable accumulated state. A Source’s try_spill always frees zero bytes — its only real lever is the pause its can_back_pressure=true advertises. A streaming Aggregate emits each group as it completes and never accumulates a spillable group table. The N/A here is about the operator’s own state, not its downstream handoff: when a streaming stage’s output rides a per-batch streaming handoff to a single consumer, that handoff registers a priority-0 consumer just like a node_buffers slot does, and its in-flight batches are spilled to disk one batch at a time if RSS crosses the soft threshold while they are in flight. So a streaming Aggregate’s group table is never a spill victim, but the batches it hands downstream can be.

Source-order barrier accounting

A Source that declares record-level sort_order is the exception to the usual pause-only Source shape: it inserts a verification barrier around each physical file before the ingest channel releases that file downstream. The barrier reuses the Source consumer’s live-byte counter. While a file is staged, that counter is the shared SortBuffer’s resident bytes plus one adjacent record retained for inversion detection, plus the charges of any verified records from the preceding file still queued downstream. During resident release, ownership moves from the sorter to an explicitly charged release total and then to the queued record’s own charge, the same exact per-attempt charge an unordered Source’s attempts carry. That last hand-off is one step on the counter, taken before the send, so the same row is neither omitted nor charged twice; a failed send drops the record together with its charge. The barrier moves the counter only by the change in its own figure, so it never overwrites the charges its queued records carry. Document punctuation does not grow with row count: admission allows only a flat file or one matching inner frame, so the barrier holds a statically bounded set of open/close events. Different physical files and different Sources never share a barrier or an authored-key comparison.

The barrier does not expose a second arbitrator victim. At the Source’s next record boundary it honors the shared consumer’s spill request, the sort buffer’s resident threshold, or the run-wide soft-pressure signal and calls the existing SortBuffer spill path. SortedRunMerger performs any bounded-fan-in cascade; each intermediate run is charged before its consumed inputs are unlinked and their exact charges are released. For final release, the barrier writes one merged spool, charges its exact completed size while the input-run charge is still live, then releases the input charge. It reads that final spool once to validate every row before emitting the file, and releases the final charge only after the second read has drained. Thus a decode or merge failure cannot leak a prefix of an unverified file, and disk accounting covers the input/output overlap at every completed-run transition.

The same Source counter remains live during release. A resident result moves from the remaining sorted-spool estimate to the bounded-channel estimate only after each record send transfers ownership, so the handoff never reports a zero-byte gap. A spilled result charges each merge reader’s 8 KiB I/O buffer plus one decoded-record estimate, along with the final writer or validation reader that overlaps it. Cleanup clears these transient charges and every outstanding stage disk charge on cancellation or read/write/merge failure.

Records spilled through this path keep their typed SourceRowId as the stable tie-break payload. Spill serialization reconstructs record-owned context, so the barrier carries the physical file’s original shared document-context handle once outside the row spool and reattaches that exact allocation to every repaired record on release. This preserves pointer identity without retaining one extra context per row and without replaying the source.

The soft-threshold poll (MemoryArbitrator::should_spill, called at batch boundaries) trips when the charged total or the process’s peak resident reading crosses the soft threshold (80 % of limit), and then takes two separate steps. Under a pausing policy (pause, both), reconcile_backpressure reads only the charged total: above the soft threshold it pauses one back-pressureable consumer that is not the Source being drained (its producer’s hot loop parks on a Condvar until resume), and below the resume watermark it resumes every paused one. The spill arm (poll_arbitration) asks the active policy for one victim among the consumers that can be paused or have reclaimable bytes and, when that victim cannot be paused, calls its try_spill, which raises its spill request for the operator to read at its next batch boundary. The poll runs no reclaim round and frees nothing itself, and a pause frees no charged bytes. A request on the walk that does not fit beside the charged total is refused with E310 only after the reclaim round described above. Three refusals run no round: a request made off the walk, a request larger than the whole limit, and Cull’s per-group decision checks. The hard-limit backstops (below) also refuse at once when the process’s peak resident reading passes the limit.

This means:

  • State that can spill lets a pipeline complete on input larger than the limit when disk space is available; state that cannot spill (see “How a pass reaches state”) still ends the run with E310 when it does not fit.
  • Performance degrades gracefully under memory pressure: while what the run holds can spill, you see slower execution (and possibly disk I/O), not failures.
  • Checked growth (reserve, try_grow, try_resize) never takes the charged total past the limit. Unchecked charges (set_bytes, add_bytes) and memory outside the ledger can pass it briefly before a poll or a backstop sees it.

Bounded-memory contract for non-fused stages

A stage runs streaming — no charged per-stage node_buffers slot — when it hands its output to a single downstream Sink and roots no window: fused Source → Transform → Sink and Merge.interleave-of-Sources chains, plus single-branch Route, non-fused Merge, streaming-strategy Aggregate, and hash-build-probe Combine probe-side feeding one Sink (see Streaming vs. Blocking Stages). The remaining boundaries — multi-branch Route fan-out, a Merge or other operator whose output forks to several consumers, Composition bodies, diamond DAGs, and every blocking strategy — materialize records into per-stage node_buffers. Each slot registers a NodeBufferConsumer with the arbitrator (priority 0 — the cheapest-to-spill victim class), so the active policy’s victim selection is fully attributed.

When a buffer crosses the soft threshold (80 % of the limit) the arbitrator runs the active policy. Under the default pause, the producer feeding the buffer is paused at its inbound channel; under spill or when no consumer can be paused, the slot spills to disk using the same LZ4 + postcard frame format as grace-hash sort partitions. Every hard-limit backstop (the inline hash build and its finished-table check, the hash and grace probe loops, the grace chunked fallback, the sort-merge emit loop, the range join’s finalize) goes through one check, MemoryArbitrator::check_hard_limit(node, surface, requester, uncharged), where uncharged is what the site is about to hold that no consumer has charged yet (0 for a check made after the fact). A hash table’s build counts its whole footprint, with two exceptions. A grace partition’s (CombineHashTable::build_from_charged) rows stay charged to the grace consumer, so only the index, chains and key cache it adds are uncharged. The inline hash join’s build rows stay charged under its build input’s reservation while the table is built (CombineHashTable::build_from_reserved): its periodic checks count the whole partial table, because the input and the table coexist, and its finished-table checks count the table less that reservation’s charge, which then moves to the table’s handle in one ledger step (ConsumerHandle::take_over), so each row’s slot is charged once. A row’s text its Source read stays charged under the Source as well as in the table’s figure (#1394). The other Combine inputs end their reservations where their rows leave them, so no backstop counts an input beside the kernel’s own charge of the same rows: the grace build’s and driver’s once the partition build and the probe loop have consumed them, the inline join’s materialised driver once its probe loop has, the range join’s once each side’s drain returns, and the sort-merge join’s in the ledger step that first charges each side’s rows to its own handle, the driver’s null-key rows included, which stay charged there until each miss is dispatched. That sort-merge step is not one of these backstops: what it adds to the input’s charge (the rows’ heap and the null-key rows) is admitted first as a checked growth of the join’s handle (try_grow), which reads only the ledger and on the walk reclaims before it refuses, so the ledger never passes the limit through the step, and a refusal is an E310 naming the join with the growth it asked for. The join’s own state reports nothing reclaimable, so such a refusal does not spill the side being taken over. The check decides which reading tripped first. When the charged total plus uncharged is over the limit it calls reclaim_before_abort: on the walk that runs the same reclaim loop a charge does (the requester elected last) and retries, off the walk it returns at once. Only the shortfall that ends that round refuses, as an E310 naming the operator and its surface, whose request is uncharged, whose floor is the charged total plus that request rounded up, and whose reclaim line is the round the walk ran. A refusal that asked for nothing reads the run held X, over memory.limit L, while <node> held <surface>. When only the process’s peak resident reading is over the limit, the check refuses at once with the process-memory form (LimitReading::ProcessMemory: the headline states that reading and the charged total instead of a full limit, and the suggested limit is the peak rounded up); no pass can lower a peak. reclaim_before_abort’s projection fits only when the charged total plus the projection fits the limit, so a projection of 0 does not fit a ledger an unchecked growth has carried past it. Checks that refuse one piece larger than the whole limit (the window index, the spill merge’s fan-in, the range join’s loaded block pair, an oversized aggregate row, a giant group) and Cull’s decision checks are not backstops and refuse without a round. The explain --code E310 diagnostic covers the full diagnostic model, including the composition-involved two-shape error model.

Every materialized slot is spill-eligible, including slots with several consumers and slots keyed by a producer output port. Pressure can therefore spill whichever live priority-0 slot the policy elects; exact (producer, producer_port) keys keep independently spilled Route/Cull branches isolated.

Every materialized slot declares an O(1) remaining-reader count when it is published. The single-reader path removes the authoritative slot directly. For several readers, dispatch remains sequential: each earlier reader borrows the same immutable backing and opens a fresh cursor, while the ledger retains the slot through its final reader. Memory backing clones one event at a time; Spilled and Mixed backing opens at most one spill file for the active scan, so file-descriptor use is O(1) per active scan rather than O(number of readers). No N-way copy or N-way cursor set is created.

A consumer that needs a full resident vector reserves that materialization before collecting the cursor. A projected overlap beyond the nonzero hard limit returns an E310 naming the consuming node and its rows collected for a full scan. A consumer that stays lazy, such as an Output writer on the envelope-reconstruction path, reads directly from the cursor without a full duplicate. The final reader reclaims the authoritative backing and its existing registration.

Document dead-letter state. Under dlq_granularity: document the run has one document dead-letter consumer, registered by the run-scoped document state. Its ledgers (above) do not spill. Every failing record of a failed document is encoded as its dead-letter row where it fails and held, behind a small header, in a per-document resident tail; no record is kept. The tails leave memory only on the arbitrator’s signals, never on a size of their own: a reclaim pass on the walk that elects the consumer while the state is between steps, which flushes them at once; the consumer’s election while the state was busy or by a round off the walk (its spill request, read before every held row and at every document decision); the soft threshold (polled every pipeline.batch_size held rows and at every decision); and a held row’s own admission when the walk’s reclaim leaves it short. Any of them flushes every tail to one chained-extent spill file in the run’s spill directory: each flush writes a document’s tail as one extent at the end of the file and links it from the document’s previous extent. A held row, its index entry (one per failed document) and, on a document’s first failure, the document’s slot are admitted in one checked growth of the consumer’s charge, which is their only charge: on the walk it first spills every other walk victim the pass elects, the document state (the requester) last, and only when that falls short does the state flush its own tails and retry once. A row that still does not fit fails the run with an E310 naming the failing node’s held failing rows. A flush is credited, in the spill quota and the run’s per-stage spill figures, to the failing node (for a pass’s flush, the node whose failure was held last); the rejecting Sink for a flush at a decision or a ledger admission; and, at the end-of-run sweep, the node that first failed the document, which the document’s failed verdict records. One flush writes every resident tail, so a node’s figure can include rows other nodes failed. Past max_spill_bytes a flush returns E320 naming the same node.

The Output that runs under the document granularity keeps each open document’s records in a bucket of its own, one NodeBufferConsumer per bucket registered under the Output’s name. A record’s bytes (what its run has not already charged) are grown through its bucket’s handle before the record is pushed, with nothing of the Output’s buckets borrowed, so the pass the growth starts can spill sibling buckets and every other walk victim, the growing bucket last; if that falls short the bucket spills itself and retries once, and a second shortfall is an E310 naming the Output’s rows held until their document is decided. A pass that elects a bucket spills its resident records as a new chunk after any it already has, recorded under the Output’s name; one that finds the Output’s buckets borrowed raises the bucket’s spill request, which its next push answers first. Until the soft-threshold poll is retired, a push while the threshold is tripped also spills the bucket. A document’s first rejection streams its chain row by row through its ledger into the dead-letter writer, in the order the rows were held. The file is removed with the state, and nothing in it is ever promoted.

Rows parked for a deferred (relaxed-key) consumer. A relaxed-key pipeline runs the steps below its relaxed aggregate at the commit, once per retraction iteration. An edge from outside that deferred region into it (a Source, Route branch, Cull port or composition-body node feeding a deferred Combine) cannot hand its rows over on the forward pass, so the producer parks a copy of them in the run’s parked-row store, keyed by the edge and the composition body it belongs to. Each edge has its own consumer, registered under the producer at its first park with the surface rows held between <producer> and <consumer> for commit (priority 0, can_back_pressure false). A park charges, before it copies any row, what the copy allocates or alone may keep alive with no other charge in this run: each row’s place and value slots, text a clone copies (unique text, governed or not), and shared text no admission in this run covers (text a computed expression built, or text another allocation authority admitted), since the copy may outlive the row that carried its only charge. Shared text this run admitted is not charged again: its admission travels with the allocation to every copy until the last one drops. The charge is a checked growth of the edge’s charge; on the walk that growth first spills what the pass elects, and when it still falls short the edge’s own resident rows spill and the growth is retried once, after which the borrowed rows are written straight to disk and no resident copy is made. A park is never refused for memory. The edge’s rows are kept as ordered segments, one per park, so arrival order survives any mix of resident and spilled segments; the edge ranks by the resident segments no open cursor shares, and any reclaim pass on the walk that elects it spills them (an election that finds the store busy raises the edge’s spill request, which its next park answers). Every retraction iteration publishes a fresh cursor over all the edge’s segments, in parking order, as the reading node’s input slot: the view adds no charge, the reader charges its own materialization as for any slot, and a spilled segment is read again from its file, its disk charge recorded once. Rows a region member parks during the commit pass, for a member of another region, form a generation of their own that the next iteration discards before it parks again; the commit walks each region after the regions whose members park rows for it. The store is released — every edge’s consumer unregistered and its spill files removed — when the commit returns, on success or error, and at the end of a walk that never reached the commit.

MergeSpilled is the one destructive spill form: its k-way merger consumes and unlinks input runs. On the first shared read, the executor folds those runs once into one ordinary re-readable spill file. It charges the replacement file before releasing the input-run charges, so the real disk-overlap peak is enforced; exceeding max_spill_bytes returns E320 SpillCapExceeded and removes both replacement and input registrations/files. Later readers reopen the folded file and do not repeat the fold.

Use clinker run --explain to predict which stages will dominate the budget before runtime — each node carries a buffer: streaming | materialized annotation. Materialized nodes charge pipeline.memory.limit as one full-stage slot and spill the whole stage; streaming nodes charge per in-flight batch and, on a single-consumer edge, spill those batches one at a time. Both classes count against the limit and can spill — the annotation tells you the granularity (whole-stage vs. per-batch), not whether a stage is exempt from the budget. Under dlq_granularity: document every Sink reports materialized: it holds a charged, spillable NodeBuffer bucket per open document, registered under the Sink, until the document’s verdict.

Reading --explain arbitration output

Alongside the buffer: class, every node in the Physical Properties stanza of --explain carries an arbitration: line giving the per-operator parameters the arbitrator would apply at runtime. The numbers are derived at plan time — --explain does no I/O, so there are no live consumers to query — but they mirror the runtime values exactly, so an author can read the spill/pause model before running the pipeline.

For a fast Source feeding a slow Aggregate (the canonical bounded-memory shape), the relevant lines read:

=== Physical Properties ===

source.orders:
  buffer: materialized
  arbitration: spill_priority=N/A, can_back_pressure=true, predicted_peak=1K, predicted_freed=0B, predicted_subtree_reclaim=1K

aggregation.dept_totals:
  buffer: materialized
  arbitration: spill_priority=30, can_back_pressure=false, predicted_peak=1K, predicted_freed=1K, predicted_subtree_reclaim=1K

The Source advertises can_back_pressure=true and spill_priority=N/A: when memory pressure rises, the arbitrator pauses the Source rather than asking it to spill (it has nothing to free). The hash Aggregate advertises the opposite — spill_priority=30, can_back_pressure=false — so it is a spill victim, ranked behind any cheaper consumer.

The three predicted_* values are the scheduler’s inputs (see Scheduling below). predicted_peak is the live volume a node is expected to hold at its peak — seeded at a file-backed Source from its path: file’s on-disk size and propagated forward. predicted_freed is what the node returns to the budget the instant it finishes draining: a blocking Aggregate holds its whole accumulated input (predicted_peak=1K) and frees it on drain (predicted_freed=1K), while a streaming Source carries the volume through but frees nothing the instant it drains (predicted_freed=0B). predicted_subtree_reclaim is the largest reclaim the node’s downstream chain eventually unlocks: the Source frees nothing itself, but launching it is the only way to reach the point where its downstream Aggregate can drain, so it inherits that Aggregate’s reclaim (predicted_subtree_reclaim=1K). Propagation of the subtree value stops at a convergence node — the Combine two independent chains feed — so each feeding chain keeps the distinct reclaim it owns up to the join rather than the shared post-join total. All three render 0B when no file-size seed reached the node — a multi-file (glob/regex/paths) or absent/unreadable Source, or any node downstream of one. The bytes are formatted in the same binary-prefix units as memory.limit (1K, 64M, 2G), and the same three values appear in --explain --format json under node_properties.<name>.predicted_peak_bytes, predicted_freed_bytes_on_complete, and predicted_subtree_reclaim_bytes.

A === Buffer Edges === section follows, listing the node_buffers slot between each pair of non-fused stages. Every slot is a priority-0, non-back-pressureable NodeBufferConsumer — the cheapest victim class — and the slot= number is the stable producer index the executor admits into. For a multi-output producer, port= completes the exact runtime slot identity. The slot carries the producer’s predicted volume (it holds the producer’s materialized output and frees that whole buffer once the consumer drains it):

=== Buffer Edges ===

edge source.orders -> aggregation.dept_totals:
  buffer: node_buffer (slot=0)
  arbitration: spill_priority=0, can_back_pressure=false, predicted_peak=1K, predicted_freed=1K (producer: source)

Reading top to bottom: under memory pressure the arbitrator first spills the inter-stage buffer (priority 0), then — if the soft threshold is still tripped — pauses the Source before it ever forces the Aggregate (priority 30) to spill. That ordering is exactly what the default pause policy (BackPressurePreferred -> Priority) encodes. Cross-reference the per-operator table to see where any operator in your own pipeline lands.

Scheduling

When a pipeline has several nodes that are simultaneously runnable — every one of their inputs is ready, so the executor could legally run any of them next — the engine picks one deterministically rather than walking topological position blindly. The common case is a single linear chain where only one node is ever runnable at a time, and there is nothing to choose. The choice matters only for a pipeline whose DAG has multiple independent subgraphs (for example, two unrelated Source → Aggregate branches that a later Combine or Merge joins): both branches’ lead nodes become runnable together.

The engine runs one node to completion before dispatching the next. When two independent chains converge — two Source → Aggregate branches a later Combine joins — both branches’ outputs must be materialized and held until the Combine consumes them, so the chain that runs second builds its working set while the first chain’s output already sits in a buffer. Running the memory-heaviest chain first therefore drains and releases its large state before the lighter chain’s output has to coexist with it, lowering the peak resident working set; running it last makes its large state coexist with the already-materialized output of every chain that finished before it. What the ranking also buys is when the frontier offers a mix of node kinds: with a blocking operator ready to drain (and reclaim its accumulated state) alongside a fresh Source about to charge a new buffer, draining first reclaims headroom before the new charge lands, and under a tight budget the engine prefers the runnable node that fits the remaining headroom over one that would overflow it.

The engine ranks the simultaneously-runnable nodes by these rules, in order:

  1. Headroom fit. A node whose predicted_peak fits within the budget’s remaining headroom is preferred over one that does not. Running a node that fits avoids tipping the live working set over the soft threshold and forcing a spill that a different ordering would have avoided. A node with an unknown peak (predicted_peak=0B — no file-size seed reached it) counts as fitting, because 0 is always within any headroom; this keeps an unestimated pipeline on its topological order rather than deprioritizing every node.

  2. Immediate-freed tiebreak. Among nodes that fit equally, the one with the larger predicted_freed runs first. Finishing a node that returns more bytes to the budget the instant it completes maximizes the headroom available to everything still waiting — the same intuition as shortest-remaining-state-first. A ready blocking operator (which reclaims its accumulated state now) therefore wins over a fresh Source (which frees nothing the instant it drains), because the immediate reclaim is the headroom-minimizing choice.

  3. Subtree-reclaim tiebreak. Among nodes that also tie on immediate freed — most importantly the fresh Sources of independent chains, which all free 0 the instant they drain — the one with the larger predicted_subtree_reclaim runs first. This front-loads the chain whose completion eventually frees the most: a Source’s value is the reclaim its downstream Aggregate will release, so the heavier chain’s Source is dispatched ahead of the lighter one even when it sorts later in topological order. Because it ranks below immediate freed, it never elects a fresh heavy Source over a ready light Aggregate (which would raise the peak), only between candidates whose immediate reclaim is equal.

  4. Stable-index tiebreak. If two nodes still tie (equal fit, equal immediate freed, equal subtree reclaim — including the all-unknown case where all are 0), the one with the lower stable node index wins. The index is each node’s position in the plan’s topological order — the exact sequence the executor walks the DAG — so this tiebreak is fully deterministic and independent of the machine, the thread schedule, and the order the runnable set happened to be assembled in.

Sinks after operators under document granularity. When any Source declares dlq_granularity: document, the plan orders every Sink after every other node, and the scheduler does not consider a Sink runnable until every other node of its pass has run. The four rules then rank the operators among themselves and the Sinks among themselves. Every place that condemns a document is an operator, so each document’s verdict is final before any Sink writes one of its records. The constraint overrides rule 2 where a Sink would otherwise run early to free its input, so every Sink’s input stays in its charged, spillable node buffer until the Sink phase. A composition body that holds a Sink is refused under document granularity (E378), because a body Sink runs inside its composition’s dispatch, where this ordering cannot reach it.

Fallback to topological order. When no node carries a volume estimate (every predicted_peak is 0B), rules 1–3 are no-ops — every node fits and every node frees the same 0 — so rule 4 alone decides, and the engine runs nodes in exactly the lowest-index / topological order it used before any volume estimates existed. This is the load-bearing guarantee: scheduling never changes record output or branching order. A pipeline’s data output is byte-identical regardless of the predictions; the estimates only steer which runnable node goes first to reclaim headroom sooner, front-load the heaviest chain, and prefer fitting nodes under pressure, never what each node computes.

Because the predictions are a pure function of the plan shape and the input files’ on-disk sizes (resolved against the pipeline file’s directory, never the process working directory), the scheduling decision is identical on every machine for an identical plan over identically-sized inputs.

Correlation Key Lifecycle & Rollback Narrowing

This page is the engineer’s reference for how a correlation key is born, carried, and unwound inside the executor. A correlation key declares a set of records from a single source as an atomic group: if any record in the group fails validation or processing, the whole group is sent to the DLQ. The mechanics here are the lineage substrate — $ck.<field> shadow columns, SourceRowId lineage, the per_source_rollback_cursors map, and Combine input snapshots — that the Retraction Protocol builds on. Where this page describes how identity is tracked and how a failure narrows the rollback, the retraction protocol describes how an aggregate refinalizes without DLQ’ing whole groups.

User-facing view: the User Guide’s “Correlation Keys” page.

Lifecycle

The engine adds a shadow column named $ck.<field> (one per correlation-key field) to every declaring source’s schema and copies the field’s value into it at ingest. From that point on, the shadow column is the authoritative group identity — if a downstream transform rewrites the user-declared correlation field, the shadow column is untouched and the group identity is preserved.

Shadow columns are an internal engine namespace. You never write $ck.<field> in YAML or CXL — the engine manages them. They are stripped from default writer output. To surface them for debugging, set include_correlation_keys: true on a Sink node:

- type: sink
  name: debug_out
  input: validate
  config:
    name: debug_out
    type: csv
    path: "./debug.csv"
    include_correlation_keys: true

A correlation key is declared per source: each source’s config: block carries an optional correlation_key: field naming the column (or list of columns) whose value identifies a record’s correlation group within that source. The engine widens each declaring source’s schema with one $ck.<field> shadow column per field and stamps the user-declared value into it at ingest. A record’s correlation group is identified by the tuple of values for that source’s listed fields; records sharing the same tuple within the same source belong to the same group. There is no pipeline-level correlation key. A source whose declared correlation_key: field names a column not present in its own schema: block is rejected at compile time with diagnostic E153.

Multi-source pipelines

Different sources can declare different correlation-key fields. The engine treats each source’s CK identity as locally consistent: a record from customers is a member of the customer-id group named in its row, and a record from orders is a member of the order-id group named in its row, regardless of whether customer_id appears in orders or vice versa. Combine and Merge nodes that join across sources negotiate which CK columns survive into the joined output via the Combine node’s propagate_ck: field (see Combine Join Strategies).

A source that declares no correlation_key: carries no $ck.* widening. Records from such a source flow through the pipeline without group identity; per-record errors DLQ on a per-record basis with no group fan-out. The orchestrator’s relaxed-aggregate retraction protocol still activates if any other source on the same DAG carries a CK field that an aggregate’s group_by omits — the retraction protocol scope is the DAG’s lattice of $ck.* columns, not any single source’s declaration.

DLQ semantics

When a record fails inside a correlation group:

  • The failing record produces a trigger DLQ entry. Its category reflects the actual failure (e.g. type_error, validation_failed).
  • Every other record from a source that contributed a trigger to the same group produces a collateral DLQ entry. Collaterals carry the category correlated.
  • Records belonging to other (clean) groups proceed normally.

A record with a null value for the correlation-key field is treated as its own per-record group: it has no peers and DLQ atomicity does not span multiple records.

A Combine output-row eval failure that the engine recovers from (under continue) produces entries under the combine_output_row category — distinct from the upstream-Transform type_coercion_failure because the entry carries the contributing-build lineage and rewinds both the driver and the matched build source’s rollback cursor. See Per-source rollback narrowing below for the cursor-rewind detail.

Every failure is one HeldFailure (executor/held_failure.rs): the row the failing evaluation is attributed to, written as its trigger, and, for a Combine output failure, the build row that contributed to it, stamped from the trigger with DlqFailureStamp::sibling. Without a key, write_failure writes it at once; under one, hold_failure_if_grouped parks it in the driver’s group cell. The build row travels inside its failure, so it is written only by the call that writes its trigger and can never name a trigger that was not written. It never makes a cell dirty, never widens its per-source narrowing, and never parks under the build record’s own key, so the build record’s own group is untouched. It did not fail, so its row stays out of the cell’s error_rows and out of the relaxed-CK retract scope: a relaxed Aggregate the build record also fed keeps its contribution. If the build record also reached a Sink under the same group, that Sink slot is spared or condemned like any other row of the group; a condemned one is not written a second time. At commit, write_held_failures writes every held failure in parking order through the same entries the keyless path writes: one trigger row per failure, then its build row. A row that failed twice, on two Route branches or against two build rows, is written twice.

The dlq_count counter sums triggers and collaterals, one trigger per failure.

Per-source rollback narrowing

When two sources contribute records to the same correlation group, a failure originating from one source does NOT collaterally DLQ records from the OTHER source. The collateral fan-out is scoped to the failing source’s records only.

Concretely, consider [src_a, src_b] → merge → tfm → out with both sources declaring correlation_key: id. A mid-stream Transform error fires on every src_b record but leaves src_a records untouched:

- type: transform
  name: tfm
  input: m
  config:
    cxl: |
      emit id = id
      emit ratio = if($source.name == "src_b") then (1 / 0) else amt

Under per-source rollback, the dirty correlation group for each id value contains:

  • One trigger DLQ entry — the src_b row that hit 1 / 0.
  • The src_a row sharing the same id is spared and reaches the output.

The engine identifies origin per record via the engine-stamped $source.name column. Within the failing source’s records, the existing CorrelationFanoutPolicy (Any / All / Primary) determines which records DLQ — the policy semantics are unchanged. Single-source pipelines see bit-identical behavior to the pre-narrowing engine because every co-grouped record shares the failing source by construction.

Records that carry no single-source attribution — synthetic aggregate emits and Combine output rows — are NOT spared by per-source narrowing. They flow through the existing collateral path because their stamp falls back to the merged-source identity which is ambiguous about origin.

The engine also surfaces a per_source_rollback_cursors map on the ExecutionReport, keyed by source name and carrying the highest source row number that cleanly exited a forward operator. The map advances per record at the clean exit of Transform / Route / Aggregate, and rewinds per contributing source on max_group_buffer overflow to the lowest row_num any group member of that source contributed, whether the member was buffered at a Sink or parked as a failure. A source whose records all DLQ lands in the map only through such an overflow rewind. The map is the replay anchor for per-source resume: a downstream rerun reads each source’s cursor as the floor for what must be reprocessed.

On max_group_buffer overflow, every record in the overflowing group still lands in DLQ (its parked failures as their own triggers, then one GroupSizeExceeded trigger over the remaining buffered rows as collaterals), but the per-source rollback cursor rewinds independently per contributing source, over parked failures as well as buffered rows. Attributing the overflow failure itself to one source would be a fiction — every contributing source shared blame proportionally — so the DLQ shape stays group-wide while the rewind narrows per source.

The relaxed-CK aggregator’s per-row lineage is keyed by SourceRowId, a typed identity made of the compiled Source node and the ordinal that Source minted for the row (executor/stream_event.rs). The Source half is load-bearing under multi-source ingest: each source numbers its rows from its own counter, so two sources that both feed the same aggregate group can contribute records with identical ordinals. Keying by Source as well keeps src_a’s row 1 distinct from src_b’s row 1 when both land in one group, so a retract that must remove both reaches each one instead of collapsing the colliding ids and stranding the second source’s contribution. The aggregator also stores each row’s source name beside its SourceRowId; retract matches on the SourceRowId alone and nothing currently reads the name.

Combine input snapshots are captured at fold start and cleared at every Combine arm’s exit (inline, IEJoin, GraceHash, SortMerge). When a Combine output-row eval fails recoverably under continue in the hash build-probe (inline) arm — a probe-key or on_miss: null_fields body failure on one driver row, or a residual-filter or matched body failure on one matched pair — the snapshot restores each contributing source’s rollback cursor to the value it held at the start of the fold (its pre-fold floor), lowering the cursor only if it had since advanced, then routes the row to the DLQ under the combine_output_row category. Only the sources that fed the failing row rewind; co-folded sources that did not contribute keep their forward progress. Under continue, the IEJoin, grace-hash, and sort-merge arms defer each recoverable output-row failure and route it through the same combine_output_row path after the kernel returns, while the snapshot is still installed. The build-side entry keeps the build record’s own row identity on every arm. See Combine Join Strategies for the per-arm execution detail.

Group buffering

The engine buffers records per correlation group until either the group completes (all source records observed) or a failure triggers a flush. The max_group_buffer: field on the pipeline-level error_handling: block caps per-group buffering across every source’s groups:

error_handling:
  max_group_buffer: 100000     # Default: 100,000

The cap counts held entries, not distinct source rows: every Sink slot a row occupies and every parked failure is one entry (CorrelationGroupBuffer::admit_entry). A Combine build row travels inside its driver’s failure and is not admitted as an entry of its own. The first admission that takes a group over the cap stamps the group’s overflow (overflowed_at), which becomes the failure stamp of its group_size_exceeded row, so that row’s timestamp and id record when the group crossed the cap rather than when it committed. The relaxed-CK archive and merge keep the earliest crossing. They carry held failures across retraction iterations with union_held_failures, which identifies a failure by its row, stage, route, message and contributing build row, never by its stamp (each re-dispatch stamps it afresh), and keeps each failure as many times as the larger side holds it: a failure observed again is held once, two failures of one driver against two build rows stay two, and a row that reached the same failing stage twice keeps both.

Crossing the cap changes nothing at admission. Later rows are still projected and buffered at their Sinks and later failures still parked, so the group keeps buffering until commit and the cap does not bound its memory today. Bounding memory at the cap is tracked separately.

At commit an overflowed group is DLQ’d entirely, as the dirty path would write it and more:

  1. Every held failure is written first, in parking order, by the same code the dirty path uses (write_held_failures): one trigger row per failure, with its own category, message, stage, route, triggering field and value, and failure stamp, followed by its contributing build row, which keeps its driver’s trigger id.
  2. The buffered rows not already written follow: the first becomes the group_size_exceeded trigger, stamped at the crossing, and the rest are correlated collaterals condemned by it. Overflow spares nothing, so neither per-source narrowing nor the fan-out policy applies. A row that is both a parked failure and buffered (an inclusive Route fan-out, a Combine driver failing on one match and succeeding on another) was written in step 1 and is skipped here.
  3. If no buffered row is left unwritten (a group of failures only), no group_size_exceeded row is written; the overflow is still counted by the clinker.correlation.group_overflows metric.

Every row goes through the ordinary DLQ push, so it counts toward dlq_count and the E315/E316 rate limits. It is not a hard error. The interaction with the per_source_rollback_cursors map on overflow is described under Per-source rollback narrowing above: the DLQ shape stays group-wide while the cursor rewind narrows per contributing source.

See also

  • The Retraction Protocol — how relaxed-CK aggregates refinalize affected groups instead of DLQ’ing them wholesale, and how the synthetic $ck.aggregate.<name> column lifts the post-aggregate retract path.
  • Operator Retraction Cost Reference — the per-operator memory/CPU footprint under retraction, plus the === Retraction === explain block and the run-time retraction metrics counters.
  • Combine Join Strategies — propagate_ck: semantics, match modes, and the per-arm snapshot/rewind behavior.
  • Memory Arbitration & Scheduling — how the RSS budget and spill thresholds interact with group buffers.

The Retraction Protocol

When an aggregate’s group_by omits a correlation-key field that is visible upstream, a single correlation group no longer maps cleanly onto a single aggregate group — one CK group can span many aggregate groups. The strict-collateral DLQ shape (roll back the whole group, including the aggregate output row) would then over-reject: a single bad source row would void an entire department’s total. The retraction protocol is the engine’s answer. It retracts only the failing records’ contributions and refinalizes the affected aggregate groups, so surviving contributions still produce a correct output row.

This page is the protocol itself, woven from three operator surfaces: the aggregate’s strict-vs-retraction path selection, the synthetic $ck.aggregate.<name> lineage column that lifts post-aggregate failures, and the buffer-mode window behavior. It builds directly on the lineage substrate in Correlation Key Lifecycle & Rollback Narrowing. For the per-operator memory/CPU footprint and the explain/metrics surfaces, see Operator Retraction Cost Reference.

User-facing view: the User Guide’s “Correlation Keys” / “Aggregate Nodes” pages.

Interactive companion: the retraction-loop explainer replays the commit loop phase by phase on three of the engine’s test pipelines — the aggregator state, the correlation buffer, the retract set and every retraction counter at each step.

Path selection: strict vs. retraction

The engine inspects each aggregate’s group_by against the upstream CK lattice (the union of $ck.* shadow columns visible at the aggregate’s input). Authors do not configure this — the engine inspects the configuration and picks the correct path:

  • group_by covers every upstream CK field — strict-collateral path. Each emitted row inherits the correlation identity of its inputs, the aggregate emits one row per group, and a DLQ trigger anywhere in the group rolls back the whole group including the aggregate output row. This is the zero-overhead default; strict aggregates short-circuit to the two-phase commit body and pay no retraction overhead.

    - type: aggregate
      name: order_totals
      input: orders
      config:
        group_by: [order_id]               # strict — covers the upstream CK
        cxl: |
          emit total = sum(amount)
    
  • group_by omits any upstream CK field — retraction protocol path. A single correlation group may span multiple aggregate groups; CK fields omitted from group_by stop being visible to downstream consumers of this aggregate’s output as user-named columns. The engine retracts only the failing records and refinalizes affected groups, so the aggregate output row reflects the surviving contributions.

    - type: aggregate
      name: dept_totals
      input: orders
      config:
        group_by: [department]             # retraction protocol is active
        cxl: |
          emit total = sum(amount)
    

On the strict path, aggregate output rows inherit the correlation meta of the records that fed them. If any input record in a correlation group fails, the surviving records in that group still flow through the aggregator and produce one aggregate row — but that aggregate row is itself DLQ’d as a collateral and never reaches the writer.

On the retraction path, the engine retracts only the failing records and refinalizes affected groups, so the aggregate output row reflects the surviving contributions. Operators downstream of a retraction-mode aggregate run only at commit time, on the post-recompute aggregate emits, and they re-run on every iteration of the commit loop below; only the last iteration’s results are written. A non-deterministic CXL builtin such as now is therefore evaluated once per iteration, and the written value is the last iteration’s.

E15Y: streaming incompatibility

The retraction protocol’s runtime constraint is enforced automatically once the engine has classified the aggregate. A retraction-mode aggregate is incompatible with strategy: streaming and is rejected with E15Y:

clinker explain --code E15Y   # retraction-mode aggregate incompatible with strategy: streaming

The reason is structural: streaming aggregates emit at group-boundary close, before the terminal correlation commit, and that early emit defeats the rollback window the retraction protocol depends on. There is nothing left to retract from once a streaming group has already emitted and been handed downstream. The engine selects the path from group_by content, so an author who writes strategy: streaming on what turns out to be a relaxed-CK aggregate gets the compile-time E15Y rather than silent incorrect behavior.

Reversible vs. BufferRequired accumulators

The cost of refinalizing a group depends on whether the accumulator can be run in reverse:

  • Reversible accumulators (sum, count, avg, weighted_avg, collect, any) carry a per-row lineage map (input_row_id → group_index) alongside accumulator state. A retract is O(retracted_rows) reverse-op calls plus one finalize_in_place. Per input row the aggregator keeps the row’s SourceRowId with its group index, plus the row’s source name, which retract does not read; --explain estimates the lineage at ~8 bytes per row. sum, avg and weighted_avg hold exact sums, which subtract exactly, so a retracted group finalizes to the bytes of a fresh fold over the surviving rows, at any memory limit. A sum whose every contribution is retracted has no value left and finalizes to null, as a group with no rows does.

  • BufferRequired accumulators (min, max) cannot be unwound by a reverse op — removing the current max, for instance, requires knowing the second-largest value, which the running accumulator never retained. They hold per-group raw contributions until commit and recompute affected groups from contributions − retracted_rows.

The full per-accumulator memory formulas live in Operator Retraction Cost Reference.

Synthetic correlation column

A retraction-mode aggregate emits one engine-managed $ck.aggregate.<name> column on its output schema, alongside the user-emitted bindings ([group_by_columns] ++ [emitted_binding_columns]). The column carries the aggregator’s per-group index at finalize and costs ~16 bytes per emitted row (the Value::Integer payload plus its slot overhead). It is hidden from default writer output, mirroring the source-CK shadow column posture, and lives outside any user-visible CXL surface — authors never write or read it.

The synthetic column is the lineage hook that lifts the post-aggregate retract path. Without it, a failure on an aggregate output row would have no way back to the source rows that produced it: the aggregate has already collapsed many source rows into one. The column lets the orchestrator’s detect phase decode the per-group index back to the contributing source row ids via the retained aggregator’s input_rows table, and the recompute phase then retracts those source rows just as it would retract a directly-failing source record — matching the upstream-failure DLQ fan-out semantic.

Where retraction triggers are sourced

Retraction handles failures on both sides of the aggregate, via two different lineage hooks:

  • Upstream of a retraction-mode aggregate (Source ingest, Transform evaluation, Combine probe, Validation): retraction is fine-grained. The failing record carries $ck.<field> shadow columns, the engine identifies its correlation group from those columns, and retract_row removes that record’s specific contribution from every affected aggregate group while leaving every other contributing record intact.

  • Downstream of a retraction-mode aggregate (a Transform that fails on an aggregate output row, an Output writer that rejects an aggregate row): the failing record carries the synthetic $ck.aggregate.<name> lineage column described above. The detect phase resolves that column to the contributing source row ids and feeds them into the same recompute pipeline as upstream failures.

Both surfaces converge on one recompute pipeline. The end-to-end demo at examples/pipelines/retract-demo/ runs both surfaces in one pipeline (a Transform failing on an aggregate output row alongside an upstream Transform error).

The commit loop

The deferred region downstream of each relaxed aggregate (its producer) does not run on the forward pass: the producer runs, keeps its aggregator state, and parks its output; Sinks outside a region buffer their rows in correlation cells keyed by the row’s $ck.* values, holding per-record failures there too. The forward-pass buffer is saved as a baseline. At commit (executor/commit/mod.rs):

  1. Detect. Cells holding failures are triggers. A source-key cell contributes its failing rows; a cell keyed by $ck.aggregate.<name> is decoded to the group index and expanded to every contributing SourceRowId. The result seeds the retract set.
  2. Recompute. This iteration’s new rows are retracted from every relaxed aggregate (a row the aggregate never saw is a tolerated “not found”), and every non-empty group is re-emitted in full; a group with no rows left is not emitted. With an empty delta this step is skipped.
  3. Dispatch. The region members re-run in topological order on the re-emitted rows; their per-record failures are held in buffer cells, and a failure that goes straight to the dead-letter queue (such as an aggregate finalize failure) is captured with its source row.
  4. Re-detect and expand. Detect runs again on the live buffer, its source rows are combined with the captured direct failures, and failures are copied to an error archive (messages only, no records). Rows not already in the retract set are the next iteration’s delta. If there are none the loop stops; otherwise the live buffer is replaced by the baseline and the loop repeats.
  5. Flush. The archive is merged back into the buffer and the cells are committed: a dirty cell dead-letters each failure as a trigger and its rows as collateral; a clean cell is written. Contributors retracted from a failed aggregate row are not dead-lettered themselves: they are simply no longer in any group.

The loop is bounded by plan nodes (composition-body nodes included) + source rows + 1 iterations, since every iteration must add at least one source row; reaching the bound is currently a panic! rather than a typed error (#1378).

Window interaction

When the pipeline has any relaxed aggregate, the planner marks every windowed Transform whose partition_by does not cover the window’s correlation-key set as needing a buffered recompute (requires_buffer_recompute), and --explain counts those windows as buffer-mode windows. What happens at run time: windows rooted at a relaxed aggregate are rebuilt from the re-emitted rows on every iteration of the commit loop, so a window after the aggregate sees the post-retract groups. Windows upstream of the aggregate are never re-evaluated, yet the mark still lifts the E150 check for Source-anchored windows that are never rerun (#1376).

Degrade fallback

The design: when retraction’s preconditions break at run time, the orchestrator degrades to dead-lettering the whole affected group, the strict-collateral shape. Today:

  • An aggregate whose recompute cannot proceed (no retained state, a failed retract, or a failed re-emit) has its output slot drained and is added to a degrade list, and degrade_fallback_count is incremented. Nothing reads the list, so the strict-collateral dead-lettering never happens: the aggregate’s groups are lost rather than dead-lettered, and a region member that needed the drained output can stop the run with an internal error (#1288).
  • A relaxed aggregate whose state spills cannot finalize in place; the run stops with an internal “spill failed” error instead of degrading (#1288).
  • There is no window degrade path.

See also

Operator Retraction Cost Reference

This is the capacity-planning reference for pipelines running the retraction protocol. An aggregate whose group_by omits any upstream CK field activates retraction automatically (see The Retraction Protocol for the path-selection rules and Correlation Key Lifecycle & Rollback Narrowing for the underlying lineage substrate). Each operator on the post-source DAG carries a different cost profile under retraction; the table below is the centerpiece — it summarizes the per-operator footprint so you can size memory and pick propagate_ck settings before pipelines hit production.

User-facing view: the User Guide’s “Correlation Keys” / “Aggregate Nodes” pages.

Per-operator cost table

OperatorRetraction cost
SourceNone at retraction time. The CK shadow columns are stamped at ingest; replay never re-reads the source file.
TransformRuns only at commit time on post-recompute aggregate emits when sitting inside a deferred region, once per iteration of the commit loop. Cost = O(rows_emitted_post_recompute) per region member per iteration, no extra state held. Non-deterministic CXL builtins (e.g. now) are evaluated once per iteration; the last iteration’s value is written.
Aggregate (strict, group_by covers upstream CK lattice)None. Strict aggregates short-circuit to today’s two-phase commit body and pay zero retraction overhead.
Aggregate (retraction-mode, Reversible bindings)Per-row lineage map (input_row_id → group_index) carried alongside accumulator state (each row’s SourceRowId with its group index; --explain estimates ~8 bytes/row) plus one synthetic $ck.aggregate.<name> shadow column on every output row at ~16 bytes/row. Retract is O(retracted_rows) reverse-op calls plus one finalize_in_place. Reversible accumulators: sum, count, avg, weighted_avg, collect, any.
Aggregate (retraction-mode, BufferRequired bindings)Per-group raw contributions held until commit, plus one synthetic $ck.aggregate.<name> shadow column on every output row at ~16 bytes/row. Memory cost = O(input_rows × Σ binding_value_size) plus the synthetic-column tail. Retract recomputes affected groups from contributions − retracted_rows. BufferRequired accumulators: min, max. A binding list with one of them puts the whole Aggregate on this path.
Combine (driver propagation)One propagated $ck.<field> slot from the driver record. No retraction state held by the combine itself; replay carries upstream deltas through.
Combine (propagate_ck: all / named: [...])Same per-row cost as driver propagation, plus the widened output schema’s $ck.<field> columns must be re-populated on replay. Cost scales with the output schema width, not retraction frequency.
WindowA window rooted at a relaxed aggregate is rebuilt from the re-emitted rows on every iteration: O(re-emitted rows) per iteration. Windows upstream of the aggregate are not re-evaluated, although the planner’s buffer-mode mark lifts the E150 check for them (#1376).
OutputHolds rows in correlation buffer cells until commit. Every iteration that continues restores the forward-pass buffer and re-runs the region, so rows are re-produced, not substituted in place; after the last iteration clean cells flush to the writer and dirty cells dead-letter. correlation_fanout_policy: all and primary currently behave like any (#1375).

Degrade fallback and metrics counters

The degrade fallback is designed to dead-letter the whole affected group when retraction’s preconditions break at run time. Today a degraded aggregate’s output is drained and the degrade list is never read, so its groups are lost rather than dead-lettered (#1288), and a relaxed aggregate whose state spills stops the run with an internal error instead of degrading (#1288). See The Retraction Protocol.

The metrics spool reports the run-time counters under its retraction object (see the User Guide’s Metrics & Monitoring page). Each is 0 on strict pipelines:

  • iterations — commit-loop iterations run (each recompute → dispatch → re-detect cycle).
  • groups_recomputed — rows re-emitted by relaxed aggregates during recompute, counting every non-empty group re-emitted, not only the changed ones.
  • partitions_dispatched — windowed-Transform member dispatches during the commit pass (one per member per iteration), not partitions.
  • degrade_fallback_count — aggregates that took the degrade path, per iteration.
  • synthetic_ck_columns_emitted_total — $ck.aggregate.<name> values written, on the forward pass and on every recompute.
  • synthetic_ck_fanout_lookups_total — aggregate-keyed trigger cells decoded back to their group, on every detect.
  • synthetic_ck_fanout_rows_expanded_total — contributing source rows those lookups produced.

Use the explain block (below) for plan-time capacity sizing, the metrics spool for post-run confirmation.

The === Retraction === explain block

Pipelines whose at least one Aggregate has a group_by that omits a correlation-key field get a === Retraction === block in the clinker run --explain text output. The engine selects the retraction-mode path automatically based on group_by content; the block is silent on every other pipeline, so strict-correlation and non-correlated --explain output stays identical to today’s text. (For the rest of the explain surface — buffer classes, arbitration parameters, the === Statistics === section — see the broader explain documentation.)

The block opens with a one-line summary —

retraction enabled — N relaxed aggregates, M buffer-mode windows, fanout policy: <policy>.

— followed by one block per retraction-mode Aggregate and one per buffer-mode window index.

Per retraction-mode Aggregate the block reports:

  • the resolved accumulator path (Reversible or BufferRequired),
  • the per-row lineage memory cost (~8 bytes/row for Reversible, n/a for BufferRequired which holds raw contributions instead),
  • the per-aggregate synthetic-CK column and its ~16-byte/output-row cost,
  • the worst-case degrade fallback when retraction’s preconditions break at runtime.

Per buffer-mode window index the block reports:

  • the source name and partition_by fields,
  • the per-row buffer cost in Value slots over the index’s arena fields,
  • the worst-case partition memory ceiling under degrade.

Group cardinality is honestly surfaced as “unknown at plan time” — the planner has no group-cardinality side-table to consult before the run. Use this per-operator cost table and the per-row figures the explain block prints for capacity planning, then confirm the live shape via clinker metrics collect after the first production run.

See also

Combine Join Strategies

Combine is the N-ary record-combining operator: every input is declared up front, the where: predicate matches records across inputs, and the cxl: body shapes the output row. This page covers the parts an engine engineer reaches for when reasoning about how a Combine executes — the strategy selection the planner performs from the predicate shape, the heap cost of materializing build sides, how reconciled document boundaries flow through every join path, the runtime mechanics of correlation-key propagation across all four execution paths, and the join-planner statistics catalog that drives build-side selection and grace-hash partitioning.

User-facing view: the User Guide’s “Combine Nodes” page.

Predicate classification

The where: expression is a CXL boolean evaluated for every candidate record pair across inputs. The planner splits a compound and-predicate into three conjunct classes, and the classification is what selects the execution strategy:

  • Equi conjunct — a cross-input equality (a.x == b.y). Drives the hash lookup or the sort-merge join.
  • Range conjunct — a cross-input ordered comparison (a.start <= b.ts and b.ts <= a.end). Handled by the block-band IEJoin path, with or without an equi conjunct on the same input pair (IEJoin for pure range, HashPartitionIEJoin when an equality is present); a pure-range predicate (no equality) with a single range conjunct over ascending-presorted inputs can take SortMerge instead.
  • Residual conjunct — any other CXL predicate (intra-input filter, function call, and so on). Applied as a post-filter after the equi/range match succeeds.

At least one cross-input equality is required for every Combine, except for pure-range predicates, which IEJoin handles without an equi conjunct.

Match selection versus the projection body

Every strategy separates two steps: match selection (the where: predicate, including any residual re-check) chooses which build records pair with a driver, and the cxl: body projects each selected pair into an output row. Under match: first this distinction is load-bearing and is contractually uniform across the four physical join paths — in-memory hash build-probe, grace-hash, sort-merge, and the block-band IEJoin. The block-band path serves both IEJoin plan variants: pure-range (IEJoin) and equi+range (HashPartitionIEJoin), where equality is an added block-pair prune axis on the same machinery, so both are bounded and behave identically at the emit boundary.

  • One candidate order: build arrival. A driver’s candidates are its matching build rows in ascending BuildSeq, the position the build input delivered each row at. first takes the lowest, all emits in that order within a driver, and collect builds its array in that order, on every strategy:

    • the hash table’s collision chains append, so a probe walks a key’s build rows oldest first;
    • grace-hash mints BuildSeq where it receives its build input and carries it beside each row through partitions, spill files, recursive repartition and fallback chunks. A probe checks that its candidates rise in BuildSeq and fails as an internal error otherwise, so a table walked out of order cannot silently pick another row;
    • sort-merge replays its window ascending over inputs declared ascending on the range key, so window order is arrival order;
    • the block-band path (both IEJoin variants) keeps the minimum build input index per driver and sorts its output by (driver order, driver index, build index).

    The pick is therefore a function of the data and the build input’s order, not of the strategy, the predicate’s shape or the memory-derived block, spill or hash layout. Grace-hash’s block-nested-loop fallback, which a reloaded partition takes when it is still over budget and either the assigner is at its MAX_HASH_BITS (12-bit) cap or a split leaves the larger child with more than 80% of the parent’s rows (SKEW_REDUCTION_THRESHOLD), still decides first, collect and on_miss once per build chunk instead of once per driver, and its output order depends on which partitions spilled; both are known defects tracked separately.

  • One verdict per driver. Each candidate’s residual has one of three outcomes (PredicateOutcome in pipeline/combine_verdict.rs): true, not true (false or null), or failed. Every strategy feeds a driver’s outcomes to one DriverScan, which folds them into the driver’s DriverVerdict under its match mode; no strategy decides “did this driver match” from its own container. first is decided by the earliest candidate that is not “not true”: a true one is selected, a failed one is the driver’s only result, and a later candidate is outside the result. all acts on every true and every failed pair as it is found. collect writes no row once a candidate failed. Only a driver with no true and no failed candidate is a miss, and every on_miss path takes the MissToken only DriverScan::finish (or a driver whose keys admit no candidate) can make, so on_miss after a failure does not compile.

    • Hash, grace-hash and sort-merge walk candidates in arrival order and stop a first driver once its scan settles.
    • The block-band path visits candidates in block order, so each driver’s scan keeps the earliest candidate by build input index, true or failed, and the driver block’s finalize reads its verdict. A held candidate is charged in the match state’s held bytes; the fixed per-driver scan slot is charged with the driver block.
    • Under fail_fast a failure outside the result (after a first driver’s deciding candidate) never aborts the run.
  • The body runs once, as a post-match projection. A body that skips the chosen build (a filter that fails, an EvalResult::Skip) or defers on a recoverable error drops only that output row. The engine does not retry a later matching build, and it does not route the driver to on_miss: the driver matched the predicate, so it is not a zero-match driver. The verdict reads predicate outcomes only, so a body result never reaches it, under any match mode.

  • A recoverable failure belongs to one pair. On every strategy a failing residual or body on one (driver, build) pair that is part of the driver’s result writes one combine_output_row failure with that pair’s build row. A failing probe key or on_miss: null_fields body is one failure with no build row. Each pair is evaluated under the driver’s row, so a failure’s detail names the same row on every strategy.

    • The materialized hash loop writes each driver’s failures before it reads the next driver, so it holds none. The streaming probe cannot: its thread has no way to reach the dead-letter output, so it holds every failure until the thread joins and the dispatcher replays them. Under match: all that set grows with a hot key’s fan-out, so each held failure (the driver row, the build row and the failure itself, priced by held_record_bytes) is charged to the Combine’s consumer as it is appended and discharged as it is written, and the thread polls the arbitrator after each driver once 10,000 rows and failures have accumulated. The overshoot past the limit is one 10,000-entry window plus the failures of the driver in flight, never the whole stream’s.
  • The residual covers every conjunct the kernel does not. The block-band kernel verifies the first two range conjuncts itself and evaluates the residual when the predicate has a third range or any conjunct that is not a range (DecomposedPredicate::residual_exceeds_ranges), so no part of where: goes unapplied.

Strategy hint

The strategy config field carries a hint; the planner has final say.

ValueBehavior
auto (default)Planner picks a strategy from the predicate shape. Pure equality: HashBuildProbe, or GraceHash when the build side’s estimated size approaches the memory limit. Equality plus range: HashPartitionIEJoin. Pure range: SortMerge when there is exactly one range conjunct on plain fields of an integer, float, date or datetime axis (not decimal, not mixed integer/float) and both inputs arrive globally sorted ascending on those fields (a declared sort_order whose first key is that field, on an input that is not partitioned per file); otherwise IEJoin. Both IEJoin variants run on the same block-band executor.
grace_hashForce grace hash join (disk-spilling partitioned hash). Applies only to pure-equi predicates; ignored on predicates carrying range conjuncts.

grace_hash is the right hint when build-side inputs may be larger than the memory budget but fit on disk after partitioning. It is a behavioral switch, not only an assertion: the hash-versus-grace choice is made once, at plan time (grace_hash_should_fire in clinker-plan/src/plan/combine.rs), and there is no runtime fallback. An in-memory HashBuildProbe build that outgrows the budget returns CombineError::MemoryLimitExceeded, which the dispatcher maps straight to PipelineError::MemoryBudgetExceeded (E310) where it calls CombineHashTable::build in executor/combine_dispatch.rs; nothing converts the join to grace hash at that point (#1337). strategy: grace_hash pins the spilling strategy regardless of the plan-time size estimate.

The choice of in-memory hash versus grace-hash for a pure-equality Combine is driven by the build-side row-count estimate (see Join-planner statistics below): a build side large enough to risk overrunning the memory limit is what tips the planner from the in-memory hash strategy to the disk-spilling grace-hash strategy.

Memory considerations

Build-side inputs are materialized in memory as hash tables keyed by the equi columns. For each non-driving input, plan for roughly 1.5–2× the raw CSV size in heap. A 50 MB product catalog typically occupies 75–100 MB of hash-table memory — the multiplier covers the per-key Value boxing, the bucket array overhead, and the per-entry chaining structure on top of the raw payload bytes.

This heap cost is the quantity the memory arbitrator charges against pipeline.memory.limit. For an in-memory HashBuildProbe Combine, exceeding the limit during build is an E310 abort; it is not flipped to grace-hash spill at runtime (#1337). See Memory Arbitration & Scheduling for the spill thresholds, the back-pressure knob, and strategy overrides.

Block-band IEJoin: bounded on both input axes and the output

Interactive companion: the range-join explainer runs this path on small inputs: the sort and slice into min/max-tagged blocks, the block-pair pruning, the IEJoin kernel step by step, the nested-loop fallback, and match:/on_miss:, checked against a brute-force join.

The block-band path that serves every range and equi+range Combine is bounded on both input axes and the output axis, so a range join whose inputs — or whose result — exceed the memory budget spills and completes rather than failing:

  • Input axes. Each side is drained into a payload-ordered, spillable sort buffer, then sliced into contiguous, key-sorted, min/max-tagged blocks. Blocks stay resident under a shared budget and spill to their own files past it; the scheduler joins one block-pair at a time, so at most a bounded working set (the resident blocks, one loaded spilled block per side, and the kernel’s O(n) sort arrays) is live. Equi+range adds an equality-hash prune axis on the same block machinery. Both deferred on_miss piles (NULL/non-orderable-keyed drivers and in-block zero-match drivers) drain to their own spillable, driver-index-ordered buffers rather than resident vectors.
  • Numeric range axis. Each range key reduces to a single order-preserving i128 so the IEJoin/block-band kernels compare, sort, and bounds-prune one integer per axis. The reduction is chosen per axis at plan time from the CONCRETE operand types: integers widen directly, finite floats take the orderable-float transform, dates take their day number (num_days_from_ce), datetimes take their canonical i128-nanosecond key (timestamp_seconds × 10⁹ + subsecond_nanos, exact and injective for the full calendar range — no microsecond rounding, no saturation outside 1677–2262), and decimals are placed on a shared 10^18 fixed-point grid (mantissa × 10^(18−scale)) — exact and scale-invariant (2.5 and 2.50 reduce identically). A mixed integer/float axis reduces both sides through the float transform, so the two encodings agree; a decimal value whose scaled magnitude overflows i128 or whose scale exceeds 18 fractional digits is not representable and aborts the run with E326 rather than dropping the row. An operand pairing that cannot reduce to one exact axis — an ambiguous Numeric operand (abs/min/max/clamp, whose runtime value may be an integer compared exactly or a float) or a non-orderable pairing — is rejected at plan time with E327 rather than routed to the axis.
  • SortMerge key encoding. The alternate pure-range strategy (SortMerge, chosen for a single-inequality predicate on pre-sorted inputs) maps each range Value to a monotone key WITHIN its own type: integers, floats (orderable transform), and dates (day number, num_days_from_ce) reduce to an i128 key, and datetimes reduce through the SAME canonical i128-nanosecond key the IEJoin axis and the memcomparable sort bytes use, so a presorted (SortMerge) and an unsorted (IEJoin) datetime run agree row for row. That within-type reduction is correct only for a same-type pair, so strategy selection restricts SortMerge to Integer/Float/Date/DateTime axes; a decimal axis (which it would map to a constant) and a mixed int/float axis (disjoint per-type encodings) route to IEJoin instead, which reduces both sides to one i128 space. The result is identical either way — only the physical strategy differs.
  • Output axis. Emitted rows accumulate in a payload-ordered sort buffer keyed on (driver order, driver index, build index) that spills on its own byte threshold, so the O(N·M) result of a high-fan-out join never has to fit in RAM. The final order is realized by that (possibly external) sort, so the output is byte-identical across memory budgets.
  • Dense block pairs. A surviving block-pair normally materializes its candidate vector before emitting. The kernel is asked for at most max_pairs = (hard limit − reserved pair peak) / 48 pairs (16 bytes per pair, times three for the vector’s growth), and stops the moment the pair’s actual match count would exceed that, without allocating the over-budget vector (iejoin/block.rs, the materialize-vs-fallback decision). A selective pair (matches far below driver × build) takes the materialized path however large its blocks are; only a genuinely dense pair, such as a single hot equality value whose block-pair is a near cross product that finer slicing cannot reduce, is instead streamed through a bounded block-nested-loop that buffers only a small tile of candidate index pairs (at most min(headroom / 16, 4096)), emitting each match immediately through the same spillable output sink. Extra residency is O(tile), so the pair completes within the budget instead of aborting. The nested-loop path is byte-identical to the materialized path (the output re-sorts, so candidate arrival order never shows) and preserves match: first/all/collect, the residual filter, and on_miss exactly.

max_output_rows, when set on the node, is a strategy-agnostic result-size runaway guard enforced at each strategy’s output-emit chokepoint — the spill-backed paths (block-band IEJoin, sort-merge) check the output sort’s cumulative total_rows() before each push, and the RAM-vec paths (hash build-probe, grace-hash) check the output vector’s cumulative length. On a breach the run fails loud with E325 rather than truncating. Because every emitted row on every strategy passes one such check, the cap covers all match modes, the deferred on_miss rows, and both the materialized and nested-loop block-pair paths. It is a result-size guard, orthogonal to the byte-based spill / RSS machinery above (which answers memory pressure by spilling or a typed budget abort, not by capping rows). Recoverable dead-lettered rows never reach the emit chokepoint, so they are not counted; under match: collect the counted rows are the per-driver output rows.

Document boundaries

A Combine forwards reconciled document boundaries to its output on every strategy — the inline hash build-probe, IEJoin, grace-hash, sort-merge, and the streaming-probe path. The boundary semantics are uniform across the strategy matrix, so a downstream operator never has to know which join algorithm ran.

Concretely:

  • A per-document Aggregate downstream of a join flushes per document. A driver source that carries several documents (a glob: over monthly files, say) produces one roll-up per driver document after the join, not one fold spanning all of them.
  • A document that spans both join inputs — the same document carried on the driver and on the build side — opens and closes exactly once downstream. The boundary is reconciled, never double-fired: the join does not emit a separate open/close for the driver-side and build-side appearances of the same document.

This reconciliation is what lets the per-document aggregation model compose with joins without special-casing the operator order.

Correlation-key propagation

Combine declares which correlation-key columns its output rows carry via the required propagate_ck field. The choice shapes both the compile-time output schema and the runtime record builder — those are the two internal surfaces an engine engineer touches when changing CK behavior.

propagate_ck valueCompile-time output schemaRuntime record builder
driverCarries only the driver input’s $ck.<field> columns.Build-side records contribute body fields only; their CK identity is consumed by the match and not copied onto the output row.
allCarries every input’s $ck.<field> columns.Copies build-side CK values onto each output row alongside the body’s emit columns. Use when the build side carries CK fields downstream operators must read.
{ named: [<field>, ...] }Carries the explicit subset, intersected with what is actually present upstream.Copies exactly the named subset. Use to project a multi-field CK down to a single field after a join.

Driver wins on a name collision. If both the driver and a build input declare $ck.<field>, the column appears once on the output schema and the runtime keeps the driver’s value.

propagate_ck is required on every Combine; a pipeline without an explicit value fails to compile.

Match-mode interaction across the strategy paths

The propagation contract holds identically across the hash build-probe, IEJoin, grace-hash, and sort-merge paths — the record builder is shared, so a build-side CK value lands on the output row the same way regardless of which algorithm produced the match. The interaction that does vary is by match mode rather than by strategy:

  • match: first / match: all — each emitted row is one driver × one build pairing, so the propagated $ck.<field> slot holds a single value (the driver’s, or the build’s, per the table above).
  • match: collect — the propagated CK slot is single-valued (it tracks the driver’s correlation-group identity), while the collected array column preserves the full lineage of every build match. The single-valued slot and the array column are distinct: the slot answers “which correlation group does this output row belong to,” the array answers “which build records were gathered.”

See the User Guide’s correlation-keys reference for the per-mode lifecycle narrative; the lifecycle and rollback-narrowing mechanics on the engine side are in Correlation Keys: Lifecycle & Rollback Narrowing.

Join-planner statistics

When the plan carries column statistics, --explain ends with a === Statistics === section listing the planner-wide statistics catalog. These are the figures that drive build-side selection and grace-hash partition-bit choice, so they belong to the join planner. Every figure is tagged with its provenance, so a metadata-derived estimate is distinguishable from a record-exact measurement.

Row counts — [file metadata] vs [exec sketch]

One line per source node, for example:

orders: ≈90 rows [file metadata] (informs combine build/probe + partition bits)
  • A [file metadata] figure is derived at plan time by dividing the input file’s on-disk byte length by an average-record-bytes constant, before any record is read. This is the same row count that drives a Combine’s build-side selection and its grace-hash partition-bit choice. A build side large enough to risk overrunning the memory limit is what tips a pure-equality Combine from the in-memory hash strategy to the disk-spilling grace-hash strategy.
  • An [exec sketch] figure is the exact count a source measured during a run, superseding the plan-time estimate.

Row counts also appear inline on each Combine’s driving and build inputs (est. 90 [file metadata] rows).

Column sketches — distinct counts, heavy hitters, membership filters

Three sketch kinds are populated by operators while records flow. All three are maintained by the grace-hash Combine over its build-side join keys, recorded under the build input’s (node, column):

  • Distinct-count estimate — product_id: 12,431 distinct [exec sketch].
  • Top-k heavy-hitter list with lower-bound counts — product_id: heavy hitters [exec sketch, lower bound]: widget=9,000, gadget=3,200, .... The list is a lower bound on frequency: a value absent from it may still be frequent, so it is only ever used to promote a key, never to exclude one.
  • Membership filter — product_id: membership filter, 119048 bits / 7 probes [exec sketch, sized from estimate]. Sized up front from the build node’s plan-time row-count estimate, built in the single build pass with no per-row buffer, and skipped entirely when no plan-time estimate is available.

Honest nulls and missing sections

A statistic that was never gathered renders as null rather than a fabricated zero. A plan over sources whose sizes cannot be read — a glob/regex multi-file source, a network source, or a missing/unreadable file — adds no Statistics section at all, and (per the membership-filter rule above) skips the membership filter that the plan-time estimate would have sized. Confirm the live shape via clinker metrics collect after the first production run, since the planner has no group-cardinality side-table to consult before the run.

Merge & Back-pressure

Merge concatenates upstream branches that share a schema into a single stream. For the engine, the interesting surface is not the YAML — it is where Merge fuses into the Source ingest loop, where the seeded-interleave path deliberately opts out of that fused channel topology, and how back-pressure propagates (or fails to) through the bounded crossbeam channels behind each mode. This page covers those mechanics.

User-facing view: the User Guide’s “Merge Nodes” page.

Fusion of interleave over Sources

When every direct predecessor of an unseeded interleave Merge is an exclusively owned Source node, the executor fuses the Merge into the source ingest loop. Exclusive ownership means that each Source has exactly one outgoing edge and that edge targets this Merge. The predecessor channels are polled directly and Merge consumption proceeds at live ingest rate, with no intermediate buffering tier between the Source readers and the Merge arm.

This fused live-channel path is what makes end-to-end back-pressure possible across the Source-to-Merge boundary (see Back-pressure semantics below). The separate Merge-to-Sink streaming predicate can avoid materializing the Merge’s own output whether or not its inputs are fused. Only when both boundaries qualify does back-pressure extend from the writer through the Merge to the Source readers; see Sink Internals.

When the predecessors are not all Sources (e.g. Transform → Merge) or any predecessor Source also feeds another consumer, fusion does not apply. Eligibility is atomic: one shared Source sends every predecessor of that Merge through the materialized path rather than partially claiming the otherwise-exclusive receivers. The Merge consumes pre-buffered predecessor outputs in round-robin order, and live back-pressure across the Merge boundary itself is unavailable in that shape (though the upstream operator’s own bounded buffer still throttles its predecessors).

Seeded interleave

Snapshot tests and benchmarks that need reproducible cross-input ordering opt into a deterministic schedule via interleave_seed::

- type: merge
  name: combined
  inputs: [east, west]
  config:
    mode: interleave
    interleave_seed: 42

A seeded interleave bypasses the fused live-channel path entirely. Instead of polling predecessor channels at ingest rate, the Merge:

  1. Pre-buffers each predecessor’s full output into a Vec.
  2. Emits records in fastrand-driven order, seeded by interleave_seed.

Output is reproducible regardless of upstream timing. The cost is that the seeded path opts out of live back-pressure across this Merge — the buffers fill to completion before emission begins, so a slow downstream consumer cannot throttle the Source readers while those Vecs are still filling.

Back-pressure semantics

How a slow consumer or a slow upstream reader propagates back through the DAG depends entirely on the Merge mode.

concat

Each Source ingest thread pushes into its own bounded crossbeam channel, capacity 1024 events per Source. Peer sources produce concurrently up to that capacity — the dispatch arm consumes from inputs[0]’s channel before turning to inputs[1]’s.

Consequences:

  • Memory. A non-leading input can hold up to one channel’s worth of buffered records (1024) before its producer blocks. Multi-input concat over N Sources may carry up to (N − 1) × 1024 records in flight even while only one input is being drained.
  • Latency. A record produced by inputs[1] while inputs[0] is still draining will not reach output until inputs[0] finishes, regardless of how fast it was produced.
  • Producer-side back-pressure. When a non-leading input’s channel fills, its reader blocks at Sender::send, propagating pressure back to the upstream file/network reader. The upstream is throttled even though it is not the currently-consumed input.

concat is the right choice when downstream consumers depend on declaration-ordered records (snapshot tests asserting byte-identical output) or when inputs represent ordered time partitions that must remain contiguous.

interleave (unseeded)

Fused with exclusively owned Source predecessors, the Merge arm polls every predecessor’s channel concurrently. Live back-pressure flows end-to-end:

  • A slow downstream operator delays Merge consumption → the predecessor channels fill → the Source reader tasks block.
  • A fast input does not wait on a slow peer — the Merge schedules whichever channel has a ready record.

When predecessors are not all Sources or any Source has another outgoing edge, fusion does not apply: the Merge consumes pre-buffered predecessor outputs in round-robin order, and live back-pressure across the Merge boundary itself is unavailable, though each upstream operator’s own bounded buffer still throttles its predecessors.

Unseeded interleave is the right choice when end-to-end latency matters and the downstream consumer is order-insensitive (an aggregator grouping on a key, or a writer that does not assert on row sequencing).

interleave (seeded)

The seeded path does not preserve live back-pressure across the Merge. It pre-buffers each predecessor’s full output into a Vec before emitting in fastrand-driven order, so a slow consumer downstream of a seeded Merge will not throttle the Source readers while the buffers are still filling.

If you need both run-to-run determinism and live back-pressure, prefer asserting on the multiset of records rather than their sequence and use unseeded interleave, or fall back to concat over deterministically-declared inputs.

Sink Internals

Sink nodes are the terminal destinations of a pipeline. Authored type: sink deserializes to PipelineNode::Sink with SinkConfig, lowers to PlanNode::Sink, and executes through executor/sink_dispatch.rs. When the planner certifies a single linear producer feeding one Sink, the executor can take a streaming handoff that wires the producer arm to a dedicated writer thread through a bounded crossbeam channel and fires Writer::write_record per record, concurrent with producer emission. Other producer shapes materialize their output before the writer fires. This page covers the topology that selects the streaming handoff, its relationship to Source-to-Merge fusion, the back-pressure chain, the counter semantics that must match the buffered arm, and the per-format nested-value contract.

User-facing view: the User Guide’s Sink Nodes page.

Use-time filesystem containment

Plan-time path validation produces a ValidatedPath, but that proof alone is not durable: an ancestor can be replaced after validation and before the operating system opens the output. Output creation therefore passes through a second, use-time boundary in clinker-exec:

  1. Resolve the runtime policy before opening the leaf. Normal output opens use the filesystem observed from the retained parent handle; explicit profile names exist only for the qualification harness and must match that observation.
  2. Walk the destination ancestors without following symbolic links or reparse points and retain the destination-parent handle.
  3. Claim a hidden sibling reservation, then create a uniquely named hidden quarantine leaf relative to that handle with owner-only Unix mode and no-follow semantics. The final leaf is not opened or truncated during staging. Every disposition, including replacement, holds the reservation so concurrent runs cannot both mutate one destination. The reservation carries the owner PID and an exclusive non-blocking fs4 lock. An existing reservation is reclaimed only after a creation-grace interval and a successful lock, which distinguishes a dead owner from a live or newly starting publisher. Lock or initialization failure immediately removes any reservation created by that attempt.
  4. Retain the boundary in a run-scoped publication ledger while single, per-source-file fan-out, and split writers produce their bytes. After the executor succeeds, preflight every quarantine and destination before the first mutation, then promote each quarantine directly through the retained handles. Replacement is one atomic rename; the old final is never moved to a backup first.
  5. If any promotion or post-rename directory synchronization fails, stop and return a typed outcome containing the exact synchronized-visible, visible-unsynchronized, and unpublished sets. Already-visible finals are not rolled back, and unvisited finals remain untouched. Deterministic fault tests pin this partial-set accounting.
  6. After each successful promotion, remove its reservation and record the exact committed final path. Cleanup failure is typed post-publication debt naming the visible final and stale reservation. Metadata sidecars are serialized before this point and join the same ledger and collision namespace as data outputs. Cross-filesystem promotion is refused; it never degrades to copying through a visible final path.

Linux uses the locked nix filesystem bindings for openat, renameat / renameat2, fstatfs, and directory synchronization. macOS uses the matching libc openat, fstat, fstatfs, renameat, and renameatx_np(RENAME_EXCL) primitives. Windows opens the drive root with reparse-point-aware CreateFileW, then walks every descendant and opens every leaf relative to the retained handle with NtCreateFile. It inspects each handle, compares volume identity, and promotes relative to the retained destination handle with NtSetInformationFile. The logical ValidatedPath remains required at the public containment boundary on every platform. A promotion that made the destination visible but could not synchronize its parent enters the explicit visible-but-unsynchronized state; the ledger reports it as an operational failure and never reduces it to a warning or claims rollback.

The set protocol is recoverable, not globally atomic: individual renames are atomic, while a multi-file commit has an observation window in which old and new entries can coexist. An uncatchable process or machine failure may leave hidden .partial and .reservation entries alongside a subset of newly visible finals. The reported/path-observed set is the recovery record; reservation liveness is established by the lock plus grace rule rather than by silently deleting a file that might belong to a live run. Output publication does not use .backup entries.

Attempt-owned publication modes

Output publication is resolved from the strict [storage.publication] block before an attempt directory or output leaf is created. The run-owned attempt uses the invocation’s existing execution ID and registers primary, per-source fan-out, split, dead-letter, and metadata-sidecar artifacts in one bounded ledger. A bounded map of compiled destination parents lets one run target more than one directory without inventing a common filesystem root. Duplicate destinations are rejected before attempt creation.

mode = "direct" is the default. Each writer receives an owner-only file in the attempt directory on its destination filesystem. Publication synchronizes that file and promotes it by same-filesystem rename; it never copies and never falls back to another mode.

mode = "local_then_publish" requires local_spool_dir on a local filesystem. The writer first produces and synchronizes an owner-only local file. The publisher then copies it in bounded 1 MiB chunks into the destination attempt directory, synchronizes the destination file, and verifies both the checked byte count and BLAKE3 digest before marking the artifact ready. The local copy is unlinked only after the destination-owned manifest state is durable. The final leaf is still reached solely by destination-local promotion; no copy writes directly to a visible final. A copy, synchronization, digest, manifest, rename, or directory-sync failure retains truthful incomplete state and never changes the selected mode.

destination_profile is explicit: local (the default), nfs_v4_1, or smb_3_1_1. A detected share under local, or a detected protocol that does not match the qualified share profile, fails before publication effects. The probe distinguishes NFS, SMB/CIFS, and other network or userspace mounts; platforms that report only an undifferentiated remote drive fail closed for the qualified NFS and SMB profiles. The remaining strict keys are failed_retention_seconds, creation_grace_seconds, max_attempt_bytes, retained_byte_limit, retained_attempt_limit, min_free_bytes, sweep_entry_limit, sweep_byte_limit, and sweep_time_limit_ms. Their fixed defaults and hard ceilings are enforced during policy resolution; only failed-attempt retention permits zero. Resolution also requires sweep_byte_limit to cover max_attempt_bytes plus the bounded 4 MiB manifest, so a valid maximum-sized attempt cannot permanently stall cleanup paging.

Free space is observed once at admission and compared with the checked attempt estimate plus min_free_bytes. That observation is advisory. It reserves no blocks or quota, proves no completion guarantee, and does not suppress a later ENOSPC or EDQUOT from a write or synchronization call. Default attempt results carry logical execution/artifact IDs, logical leaves, and exact published, visible-unsynchronized, or unpublished states. Physical paths are available only through an explicit opt-in intended for sanitized diagnostics.

Aggregate retained-attempt admission acquires handle-relative lock files under each root’s internal .clinker-attempts namespace in canonical root order and holds them through inventory, eligible cleanup, limit checks, and creation of every execution root. The lock never occupies an author-addressable final leaf, and namespace paging recognizes it as internal metadata. Each new root manifest durably carries the admitted byte estimate before those locks are released. Inventory sums manifest-owned regular files across roots and charges at least one per-execution reservation until every artifact size is exact; simultaneous local-spool and destination quarantine copies still count physically. Missing or uninspectable ownership evidence is conservative debt, never a zero-byte assumption. Namespace enumeration is bounded by the publication policy’s fixed maximum, not the current desired retained count. A configuration downgrade therefore still returns physical attempts through advancing continuation tokens while also reporting that the aggregate count exceeds current policy.

The operator query recompiles with the same default anchor as run (the pipeline file’s directory when no base is explicit) and replays bounded file-source discovery before rendering per-source output paths. Execution ID used in a path template is a separate typed input from an exact purge selector, which lets expired cleanup reconstruct an execution-scoped root. Continuations cross the CLI as their canonical raw bytes; structured argument arrays are the authoritative automation surface and text commands apply platform quoting.

Run-owned manifests also replicate a bounded historical-root receipt into the stable pipeline root. The receipt is bound to the compiled-plan hash, stores only typed logical source name/path pairs needed by {source_file} and {source_path}, and records sorted path-free identifiers for the execution’s output and spool roots. It contains no direct deletion path. When live source discovery no longer finds a retained failure’s inputs, operator compilation re-renders the authored templates from the receipt, validates the resulting paths, and requires their identifiers to match exactly. The stable replica is itself ordinary manifest ownership: successful publication removes it, while bounded purge removes it only after the same ownership checks as other roots.

Remote filesystem qualification

Normal CLI output is admitted from the filesystem type observed through the retained parent handle, including NFS and SMB shares. The profile strings below are qualification labels, not pipeline settings and not runtime admission tokens:

  • linux-nfsv4.1-loopback-ci: a disposable GitHub-hosted ubuntu-24.04 VM, Linux kernel NFS client/server, NFSv4.1 over TCP, a hard mount, and remote locking without a local-only lock mode.
  • linux-smb3.1.1-loopback-ci: the same runner class, Linux kernel CIFS client, Samba server, SMB3.1.1, cache=strict, remote byte-range locking without nobrl, and strict synchronization without nostrictsync. The loopback client disables only client-side permission checks with noperm; Samba still authorizes I/O as the configured guest identity.

The dedicated CI matrix provisions each server and mount inside its runner, places the exported root on a mounted 64 MiB temporary filesystem, and executes both publication modes. Success pauses after the exact Complete manifest and final are synchronized but before normal attempt cleanup. The harness reopens both through the mounted client, releases cleanup, and then proves the final is still present while the successful attempt root and manifest are absent. This barrier is installed only through a programmatic Linux qualification API; YAML, the ordinary CLI, and environment configuration cannot enable it, and success still creates no receipt or sidecar.

The same local, identity-bound control pauses before copy, file sync, rename, and parent-directory sync. For every applicable mode/stage pair, the harness withdraws the exact NFS export or stops the exact Samba PID, force-lazy-detaches the client mount, releases the operation, observes a bounded non-success, restores and remounts the exact profile, and reopens its retained manifest. SMB remount uses a bounded retry while detached kernel client state is released. The harness then exercises bounded list, inspect, purge preview, and purge execution. A separate attempt fills the mounted bounded backing until the operating system returns ENOSPC; the final must remain absent and operator cleanup must remove the staging attempt. Deterministic EDQUOT coverage is recorded only as seam_covered unless a real quota is separately provisioned and observed.

Evidence uses clinker.filesystem-matrix-evidence/3. It records the runner and kernel, exact package and protocol observations, mount and lock behavior, the six unchanged edge outcomes, six lifecycle classes, ordered success and interruption readbacks, real capacity behavior, recovery, persistence, operator cleanup, and environment teardown. Teardown and bounded evidence/log upload run unconditionally for passing, failing, timed-out, and interrupted matrix cells. Legacy schema 1 or any missing, unknown, or truncated proof is ineligible.

The byte-range/OFD lock observation remains separate from publication admission. A dedicated production admission-lock section records independent test-binary processes calling RunAttemptPublication::create on the mounted profile with opposite multi-root order. Both the retained-count and retained-byte scenarios require bounded completion, exactly one admission, one rejection, and readback of the same retained execution from every mounted root.

A profile is support-eligible only when that exact cell writes status: passed, support_eligible: true, and successful client-mount, bounded-backing, service, and workspace teardown evidence. Missing packages or administrative capability, mount/provision failure, incomplete observations, a semantic failure, or cleanup failure leaves status: incomplete and cannot be interpreted as support. These loopback results do not certify a corporate share, vendor NAS device, Windows/macOS client, clustered server, or different server/mount configuration. They prove the implementation against controlled representatives. Operators of other shares must validate their server and mount semantics; runtime detection does not turn CI evidence into a vendor support claim.

Executor spill should remain on local storage for predictable bounded-memory performance. Local working data is distinct from destination quarantine: completed bytes still need to be streamed into a destination-local hidden file, synchronized there, and promoted on that same share. A local working copy reduces random network I/O; it cannot make a cross-filesystem rename atomic.

CSV preparation and writer ownership

The registry constructs CSV writers with the run’s finite WriterResources. Direct, split, per-source-file, correlation-deferred and combined split/fan-out routes retain the same admitted encoder path. CsvEncoderConfig owns shared policy; each encoder owns its schema mapping. CsvHeaderCapture shares only successfully committed header names. Shared backings remain charged until their final alias is destroyed.

CsvEncoder prepares into private operation storage. The first body record and automatic header form one operation; explicit document begin/end and finalizing operations have their own boundaries. A preparation failure writes none of that operation and leaves committed format state intact. Successful preparation is not a delivered record: delivery and storage completion precede encoder commit. Once delivery starts, a destination may accept a prefix before failure. The writer then rejects further operations, including flush, instead of retrying bytes. This boundary is distinct from quarantine-file publication above.

With repeat_header: true, rotated CSV writers replay the captured names. With it disabled, only the first file emits the automatic header; a failed writer construction does not consume that first-file policy. include_header: false and reconstructed envelopes do not invent an automatic header. Header capture commits after successful delivery and cleanup, not when header bytes are merely prepared.

FormatWriterHandle owns each boxed CSV writer, including outer counted and splitting wrappers. WriterFactory owns the concrete split-construction closure. Their external leases outlive box deallocation; internal buffers and captured state retain independent allocation owners. Explicit legacy handles preserve unchanged codecs outside CSV/JSON/XML without claiming admission. The extension seam describes the constructor contracts, and prepared storage describes spill and cleanup debt.

Memory, disk, allocation, descriptor and temporary-storage errors remain typed resource failures and are fatal even under strategy: continue. Cancellation is interrupted work, not a Sink error. A real failure the walk reached retains its classification when shutdown also occurs; a reader failure the walk never reached because the cancellation stopped it first is logged and the run stays cancelled. The CSV cell encoder preserves its original failure across the library writer’s drop-time flush. Malformed UTF-8 in source headers and body cells is an input-data failure, including schema discovery; unsupported authored encodings fail configuration admission instead.

The executable contract is covered by clinker-format’s encoding_contract and writer_preparation, the executor’s encoding_runtime_contract and writer_resources, and the CLI’s encoding_cli and sink_surface tests. The runtime tests observe actual spill completion and bytes; the CLI checks exact output, summary counters and diagnostic categories. Destination-prefix and cancellation fault guarantees come from the injected integration tests, not from successful CLI fixtures.

Native JSON/XML prepared output

JsonEncoder and XmlEncoder implement the same FormatEncoder transaction: prepare borrowed input into an OperationStage, seal, deliver completely, then commit pending framing/count/cache state. Begin-document, record, end-document and finalization each have their own operation. Memory, allocation and storage refusal before delivery leave committed state and destination unchanged. Readback, destination or cancellation failure after delivery begins can leave a prefix and poisons continuation. Drop never finalizes or retries. Draining bytes through flush_bytes never closes syntax.

The registry uses admitted immutable configs, FormatWriterHandle and WriterFactory for ordinary, split and per-source destinations. CSV/JSON/XML prepared writers directly wrap the destination, without an extra BufWriter: a buffered outer writer could count undelivered bytes or retry them on drop after poisoning. Other codecs retain their existing wrapper behavior. Sink byte metrics count actual accepted destination bytes; records count only complete delivered record operations. Cancellation is interruption, with zero Sink errors unless a real failure also occurred.

Ordinary NDJSON always emits one compact record and LF, regardless of pretty. Reconstructed JSON keeps its existing document grammar and pretty behavior; XML keeps native attributes, text, repeated elements and configured wrappers. No XML declaration is emitted. Empty library finalization and explicitly opened empty envelopes are supported, but a CLI run with no native body rows never opens a writer and publishes an empty staged file. Unsupported envelope/routing combinations fail configuration validation before execution.

Native source schema discovery and envelope pre-scan preserve input errors as format failures. The UTF-8 adapter marks its own encoding failures structurally across std::io::Error; unrelated transport errors are not relabeled by that adapter. Each physical open revalidates BOM/declaration policy. XML selection resets matched depth at the selected record’s closing event, so metadata and repeated containers do not inflate body cardinality. File dataset identity and declared-column DIRECT lineage remain unchanged; structural element names do not invent data-column influence edges.

Fault tests in writer_preparation and writer_resources cover all native operation boundaries, storage failures, exact destination prefixes, cancellation, cache rollback and telemetry admission loss. encoding_cli and encoding_runtime_contract cover literal files, publication, exit codes, Source read counts, and what a later file’s failure leaves: nothing published, and in the failed attempt only the rows a streaming Sink had already written. See native ownership for configuration, schema lifetime and the remaining reader allowance.

Streaming vs. buffered

When a single Sink sits directly downstream of an eligible linear producer, a bounded crossbeam channel connects the producer arm to the writer thread, and Writer::write_record fires per record as the producer emits. For a Merge.interleave whose direct predecessors are exclusively owned Sources, this combines with Source-to-Merge receiver fusion to form an end-to-end live path. Each Source must have exactly one outgoing edge, targeting that Merge; sharing any predecessor rejects receiver fusion for the whole Merge.

A shared Source does not necessarily materialize the Merge’s output. After the Source inputs materialize and the non-fused Merge reads those slots, an otherwise-eligible Merge with one downstream Sink can still hand its result to the writer without admitting a node_buffers[merge] slot. Explain therefore reports the shared Sources as materialized while the Merge may remain streaming; that label does not claim live back-pressure across the Source-to-Merge boundary.

When the producer-to-Sink edge is not certified for streaming, the producer’s output materializes before the Sink arm invokes the writer. With a fused Merge.interleave, that extra slot would break the live back-pressure chain at the Merge output. The streaming handoff avoids that slot. For a non-fused Merge, it still avoids materializing the Merge’s own output, but the already-materialized Merge inputs mean back-pressure cannot extend through to the Source readers.

The streaming path is selected automatically — there is no opt-in setting. Pipelines that don’t match the topology keep the buffered path.

Topology

- type: source
  name: src_a
  config: { type: csv, path: a.csv, schema: ... }
- type: source
  name: src_b
  config: { type: csv, path: b.csv, schema: ... }
- type: merge
  name: merged
  inputs: [src_a, src_b]
  config:
    mode: interleave        # required
- type: sink
  name: out
  input: merged
  config:
    name: out
    type: csv
    path: out.csv

Eligibility

Every condition must hold for the producer-to-Sink streaming handoff to engage. Source exclusivity is a separate condition for the Source-to-Merge boundary: if it fails, the Merge inputs materialize even though an eligible Merge-to-Sink handoff may still stream.

  • The Sink has exactly one incoming edge.
  • Its producer is a supported linear producer: a Merge, fused Source-to-Transform, single-branch Route, streaming Aggregate, or an eligible streaming-output Combine strategy.
  • The producer has no other downstream consumer besides this Sink, roots no node-anchored window arena, and satisfies its producer-specific streaming requirements.
  • The Sink is not in the init-phase ancestor closure.
  • The SinkConfig has no split: block — splitting writers manage their own file rotation lifecycle.
  • The writer is registered in the single-file writer registry (not fan_out_per_source_file).
  • No Source in the pipeline declares a correlation key or document-level DLQ, and no Sink reconstructs envelopes. Those paths own deferred or document-scoped writer lifecycles that are incompatible with the per-record writer thread.

For a Merge to receive directly from live Source channels as well, it must be an unseeded interleave and every direct predecessor must be a Source exclusively owned by that Merge. Eligibility is atomic, so one shared predecessor rejects receiver fusion for all of that Merge’s Sources (see Merge & Back-pressure).

Back-pressure flow

Across a certified producer-to-Sink handoff, back-pressure flows toward the producer. When the upstream boundaries are also streaming or fused, the chain continues to the Source reader:

writer slow → bounded crossbeam Sender::send blocks
             → producer arm blocks
             → Source channel fills (when the upstream boundary is fused)
             → Source ingest thread blocks on send

The bounded handoff channel between the producer and Sink (256 events) limits that edge’s in-flight data. With a fused Source-to-Merge boundary, it joins the existing bounded Source channels into a single pace-bound chain from the underlying Write sink back to the source reader. A slow file system, a saturated network sink, or a deliberately paced writer then slows the upstream readers rather than accumulating the producer’s whole output in a pipeline-internal Vec. When an earlier boundary is materialized, back-pressure stops at that boundary.

Counter semantics

Counter behavior under the streaming path matches the buffered Sink arm exactly:

  • records_written increments once per successful Writer::write_record call.
  • ok_count counts distinct source row_nums reaching the Sink.
  • dlq_count counts the Sink’s own dead letters, its CSV join_values on_conflict: error collisions under strategy: continue, the same way in both arms; the arms differ only in when a collision is pushed. The buffered arm pushes each collision at the record that collided, through the run’s dead-letter funnel: it is counted, written into its bucket’s staged DLQ file through the walk’s DLQ writer, and checked against the rate ceilings (E315/E316) before the next record. The streaming arm’s writer thread cannot reach the run-scoped counters, so it keeps its collisions in a pending list, capped at 65,536 entries (the cap fails the run with an internal error), and the walk pushes them through the same funnel, in arrival order, when it joins the thread at the end of its producer’s turn.

Stage metrics (SchemaScan, Write, Projection) accumulate into the same fields the buffered path uses. The walk folds the streaming task’s per-task accounting back into the run-wide totals when it joins the Sink’s thread at the end of its producer’s turn, so a streaming run and a buffered run over the same input produce identical counter output.

End of input and failure

The channel carries the producer’s events and then, only once the producer’s dispatch has returned Ok, one HopMessage::End that the walk sends after the producer’s turn. A streaming Sink closes its output, and with it any document syntax (a JSON array’s ], an XML wrapper’s end tag), only when it reads that End. A channel that closes without End (the producer failed, or the run was cancelled) is an incomplete input and finishes nothing: the Sink abandons its staged file unpublished, keeping the rows it had already written, without closing framing, and records the Sink as failed (clinker.sink.failed, at least one Sink error) or, on a cancelled run, interrupted. A failed run publishes nothing; its retained attempt holds exactly what each Sink had written.

A producer delivers every row it emitted before it returns an error, so the Sink meets its rows in data order. The walk joins the Sink’s writer thread at the end of its producer’s turn and settles the turn with settle_hop: a Sink that failed on a row its producer emitted earlier is the run’s error, and the producer’s later failure is logged on the walk thread with node and upstream. A Sink’s own failure stops the walk only when its producer’s turn failed too; beside a producer that finished, the walk goes on, so every other Sink still reports its own failure. A Sink whose input closed without End beside a producer that finished is an engine defect (PipelineError::Internal).

Memory, telemetry, and lineage

A Sink does not retain an unbounded private collection. Incremental paths hold at most the bounded handoff channel described above, plus, on the streaming arm, the capped collision list described under Counter semantics. Dead-letter rows the Sink pushes are written through the run’s DLQ writer, one fixed buffer per open DLQ file, and are not Sink state. Materialized producers charge their node buffer to the run-scoped memory authority, and an authored Sink sort_order uses the shared stable resident/spill sorter through the planning-owned PhysicalWriterBoundary. Document-DLQ, envelope, per-source, and correlation-deferred modes apply that boundary at their actual population grain; they do not create a second memory budget.

CSV adds admitted per-cell workspace, retained policy/header state and prepared operation storage under that same authority. JSON/XML add admitted immutable configuration, schema caches and prepared operation storage. Raw parser buffers, JSON parser intermediates and unchanged downstream copies remain outside this allocation guarantee. An explicit spill root permits prepared output to spill, but does not remove the finite memory needed for cell rendering and metadata. See CSV decoding and document ownership for the remaining Source/Combine residency boundary.

An XML source feeding a Sink still follows the XML reader’s two-pass contract. The envelope pre-scan and body stream each open the finite source independently; the pre-scan retains only planner-attributed $doc.* subtrees and charges them incrementally against max_index_bytes (64 MB by default). Unreferenced XML is event-walked and discarded, body records stream one at a time, and a path-backed file that changes between the two opens fails instead of combining metadata and body from different bytes.

Telemetry is attached to the runtime owner that establishes each real Sink work unit: the synchronous dispatcher, streaming writer thread, or deferred correlation commit. The closed metric set is clinker.sink.started, exactly one of clinker.sink.completed, clinker.sink.failed, or clinker.sink.interrupted, plus clinker.sink.records, clinker.sink.errors, clinker.sink.bytes, and clinker.sink.truncations; each terminal work unit emits one complete clinker.sink span after the outcome is known. Admission loss is behavior-neutral: a full telemetry arena may drop the optional span but cannot change writer bytes or run status.

Fixed-width truncation: warn is reported through a run-scoped truncation ledger rather than the telemetry arena. build_format_writer wraps every Sink writer it builds, so no Sink path can omit the report; when a writer drops, its schema-sized account (per-column counts, longest length, first eight delivered record numbers) settles into the ledger, shifted past the records earlier writers of the same Sink delivered. The clinker.sink.truncations counter of a work unit is that Sink’s ledger delta across the unit, and the end-of-run W367 advisories render from the same ledger after the streaming writer threads join, so the metric and the warning report one count.

Preparation adds the existing fixed-cardinality admission, stage, spill and cleanup observations described in prepared-output telemetry. Those counters describe resource work and never substitute for delivered Sink records or accepted destination bytes. No new lineage edge or dataset is needed for allocation ownership or temporary preparation: neither changes field dependencies, row routing or the external Sink identity.

Lineage keeps the terminal role as an OpenLineage output dataset. PlanNode::Sink resolves the physical or catalog dataset identity, Sink mapping contributes direct column edges, and upstream filters and authored terminal ordering contribute indirect influence. Composition-scoped Sinks keep distinct external identities. The Rust/YAML node name changed to Sink; the lineage role and its output-dataset vocabulary did not.

Writer handling of structured payloads

CSV, fixed-width, EDIFACT, X12, and HL7 writers refuse records carrying a Value::Map payload at any column slot, raising:

FormatError::UnserializableMapValue { format, column }

JSON serializes Value::Map natively as a nested object. XML also accepts a map at an element field and recursively maps ordinary keys to child elements, unescaped @... keys to attributes, #text to text, and arrays to repeated children. Both recursive writers validate the shared key grammar, decoded-key collisions, and the 64-container depth cap before any record bytes reach the sink.

Top-level arrays have a separate compile/runtime contract. Planning resolves every reachable multiple: true source column, applies Sink exclude:, mapping:, and include_unmapped: rules, unions the Sink’s own schema, and retains the resulting output-facing names on the compiled Sink config. CSV joins arrays and XML repeats elements only for those names. Either writer raises FormatError::UnserializableArrayValue when an array reaches any other column. The derived set is schema-bounded; CSV resolves it to one boolean per emitted column and XML stores one boolean per schema leaf, so the record path performs no name-set allocation or input-cardinality growth. Buffered, fused streaming, split, correlation, and document paths all construct writers through the same factory boundary.

The engine-stamped $widened sidecar is handled at projection: it is expanded or stripped rather than exposed as author XML. The contract is the same on the streaming and buffered paths. See Schema Drift & the $widened Sidecar for that lifecycle.

Schema Drift & the $widened Sidecar

User-facing view: the User Guide’s “Auto-Widen & Schema Drift” page.

This page is the engine-internals reference for how Clinker absorbs input columns the source’s declared schema: block does not name, carries them through the DAG, and either expands them back at the sink or refuses them at a writer that cannot serialize them. The mechanism is an on-schema sidecar column named $widened, stamped by the engine and propagated by the same machinery that carries every user-declared column. The depth here is the sidecar’s data model (FieldMetadata::WidenedSidecar, Value::Map), the per-node-type propagation rules, the writer/DLQ rejection paths, and the structural reasons the design is on-schema rather than off-schema. The user-facing page documents the three policy modes and the YAML knobs; this page documents why the absorber is shaped the way it is.

The three modes (context)

The per-source on_unmapped policy selects one of three behaviors for input fields absent from the declared schema. The engine-wide default is auto_widen.

  • auto_widen (default) — per-record undeclared fields are absorbed into a Value::Map payload carried by an engine-stamped $widened sidecar column appended to the source’s schema. The payload propagates downstream and the sink expands it back to top-level columns when include_unmapped: true on the Sink node (the default). Pattern precedent: Databricks Auto Loader’s _rescued_data sidecar and ClickHouse’s JSON column type.
  • drop — undeclared input fields are silently stripped at read time. No sidecar; the source’s plan-time schema equals the declared schema:.
  • reject — any input record carrying a key not in the declared schema fails the source with a FormatError::UndeclaredField diagnostic naming the offending field.

Everything below concerns auto_widen, the only mode that materializes a sidecar.

The $widened sidecar absorber

auto_widen is implemented as an on-schema sidecar: the engine appends a single column named $widened to the source’s schema, marked with FieldMetadata::WidenedSidecar. Each record’s undeclared input fields are stored as the sidecar’s Value::Map payload — keyed by input field name, valued by the read scalar.

The on-schema design is deliberate, and the reason is a silent-loss bug class. An off-schema sidecar — a parallel data structure living outside Schema — would be dropped by any code path that reconstructs a Record from schema.columns() and a value vector. The DAG has many such reconstruction points (projection, sort, spill round-trips, combine collect-array assembly), and each one would carry a standing obligation to “remember to copy the side-channel.” The on-schema slot instead inherits the exact same serialization, span propagation, sort/spill, and projection machinery as a user-declared column: there is no separate copy obligation on any consumer, because to every consumer $widened is a column.

The trade-off the on-schema design accepts is that the sidecar occupies a real schema slot the typechecker can see by name. CXL expressions cannot read or write the sidecar — the typechecker is blind to its contents (the Value::Map interior is never type-resolved into addressable fields), and the parser rejects a literal $widened reference at the system-variable layer. The net effect is that the sidecar rides through every structural transform automatically while remaining unaddressable from user CXL.

Propagation through the DAG

The $widened sidecar follows these rules through downstream nodes. The table is the propagation contract each node type’s executor implements.

Node typeSidecar behavior
TransformInherits unchanged from input (transforms are row-preserving).
AggregateOutput’s $widened slot is Value::Null — per-row payloads have no canonical aggregation. Users who need an unmapped field at aggregate output must add it to group_by or emit it explicitly via an aggregate function.
CombineDriver’s sidecar rides through; build-side sidecars are dropped (mirrors propagate_ck: Driver). Build-side iter_user_fields() filters every engine-stamped column from match: collect array payloads, so build $widened cannot leak into the collect array. Users can lift a build-side unmapped field via <build_qualifier>.<field> in the combine body’s CXL.
Route / MergeRow-preserving — sidecar passes through. Merge requires every input source to share the same on_unmapped policy; mixing fails compile with E315 (see below).
CompositionBody inherits the parent’s sidecar via the synthetic input port; whatever the body’s terminal node carries flows back to the parent. The body’s terminal-node propagation rule applies (e.g. an Aggregate terminal yields Value::Null at the parent boundary, a match: first Combine terminal carries the driver’s payload).
OutputSidecar expands to top-level columns when include_unmapped: true (the default). Set include_unmapped: false to strip the sidecar (and every other unmapped input field) so only explicitly-emitted columns reach the writer.

Two of these rows encode load-bearing internal mechanics worth restating:

  • Combine propagate_ck: Driver mirroring. The sidecar follows the same provenance rule as correlation keys: only the driver (probe) side’s payload survives the join, build-side payloads are dropped. The build-side iter_user_fields() iterator is the single filter that excludes every engine-stamped column — $widened and the $ck.* lattice alike — from the match: collect array, so a build record’s sidecar can never appear as an element of a collect array even though build user fields can. The escape hatch for a genuinely needed build-side unmapped field is to lift it explicitly through the combine body CXL via the build qualifier.
  • Aggregate null-out. There is no canonical way to fold N per-row Value::Map payloads into one, so the Aggregate output slot is deliberately Value::Null rather than (say) the first row’s payload or a merged map. The explicit-emit path (group_by membership or an aggregate function) is the supported way to carry a specific unmapped field across an aggregation boundary.

Output controls

When include_unmapped: true (the default), fields the source absorbed into $widened are expanded back to top-level columns at the sink. The expansion happens at the projection layer, before the writer sees the record, so the literal $widened slot is stripped during expansion and the writer never sees a Value::Map for a well-formed pass-through. Setting include_unmapped: false strips the sidecar (and every other input field not explicitly emitted upstream) so the writer sees only user-declared columns.

include_unmapped composes independently with include_correlation_keys; the two flags are orthogonal — include_correlation_keys does not surface $widened. Because expansion is a projection-layer operation, a CSV source under auto_widen feeding a JSON output under include_unmapped: true produces JSON objects whose top-level keys include both declared columns and absorbed input columns, with no sidecar key remaining.

Writer handling of Value::Map payloads

CSV and fixed-width writers refuse a user-visible Value::Map column, raising FormatError::UnserializableMapValue { format, column }. JSON writes a map as a native object; XML recursively maps it to elements, attributes, and text.

The $widened map is different from a user-visible nested value: Output projection expands it to top-level fields under include_unmapped: true, or strips it under include_unmapped: false. It is never passed through as the author’s JSON/XML nested structure. That projection rule keeps schema-drift handling separate from the recursive writer vocabulary.

DLQ filtering

The dead-letter-queue writer applies the same exclusion at its own layer rather than relying on the main-path projection. dlq::dlq_user_columns strips every column tagged FieldMetadata::WidenedSidecar, so the DLQ CSV header never contains a $widened column even when a DLQ entry’s original_record still carries the full auto_widen schema shape. Correlation-lattice columns ($ck.*) are deliberately retained in DLQ output for collateral debugging — the DLQ filter excludes only the unserializable sidecar, not the engine-stamped provenance columns.

E315 — Merge inputs must agree on policy

Merge concatenates streams positionally against the merge node’s output_schema (taken from the first input). Every input must agree on column shape — same column names, same on_unmapped policy, same correlation_key set. The $widened agreement is a special case of that rule: if one upstream source uses auto_widen (and therefore carries the sidecar column) while another uses drop or reject (and does not), the two input schemas disagree on the presence of the $widened slot, and the positional concatenation has no coherent column to align. Compile fails:

E315: merge "merged": input schemas disagree on the `$widened` auto_widen sidecar column.

The remediation is to set every merge upstream source to the same on_unmapped policy; for sources that should explicitly omit the sidecar, declare on_unmapped: { mode: drop } (or reject) on each so the absent-sidecar shape is uniform across inputs.

Fixed-width sources are structurally inert

Fixed-width sources are positional: the schema is constructed from width / start..end byte ranges, and bytes outside the declared ranges are invisible to the reader. There is no notion of an “undeclared field” to absorb — a byte either falls inside a declared range (and becomes a declared column) or is never read. auto_widen therefore can never populate the sidecar for a fixed-width source; the $widened slot stays Value::Null for every record.

Because the policy is silently inert rather than wrong, the executor emits a tracing::info diagnostic at source-reader construction time when auto_widen is the policy on a fixed-width source, naming the source. The diagnostic fires once per reader instance — a source used as a combine build-side input across multiple combines may produce one log line per combine. To avoid the noise, switch to on_unmapped: drop (or reject) for explicit scalar semantics, or accept the empty sidecar.

Guess multiplicity authoring proof

clinker guess --field <source>.<column> can review a directly authored, single-record CSV, JSON, or XML column that is not already multiple. This is an authoring path, not runtime schema drift: it never writes $widened, resolved layout names, or any other system field into author vocabulary.

The read pass clones the selected source schema and marks candidate columns multiple only in that finite probe. The ordinary format reader then returns the same ordered Value::Array shape runtime ingestion would use after an author declares multiple: true. XML sibling and JSON array counts are inspected per record, so two singleton records cannot combine into proof. Null, empty, and singleton evidence is counted but abstains.

CSV retains one fixed-size state table rather than cells. Delimiter candidates come from observed non-alphanumeric, non-whitespace characters and are capped at 16 distinct values. Escape handling is active only after an observed escape-before-delimiter or escape-before-escape sequence. Each live interpretation splits and re-encodes through the source Charset; byte inequality removes it. Exhaustive check/write requires one survivor across every non-null cell and at least one multi-token cell. No survivor is unconfirmed, while several survivors or a candidate-bound overflow is review-only.

A conclusive edit goes through the same guarded publication path as numeric concretization: frozen input hashes, stable lock, staged sibling file, direct owner re-resolution, exact-byte comparisons, typed semantic proof, atomic replacement, and parent-directory sync. The typed mutation sets multiple and, only for CSV, one complete split_values value. An already-multiple column, conflicting split declaration, indirect owner, sibling semantic change, interruption, or competing writer produces no publication. Reports contain paths, counters, proof states, and proposed syntax, but never sampled field values.

This authoring-only work changes neither row selection nor field values in an executed plan, so it adds no lineage edge. It reuses the fixed-cardinality Guess lifecycle signals and retains only bounded counters/interpretation state; no new execution span, metric label, or memory consumer is introduced.

Staging, Crash Durability & Locks

User-facing view: the User Guide’s “Storage & Spill Location” page.

This page is the engine-internals reference for the durability and concurrency mechanics behind Clinker’s storage subsystem: how a matched source file is copied to a local staging volume without ever leaving a corrupt or half-trusted artifact behind, how concurrent clinker invocations sharing one staging or spill volume coordinate through advisory locks, and how a startup crash purge reclaims the artifacts a SIGKILL-ed run could not clean up. The depth here is the staging copy protocol (single-pass copy + hash, atomic publish via rename, parent-directory fsync, verify, manifest commit), the per-source reader-writer lock semantics, the orphan-detection liveness gates, the file-permission model, and the filesystem-journal reasoning behind the directory fsync. The user-facing page documents the [storage] config block, the spill dir, disk cap, compression, and observability surfaces; those are out of scope here.

Prepared output storage

CSV, JSON, XML, fixed-width and SWIFT output in the CLI and executor seals complete operation bytes in finite memory or executor-backed raw temporary storage. EDIFACT, X12 and HL7 retain their existing writer paths. Preparation is distinct from source file staging and output publication described below.

An executor provider with no resolved spill root is memory-only. It never silently uses the operating system temporary directory. With a resolved root, each operation reserves its future descriptor slot and file/path metadata before accepting bytes, plus a 16 KiB progress buffer retained through readback. Memory is held in independently admitted 16 KiB chunks. A spill request, memory-budget refusal or the 64 KiB residency threshold moves those bytes to a raw file. The threshold is not an operation-size ceiling: disk staging can continue within the supplied quota. Standalone memory-only staging is limited by its explicit budget, not by that spill threshold.

Every temporary write admits at most 8 KiB of prospective disk growth first. Short writes retain only the bytes actually written and release unused quota; memory remains charged until its copy has finished spilling. Sealing rewinds storage and establishes its exact byte length. The existing stage box becomes readback storage, and the retained progress buffer streams bytes to the destination without a second operation-sized allocation. A truncated readback fails instead of delivering an apparently complete operation.

OperationStage is a sealed owner, returned directly by ResourceAuthority::create_stage. Storage extensions implement StageStorage and construct it with StorageStage::create; they cannot extract the boxed backend or release its metadata grant before deallocation. The grant stays with both writable and sealed storage through successful delivery, failure, cancellation, early drop and unwinding.

Admission, allocation, quota and cancellation failures retain a typed inline ResourceError. StageStorage::failure and resource_error recover that evidence from the nonallocating standard-I/O sentinel at write, flush and readback boundaries. Direct storage callers must inspect this channel instead of interpreting the sentinel as an unclassified data error. Failed memory storage cannot accept more writes or seal its prefix. Metadata/progress/box failure notifications run while the original owners and grants are live. No error needs a heap-allocated diagnostic after resource denial.

Before delivery, errors leave the destination and committed encoder state untouched. Generic destination I/O may accept a prefix before failing; after delivery begins, PreparedWriter poisons continuation on any error and never retries or commits pending state. Even zero-byte destination acceptance is a delivery failure. Storage is closed and removed before encoder state commits; cleanup failure after byte delivery also prevents commit and poisons the writer. Drop performs cleanup but does not finalize or retry destination writes.

CSV prepares the first automatic header and body row together; explicit document start and end are separate operations. On preparation failure, the CSV adapter disables stage writes before dropping the library writer, whose destructor can flush buffered bytes. That discarded flush cannot replace the original representability or structured-value error with a later resource failure. Actual resource-denied I/O still recovers its typed stage evidence.

Fixed-width prepares document start, each body record and document end separately. Pending warnings and document counters commit only after sealed bytes have been delivered and storage released. The entire header or footer is staged from borrowed scalar fields; an unsupported late field cannot leak a prefix. SWIFT stages first-record service blocks with that record’s body, retaining a document-sourced trailer only on successful delivery. Finalization stages the close marker and trailer once; flush_bytes never finalizes.

Ordinary and split fixed-width destinations, and ordinary SWIFT destinations, receive prepared bytes directly without an outer BufWriter. Accepted-byte counters therefore observe the destination’s writes, and no outer buffer can retry a failed prefix during drop. All SWIFT splitting remains rejected by planning. These delivery rules do not turn a generic Write destination into an atomic publication boundary.

Unlink failure transfers the exact disk charge, path grant and descriptor slot to an already-admitted cleanup-debt slot. Retry runs outside admission locks; failure never reports reclaimed bytes. An uncertain native close retains conservative debt and is never retried by raw handle, since the handle could have been reused. Debt can outlive the provider in the run arbitrator, keeping its memory consumer live. Confirmed removal releases charges exactly once; unresolved debt remains visible at run teardown. Cleanup still runs when the run is cancelled. The memory contract documents the fixed lifecycle counters and their limits.

How a file is staged

When storage.staging is enabled and a source path matches a configured pattern, the source is copied to a local volume before the pipeline reads it. Each matched source maps to a stable, content-addressed set of files directly under the staging dir, deterministic across runs of the same source:

  • <source-id>.staged — the local copy the reader opens.
  • <source-id>.manifest.json — a sidecar recording the source’s identity: its path, modification time, size, the BLAKE3 content hash, and the stage time.
  • <source-id>.lock — a small advisory-lock file that serializes concurrent invocations staging the same source. It carries no data and persists between runs as the per-source coordination point, alongside the cached copy it guards.

<source-id> is derived from the source’s canonical path, so the same source always resolves to the same staged file. That stability is what makes the reuse cache work — a later run can find the prior copy — and it is why the layout is stable rather than per-run UUIDs.

The copy is built to survive a crash at any point without leaving a corrupt or partial file a later run might trust. The five steps are ordered so that the only trustworthy state is one a crash cannot fabricate:

  1. Single-pass copy + hash. The source is read once in ~1 MiB chunks; each chunk is fed to both the BLAKE3 hasher and the destination file in the same pass. The copy never holds the whole file in memory, so it stays a memory-budget no-op regardless of file size.

  2. Atomic publish. Bytes are written to a <source-id>.<run>.partial temp file, flushed and fsync’d, then renamed to <source-id>.staged. A rename is an atomic replace on Linux, macOS, and Windows (Windows 10 1607+), so a reader scanning for .staged files sees either nothing or the complete file — never a half-written one. The <run> segment in the partial name keeps any two in-flight copies of one source on distinct temp files, and the per-source lock (see the staging cache) ensures only one of them ever runs at a time.

  3. Durable rename. On Linux/macOS the parent directory is fsync’d after the rename, because on ext4/xfs a rename is only crash-durable once the directory entry itself is flushed. On Windows the NTFS journal makes the rename durable, so there is no separate directory flush to do. (See Crash durability and the parent-directory fsync below for the full filesystem-journal reasoning.)

  4. Verify. With verify = blake3 (the default) the source is independently re-read and hashed, and the two digests must match. A size check cannot catch a soft-mount that silently truncated the read; two content digests can. A mismatch removes the published copy and fails the run with a distinct “staged copy is corrupt” diagnostic (E335) — not a generic I/O error.

  5. Commit the manifest. The identity manifest is written with the same atomic temp-file + rename discipline. The manifest is the commit marker: a .staged file is only trustworthy once its manifest exists. A crash between the copy and the manifest leaves a .staged with no manifest, which the next run’s crash purge reaps as an orphan rather than half-trusting.

If the copy fails partway, the .partial is removed before the error propagates. The invariant that closes the protocol: a complete .staged paired with a committed .manifest.json is the only shape any later run will trust, and that pairing cannot exist unless every step above completed.

The staging cache

Because staged copies live at stable paths, a copy from a prior run is still on disk when the next run starts (unless cleanup removed it). The on_existing policy decides what happens when that prior copy is found — overwrite always re-stages, error refuses, and reuse reuses the prior copy only when it is still fresh (the source’s current modification time and size both match what the manifest recorded). A fresh reuse match skips the copy entirely: no bytes read off the share, nothing charged against the disk cap. The freshness check is mtime + size, not a re-hash, so it is a cheap stat rather than a full read of the source.

The internals that matter here are how this stays correct under concurrent invocations. Under the partition-and-run model — several clinker processes over a partitioned input sharing one staging volume — independent runs may stage, reuse, or clean up the same shared source at the same time. The per-source <source-id>.lock file is a reader-writer (shared/exclusive) advisory lock that keeps every such overlap safe on Linux, macOS, and Windows:

  • Exactly one copy. A run that needs to copy takes the lock exclusively for its copy-and-publish. The first run to take it copies and publishes; every other run blocks, then acquires the lock, finds the now-fresh .staged, and reuses it. So a source is copied exactly once no matter how many invocations race for it.

  • A reader is never yanked. A run reading a staged copy holds the lock in shared mode for as long as it has the file open, and keeps it held across the moment it decides to reuse a copy and the moment it opens that copy — so the file it chose cannot be deleted or replaced in between. Any number of readers share the lock at once, so concurrent runs all read the same copy in parallel.

  • Cleanup and overwrite wait for readers. Removing or re-copying a staged pair takes the lock exclusively, which a live reader’s shared lock blocks. Cleanup probes the lock without waiting (a try-lock): if a concurrent run is still reading the copy, cleanup leaves it in place — the last run to release it, or a later crash purge, reaps it. An overwrite re-stage instead waits for in-flight readers to finish, then publishes atomically.

The reader-vs-writer distinction in lock-acquisition discipline is the whole safety argument: a copy/cleanup/overwrite mutates the staged pair and must be exclusive; a reuse-and-read only observes it and can be shared; and because the shared lock is held across the choose-then-open gap, no exclusive holder can slip a delete or replace into that window.

Windows share-mode interoperation

On POSIX an unlinked-but-open file stays readable, so a concurrent delete or atomic-rename replace coexists naturally with an open reader. Windows has no such default. To match the POSIX behavior, the staged copy is opened on Windows with a share mode that permits a concurrent delete or atomic-rename replace (FILE_SHARE_DELETE), so an open reader and a concurrent publish/cleanup interoperate there exactly as they do on POSIX. The net guarantee across any mix of concurrent runs sharing a staged source: a reader always sees a complete, coherent .staged file and no run fails spuriously.

Crash purge of orphaned artifacts

A clean (or panicking) run runs its configured cleanup. But a SIGKILL, the Linux OOM-killer, or a power loss kills the process before any cleanup runs, leaking its staged artifacts under the staging root. To stop that from accumulating across crashes, every run performs an idempotent crash purge at startup, before it stages anything. It reaps four orphan shapes left under the staging root:

  • a *.partial — an interrupted copy. Reaped only when its owning run is dead (see the liveness gate below), so a concurrent sibling’s in-flight copy is never reaped;
  • a *.staged with no matching manifest — a copy that crashed before it could commit its manifest;
  • a *.manifest.json with no matching staged file;
  • a *.lock whose source has no surviving cache entry — a coordination lock left by a source that is no longer staged (not necessarily from a crash), reclaimed under the liveness and age gates below.

A clean pair (a .staged with its committed .manifest.json) is the reuse cache and is kept — the purge never removes a complete, trustworthy copy — and the source’s .lock is kept alongside it so a later reuse run has a lock to take.

Liveness gate: try-lock vs exclusive, and the creation grace window

The purge must distinguish a crash corpse from a live sibling’s in-flight work, because several invocations can share one staging volume. It tells them apart the same way the spill-directory purge does — it asks the operating system “is anyone still staging this?” rather than guessing — and a reap proceeds only when both of two gates pass:

  1. Acquirable under a try-lock. A .partial is reaped only when the source’s .lock is acquirable under a non-blocking try-lock. If the try-lock succeeds, no live process holds the lock, so the owning run is gone. If the lock is still held, a concurrent live run owns the work and the artifact is kept.

  2. Aged past a short creation grace window. Even with an acquirable lock, the artifact must have aged past a short creation grace window. This covers the race where a sibling has just started a copy — created the .partial — but has not yet taken the lock. A partial too young to have been locked yet is kept regardless of the try-lock result.

The actual removal is then performed while the purge holds the lock exclusively, so the reap itself cannot race a sibling mid-acquire.

A .lock whose source has no surviving cache entry (no .staged and no .manifest.json) is itself reclaimed under the same two-gate discipline: removed only when it is acquirable under a try-lock and has aged past the creation grace window, with the removal performed under the exclusive lock. A held lock, a lock still guarding a cached copy, and a freshly created lock are all kept. This bounds what would otherwise be unbounded growth of one zero-byte lock file per distinct source ever staged — relevant for a long-lived persistent cache (on_existing = reuse, cleanup = never) — while never removing a coordination point a live or cached source still needs. The net effect: a concurrent purge can never delete a running sibling’s work, and a persistent staging root does not accumulate one orphan lock per source that has ever passed through it.

File permissions

Staged copies hold verbatim source records — potentially PII, credentials, or financial data — and on a shared staging volume they must not be readable by other users. On Unix each staged file and its manifest are created with mode 0o600 (owner-only). On Windows there is no portable mode bit; staged files inherit the staging directory’s ACL, so the directory’s ACL must be restricted if the volume is shared. The asymmetry is deliberate: the Unix mode bit is enforced per-file at creation, whereas the Windows ACL model pushes the responsibility up to the directory the operator provisions.

Crash durability and the parent-directory fsync

The atomic-rename guarantee that underpins both staging publish (step 2) and the manifest commit (step 5) only holds across a crash if the rename is durable. On POSIX filesystems (ext4, xfs) a rename’s directory entry can still be in the page cache after rename returns: the inode’s data is durable, but the directory entry that names it may not be, so a power loss between the rename and the next implicit flush could lose the entry and leave the file unreachable. To close that gap, Clinker fsyncs the parent directory after the rename, forcing the directory entry to stable storage before the protocol proceeds.

Windows is intentionally exempt. The NTFS metadata journal already makes the rename crash-durable — the same semantics a MOVEFILE_WRITE_THROUGH rename requests — so the directory entry is journaled atomically with the rename and survives a crash without a separate flush. Windows also offers no directory-fsync equivalent to call, so there is nothing to do there. The single cross-platform rule is therefore: durable rename means rename + parent-dir fsync on ext4/xfs, and rename alone on NTFS, with the journal standing in for the directory flush.

Compiler Phases & Type Unification

User-facing view: the User Guide’s “CXL Overview” and “Types & Literals” pages.

This page is the engine-internals reference for how Clinker Expression Language (CXL) source becomes typed planner artifacts and, later, evaluated records. The shared front end parses, resolves, and typechecks a program. Planner consumers may then analyze the typed tree and extract or lower aggregate behavior before runtime evaluation. CXL is a per-record ETL expression language — not SQL — so type errors are reported before records flow. This page covers those boundaries, the miette diagnostic surface, and unification over CXL’s ten runtime value types.

clinker-plan is the sole execution-admission authority. It supplies the bound row schema, compiles CXL, and decides whether the resulting pipeline may run. The CXL crate owns expression semantics but does not parse pipeline YAML or schedule operators. clinker-schema may report advisory findings from bounded discovery and heuristic field extraction; those warnings do not replace canonical planner parsing or admit a rejected pipeline (D-17).

One lower-layer dependency is intentionally narrow: clinker-format may use only CXL’s logical Type and document DocPath/DocIndex vocabulary. D-20 does not permit it to depend on the parser, resolver, evaluator, planner, or other analyzers. A neutral lower-vocabulary extraction may replace this edge later, but no broader dependency is approved now.

Compilation and evaluation pipeline

CXL catches type errors before data processing begins. Its front-end phases are ordered, and a failure short-circuits the remaining work: a parse error never reaches the resolver, and a type error never reaches planner analysis or runtime evaluation.

  1. Parse — tokenize and build an AST from CXL source text. The lexer turns raw source into a token stream; the parser assembles those tokens into an abstract syntax tree of statements (emit, let, filter, distinct) and the expressions inside them. This is the phase that rejects the symbolic boolean operators: &&, ||, and ! are syntax errors in CXL — the language uses the and / or / not keywords — and that rejection happens here, at parse time, before any name or type is known.

  2. Resolve — bind field references, validate method names, and check arity. With the AST in hand, the resolver binds each field reference to a column in the input schema, confirms every method call names a real method, and checks that each call site supplies the right number of arguments. Name and arity errors are structural — they do not depend on types — so they are settled here, ahead of type inference, which lets the typechecker assume every reference resolves and every call is well-formed.

  3. Typecheck — infer types, validate operator compatibility, and check method receiver types. The typechecker walks the resolved tree, infers a type for every expression, and applies the unification rules below at each point two types meet (a binary operator, a method receiver, a conditional’s branches). It rejects incompatible combinations — applying + to a String and an Int, for instance — and emits a span-annotated diagnostic that names both operand types and suggests a coercion. The output of this phase is a TypedProgram: the AST annotated with the inferred type of every node, ready to evaluate without further inference.

  4. Analyze and extract — planner consumers inspect the TypedProgram for execution properties and, where the node kind requires it, extract compiled aggregates or other lowered artifacts. This is not one universal AST rewrite: individual planning paths invoke the analyses they need.

Runtime evaluation then executes the typed or extracted artifact against records. Statements execute top to bottom; later statements can reference fields produced by earlier emit or let statements, and a false filter excludes the record. Evaluation performs no type inference.

Array literals, map literals, and array comprehensions are ordinary expression nodes throughout this pipeline, including inside aggregate residuals. Every AST walker must recurse into item/value expressions, computed map keys, comprehension sources, and predicates; otherwise schema binding, dependency analysis, semantic identity, or lineage can silently miss an input. Runtime construction preserves author order, rejects duplicate logical keys after canonical escape decoding, and shares a per-record 10 MiB allocation budget and 64-container depth cap across nested constructors. The aggregate residual evaluator enforces the same rules.

Expression parsing uses an explicit stack of Pratt continuations, so a nested constructor, operator, call argument, or subscript does not retain a native parser call frame. The existing limit remains 256 simultaneously active expression contexts: a root scalar counts as one, and each nested child expression adds one. Thus 255 containers around a scalar parse successfully; the next child returns the nesting diagnostic. This parser limit is separate from the 64-container runtime value limit. Continuations are discarded on a parse error, so recovery starts the next statement with a fresh depth budget. The limit does not measure the final AST depth of iterative postfix or left-associative operator chains. Statement-level emit each nesting keeps its independent 32-level limit.

The phase split is what makes CXL’s compile-time guarantee meaningful: a cxl check transform.cxl runs Parse → Resolve → Typecheck and reports any error with a span before a single record is read, e.g.

error[typecheck]: cannot apply '+' to String and Int (at transform.cxl:12)
  help: convert one operand — use .to_int() or .to_string()

Because the typecheck phase produces a fully typed program, that class of type mismatch is eliminated before evaluation rather than merely detected earlier.

The type lattice

CXL has 10 value types, and unification operates over them plus two compile-time-only constructs (Numeric and Any) and the Nullable(T) wrapper. The concrete value types and their Rust backings:

TypeRust backingDescription
NullValue::NullMissing or absent value
Boolbooltrue or false
Integeri6464-bit signed integer
Floatf6464-bit double-precision float
Decimalrust_decimal::DecimalExact base-10 fixed-point number (16 bytes) for monetary/financial data
StringFieldStrUTF-8 text
DateNaiveDateCalendar date without timezone
DateTimeNaiveDateTimeDate and time without timezone
ArrayOwnedValuesOrdered collection of values
MapOwnedMapKey-value pairs

Two further type-level constructs appear only at compile time, never as a runtime Value:

  • Numeric — an inference-only union accepting either Int or Float. Unification resolves it when enough context supplies a concrete numeric type. It may not survive into a compiled source schema: an unresolved authored type: numeric is rejected with E158, so source authors must declare int or float.
  • Any — an unconstrained type with no type constraints, the declared type for a column whose type is unknown. It unifies away to whatever it meets.

And the Nullable(T) wrapper marks a type whose value may be null. Nullability is tracked through unification rather than discarded, so a nullable operand propagates its nullability into the result.

Type unification rules

When two types meet in an expression — the two operands of a binary operator, the receiver and a method’s expected type, the branches of a conditional — the typechecker unifies them to a single result type. The algorithm is a small, ordered set of rules; each is tried against the pair of types until one applies:

  1. Identity. Same types unify to themselves: Int + Int produces Int. This is the base case — when both sides already agree, the result is that shared type.

  2. Any absorbs. Any unifies with anything: Any + T produces T. An Any operand imposes no constraint, so the result takes the other operand’s type. (When both are Any, identity covers it.)

  3. Numeric resolves to the concrete type. Numeric + Int produces Int; Numeric + Float produces Float. The Numeric union collapses to whichever concrete numeric type it meets, rather than staying an unresolved union in the result.

  4. Int promotes to Float. Int + Float produces Float. When the two concrete numeric types differ, the result is the wider one — integer arithmetic against a float yields a float, matching the runtime promotion the evaluator performs.

4a. Int widens into Decimal, but Float does not. Decimal + Int produces Decimal — an integer literal or column joins exact decimal arithmetic without loss, so amount + 1 typechecks as Decimal. Decimal deliberately does not unify with Float or Numeric (which admits Float): mixing an exact base-10 value with a binary float is a hard type error that requires an explicit conversion. This is what preserves the decimal type’s exactness guarantee — a lossy float can never silently contaminate a decimal computation. The same rule governs comparisons: decimal > float is rejected, decimal > int is fine. It also governs the branches of a conditional: an if, a match or a ?? whose branches (nullability stripped) include both Decimal and Float is an E200 that names the construct, the first decimal and first float branch (by field name when the branch is a bare field, else by position) and one fix: when the float branch is a bare reference to a Source column the schema declares (not a lexical name or a let binding), the column’s Source schema type (type: decimal in place of type: float, or the nullable forms), because the reader parses the column’s text exactly; otherwise .to_float() on the decimal branch. No fix converts a float to a decimal. The typechecker has no column provenance, so a column an upstream Transform computes also gets the schema fix. The join is typed Any so it raises one diagnostic. The fix is in the message because the pipeline’s E200 joins messages and drops help. Decimal against Numeric stays permissive in operators and joins alike — Numeric may be an integer at run time, as for a decimal.clamp(lo, hi) result — and the aggregates’ run-time rule (a sum, avg or weighted_avg group holding a decimal and a float fails) covers what typecheck cannot see.

  1. Null wraps. Null + T produces Nullable(T). Any operation involving the Null type produces a nullable result: meeting Null cannot guarantee a non-null outcome, so the result type carries the Nullable marker. (Runtime behavior matches — e.g. null + 5 evaluates to null — and the type reflects that the result may be absent.)

  2. Nullable propagates. Nullable(A) + B produces Nullable(unified(A, B)). When a nullable type meets any other type, unification recurses on the inner type A against B, then re-wraps the result in Nullable. Nullability is sticky: it survives the unification and re-wraps whatever the inner types unify to, so a nullable operand anywhere in an expression makes the whole result nullable.

  3. Incompatible types fail. When no rule above applies — String + Int, for instance — unification fails and the typecheck phase emits a span-annotated type error naming both operand types and suggesting a coercion.

The ordering matters: Any and Numeric are resolved before the promotion and nullability rules, so by the time rules 4–6 run, both sides are concrete (or nullable-wrapped concrete) types. Rule 6’s recursion is the only point the algorithm re-enters itself, and it always recurses on strictly-inner types, so unification terminates.

These rules let the typechecker hand later planner phases a resolved TypedProgram: every binary operator, method receiver, and conditional has a single inferred result type, computed once before the per-record evaluator runs.