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