Repository navigation
feat(layout): add push-based pipelines and the scan that runs them - #10391
Open
joseph-isaacs wants to merge 5 commits into
Open
joseph-isaacs wants to merge 5 commits into
joseph-isaacs wants to merge 5 commits into
Conversation
Merging this PR will not alter performance
|
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
joseph-isaacs
force-pushed
the
ji/exec-2-exec-node
branch
from
October 8, 2026 20:16
75184f0 to
8e49f35
Compare
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
`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
…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
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
Pulled out of the V2 scan work in #10349: the executor a
vortex-layoutphysical plan runs on, so aPlanRefbuilt 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 aFilterplan 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 moreOperatorstages, and the ports it writes. Both traits have one method,compute(input, cx) -> Step, withStepone ofMore,Last,Consumed,Blocked, orFinished. 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, andclosed.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.
Scancompiles a plan per split, runs several splits at once, and returns each read to its owner asTurn::Read;Scan::deliverhands the bytes back. Nothing in a scan awaits.Plan hooks.
PlanVTable::compilecompiles a plan into a chain (a source and stages) through aCompiler, andPlanVTable::reachreports what it reads, for sharing.Synthetic plans.
plan::pipeline::synthetichasSyntheticPlan, compiled by a closure,RowSource, andProbe, 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-layoutgains asmallvecdependency.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