1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
//! Pull-based row pipeline.
//!
//! [`RowSource`] is the fallible row cursor; [`PullExecutor::open_compiled`]
//! and [`MutablePullExecutor::open_compiled`] return a `Box<dyn RowSource + 'a>`
//! representing a streaming query plan execution.
//!
//! ## Architecture
//!
//! The streaming-listed operators have real per-operator
//! [`RowSource`] implementations that pull from their upstream one
//! row at a time:
//!
//! * [`ArgumentSource`]
//! * [`NodeScanSource`]
//! * [`NodeByLabelScanSource`]
//! * catalog-backed index scans for range/text/point predicates
//! * [`ExpandSource`] (single-hop)
//! * [`VariableLengthExpandSource`]
//! * [`FilterSource`]
//! * [`ProjectionSource`]
//! * [`DistinctSource`]
//! * [`UnwindSource`]
//! * [`LimitSource`]
//! * [`SortSource`] (buffers internally, yields lazily)
//! * [`HashAggregationSource`] (buffers internally, yields lazily)
//! * [`OptionalMatchSource`] (streams outer input, buffers inner once)
//! * [`PathBuildSource`]
//!
//! Blocking internals such as sort, aggregation, and shortest-path
//! filtering still allocate where the Cypher semantics require a
//! complete input set. Deduping operators keep only their seen-key
//! state and stream rows as soon as a new key appears.
//!
//! Hydration happens once at the top of the pipeline — operator
//! sources yield raw rows so intermediate evaluations work on
//! storage-borrowed values, and the topmost [`HydratingSource`]
//! converts node / relationship references to their full hydrated
//! map form before the row leaves the cursor.
//!
//! ## Layout
//! - `source` — the [`RowSource`] trait, [`drain`],
//! [`BufferedRowSource`], and [`ArgumentSource`].
//! - `context` — [`StreamCtx`], the shared storage / params handle.
//! - `hydration` — [`HydratingSource`] and [`hydrate_value`].
//! - `traits` — the read-side plan walker (`is_streaming_op`,
//! `subtree_is_fully_streaming`, `build_streaming`,
//! `compiled_to_streaming`, `write_op_input`), [`PullExecutor`],
//! and [`collect_compiled`].
//! - `mutable` — [`MutablePullExecutor`] and the mutable cursor
//! machinery ([`StreamingWriteCursor`], [`MutableUnionSource`]).
//! - `shape` — [`StreamShape`] and [`classify_stream`].
//! - `columns` — [`plan_result_columns`] / [`compiled_result_columns`].
//! - `scan` — node and index scan operator sources ([`NodeScanSource`],
//! [`NodeByLabelScanSource`], [`NodeByPropertyScanSource`]).
//! - `expand` — single-hop and variable-length expansion
//! ([`ExpandSource`], [`VariableLengthExpandSource`]).
//! - `filter` — predicate filter ([`FilterSource`]).
//! - `projection` — projection / unwind / distinct
//! ([`ProjectionSource`], [`UnwindSource`], [`DistinctSource`]).
//! - `sort` — sort and limit ([`SortSource`], [`LimitSource`]).
//! - `aggregate` — hash aggregation and the streamable fold-only fast
//! path ([`HashAggregationSource`], [`StreamableAggKind`],
//! [`AggState`]).
//! - `optional` — outer OPTIONAL MATCH ([`OptionalMatchSource`]).
//! - `path` — path construction including SHORTEST PATH filtering
//! ([`PathBuildSource`]).
//! - `union` — read-side UNION ([`UnionSource`]).
// Public surface — these names appear in `lora_executor`'s public
// API via the explicit `pub use pull::{...}` list in `lib.rs`.
pub use ;
pub use MutablePullExecutor;
pub use ;
pub use ;
pub use ;
// Crate-internal re-exports used by the buffered executor in
// `crate::executor` for the streaming aggregate fast-path and for
// the `StreamingWriteCursor` plan-shape probes.
pub use ;
pub use StreamCtx;
pub use ;
pub use ArgumentSource;
pub use ;