Skip to content

feat(layout): add push-based pipelines and the scan that runs them - #10391

Open
joseph-isaacs wants to merge 5 commits into
developfrom
ji/exec-2-exec-node
Open

joseph-isaacs wants to merge 5 commits into
developfrom
ji/exec-2-exec-node

Conversation

@joseph-isaacs

@joseph-isaacs joseph-isaacs commented Oct 8, 2026 •

Copy link
Copy Markdown
Contributor

Summary

Pulled out of the V2 scan work in #10349: the executor a vortex-layout physical plan runs on, so a PlanRef built in memory can be run over row ranges. This PR adds the pipeline runtime, the plan hooks, and synthetic plans that exercise it; no layout operator implements the hooks yet. The operators follow in the PRs stacked on this one. #10389 adds a Filter plan and is independent of this PR.

The design is in vortex-layout/docs/pipeline-exec-design.md, added later in the stack.

Changes

Operators and sources. A pipeline is one Source, zero or more Operator stages, and the ports it writes. Both traits have one method, compute(input, cx) -> Step, with Step one of More, Last, Consumed, Blocked, or Finished. Only a source blocks, and only on a segment read it requested, an empty inlet, or a full outlet, which the driver detects.

Ports. Every port has one writer and one reader and lives in an arena the scan owns. The reading source declares each inlet's capacity in batches; an empty port always has room. A source reads an inlet through peek_mut, take, and closed.

Driver and scheduler. The driver asks a source for a batch, pushes it through the stages into the outlets, and asks again until the source blocks or finishes. The scheduler runs the most recently woken pipeline first and wakes a blocked one only when its condition clears, checking only the ports a run touched.

Scan. Scan compiles a plan per split, runs several splits at once, and returns each read to its owner as Turn::Read; Scan::deliver hands the bytes back. Nothing in a scan awaits.

Plan hooks. PlanVTable::compile compiles a plan into a chain (a source and stages) through a Compiler, and PlanVTable::reach reports what it reads, for sharing.

Synthetic plans. plan::pipeline::synthetic has SyntheticPlan, compiled by a closure, RowSource, and Probe, which records its inlets each time it runs. Tests drive them through a real scan and check order, wakeups only on change, backpressure from a full inlet, and several splits in one scan.

vortex-layout gains a smallvec dependency.

Checks run: cargo clippy -p vortex-layout --all-targets --all-features -- -D warnings, cargo test -p vortex-layout --lib, cargo +nightly-2026-09-10 fmt -p vortex-layout.

🤖 Generated with Claude Code

https://claude.ai/code/session_01NWYmfHTFtaEdhdyd3u5hv3

@codspeed

codspeed Bot commented Oct 8, 2026 •

Copy link
Copy Markdown

Merging this PR will not alter performance

⚠️ Unknown Walltime execution environment detected

Using the Walltime instrument on standard Hosted Runners will lead to inconsistent data.

For the most accurate results, we recommend using CodSpeed Macro Runners: bare-metal machines fine-tuned for performance measurement consistency.

⚠️ 11 benchmarks spent significant time in system calls

System calls cannot be consistently instrumented, so they are not included in the measure, which understates the real cost. Please switch to the Walltime instrument to accurately measure system calls.

Measurement and system calls

✅ 2183 untouched benchmarks
⏩ 359 skipped benchmarks1


Comparing ji/exec-2-exec-node (06f2233) with develop (731a231)

Open in CodSpeed

Footnotes

  1. 359 benchmarks were skipped, so the baseline results were used instead. If they were deleted from the codebase, click here and archive them to remove them from the performance reports. ↩

@joseph-isaacs
joseph-isaacs added this pull request to stack #10395 October 8, 2026 16:01
Add `plan::exec`, a push-based executor interface for physical plans.

A node is built by `PlanVTable::exec` for a range of its plan's rows and a
selection over them. It has input ports, one output, and two entry
points: `start`, called once to spawn children and request segments, and
`compute`, called whenever something changed for the node and its
readiness rule holds. A port is an ordered queue of arrays fed by one
child plus a closed flag, owned by the graph and read by the node through
`StepCx::input`. Every node emits its selected rows in row order as
non-empty arrays and returns `Done` to close its own output; there are
no close events, no empty arrays, and no row ranges on the data.

`Ready` is the readiness rule: `Any` for nodes that stream input through,
`Closed(ports)` for nodes that need one input whole before using the
others, as a dictionary needs its values before its codes, and
`AllClosed` for nodes that need everything. Chains of single-port `Any`
nodes are pipelines; the other rules are the barriers between them.

