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, NodeByIdSeekExec,
13    NodeByLabelScanExec, NodeByPointScanExec, NodeByPropertyRangeScanExec, NodeByPropertyScanExec,
14    NodeByTextScanExec, NodeScanExec, OptionalMatchExec, PathBuildExec, PhysicalNodeId, PhysicalOp,
15    PhysicalPlan, ProjectionExec, RelByIdSeekExec, RelByPointScanExec, RelByPropertyRangeScanExec,
16    RelByTextScanExec, SortExec, UnwindExec,
17};
18use lora_compiler::CompiledQuery;
19use lora_store::GraphStorage;
20
21use crate::errors::{ExecResult, ExecutorError};
22use crate::eval::clear_eval_error;
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::NodeByIdSeek(_)
121        | PhysicalOp::RelByIdSeek(_)
122        | PhysicalOp::NodeByPropertyRangeScan(_)
123        | PhysicalOp::NodeByTextScan(_)
124        | PhysicalOp::NodeByPointScan(_)
125        | PhysicalOp::RelByPropertyRangeScan(_)
126        | PhysicalOp::RelByTextScan(_)
127        | PhysicalOp::RelByPointScan(_)
128        | PhysicalOp::Filter(_)
129        | PhysicalOp::Unwind(_)
130        | PhysicalOp::Limit(_)
131        // Sort is internally O(N) but exposed as a `RowSource`:
132        // it drains its input on the first pull, sorts in place,
133        // then yields lazily. This lets a write op (CREATE / SET /
134        // DELETE) above an ORDER BY stream its writes one row at
135        // a time instead of forcing the whole subtree to
136        // materialize before the first write.
137        | PhysicalOp::Sort(_)
138        | PhysicalOp::HashAggregation(_)
139        | PhysicalOp::OptionalMatch(_)
140        | PhysicalOp::PathBuild(_)
141        | PhysicalOp::CallSubquery(_)
142        // Projection (both `DISTINCT` and non-`DISTINCT`). The
143        // `DISTINCT` form drains + dedups internally and yields
144        // lazily via `DistinctSource`.
145        | PhysicalOp::Projection(_) => true,
146        // Single-hop expands are fully per-edge. Variable-length expands still
147        // allocate the current source row's BFS result, then yield lazily.
148        PhysicalOp::Expand(_) => true,
149        _ => false,
150    }
151}
152
153/// If `node_id` is a streamable write operator
154/// (Create / Set / Delete / Remove / Merge), return its input
155/// `PhysicalNodeId`. Used by [`MutablePullExecutor::open_compiled`]
156/// to detect plans that can be driven by [`StreamingWriteCursor`].
157pub(crate) fn write_op_input(
158    plan: &PhysicalPlan,
159    node_id: PhysicalNodeId,
160) -> Option<PhysicalNodeId> {
161    match &plan.nodes[node_id] {
162        PhysicalOp::Create(o) => Some(o.input),
163        PhysicalOp::Set(o) => Some(o.input),
164        PhysicalOp::Delete(o) => Some(o.input),
165        PhysicalOp::Remove(o) => Some(o.input),
166        PhysicalOp::Merge(o) => Some(o.input),
167        _ => None,
168    }
169}
170
171/// True if every operator in the subtree rooted at `node_id` is
172/// covered by [`is_streaming_op`] (and therefore by
173/// [`build_streaming`] without falling back to buffered execution).
174///
175/// Used by the mutable executor to decide whether write operators
176/// can pull their input row-by-row instead of materializing it.
177pub(crate) fn subtree_is_fully_streaming(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
178    let op = &plan.nodes[node_id];
179    if !is_streaming_op(op) {
180        return false;
181    }
182    let child = match op {
183        PhysicalOp::Argument(_) => return true,
184        PhysicalOp::NodeScan(o) => o.input,
185        PhysicalOp::NodeByLabelScan(o) => o.input,
186        PhysicalOp::NodeByPropertyScan(o) => o.input,
187        PhysicalOp::NodeByIdSeek(o) => o.input,
188        PhysicalOp::RelByIdSeek(o) => o.input,
189        PhysicalOp::NodeByPropertyRangeScan(o) => o.input,
190        PhysicalOp::NodeByTextScan(o) => o.input,
191        PhysicalOp::NodeByPointScan(o) => o.input,
192        PhysicalOp::RelByPropertyRangeScan(o) => o.input,
193        PhysicalOp::RelByTextScan(o) => o.input,
194        PhysicalOp::RelByPointScan(o) => o.input,
195        PhysicalOp::Filter(o) => Some(o.input),
196        PhysicalOp::Unwind(o) => Some(o.input),
197        PhysicalOp::Limit(o) => Some(o.input),
198        PhysicalOp::Expand(o) => Some(o.input),
199        PhysicalOp::Projection(o) => Some(o.input),
200        PhysicalOp::Sort(o) => Some(o.input),
201        PhysicalOp::HashAggregation(o) => Some(o.input),
202        PhysicalOp::OptionalMatch(o) => Some(o.input),
203        // A `CALL { ... }` that writes runs on the mutable executor, never
204        // on the read-only pull pipeline.
205        PhysicalOp::CallSubquery(o) if subtree_has_write(plan, o.inner) => return false,
206        PhysicalOp::CallSubquery(o) => Some(o.input),
207        PhysicalOp::PathBuild(o) => Some(o.input),
208        // Already filtered by is_streaming_op above.
209        _ => return false,
210    };
211    match child {
212        None => true,
213        Some(c) => subtree_is_fully_streaming(plan, c),
214    }
215}
216
217/// True if any operator in the subtree rooted at `node_id`, including
218/// nested `OPTIONAL MATCH` and `CALL { ... }` bodies, writes.
219pub(crate) fn subtree_has_write(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
220    let op = &plan.nodes[node_id];
221    let (first, second) = match op {
222        PhysicalOp::Create(_)
223        | PhysicalOp::Merge(_)
224        | PhysicalOp::Delete(_)
225        | PhysicalOp::Set(_)
226        | PhysicalOp::Remove(_)
227        | PhysicalOp::Foreach(_) => return true,
228        PhysicalOp::Argument(_) => (None, None),
229        PhysicalOp::NodeScan(o) => (o.input, None),
230        PhysicalOp::NodeByLabelScan(o) => (o.input, None),
231        PhysicalOp::NodeByPropertyScan(o) => (o.input, None),
232        PhysicalOp::NodeByIdSeek(o) => (o.input, None),
233        PhysicalOp::RelByIdSeek(o) => (o.input, None),
234        PhysicalOp::NodeByPropertyRangeScan(o) => (o.input, None),
235        PhysicalOp::NodeByTextScan(o) => (o.input, None),
236        PhysicalOp::NodeByPointScan(o) => (o.input, None),
237        PhysicalOp::RelByPropertyRangeScan(o) => (o.input, None),
238        PhysicalOp::RelByTextScan(o) => (o.input, None),
239        PhysicalOp::RelByPointScan(o) => (o.input, None),
240        PhysicalOp::Expand(o) => (Some(o.input), None),
241        PhysicalOp::Filter(o) => (Some(o.input), None),
242        PhysicalOp::Projection(o) => (Some(o.input), None),
243        PhysicalOp::Unwind(o) => (Some(o.input), None),
244        PhysicalOp::HashAggregation(o) => (Some(o.input), None),
245        PhysicalOp::Sort(o) => (Some(o.input), None),
246        PhysicalOp::Limit(o) => (Some(o.input), None),
247        PhysicalOp::PathBuild(o) => (Some(o.input), None),
248        PhysicalOp::OptionalMatch(o) => (Some(o.input), Some(o.inner)),
249        PhysicalOp::CallSubquery(o) => (Some(o.input), Some(o.inner)),
250    };
251    first.is_some_and(|c| subtree_has_write(plan, c))
252        || second.is_some_and(|c| subtree_has_write(plan, c))
253}
254
255pub(crate) fn build_streaming<'a, S: GraphStorage + 'a>(
256    plan: &'a PhysicalPlan,
257    node_id: PhysicalNodeId,
258    storage: &'a S,
259    params: Arc<BTreeMap<String, LoraValue>>,
260) -> ExecResult<Box<dyn RowSource + 'a>> {
261    build_streaming_dispatch(plan, node_id, storage, params, None)
262}
263
264/// Like [`build_streaming`], but seeds the bottom `Argument` source
265/// with `seed` instead of an empty row. Used by `CallSubquerySource`
266/// to drive the inner sub-plan once per outer row.
267pub(crate) fn build_streaming_seeded<'a, S: GraphStorage + 'a>(
268    plan: &'a PhysicalPlan,
269    node_id: PhysicalNodeId,
270    storage: &'a S,
271    params: Arc<BTreeMap<String, LoraValue>>,
272    seed: Row,
273) -> ExecResult<Box<dyn RowSource + 'a>> {
274    build_streaming_dispatch(plan, node_id, storage, params, Some(seed))
275}
276
277fn build_streaming_dispatch<'a, S: GraphStorage + 'a>(
278    plan: &'a PhysicalPlan,
279    node_id: PhysicalNodeId,
280    storage: &'a S,
281    params: Arc<BTreeMap<String, LoraValue>>,
282    seed: Option<Row>,
283) -> ExecResult<Box<dyn RowSource + 'a>> {
284    build_streaming_inner(plan, node_id, storage, params, seed).map(|src| {
285        super::source::DeadlineSource::wrap(
286            wrap_metered(node_id, src),
287            crate::cancel::active_deadline(),
288        )
289    })
290}
291
292fn build_streaming_inner<'a, S: GraphStorage + 'a>(
293    plan: &'a PhysicalPlan,
294    node_id: PhysicalNodeId,
295    storage: &'a S,
296    params: Arc<BTreeMap<String, LoraValue>>,
297    seed: Option<Row>,
298) -> ExecResult<Box<dyn RowSource + 'a>> {
299    let op = &plan.nodes[node_id];
300
301    if !is_streaming_op(op) {
302        return build_buffered_subtree(plan, node_id, storage, &params);
303    }
304
305    match op {
306        PhysicalOp::Argument(_) => match seed {
307            Some(seed_row) => Ok(Box::new(BufferedRowSource::new(vec![seed_row]))),
308            None => Ok(Box::new(ArgumentSource::new())),
309        },
310
311        PhysicalOp::NodeScan(NodeScanExec { input, var }) => {
312            let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
313            Ok(Box::new(NodeScanSource::new(upstream, storage, *var)))
314        }
315
316        PhysicalOp::NodeByLabelScan(NodeByLabelScanExec { input, var, labels }) => {
317            let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
318            Ok(Box::new(NodeByLabelScanSource::new(
319                upstream, storage, *var, labels,
320            )))
321        }
322
323        PhysicalOp::NodeByPropertyScan(NodeByPropertyScanExec {
324            input,
325            var,
326            labels,
327            key,
328            value,
329            in_list,
330        }) => {
331            let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
332            let ctx = StreamCtx::new(storage, params);
333            Ok(Box::new(NodeByPropertyScanSource::new(
334                upstream, ctx, *var, labels, key, value, *in_list,
335            )))
336        }
337
338        PhysicalOp::NodeByIdSeek(op @ NodeByIdSeekExec { input, .. }) => {
339            let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
340            let ctx = StreamCtx::new(storage, params);
341            Ok(Box::new(BufferedIndexScanSource::node_id(
342                upstream, ctx, op,
343            )))
344        }
345
346        PhysicalOp::RelByIdSeek(op @ RelByIdSeekExec { input, .. }) => {
347            let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
348            let ctx = StreamCtx::new(storage, params);
349            Ok(Box::new(BufferedIndexScanSource::rel_id(upstream, ctx, op)))
350        }
351
352        PhysicalOp::NodeByPropertyRangeScan(op @ NodeByPropertyRangeScanExec { input, .. }) => {
353            let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
354            let ctx = StreamCtx::new(storage, params);
355            if op.order.is_some() {
356                return Ok(Box::new(super::scan::OrderedRangeScanSource::new(
357                    upstream, ctx, op,
358                )));
359            }
360            Ok(Box::new(BufferedIndexScanSource::node_range(
361                upstream, ctx, op,
362            )))
363        }
364
365        PhysicalOp::NodeByTextScan(op @ NodeByTextScanExec { input, .. }) => {
366            let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
367            let ctx = StreamCtx::new(storage, params);
368            Ok(Box::new(BufferedIndexScanSource::node_text(
369                upstream, ctx, op,
370            )))
371        }
372
373        PhysicalOp::NodeByPointScan(op @ NodeByPointScanExec { input, .. }) => {
374            let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
375            let ctx = StreamCtx::new(storage, params);
376            Ok(Box::new(BufferedIndexScanSource::node_point(
377                upstream, ctx, op,
378            )))
379        }
380
381        PhysicalOp::RelByPropertyRangeScan(op @ RelByPropertyRangeScanExec { input, .. }) => {
382            let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
383            let ctx = StreamCtx::new(storage, params);
384            Ok(Box::new(BufferedIndexScanSource::rel_range(
385                upstream, ctx, op,
386            )))
387        }
388
389        PhysicalOp::RelByTextScan(op @ RelByTextScanExec { input, .. }) => {
390            let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
391            let ctx = StreamCtx::new(storage, params);
392            Ok(Box::new(BufferedIndexScanSource::rel_text(
393                upstream, ctx, op,
394            )))
395        }
396
397        PhysicalOp::RelByPointScan(op @ RelByPointScanExec { input, .. }) => {
398            let upstream = open_input(plan, *input, storage, params.clone(), seed.clone())?;
399            let ctx = StreamCtx::new(storage, params);
400            Ok(Box::new(BufferedIndexScanSource::rel_point(
401                upstream, ctx, op,
402            )))
403        }
404
405        PhysicalOp::Expand(ExpandExec {
406            input,
407            src,
408            rel,
409            dst,
410            types,
411            direction,
412            rel_properties,
413            range,
414        }) => {
415            let upstream =
416                build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
417            let ctx = StreamCtx::new(storage, params);
418            match range.as_ref() {
419                Some(range) => Ok(Box::new(VariableLengthExpandSource::new(
420                    upstream, ctx, *src, *rel, *dst, types, *direction, range,
421                ))),
422                None => Ok(Box::new(ExpandSource::new(
423                    upstream,
424                    ctx,
425                    *src,
426                    *rel,
427                    *dst,
428                    types,
429                    *direction,
430                    rel_properties.as_ref(),
431                ))),
432            }
433        }
434
435        PhysicalOp::Filter(FilterExec { input, predicate }) => {
436            let upstream =
437                build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
438            let ctx = StreamCtx::new(storage, params);
439            Ok(Box::new(FilterSource::new(upstream, ctx, predicate)))
440        }
441
442        PhysicalOp::Projection(ProjectionExec {
443            input,
444            distinct,
445            items,
446            include_existing,
447        }) => {
448            let upstream =
449                build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
450            let ctx = StreamCtx::new(storage, params);
451            let proj: Box<dyn RowSource + 'a> = Box::new(ProjectionSource::new(
452                upstream,
453                ctx,
454                items,
455                *include_existing,
456            ));
457            if *distinct {
458                Ok(Box::new(DistinctSource::new(proj)))
459            } else {
460                Ok(proj)
461            }
462        }
463
464        PhysicalOp::Unwind(UnwindExec { input, expr, alias }) => {
465            let upstream =
466                build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
467            let ctx = StreamCtx::new(storage, params);
468            Ok(Box::new(UnwindSource::new(upstream, ctx, expr, *alias)))
469        }
470
471        PhysicalOp::Limit(LimitExec { input, skip, limit }) => {
472            let upstream =
473                build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
474            // Skip / limit expressions are evaluated against an
475            // empty row (matching the buffered executor semantics).
476            let ctx = StreamCtx::new(storage, params);
477            let eval_ctx = ctx.eval_ctx();
478            let skip_n = match skip.as_ref() {
479                Some(e) => crate::executor::eval_row_count("SKIP", e, &eval_ctx)?,
480                None => 0,
481            };
482            let limit_n = match limit.as_ref() {
483                Some(e) => Some(crate::executor::eval_row_count("LIMIT", e, &eval_ctx)?),
484                None => None,
485            };
486            Ok(Box::new(LimitSource::new(upstream, skip_n, limit_n)))
487        }
488
489        PhysicalOp::Sort(SortExec {
490            input,
491            items,
492            top_k,
493            limit,
494        }) => {
495            let upstream =
496                build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
497            let ctx = StreamCtx::new(storage, params);
498            let bound = crate::executor::sort_row_bound(*top_k, limit.as_ref(), &ctx.eval_ctx());
499            Ok(Box::new(SortSource::new_with_top_k(
500                upstream, ctx, items, bound,
501            )))
502        }
503
504        PhysicalOp::HashAggregation(
505            agg @ HashAggregationExec {
506                input,
507                group_by,
508                aggregates,
509            },
510        ) => {
511            // `MATCH (n:L) RETURN count(n)`: answer from the label count.
512            // A seeded run (CALL / OPTIONAL MATCH body) may bind the
513            // scanned variable, so it always scans.
514            if seed.is_none() {
515                if let Some(rows) =
516                    crate::executor::count_all_scan_aggregation_rows(storage, plan, agg)
517                {
518                    return Ok(Box::new(BufferedRowSource::new(rows)));
519                }
520            }
521            let upstream =
522                build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
523            let ctx = StreamCtx::new(storage, params);
524            Ok(Box::new(HashAggregationSource::new(
525                upstream, ctx, group_by, aggregates,
526            )))
527        }
528
529        PhysicalOp::OptionalMatch(OptionalMatchExec {
530            input,
531            inner,
532            new_vars,
533        }) => {
534            let upstream =
535                build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
536            let ctx = StreamCtx::new(storage, params);
537            Ok(Box::new(OptionalMatchSource::new(
538                upstream, ctx, plan, *inner, new_vars,
539            )))
540        }
541
542        PhysicalOp::CallSubquery(CallSubqueryExec {
543            input,
544            inner,
545            new_vars,
546        }) => {
547            let upstream = build_streaming_dispatch(plan, *input, storage, params.clone(), seed)?;
548            Ok(Box::new(CallSubquerySource::new(
549                upstream, plan, *inner, storage, params, new_vars,
550            )))
551        }
552
553        PhysicalOp::PathBuild(PathBuildExec {
554            input,
555            output,
556            node_vars,
557            rel_vars,
558            shortest_path_all,
559        }) => {
560            let upstream =
561                build_streaming_dispatch(plan, *input, storage, params.clone(), seed.clone())?;
562            let ctx = StreamCtx::new(storage, params);
563            Ok(Box::new(PathBuildSource::new(
564                upstream,
565                ctx,
566                *output,
567                node_vars,
568                rel_vars,
569                *shortest_path_all,
570            )))
571        }
572
573        // Already filtered out by `is_streaming_op`; keep this fallible in case
574        // planner/executor capabilities drift.
575        _ => Err(ExecutorError::RuntimeError(format!(
576            "non-streaming op reached streaming branch: {op:?}"
577        ))),
578    }
579}
580
581/// Open an upstream input source. `Option<PhysicalNodeId>` parents
582/// (NodeScan / NodeByLabelScan) treat `None` as "start from a single
583/// empty row".
584fn open_input<'a, S: GraphStorage + 'a>(
585    plan: &'a PhysicalPlan,
586    input: Option<PhysicalNodeId>,
587    storage: &'a S,
588    params: Arc<BTreeMap<String, LoraValue>>,
589    seed: Option<Row>,
590) -> ExecResult<Box<dyn RowSource + 'a>> {
591    match input {
592        Some(input) => build_streaming_dispatch(plan, input, storage, params, seed),
593        None => match seed {
594            Some(seed_row) => Ok(Box::new(BufferedRowSource::new(vec![seed_row]))),
595            None => Ok(Box::new(ArgumentSource::new())),
596        },
597    }
598}
599
600/// Materialized fallback: drain the subtree through the existing
601/// `Executor` and present the result as a [`BufferedRowSource`]. This
602/// remains the leaf path for operators that have no cursor-shaped
603/// source yet (most notably variable-length expansion inside a larger
604/// streaming tree) and for write operators in the read-only pull
605/// executor.
606fn build_buffered_subtree<'a, S: GraphStorage + 'a>(
607    plan: &'a PhysicalPlan,
608    node_id: PhysicalNodeId,
609    storage: &'a S,
610    params: &Arc<BTreeMap<String, LoraValue>>,
611) -> ExecResult<Box<dyn RowSource + 'a>> {
612    // The `Executor` consumes its `ExecutionContext` so we must
613    // clone the params map for the fallback. In practice this is
614    // small (typically empty or a handful of named parameters).
615    let executor = Executor::with_deadline(
616        ExecutionContext {
617            storage,
618            params: (**params).clone(),
619        },
620        crate::cancel::active_deadline(),
621    );
622    let rows = executor.execute_subtree(plan, node_id)?;
623    Ok(Box::new(BufferedRowSource::new(rows)))
624}
625
626// ---------------------------------------------------------------------------
627// Public entry points
628// ---------------------------------------------------------------------------
629
630/// Pull-based read-only executor.
631pub struct PullExecutor<'a, S: GraphStorage> {
632    storage: &'a S,
633    params: BTreeMap<String, LoraValue>,
634}
635
636impl<'a, S: GraphStorage> PullExecutor<'a, S> {
637    pub fn new(storage: &'a S, params: BTreeMap<String, LoraValue>) -> Self {
638        Self { storage, params }
639    }
640
641    /// Open a streaming cursor for a compiled query.
642    ///
643    /// Both no-UNION and UNION-bearing plans go through
644    /// [`compiled_to_streaming`]: UNION drains its branches via
645    /// [`UnionSource`] (memory unchanged from the previous buffered
646    /// path; UNION is inherently O(N) before dedup), but the
647    /// consumer side is now streaming so any downstream pipeline
648    /// composes uniformly.
649    pub fn open_compiled(self, compiled: &'a CompiledQuery) -> ExecResult<Box<dyn RowSource + 'a>>
650    where
651        S: 'a,
652    {
653        clear_eval_error();
654        compiled_to_streaming(compiled, self.storage, self.params)
655    }
656}
657
658/// Drain a freshly opened cursor into a `Vec<Row>`. Convenience for
659/// callers that want the streaming entry point but a buffered result.
660pub fn collect_compiled<'a, S: GraphStorage + 'a>(
661    storage: &'a S,
662    params: BTreeMap<String, LoraValue>,
663    compiled: &'a CompiledQuery,
664) -> ExecResult<Vec<Row>> {
665    let mut cursor = PullExecutor::new(storage, params).open_compiled(compiled)?;
666    drain(cursor.as_mut())
667}
668
669/// [`collect_compiled`] bounded by a cooperative `deadline`.
670///
671/// The deadline is made active for the whole open-and-drain, so every
672/// source in the pipeline (including ones built lazily mid-query by
673/// `OPTIONAL MATCH` or `CALL {}`, and buffered sub-executors) checks it.
674/// Because the pipeline is pulled, a `LIMIT` still stops the scan early.
675pub fn collect_compiled_with_deadline<'a, S: GraphStorage + 'a>(
676    storage: &'a S,
677    params: BTreeMap<String, LoraValue>,
678    compiled: &'a CompiledQuery,
679    deadline: Option<web_time::Instant>,
680) -> ExecResult<Vec<Row>> {
681    let _deadline_scope = crate::cancel::DeadlineScope::enter(deadline);
682    if let Some(deadline) = deadline {
683        if crate::cancel::deadline_reached(deadline) {
684            return Err(ExecutorError::QueryTimeout);
685        }
686    }
687    let mut cursor = PullExecutor::new(storage, params).open_compiled(compiled)?;
688    drain(cursor.as_mut())
689}