Skip to main content

lora_executor/pull/
traits.rs

1//! Pull-pipeline trait, plan walker, and public entry points.
2//!
3//! This file owns:
4//! - The [`RowSource`] cursor trait, [`drain`] helper, and the shared
5//!   [`StreamCtx`] that every operator source borrows storage and bound
6//!   parameters from.
7//! - The buffered fallback ([`BufferedRowSource`]) and the leaf
8//!   [`ArgumentSource`].
9//! - The top-of-pipeline [`HydratingSource`] and its
10//!   [`hydrate_value`] helper.
11//! - The plan walker (`is_streaming_op`, `subtree_is_fully_streaming`,
12//!   `build_streaming`, `compiled_to_streaming`, `write_op_input`,
13//!   `open_input`, `build_buffered_subtree`).
14//! - The public [`PullExecutor`] / [`MutablePullExecutor`] entry points
15//!   plus the mutable cursor machinery ([`StreamingWriteCursor`],
16//!   [`MutableUnionSource`], [`StoragePtr`]).
17//! - [`collect_compiled`], [`StreamShape`] / [`classify_stream`], and
18//!   [`plan_result_columns`] / [`compiled_result_columns`].
19
20use std::collections::{BTreeMap, BTreeSet};
21use std::mem::ManuallyDrop;
22use std::sync::Arc;
23
24use lora_compiler::physical::{
25    ExpandExec, FilterExec, HashAggregationExec, LimitExec, NodeByLabelScanExec,
26    NodeByPropertyScanExec, NodeScanExec, OptionalMatchExec, PathBuildExec, PhysicalNodeId,
27    PhysicalOp, PhysicalPlan, ProjectionExec, SortExec, UnwindExec,
28};
29use lora_compiler::CompiledQuery;
30use lora_store::{GraphStorage, GraphStorageMut};
31
32use crate::errors::{ExecResult, ExecutorError};
33use crate::eval::{clear_eval_error, eval_expr, EvalContext};
34use crate::executor::{
35    hydrate_node_record, hydrate_relationship_record, ExecutionContext, Executor, GroupValueKey,
36    MutableExecutionContext, MutableExecutor,
37};
38use crate::value::{LoraValue, Row};
39
40use super::aggregate::HashAggregationSource;
41use super::expand::{ExpandSource, VariableLengthExpandSource};
42use super::filter::FilterSource;
43use super::optional::OptionalMatchSource;
44use super::path::PathBuildSource;
45use super::projection::{DistinctSource, ProjectionSource, UnwindSource};
46use super::scan::{NodeByLabelScanSource, NodeByPropertyScanSource, NodeScanSource};
47use super::sort::{LimitSource, SortSource};
48use super::union::UnionSource;
49
50/// Fallible pull-based row cursor.
51///
52/// Each call to [`RowSource::next_row`] returns the next row,
53/// `Ok(None)` when the cursor is exhausted, or an error if execution
54/// fails. The cursor stays in a valid state after an error — callers
55/// may drop it without observing additional side effects.
56pub trait RowSource {
57    /// Pull the next row.
58    fn next_row(&mut self) -> ExecResult<Option<Row>>;
59}
60
61/// Drain a row source into a `Vec<Row>`, propagating the first error.
62pub fn drain<S: RowSource + ?Sized>(source: &mut S) -> ExecResult<Vec<Row>> {
63    let mut out = Vec::new();
64    while let Some(row) = source.next_row()? {
65        out.push(row);
66    }
67    Ok(out)
68}
69
70// ---------------------------------------------------------------------------
71// Shared streaming context
72// ---------------------------------------------------------------------------
73
74/// Storage + bound parameters shared by every operator source in a
75/// pull pipeline. `Clone` is one pointer-copy plus an `Arc::clone`
76/// (params), so passing it by value down the build tree is
77/// effectively free, while consolidating "the two pieces every
78/// expression-evaluating source needs" into one field.
79#[derive(Clone)]
80pub(super) struct StreamCtx<'a, S: GraphStorage> {
81    pub storage: &'a S,
82    pub params: Arc<BTreeMap<String, LoraValue>>,
83}
84
85impl<'a, S: GraphStorage> StreamCtx<'a, S> {
86    pub(super) fn new(storage: &'a S, params: Arc<BTreeMap<String, LoraValue>>) -> Self {
87        Self { storage, params }
88    }
89
90    /// Build a borrowing [`EvalContext`] for use inside an
91    /// operator's `next_row` method. Cheap — two pointer reads.
92    pub(super) fn eval_ctx<'b>(&'b self) -> EvalContext<'b, S> {
93        EvalContext {
94            storage: self.storage,
95            params: &self.params,
96        }
97    }
98}
99
100// ---------------------------------------------------------------------------
101// Buffered fallback
102// ---------------------------------------------------------------------------
103
104/// Buffered cursor backed by a pre-computed `Vec<Row>`. Used both as
105/// a simple "rows already collected" adapter and as the leaf fallback
106/// for operators whose internals still require full materialization.
107pub struct BufferedRowSource {
108    iter: std::vec::IntoIter<Row>,
109}
110
111impl BufferedRowSource {
112    pub fn new(rows: Vec<Row>) -> Self {
113        Self {
114            iter: rows.into_iter(),
115        }
116    }
117}
118
119impl RowSource for BufferedRowSource {
120    fn next_row(&mut self) -> ExecResult<Option<Row>> {
121        Ok(self.iter.next())
122    }
123}
124
125// ---------------------------------------------------------------------------
126// Leaf "yield one empty row" source
127// ---------------------------------------------------------------------------
128
129/// Yields a single empty row exactly once. The bottom of every plan
130/// chain that doesn't start with an explicit input.
131pub struct ArgumentSource {
132    yielded: bool,
133}
134
135impl ArgumentSource {
136    pub fn new() -> Self {
137        Self { yielded: false }
138    }
139}
140
141impl Default for ArgumentSource {
142    fn default() -> Self {
143        Self::new()
144    }
145}
146
147impl RowSource for ArgumentSource {
148    fn next_row(&mut self) -> ExecResult<Option<Row>> {
149        if self.yielded {
150            Ok(None)
151        } else {
152            self.yielded = true;
153            Ok(Some(Row::new()))
154        }
155    }
156}
157
158// ---------------------------------------------------------------------------
159// Top-of-pipeline hydration
160// ---------------------------------------------------------------------------
161
162/// Top-of-pipeline hydration. Replaces node / relationship id
163/// references in each emitted row with their full hydrated map form,
164/// matching the buffered executor's post-execution hydration step.
165pub struct HydratingSource<'a, S: GraphStorage> {
166    upstream: Box<dyn RowSource + 'a>,
167    storage: &'a S,
168}
169
170impl<'a, S: GraphStorage> HydratingSource<'a, S> {
171    pub(super) fn new(upstream: Box<dyn RowSource + 'a>, storage: &'a S) -> Self {
172        Self { upstream, storage }
173    }
174}
175
176impl<'a, S: GraphStorage> RowSource for HydratingSource<'a, S> {
177    fn next_row(&mut self) -> ExecResult<Option<Row>> {
178        match self.upstream.next_row()? {
179            None => Ok(None),
180            Some(row) => {
181                let mut out = Row::new();
182                for (var, name, value) in row.into_iter_named() {
183                    out.insert_named(var, name, hydrate_value(value, self.storage));
184                }
185                Ok(Some(out))
186            }
187        }
188    }
189}
190
191pub(super) fn hydrate_value<S: GraphStorage>(value: LoraValue, storage: &S) -> LoraValue {
192    match value {
193        LoraValue::Node(id) => storage
194            .with_node(id, hydrate_node_record)
195            .unwrap_or(LoraValue::Null),
196        LoraValue::Relationship(id) => storage
197            .with_relationship(id, hydrate_relationship_record)
198            .unwrap_or(LoraValue::Null),
199        LoraValue::List(values) => LoraValue::List(
200            values
201                .into_iter()
202                .map(|v| hydrate_value(v, storage))
203                .collect(),
204        ),
205        LoraValue::Map(map) => LoraValue::Map(
206            map.into_iter()
207                .map(|(k, v)| (k, hydrate_value(v, storage)))
208                .collect(),
209        ),
210        other => other,
211    }
212}
213
214// ---------------------------------------------------------------------------
215// Compiled-query → streaming entry helpers
216// ---------------------------------------------------------------------------
217
218/// Build a streaming `RowSource` for an entire compiled query,
219/// handling both the no-UNION and UNION cases. Replaces the
220/// "UNION-bearing → BufferedRowSource" fallback that previously
221/// sat in `PullExecutor::open_compiled`.
222///
223/// For non-UNION plans this is a thin wrapper around
224/// [`build_streaming`] + [`HydratingSource`]. For UNION plans, we
225/// build a streaming chain per branch (each ending in its own
226/// `HydratingSource` so its node / relationship references are
227/// resolved against the same view of storage), then combine them
228/// through [`UnionSource`].
229pub(super) fn compiled_to_streaming<'a, S: GraphStorage + 'a>(
230    compiled: &'a CompiledQuery,
231    storage: &'a S,
232    params: BTreeMap<String, LoraValue>,
233) -> ExecResult<Box<dyn RowSource + 'a>> {
234    let params = Arc::new(params);
235
236    if compiled.unions.is_empty() {
237        let plan = &compiled.physical;
238        let inner = build_streaming(plan, plan.root, storage, params)?;
239        return Ok(Box::new(HydratingSource::new(inner, storage)));
240    }
241
242    let mut branches: Vec<Box<dyn RowSource + 'a>> = Vec::with_capacity(compiled.unions.len() + 1);
243
244    let head_inner = build_streaming(
245        &compiled.physical,
246        compiled.physical.root,
247        storage,
248        params.clone(),
249    )?;
250    branches.push(Box::new(HydratingSource::new(head_inner, storage)));
251
252    let mut needs_dedup = false;
253    for branch in &compiled.unions {
254        let inner = build_streaming(
255            &branch.physical,
256            branch.physical.root,
257            storage,
258            params.clone(),
259        )?;
260        branches.push(Box::new(HydratingSource::new(inner, storage)));
261        if !branch.all {
262            needs_dedup = true;
263        }
264    }
265
266    Ok(Box::new(UnionSource::new(branches, needs_dedup)))
267}
268
269// ---------------------------------------------------------------------------
270// Plan walker
271// ---------------------------------------------------------------------------
272
273/// True iff this op has a per-operator streaming source. Operators
274/// that aren't on this list fall back to a single materialized
275/// [`Executor::execute_subtree`] call wrapped as a [`BufferedRowSource`].
276pub(super) fn is_streaming_op(op: &PhysicalOp) -> bool {
277    match op {
278        PhysicalOp::Argument(_)
279        | PhysicalOp::NodeScan(_)
280        | PhysicalOp::NodeByLabelScan(_)
281        | PhysicalOp::NodeByPropertyScan(_)
282        | PhysicalOp::Filter(_)
283        | PhysicalOp::Unwind(_)
284        | PhysicalOp::Limit(_)
285        // Sort is internally O(N) but exposed as a `RowSource`:
286        // it drains its input on the first pull, sorts in place,
287        // then yields lazily. This lets a write op (CREATE / SET /
288        // DELETE) above an ORDER BY stream its writes one row at
289        // a time instead of forcing the whole subtree to
290        // materialize before the first write.
291        | PhysicalOp::Sort(_)
292        | PhysicalOp::HashAggregation(_)
293        | PhysicalOp::OptionalMatch(_)
294        | PhysicalOp::PathBuild(_)
295        // Projection (both `DISTINCT` and non-`DISTINCT`). The
296        // `DISTINCT` form drains + dedups internally and yields
297        // lazily via `DistinctSource`.
298        | PhysicalOp::Projection(_) => true,
299        // Single-hop expands are fully per-edge. Variable-length expands still
300        // allocate the current source row's BFS result, then yield lazily.
301        PhysicalOp::Expand(_) => true,
302        _ => false,
303    }
304}
305
306/// If `node_id` is a streamable write operator
307/// (Create / Set / Delete / Remove / Merge), return its input
308/// `PhysicalNodeId`. Used by [`MutablePullExecutor::open_compiled`]
309/// to detect plans that can be driven by [`StreamingWriteCursor`].
310pub(super) fn write_op_input(
311    plan: &PhysicalPlan,
312    node_id: PhysicalNodeId,
313) -> Option<PhysicalNodeId> {
314    match &plan.nodes[node_id] {
315        PhysicalOp::Create(o) => Some(o.input),
316        PhysicalOp::Set(o) => Some(o.input),
317        PhysicalOp::Delete(o) => Some(o.input),
318        PhysicalOp::Remove(o) => Some(o.input),
319        PhysicalOp::Merge(o) => Some(o.input),
320        _ => None,
321    }
322}
323
324/// True if every operator in the subtree rooted at `node_id` is
325/// covered by [`is_streaming_op`] (and therefore by
326/// [`build_streaming`] without falling back to buffered execution).
327///
328/// Used by the mutable executor to decide whether write operators
329/// can pull their input row-by-row instead of materializing it.
330pub(crate) fn subtree_is_fully_streaming(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
331    let op = &plan.nodes[node_id];
332    if !is_streaming_op(op) {
333        return false;
334    }
335    let child = match op {
336        PhysicalOp::Argument(_) => return true,
337        PhysicalOp::NodeScan(o) => o.input,
338        PhysicalOp::NodeByLabelScan(o) => o.input,
339        PhysicalOp::NodeByPropertyScan(o) => o.input,
340        PhysicalOp::Filter(o) => Some(o.input),
341        PhysicalOp::Unwind(o) => Some(o.input),
342        PhysicalOp::Limit(o) => Some(o.input),
343        PhysicalOp::Expand(o) => Some(o.input),
344        PhysicalOp::Projection(o) => Some(o.input),
345        PhysicalOp::Sort(o) => Some(o.input),
346        PhysicalOp::HashAggregation(o) => Some(o.input),
347        PhysicalOp::OptionalMatch(o) => Some(o.input),
348        PhysicalOp::PathBuild(o) => Some(o.input),
349        // Already filtered by is_streaming_op above.
350        _ => return false,
351    };
352    match child {
353        None => true,
354        Some(c) => subtree_is_fully_streaming(plan, c),
355    }
356}
357
358pub(crate) fn build_streaming<'a, S: GraphStorage + 'a>(
359    plan: &'a PhysicalPlan,
360    node_id: PhysicalNodeId,
361    storage: &'a S,
362    params: Arc<BTreeMap<String, LoraValue>>,
363) -> ExecResult<Box<dyn RowSource + 'a>> {
364    let op = &plan.nodes[node_id];
365
366    if !is_streaming_op(op) {
367        return build_buffered_subtree(plan, node_id, storage, &params);
368    }
369
370    match op {
371        PhysicalOp::Argument(_) => Ok(Box::new(ArgumentSource::new())),
372
373        PhysicalOp::NodeScan(NodeScanExec { input, var }) => {
374            let upstream = open_input(plan, *input, storage, params.clone())?;
375            Ok(Box::new(NodeScanSource::new(upstream, storage, *var)))
376        }
377
378        PhysicalOp::NodeByLabelScan(NodeByLabelScanExec { input, var, labels }) => {
379            let upstream = open_input(plan, *input, storage, params.clone())?;
380            Ok(Box::new(NodeByLabelScanSource::new(
381                upstream, storage, *var, labels,
382            )))
383        }
384
385        PhysicalOp::NodeByPropertyScan(NodeByPropertyScanExec {
386            input,
387            var,
388            labels,
389            key,
390            value,
391        }) => {
392            let upstream = open_input(plan, *input, storage, params.clone())?;
393            let ctx = StreamCtx::new(storage, params);
394            Ok(Box::new(NodeByPropertyScanSource::new(
395                upstream, ctx, *var, labels, key, value,
396            )))
397        }
398
399        PhysicalOp::Expand(ExpandExec {
400            input,
401            src,
402            rel,
403            dst,
404            types,
405            direction,
406            rel_properties,
407            range,
408        }) => {
409            let upstream = build_streaming(plan, *input, storage, params.clone())?;
410            let ctx = StreamCtx::new(storage, params);
411            match range.as_ref() {
412                Some(range) => Ok(Box::new(VariableLengthExpandSource::new(
413                    upstream, ctx, *src, *rel, *dst, types, *direction, range,
414                ))),
415                None => Ok(Box::new(ExpandSource::new(
416                    upstream,
417                    ctx,
418                    *src,
419                    *rel,
420                    *dst,
421                    types,
422                    *direction,
423                    rel_properties.as_ref(),
424                ))),
425            }
426        }
427
428        PhysicalOp::Filter(FilterExec { input, predicate }) => {
429            let upstream = build_streaming(plan, *input, storage, params.clone())?;
430            let ctx = StreamCtx::new(storage, params);
431            Ok(Box::new(FilterSource::new(upstream, ctx, predicate)))
432        }
433
434        PhysicalOp::Projection(ProjectionExec {
435            input,
436            distinct,
437            items,
438            include_existing,
439        }) => {
440            let upstream = build_streaming(plan, *input, storage, params.clone())?;
441            let ctx = StreamCtx::new(storage, params);
442            let proj: Box<dyn RowSource + 'a> = Box::new(ProjectionSource::new(
443                upstream,
444                ctx,
445                items,
446                *include_existing,
447            ));
448            if *distinct {
449                Ok(Box::new(DistinctSource::new(proj)))
450            } else {
451                Ok(proj)
452            }
453        }
454
455        PhysicalOp::Unwind(UnwindExec { input, expr, alias }) => {
456            let upstream = build_streaming(plan, *input, storage, params.clone())?;
457            let ctx = StreamCtx::new(storage, params);
458            Ok(Box::new(UnwindSource::new(upstream, ctx, expr, *alias)))
459        }
460
461        PhysicalOp::Limit(LimitExec { input, skip, limit }) => {
462            let upstream = build_streaming(plan, *input, storage, params.clone())?;
463            // Skip / limit expressions are evaluated against an
464            // empty row (matching the buffered executor semantics).
465            let ctx = StreamCtx::new(storage, params);
466            let eval_ctx = ctx.eval_ctx();
467            let scratch = Row::new();
468            let skip_n = skip
469                .as_ref()
470                .and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
471                .unwrap_or(0)
472                .max(0) as usize;
473            let limit_n = limit
474                .as_ref()
475                .and_then(|e| eval_expr(e, &scratch, &eval_ctx).as_i64())
476                .map(|n| n.max(0) as usize);
477            Ok(Box::new(LimitSource::new(upstream, skip_n, limit_n)))
478        }
479
480        PhysicalOp::Sort(SortExec { input, items }) => {
481            let upstream = build_streaming(plan, *input, storage, params.clone())?;
482            let ctx = StreamCtx::new(storage, params);
483            Ok(Box::new(SortSource::new(upstream, ctx, items)))
484        }
485
486        PhysicalOp::HashAggregation(HashAggregationExec {
487            input,
488            group_by,
489            aggregates,
490        }) => {
491            let upstream = build_streaming(plan, *input, storage, params.clone())?;
492            let ctx = StreamCtx::new(storage, params);
493            Ok(Box::new(HashAggregationSource::new(
494                upstream, ctx, group_by, aggregates,
495            )))
496        }
497
498        PhysicalOp::OptionalMatch(OptionalMatchExec {
499            input,
500            inner,
501            new_vars,
502        }) => {
503            let upstream = build_streaming(plan, *input, storage, params.clone())?;
504            let ctx = StreamCtx::new(storage, params);
505            Ok(Box::new(OptionalMatchSource::new(
506                upstream, ctx, plan, *inner, new_vars,
507            )))
508        }
509
510        PhysicalOp::PathBuild(PathBuildExec {
511            input,
512            output,
513            node_vars,
514            rel_vars,
515            shortest_path_all,
516        }) => {
517            let upstream = build_streaming(plan, *input, storage, params.clone())?;
518            let ctx = StreamCtx::new(storage, params);
519            Ok(Box::new(PathBuildSource::new(
520                upstream,
521                ctx,
522                *output,
523                node_vars,
524                rel_vars,
525                *shortest_path_all,
526            )))
527        }
528
529        // Already filtered out by `is_streaming_op`.
530        _ => unreachable!("non-streaming op reached streaming branch: {op:?}"),
531    }
532}
533
534/// Open an upstream input source. `Option<PhysicalNodeId>` parents
535/// (NodeScan / NodeByLabelScan) treat `None` as "start from a single
536/// empty row".
537fn open_input<'a, S: GraphStorage + 'a>(
538    plan: &'a PhysicalPlan,
539    input: Option<PhysicalNodeId>,
540    storage: &'a S,
541    params: Arc<BTreeMap<String, LoraValue>>,
542) -> ExecResult<Box<dyn RowSource + 'a>> {
543    match input {
544        Some(input) => build_streaming(plan, input, storage, params),
545        None => Ok(Box::new(ArgumentSource::new())),
546    }
547}
548
549/// Materialized fallback: drain the subtree through the existing
550/// `Executor` and present the result as a [`BufferedRowSource`]. This
551/// remains the leaf path for operators that have no cursor-shaped
552/// source yet (most notably variable-length expansion inside a larger
553/// streaming tree) and for write operators in the read-only pull
554/// executor.
555fn build_buffered_subtree<'a, S: GraphStorage + 'a>(
556    plan: &'a PhysicalPlan,
557    node_id: PhysicalNodeId,
558    storage: &'a S,
559    params: &Arc<BTreeMap<String, LoraValue>>,
560) -> ExecResult<Box<dyn RowSource + 'a>> {
561    // The `Executor` consumes its `ExecutionContext` so we must
562    // clone the params map for the fallback. In practice this is
563    // small (typically empty or a handful of named parameters).
564    let executor = Executor::new(ExecutionContext {
565        storage,
566        params: (**params).clone(),
567    });
568    let rows = executor.execute_subtree(plan, node_id)?;
569    Ok(Box::new(BufferedRowSource::new(rows)))
570}
571
572// ---------------------------------------------------------------------------
573// Public entry points
574// ---------------------------------------------------------------------------
575
576/// Pull-based read-only executor.
577pub struct PullExecutor<'a, S: GraphStorage> {
578    storage: &'a S,
579    params: BTreeMap<String, LoraValue>,
580}
581
582impl<'a, S: GraphStorage> PullExecutor<'a, S> {
583    pub fn new(storage: &'a S, params: BTreeMap<String, LoraValue>) -> Self {
584        Self { storage, params }
585    }
586
587    /// Open a streaming cursor for a compiled query.
588    ///
589    /// Both no-UNION and UNION-bearing plans go through
590    /// [`compiled_to_streaming`]: UNION drains its branches via
591    /// [`UnionSource`] (memory unchanged from the previous buffered
592    /// path; UNION is inherently O(N) before dedup), but the
593    /// consumer side is now streaming so any downstream pipeline
594    /// composes uniformly.
595    pub fn open_compiled(self, compiled: &'a CompiledQuery) -> ExecResult<Box<dyn RowSource + 'a>>
596    where
597        S: 'a,
598    {
599        clear_eval_error();
600        compiled_to_streaming(compiled, self.storage, self.params)
601    }
602}
603
604/// Pull-based read-write executor. Wraps the existing
605/// [`MutableExecutor`] under the same row-cursor API. Mutations are
606/// applied during `open_compiled`; the returned cursor yields the
607/// resulting rows lazily.
608pub struct MutablePullExecutor<'a, S: GraphStorageMut> {
609    storage: &'a mut S,
610    params: BTreeMap<String, LoraValue>,
611}
612
613impl<'a, S: GraphStorageMut + GraphStorage> MutablePullExecutor<'a, S> {
614    pub fn new(storage: &'a mut S, params: BTreeMap<String, LoraValue>) -> Self {
615        Self { storage, params }
616    }
617
618    /// Open a cursor for a compiled write query.
619    ///
620    /// Fast path: when a branch root is one of `Create` / `Set` /
621    /// `Delete` / `Remove` / `Merge` and its input subtree is fully
622    /// streamable, returns a [`StreamingWriteCursor`] that pulls input
623    /// row-by-row and applies the per-row write through
624    /// [`MutableExecutor::apply_write_op`]. `UNION ALL` plans stream
625    /// one branch at a time. Plain `UNION` drains branches first so
626    /// rows can be deduplicated by name.
627    ///
628    /// Fallback: a branch that is not streamable materializes through
629    /// [`MutableExecutor::execute_rows`] and wraps the result in a
630    /// [`BufferedRowSource`].
631    pub fn open_compiled(self, compiled: &'a CompiledQuery) -> ExecResult<Box<dyn RowSource + 'a>>
632    where
633        S: 'a,
634    {
635        if compiled.unions.is_empty() {
636            return open_mutable_plan_cursor(self.storage, &compiled.physical, self.params);
637        }
638
639        MutableUnionSource::open(self.storage, compiled, self.params)
640            .map(|source| Box::new(source) as Box<dyn RowSource + 'a>)
641    }
642}
643
644fn open_mutable_plan_cursor<'a, S: GraphStorageMut + GraphStorage + 'a>(
645    storage: &'a mut S,
646    plan: &'a PhysicalPlan,
647    params: BTreeMap<String, LoraValue>,
648) -> ExecResult<Box<dyn RowSource + 'a>> {
649    if let Some(input) = write_op_input(plan, plan.root) {
650        if subtree_is_fully_streaming(plan, input) {
651            return StreamingWriteCursor::open(storage, plan, plan.root, params)
652                .map(|c| Box::new(c) as Box<dyn RowSource + 'a>);
653        }
654    }
655
656    let mut executor = MutableExecutor::new(MutableExecutionContext { storage, params });
657    let rows = executor.execute_rows(plan)?;
658    Ok(Box::new(BufferedRowSource::new(rows)))
659}
660
661#[derive(Clone, Copy)]
662struct StoragePtr<S> {
663    ptr: *mut S,
664}
665
666impl<S> StoragePtr<S> {
667    fn from_mut(storage: &mut S) -> Self {
668        Self {
669            ptr: storage as *mut S,
670        }
671    }
672
673    unsafe fn as_ref<'a>(&self) -> &'a S {
674        unsafe { &*self.ptr }
675    }
676
677    unsafe fn as_mut<'a>(&self) -> &'a mut S {
678        unsafe { &mut *self.ptr }
679    }
680}
681
682/// Mutable UNION cursor. `UNION ALL` streams one branch at a time
683/// against the same staged graph. Plain `UNION` streams branch-by-branch
684/// while retaining only a seen-key set for deduplication.
685pub struct MutableUnionSource<'a, S: GraphStorageMut + GraphStorage + 'a> {
686    storage_ptr: StoragePtr<S>,
687    compiled: &'a CompiledQuery,
688    params: BTreeMap<String, LoraValue>,
689    branch_idx: usize,
690    current: Option<Box<dyn RowSource + 'a>>,
691    needs_dedup: bool,
692    seen: BTreeSet<Vec<(String, GroupValueKey)>>,
693    _phantom: std::marker::PhantomData<&'a mut S>,
694}
695
696impl<'a, S: GraphStorageMut + GraphStorage + 'a> MutableUnionSource<'a, S> {
697    fn open(
698        storage: &'a mut S,
699        compiled: &'a CompiledQuery,
700        params: BTreeMap<String, LoraValue>,
701    ) -> ExecResult<Self> {
702        let needs_dedup = compiled.unions.iter().any(|branch| !branch.all);
703        Ok(Self {
704            storage_ptr: StoragePtr::from_mut(storage),
705            compiled,
706            params,
707            branch_idx: 0,
708            current: None,
709            needs_dedup,
710            seen: BTreeSet::new(),
711            _phantom: std::marker::PhantomData,
712        })
713    }
714
715    fn branch_count(&self) -> usize {
716        self.compiled.unions.len() + 1
717    }
718
719    fn branch_plan(&self, idx: usize) -> &'a PhysicalPlan {
720        if idx == 0 {
721            &self.compiled.physical
722        } else {
723            &self.compiled.unions[idx - 1].physical
724        }
725    }
726
727    fn open_branch(&mut self, idx: usize) -> ExecResult<Box<dyn RowSource + 'a>> {
728        let plan = self.branch_plan(idx);
729        // SAFETY: MutableUnionSource keeps at most one branch cursor
730        // alive at a time. `current` is dropped before advancing to
731        // the next branch, so each mutable reborrow is temporally
732        // disjoint.
733        let storage = unsafe { self.storage_ptr.as_mut() };
734        open_mutable_plan_cursor(storage, plan, self.params.clone())
735    }
736}
737
738impl<'a, S: GraphStorageMut + GraphStorage + 'a> RowSource for MutableUnionSource<'a, S> {
739    fn next_row(&mut self) -> ExecResult<Option<Row>> {
740        loop {
741            if self.branch_idx >= self.branch_count() {
742                return Ok(None);
743            }
744
745            if self.current.is_none() {
746                self.current = Some(self.open_branch(self.branch_idx)?);
747            }
748
749            match self
750                .current
751                .as_mut()
752                .expect("current branch initialized above")
753                .next_row()?
754            {
755                Some(row) => {
756                    if self.needs_dedup {
757                        let key = row
758                            .iter_named()
759                            .map(|(_, name, val)| {
760                                (name.into_owned(), GroupValueKey::from_value(val))
761                            })
762                            .collect();
763                        if !self.seen.insert(key) {
764                            continue;
765                        }
766                    }
767                    return Ok(Some(row));
768                }
769                None => {
770                    self.current.take();
771                    self.branch_idx += 1;
772                }
773            }
774        }
775    }
776}
777
778/// Streaming write cursor for plans whose root is one of
779/// `Create` / `Set` / `Delete` / `Remove` / `Merge` and whose input
780/// subtree is fully streamable.
781///
782/// # Layout invariant
783///
784/// The cursor owns a raw alias of the original `&'a mut S`.
785/// Its `upstream` was constructed using a `&'a S` reborrow derived
786/// from `storage_ptr` via unsafe lifetime extension. This is sound
787/// because the existing read-side `RowSource` impls (see
788/// `NodeScanSource::cur_ids`, `ExpandSource::cur_edges`, etc.)
789/// materialize their iteration state into owned `Vec`s at
790/// construction or first call, so no live `&S` borrow into storage
791/// persists across `next_row` calls. Read-only access happens
792/// transiently inside each `upstream.next_row` call; mutable access
793/// happens between calls inside [`MutableExecutor::apply_write_op`].
794/// The borrows never overlap in time.
795///
796/// # Drop order
797///
798/// `upstream` must drop before any caller may regain `&mut S` access
799/// to the underlying storage. The explicit `Drop` impl enforces
800/// that order — `ManuallyDrop` lets us force the sequence.
801pub struct StreamingWriteCursor<'a, S: GraphStorageMut + GraphStorage + 'a> {
802    /// SAFETY: borrows from `*storage_ptr`. Must drop first.
803    upstream: ManuallyDrop<Box<dyn RowSource + 'a>>,
804    /// Raw alias of the `&'a mut S` handed in at construction. Used
805    /// as `&S` by `upstream` and as `&mut S` inside this cursor's `next_row`.
806    storage_ptr: StoragePtr<S>,
807    /// Physical plan — kept alive for the per-row op borrow.
808    plan: &'a PhysicalPlan,
809    /// Index into `plan.nodes` of the write operator.
810    /// We re-fetch the op per call so this struct doesn't need to
811    /// be parameterized by the specific op type.
812    write_op_node: PhysicalNodeId,
813    /// Parameters; cloned per row into a fresh `MutableExecutor`.
814    /// In typical bulk-write workloads this is empty or tiny.
815    params: BTreeMap<String, LoraValue>,
816    _phantom: std::marker::PhantomData<&'a mut S>,
817}
818
819impl<'a, S: GraphStorageMut + GraphStorage + 'a> StreamingWriteCursor<'a, S> {
820    /// Build a cursor. Caller must already have verified that
821    /// `plan.nodes[write_op_node]` is a streamable write op via
822    /// [`write_op_input`] and [`subtree_is_fully_streaming`].
823    pub(crate) fn open(
824        storage: &'a mut S,
825        plan: &'a PhysicalPlan,
826        write_op_node: PhysicalNodeId,
827        params: BTreeMap<String, LoraValue>,
828    ) -> ExecResult<Self> {
829        let input = match write_op_input(plan, write_op_node) {
830            Some(i) => i,
831            None => {
832                return Err(ExecutorError::RuntimeError(format!(
833                    "StreamingWriteCursor::open called with non-write node {write_op_node:?}"
834                )));
835            }
836        };
837        let storage_ptr = StoragePtr::from_mut(storage);
838
839        // SAFETY: see struct-level comment.
840        let storage_ref: &'a S = unsafe { storage_ptr.as_ref() };
841        let upstream = build_streaming(plan, input, storage_ref, Arc::new(params.clone()))?;
842
843        Ok(Self {
844            upstream: ManuallyDrop::new(upstream),
845            storage_ptr,
846            plan,
847            write_op_node,
848            params,
849            _phantom: std::marker::PhantomData,
850        })
851    }
852}
853
854impl<'a, S: GraphStorageMut + GraphStorage + 'a> RowSource for StreamingWriteCursor<'a, S> {
855    fn next_row(&mut self) -> ExecResult<Option<Row>> {
856        let mut row = match self.upstream.next_row()? {
857            Some(r) => r,
858            None => return Ok(None),
859        };
860
861        // SAFETY: upstream's `next_row` has returned, so its
862        // dormant `&S` borrow is not in active use right now. We
863        // reborrow `&mut S` for the per-row write and drop the
864        // borrow before the next pull.
865        let storage_mut: &mut S = unsafe { self.storage_ptr.as_mut() };
866        let mut exec = MutableExecutor::new(MutableExecutionContext {
867            storage: storage_mut,
868            params: self.params.clone(),
869        });
870        let op = &self.plan.nodes[self.write_op_node];
871        exec.apply_write_op(op, &mut row)?;
872        let row = exec.hydrate_row(row);
873        Ok(Some(row))
874    }
875}
876
877impl<'a, S: GraphStorageMut + GraphStorage + 'a> Drop for StreamingWriteCursor<'a, S> {
878    fn drop(&mut self) {
879        // SAFETY: drop `upstream` first to release its borrow into
880        // `*storage_ptr`. Subsequent fields drop via the normal
881        // field-drop sequence and don't touch storage.
882        unsafe {
883            ManuallyDrop::drop(&mut self.upstream);
884        }
885    }
886}
887
888/// Drain a freshly opened cursor into a `Vec<Row>`. Convenience for
889/// callers that want the streaming entry point but a buffered result.
890pub fn collect_compiled<'a, S: GraphStorage + 'a>(
891    storage: &'a S,
892    params: BTreeMap<String, LoraValue>,
893    compiled: &'a CompiledQuery,
894) -> ExecResult<Vec<Row>> {
895    let mut cursor = PullExecutor::new(storage, params).open_compiled(compiled)?;
896    drain(cursor.as_mut())
897}
898
899// ---------------------------------------------------------------------------
900// Stream classification
901// ---------------------------------------------------------------------------
902
903/// Classification of a compiled query, used by the database layer to
904/// decide whether `db.stream` needs a hidden staged transaction.
905#[derive(Debug, Clone, Copy, PartialEq, Eq)]
906pub enum StreamShape {
907    /// No mutating operator anywhere in the plan or any of its
908    /// UNION branches. Safe to stream against the live store.
909    ReadOnly,
910    /// Has at least one mutating operator (Create / Merge / Delete /
911    /// Set / Remove). The host should run this against a staged
912    /// graph and only publish on cursor exhaustion.
913    Mutating,
914}
915
916impl StreamShape {
917    pub fn is_mutating(self) -> bool {
918        matches!(self, StreamShape::Mutating)
919    }
920}
921
922fn plan_is_mutating(plan: &PhysicalPlan) -> bool {
923    plan.nodes.iter().any(|op| {
924        matches!(
925            op,
926            PhysicalOp::Create(_)
927                | PhysicalOp::Merge(_)
928                | PhysicalOp::Delete(_)
929                | PhysicalOp::Set(_)
930                | PhysicalOp::Remove(_)
931        )
932    })
933}
934
935/// Classify a compiled query for streaming. Treats any UNION branch
936/// the same as the head: a single mutating op anywhere across the
937/// compiled query promotes the whole query to `Mutating`.
938pub fn classify_stream(compiled: &CompiledQuery) -> StreamShape {
939    if plan_is_mutating(&compiled.physical)
940        || compiled
941            .unions
942            .iter()
943            .any(|b| plan_is_mutating(&b.physical))
944    {
945        StreamShape::Mutating
946    } else {
947        StreamShape::ReadOnly
948    }
949}
950
951// ---------------------------------------------------------------------------
952// Plan-derived result columns
953// ---------------------------------------------------------------------------
954
955/// Result column names derived from the compiled plan.
956///
957/// Walks the plan from `root` looking for the topmost projection-shaped
958/// node (Projection, HashAggregation). Other operators that wrap a
959/// projection (Limit, Sort, PathBuild, OptionalMatch, Filter, Unwind,
960/// Create/Merge/Set/Delete/Remove) defer to their input. Returns an
961/// empty `Vec` for plans that have no named output (e.g. a bare
962/// scan-only plan), preserving the previous "infer from first row"
963/// behaviour for those cases.
964pub fn plan_result_columns(plan: &PhysicalPlan) -> Vec<String> {
965    plan_columns_at(plan, plan.root).unwrap_or_default()
966}
967
968fn plan_columns_at(plan: &PhysicalPlan, node: PhysicalNodeId) -> Option<Vec<String>> {
969    match &plan.nodes[node] {
970        PhysicalOp::Projection(p) => Some(p.items.iter().map(|i| i.name.clone()).collect()),
971        PhysicalOp::HashAggregation(p) => Some(
972            p.group_by
973                .iter()
974                .chain(p.aggregates.iter())
975                .map(|i| i.name.clone())
976                .collect(),
977        ),
978        PhysicalOp::Limit(p) => plan_columns_at(plan, p.input),
979        PhysicalOp::Sort(p) => plan_columns_at(plan, p.input),
980        PhysicalOp::PathBuild(p) => plan_columns_at(plan, p.input),
981        PhysicalOp::OptionalMatch(p) => plan_columns_at(plan, p.input),
982        PhysicalOp::Filter(p) => plan_columns_at(plan, p.input),
983        PhysicalOp::Unwind(p) => plan_columns_at(plan, p.input),
984        PhysicalOp::Create(p) => plan_columns_at(plan, p.input),
985        PhysicalOp::Merge(p) => plan_columns_at(plan, p.input),
986        PhysicalOp::Delete(p) => plan_columns_at(plan, p.input),
987        PhysicalOp::Set(p) => plan_columns_at(plan, p.input),
988        PhysicalOp::Remove(p) => plan_columns_at(plan, p.input),
989        PhysicalOp::Argument(_)
990        | PhysicalOp::NodeScan(_)
991        | PhysicalOp::NodeByLabelScan(_)
992        | PhysicalOp::NodeByPropertyScan(_)
993        | PhysicalOp::Expand(_) => None,
994    }
995}
996
997/// Result column names for a compiled query (head plan; UNION branches
998/// must produce the same shape so the head's columns are authoritative).
999pub fn compiled_result_columns(compiled: &CompiledQuery) -> Vec<String> {
1000    plan_result_columns(&compiled.physical)
1001}