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