Skip to main content

lora_executor/pull/
traits.rs

1//! Pull-pipeline plan walker and executor entry points.
2//!
3//! Cursor/source basics, context, hydration, stream-shape classification,
4//! and result-column inference live in sibling modules. This file keeps the
5//! code that turns physical plans into row cursors plus the public read/write
6//! pull executors.
7
8use std::collections::BTreeMap;
9use std::sync::Arc;
10
11use lora_compiler::physical::{
12    ExpandExec, FilterExec, HashAggregationExec, LimitExec, NodeByLabelScanExec,
13    NodeByPropertyScanExec, NodeScanExec, OptionalMatchExec, PathBuildExec, PhysicalNodeId,
14    PhysicalOp, PhysicalPlan, ProjectionExec, SortExec, UnwindExec,
15};
16use lora_compiler::CompiledQuery;
17use lora_store::GraphStorage;
18
19use crate::errors::ExecResult;
20use crate::eval::{clear_eval_error, eval_expr};
21use crate::executor::{ExecutionContext, Executor};
22use crate::profile::wrap_metered;
23use crate::value::{LoraValue, Row};
24
25use super::aggregate::HashAggregationSource;
26use super::expand::{ExpandSource, VariableLengthExpandSource};
27use super::filter::FilterSource;
28use super::optional::OptionalMatchSource;
29use super::path::PathBuildSource;
30use super::projection::{DistinctSource, ProjectionSource, UnwindSource};
31use super::scan::{NodeByLabelScanSource, NodeByPropertyScanSource, NodeScanSource};
32use super::sort::{LimitSource, SortSource};
33use super::union::UnionSource;
34use super::{drain, ArgumentSource, BufferedRowSource, HydratingSource, RowSource, StreamCtx};
35
36// ---------------------------------------------------------------------------
37// Compiled-query → streaming entry helpers
38// ---------------------------------------------------------------------------
39
40/// Build a streaming `RowSource` for an entire compiled query,
41/// handling both the no-UNION and UNION cases. Replaces the
42/// "UNION-bearing → BufferedRowSource" fallback that previously
43/// sat in `PullExecutor::open_compiled`.
44///
45/// For non-UNION plans this is a thin wrapper around
46/// [`build_streaming`] + [`HydratingSource`]. For UNION plans, we
47/// build a streaming chain per branch (each ending in its own
48/// `HydratingSource` so its node / relationship references are
49/// resolved against the same view of storage), then combine them
50/// through [`UnionSource`].
51pub(super) fn compiled_to_streaming<'a, S: GraphStorage + 'a>(
52    compiled: &'a CompiledQuery,
53    storage: &'a S,
54    params: BTreeMap<String, LoraValue>,
55) -> ExecResult<Box<dyn RowSource + 'a>> {
56    let params = Arc::new(params);
57
58    if compiled.unions.is_empty() {
59        let plan = &compiled.physical;
60        let inner = build_streaming(plan, plan.root, storage, params)?;
61        return Ok(Box::new(HydratingSource::new(inner, storage)));
62    }
63
64    let mut branches: Vec<Box<dyn RowSource + 'a>> = Vec::with_capacity(compiled.unions.len() + 1);
65
66    let head_inner = build_streaming(
67        &compiled.physical,
68        compiled.physical.root,
69        storage,
70        params.clone(),
71    )?;
72    branches.push(Box::new(HydratingSource::new(head_inner, storage)));
73
74    let mut needs_dedup = false;
75    for branch in &compiled.unions {
76        let inner = build_streaming(
77            &branch.physical,
78            branch.physical.root,
79            storage,
80            params.clone(),
81        )?;
82        branches.push(Box::new(HydratingSource::new(inner, storage)));
83        if !branch.all {
84            needs_dedup = true;
85        }
86    }
87
88    Ok(Box::new(UnionSource::new(branches, needs_dedup)))
89}
90
91// ---------------------------------------------------------------------------
92// Plan walker
93// ---------------------------------------------------------------------------
94
95/// True iff this op has a per-operator streaming source. Operators
96/// that aren't on this list fall back to a single materialized
97/// [`Executor::execute_subtree`] call wrapped as a [`BufferedRowSource`].
98pub(super) fn is_streaming_op(op: &PhysicalOp) -> bool {
99    match op {
100        PhysicalOp::Argument(_)
101        | PhysicalOp::NodeScan(_)
102        | PhysicalOp::NodeByLabelScan(_)
103        | PhysicalOp::NodeByPropertyScan(_)
104        | PhysicalOp::Filter(_)
105        | PhysicalOp::Unwind(_)
106        | PhysicalOp::Limit(_)
107        // Sort is internally O(N) but exposed as a `RowSource`:
108        // it drains its input on the first pull, sorts in place,
109        // then yields lazily. This lets a write op (CREATE / SET /
110        // DELETE) above an ORDER BY stream its writes one row at
111        // a time instead of forcing the whole subtree to
112        // materialize before the first write.
113        | PhysicalOp::Sort(_)
114        | PhysicalOp::HashAggregation(_)
115        | PhysicalOp::OptionalMatch(_)
116        | PhysicalOp::PathBuild(_)
117        // Projection (both `DISTINCT` and non-`DISTINCT`). The
118        // `DISTINCT` form drains + dedups internally and yields
119        // lazily via `DistinctSource`.
120        | PhysicalOp::Projection(_) => true,
121        // Single-hop expands are fully per-edge. Variable-length expands still
122        // allocate the current source row's BFS result, then yield lazily.
123        PhysicalOp::Expand(_) => true,
124        _ => false,
125    }
126}
127
128/// If `node_id` is a streamable write operator
129/// (Create / Set / Delete / Remove / Merge), return its input
130/// `PhysicalNodeId`. Used by [`MutablePullExecutor::open_compiled`]
131/// to detect plans that can be driven by [`StreamingWriteCursor`].
132pub(crate) fn write_op_input(
133    plan: &PhysicalPlan,
134    node_id: PhysicalNodeId,
135) -> Option<PhysicalNodeId> {
136    match &plan.nodes[node_id] {
137        PhysicalOp::Create(o) => Some(o.input),
138        PhysicalOp::Set(o) => Some(o.input),
139        PhysicalOp::Delete(o) => Some(o.input),
140        PhysicalOp::Remove(o) => Some(o.input),
141        PhysicalOp::Merge(o) => Some(o.input),
142        _ => None,
143    }
144}
145
146/// True if every operator in the subtree rooted at `node_id` is
147/// covered by [`is_streaming_op`] (and therefore by
148/// [`build_streaming`] without falling back to buffered execution).
149///
150/// Used by the mutable executor to decide whether write operators
151/// can pull their input row-by-row instead of materializing it.
152pub(crate) fn subtree_is_fully_streaming(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
153    let op = &plan.nodes[node_id];
154    if !is_streaming_op(op) {
155        return false;
156    }
157    let child = match op {
158        PhysicalOp::Argument(_) => return true,
159        PhysicalOp::NodeScan(o) => o.input,
160        PhysicalOp::NodeByLabelScan(o) => o.input,
161        PhysicalOp::NodeByPropertyScan(o) => o.input,
162        PhysicalOp::Filter(o) => Some(o.input),
163        PhysicalOp::Unwind(o) => Some(o.input),
164        PhysicalOp::Limit(o) => Some(o.input),
165        PhysicalOp::Expand(o) => Some(o.input),
166        PhysicalOp::Projection(o) => Some(o.input),
167        PhysicalOp::Sort(o) => Some(o.input),
168        PhysicalOp::HashAggregation(o) => Some(o.input),
169        PhysicalOp::OptionalMatch(o) => Some(o.input),
170        PhysicalOp::PathBuild(o) => Some(o.input),
171        // Already filtered by is_streaming_op above.
172        _ => return false,
173    };
174    match child {
175        None => true,
176        Some(c) => subtree_is_fully_streaming(plan, c),
177    }
178}
179
180pub(crate) fn build_streaming<'a, S: GraphStorage + 'a>(
181    plan: &'a PhysicalPlan,
182    node_id: PhysicalNodeId,
183    storage: &'a S,
184    params: Arc<BTreeMap<String, LoraValue>>,
185) -> ExecResult<Box<dyn RowSource + 'a>> {
186    build_streaming_inner(plan, node_id, storage, params).map(|src| wrap_metered(node_id, src))
187}
188
189fn build_streaming_inner<'a, S: GraphStorage + 'a>(
190    plan: &'a PhysicalPlan,
191    node_id: PhysicalNodeId,
192    storage: &'a S,
193    params: Arc<BTreeMap<String, LoraValue>>,
194) -> ExecResult<Box<dyn RowSource + 'a>> {
195    let op = &plan.nodes[node_id];
196
197    if !is_streaming_op(op) {
198        return build_buffered_subtree(plan, node_id, storage, &params);
199    }
200
201    match op {
202        PhysicalOp::Argument(_) => Ok(Box::new(ArgumentSource::new())),
203
204        PhysicalOp::NodeScan(NodeScanExec { input, var }) => {
205            let upstream = open_input(plan, *input, storage, params.clone())?;
206            Ok(Box::new(NodeScanSource::new(upstream, storage, *var)))
207        }
208
209        PhysicalOp::NodeByLabelScan(NodeByLabelScanExec { input, var, labels }) => {
210            let upstream = open_input(plan, *input, storage, params.clone())?;
211            Ok(Box::new(NodeByLabelScanSource::new(
212                upstream, storage, *var, labels,
213            )))
214        }
215
216        PhysicalOp::NodeByPropertyScan(NodeByPropertyScanExec {
217            input,
218            var,
219            labels,
220            key,
221            value,
222        }) => {
223            let upstream = open_input(plan, *input, storage, params.clone())?;
224            let ctx = StreamCtx::new(storage, params);
225            Ok(Box::new(NodeByPropertyScanSource::new(
226                upstream, ctx, *var, labels, key, value,
227            )))
228        }
229
230        PhysicalOp::Expand(ExpandExec {
231            input,
232            src,
233            rel,
234            dst,
235            types,
236            direction,
237            rel_properties,
238            range,
239        }) => {
240            let upstream = build_streaming(plan, *input, storage, params.clone())?;
241            let ctx = StreamCtx::new(storage, params);
242            match range.as_ref() {
243                Some(range) => Ok(Box::new(VariableLengthExpandSource::new(
244                    upstream, ctx, *src, *rel, *dst, types, *direction, range,
245                ))),
246                None => Ok(Box::new(ExpandSource::new(
247                    upstream,
248                    ctx,
249                    *src,
250                    *rel,
251                    *dst,
252                    types,
253                    *direction,
254                    rel_properties.as_ref(),
255                ))),
256            }
257        }
258
259        PhysicalOp::Filter(FilterExec { input, predicate }) => {
260            let upstream = build_streaming(plan, *input, storage, params.clone())?;
261            let ctx = StreamCtx::new(storage, params);
262            Ok(Box::new(FilterSource::new(upstream, ctx, predicate)))
263        }
264
265        PhysicalOp::Projection(ProjectionExec {
266            input,
267            distinct,
268            items,
269            include_existing,
270        }) => {
271            let upstream = build_streaming(plan, *input, storage, params.clone())?;
272            let ctx = StreamCtx::new(storage, params);
273            let proj: Box<dyn RowSource + 'a> = Box::new(ProjectionSource::new(
274                upstream,
275                ctx,
276                items,
277                *include_existing,
278            ));
279            if *distinct {
280                Ok(Box::new(DistinctSource::new(proj)))
281            } else {
282                Ok(proj)
283            }
284        }
285
286        PhysicalOp::Unwind(UnwindExec { input, expr, alias }) => {
287            let upstream = build_streaming(plan, *input, storage, params.clone())?;
288            let ctx = StreamCtx::new(storage, params);
289            Ok(Box::new(UnwindSource::new(upstream, ctx, expr, *alias)))
290        }
291
292        PhysicalOp::Limit(LimitExec { input, skip, limit }) => {
293            let upstream = build_streaming(plan, *input, storage, params.clone())?;
294            // Skip / limit expressions are evaluated against an
295            // empty row (matching the buffered executor semantics).
296            let ctx = StreamCtx::new(storage, params);
297            let eval_ctx = ctx.eval_ctx();
298            let scratch = Row::new();
299            let skip_n = skip
300                .as_ref()
301                .and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
302                .unwrap_or(0)
303                .max(0) as usize;
304            let limit_n = limit
305                .as_ref()
306                .and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
307                .map(|n| n.max(0) as usize);
308            Ok(Box::new(LimitSource::new(upstream, skip_n, limit_n)))
309        }
310
311        PhysicalOp::Sort(SortExec { input, items }) => {
312            let upstream = build_streaming(plan, *input, storage, params.clone())?;
313            let ctx = StreamCtx::new(storage, params);
314            Ok(Box::new(SortSource::new(upstream, ctx, items)))
315        }
316
317        PhysicalOp::HashAggregation(HashAggregationExec {
318            input,
319            group_by,
320            aggregates,
321        }) => {
322            let upstream = build_streaming(plan, *input, storage, params.clone())?;
323            let ctx = StreamCtx::new(storage, params);
324            Ok(Box::new(HashAggregationSource::new(
325                upstream, ctx, group_by, aggregates,
326            )))
327        }
328
329        PhysicalOp::OptionalMatch(OptionalMatchExec {
330            input,
331            inner,
332            new_vars,
333        }) => {
334            let upstream = build_streaming(plan, *input, storage, params.clone())?;
335            let ctx = StreamCtx::new(storage, params);
336            Ok(Box::new(OptionalMatchSource::new(
337                upstream, ctx, plan, *inner, new_vars,
338            )))
339        }
340
341        PhysicalOp::PathBuild(PathBuildExec {
342            input,
343            output,
344            node_vars,
345            rel_vars,
346            shortest_path_all,
347        }) => {
348            let upstream = build_streaming(plan, *input, storage, params.clone())?;
349            let ctx = StreamCtx::new(storage, params);
350            Ok(Box::new(PathBuildSource::new(
351                upstream,
352                ctx,
353                *output,
354                node_vars,
355                rel_vars,
356                *shortest_path_all,
357            )))
358        }
359
360        // Already filtered out by `is_streaming_op`.
361        _ => unreachable!("non-streaming op reached streaming branch: {op:?}"),
362    }
363}
364
365/// Open an upstream input source. `Option<PhysicalNodeId>` parents
366/// (NodeScan / NodeByLabelScan) treat `None` as "start from a single
367/// empty row".
368fn open_input<'a, S: GraphStorage + 'a>(
369    plan: &'a PhysicalPlan,
370    input: Option<PhysicalNodeId>,
371    storage: &'a S,
372    params: Arc<BTreeMap<String, LoraValue>>,
373) -> ExecResult<Box<dyn RowSource + 'a>> {
374    match input {
375        Some(input) => build_streaming(plan, input, storage, params),
376        None => Ok(Box::new(ArgumentSource::new())),
377    }
378}
379
380/// Materialized fallback: drain the subtree through the existing
381/// `Executor` and present the result as a [`BufferedRowSource`]. This
382/// remains the leaf path for operators that have no cursor-shaped
383/// source yet (most notably variable-length expansion inside a larger
384/// streaming tree) and for write operators in the read-only pull
385/// executor.
386fn build_buffered_subtree<'a, S: GraphStorage + 'a>(
387    plan: &'a PhysicalPlan,
388    node_id: PhysicalNodeId,
389    storage: &'a S,
390    params: &Arc<BTreeMap<String, LoraValue>>,
391) -> ExecResult<Box<dyn RowSource + 'a>> {
392    // The `Executor` consumes its `ExecutionContext` so we must
393    // clone the params map for the fallback. In practice this is
394    // small (typically empty or a handful of named parameters).
395    let executor = Executor::new(ExecutionContext {
396        storage,
397        params: (**params).clone(),
398    });
399    let rows = executor.execute_subtree(plan, node_id)?;
400    Ok(Box::new(BufferedRowSource::new(rows)))
401}
402
403// ---------------------------------------------------------------------------
404// Public entry points
405// ---------------------------------------------------------------------------
406
407/// Pull-based read-only executor.
408pub struct PullExecutor<'a, S: GraphStorage> {
409    storage: &'a S,
410    params: BTreeMap<String, LoraValue>,
411}
412
413impl<'a, S: GraphStorage> PullExecutor<'a, S> {
414    pub fn new(storage: &'a S, params: BTreeMap<String, LoraValue>) -> Self {
415        Self { storage, params }
416    }
417
418    /// Open a streaming cursor for a compiled query.
419    ///
420    /// Both no-UNION and UNION-bearing plans go through
421    /// [`compiled_to_streaming`]: UNION drains its branches via
422    /// [`UnionSource`] (memory unchanged from the previous buffered
423    /// path; UNION is inherently O(N) before dedup), but the
424    /// consumer side is now streaming so any downstream pipeline
425    /// composes uniformly.
426    pub fn open_compiled(self, compiled: &'a CompiledQuery) -> ExecResult<Box<dyn RowSource + 'a>>
427    where
428        S: 'a,
429    {
430        clear_eval_error();
431        compiled_to_streaming(compiled, self.storage, self.params)
432    }
433}
434
435/// Drain a freshly opened cursor into a `Vec<Row>`. Convenience for
436/// callers that want the streaming entry point but a buffered result.
437pub fn collect_compiled<'a, S: GraphStorage + 'a>(
438    storage: &'a S,
439    params: BTreeMap<String, LoraValue>,
440    compiled: &'a CompiledQuery,
441) -> ExecResult<Vec<Row>> {
442    let mut cursor = PullExecutor::new(storage, params).open_compiled(compiled)?;
443    drain(cursor.as_mut())
444}