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