`ExecGraph` is the flat node arena. It wires ports, routes segment
deliveries to the leaves that requested them, runs ready nodes last-in
first-out so data drains toward the root, and hands its owner the root's
arrays and batches of `IoRequest`s. IO may complete in any order; a node
whose children are backed by independent reads holds the early ones in
their ports and emits the ordered prefix.

`plan::exec::synthetic` has plans whose nodes are built by hand, so the
graph can be exercised without layouts. The scheduling tests drive
sources that cut their rows into pieces, request in-memory segments, and
yield, under a probe node that records when its readiness rule let it
run.

Signed-off-by: Joe Isaacs <joe.isaacs@live.co.uk>
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NWYmfHTFtaEdhdyd3u5hv3
Replaces the `ExecNode` graph with `plan::pipeline`, the executor a plan
now runs on:

- `Operator` and `Source`: one `compute` method returning a `Step`
  (`More`, `Last`, `Consumed`, `Blocked`, `Finished`). Only a source
  blocks, on a read, an empty inlet, or a full outlet.
- Ports with one writer and one reader, in an arena the scan owns; the
  reading source declares each inlet's capacity, and `Inlet` exposes
  `peek_mut`, `take`, and `closed`.
- A driver that pushes each batch through a pipeline's stages into its
  outlets, and a scheduler that wakes a blocked pipeline only when its
  condition clears, checking only the ports a run touched.
- `Scan`, which compiles a plan per split, runs several splits at once,
  and returns reads to its owner as `Turn::Read`.
- `PlanVTable::compile` and `PlanVTable::reach` replace
  `PlanVTable::exec`: the hooks a plan implements to compile into a
  chain and to report the segments it reads. Synthetic plans exercise
  the runtime without layouts.

Signed-off-by: Claude <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NWYmfHTFtaEdhdyd3u5hv3
@joseph-isaacs joseph-isaacs changed the title feat(layout): add the ExecNode trait and push-based exec graph feat(layout): add push-based pipelines and the scan that runs them Oct 9, 2026
`Cx::spawn` let a source ask for a plan to be compiled into a new inlet
while the scan ran. Every plan is now compiled up front: what a stage
reads is decided by its plan and its masks before it runs, so the
runtime never plans mid-run.

Signed-off-by: Claude <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NWYmfHTFtaEdhdyd3u5hv3
@codecov

codecov Bot commented Oct 9, 2026

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 81.76539% with 157 lines in your changes missing coverage. Please review.
✅ Project coverage is 83.41%. Comparing base (1557fb0) to head (94d7196).
⚠️ Report is 37 commits behind head on develop.

Files with missing lines Patch % Lines
vortex-layout/src/plan/pipeline/scan.rs 86.24% 52 Missing ⚠️
vortex-layout/src/plan/pipeline/compile.rs 64.70% 24 Missing ⚠️
vortex-layout/src/plan/typed.rs 48.83% 22 Missing ⚠️
vortex-layout/src/plan/pipeline/synthetic.rs 88.00% 18 Missing ⚠️
vortex-layout/src/plan/vtable.rs 0.00% 18 Missing ⚠️
vortex-layout/src/plan/pipeline/mod.rs 48.00% 13 Missing ⚠️
vortex-layout/src/plan/pipeline/port.rs 90.90% 9 Missing ⚠️
...ortex-layout/src/plan/pipeline/scheduling_tests.rs 98.75% 1 Missing ⚠️

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

claude added 2 commits October 9, 2026 22:05
…et a reader waits on

Ports allocate their buffer once, at their capacity, and keep it when a freed
slot is reused. The default capacity rises from 2 to 8 batches, which halves
scheduler round trips on narrow pipelines. A stage's output port is bounded
like any other, so it needs no growth either; only a share's port, whose
reader may be compiled later, can grow past the default.

A reader blocked on inlet i is woken only by a batch or close on that inlet,
and the inlets a run touched are tracked with a flag per inlet, so a source
reading thousands of inlets pays per inlet it read, not per wake. Pipelines
and ports live in a slab that never moves its entries.

Signed-off-by: Claude <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NWYmfHTFtaEdhdyd3u5hv3
A source could ask for one segment at a time: the driver issued its read and
did not compute the source until the bytes came back, so a source reading
several segments waited for each in turn. Now the driver asks the source for
segments until it has none to ask for and issues every read at once. The
source is computed as each arrives and takes the bytes tagged with their
segment through Cx::take_segment, oldest arrival first.

Signed-off-by: Claude <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NWYmfHTFtaEdhdyd3u5hv3

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants