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