Skip to main content

lora_executor/executor/
immutable.rs

1//! Read-only buffered executor: lower a [`PhysicalPlan`] into a fully
2//! materialized `Vec<Row>` without touching the store.
3//!
4//! [`Executor`] mirrors the operator set of the streaming pipeline in
5//! `crate::pull` so write operators can fall back to it for subtrees
6//! that are not fully streamable. Aggregation and DISTINCT projection
7//! reuse the streaming `StreamableAggSpec` / `AggState` machinery
8//! (re-exported as `crate::pull::*` for that purpose) on the
9//! fold-only fast path; everything else materializes.
10
11use crate::errors::{value_kind, ExecResult, ExecutorError};
12use crate::eval::{clear_eval_error, eval_expr, eval_expr_result, EvalContext};
13use crate::value::{LoraValue, Row};
14use crate::{project_rows, ExecuteOptions, QueryResult};
15
16use lora_analyzer::{ResolvedExpr, ResolvedProjection};
17use lora_ast::RangeLiteral;
18use lora_compiler::physical::*;
19use lora_compiler::CompiledQuery;
20use lora_store::{GraphStorage, Properties};
21
22use std::cmp::Ordering;
23use std::collections::BTreeMap;
24use std::time::Instant;
25use tracing::{error, trace};
26
27use super::helpers::{
28    build_path_value, check_deadline_at, compare_sort_item, compute_aggregate_expr, dedup_rows,
29    filter_rows_checked, filter_shortest_paths, hydrate_node_record, hydrate_relationship_record,
30    indexed_node_property_candidates, label_group_candidates_prefiltered,
31    node_matches_label_groups, node_matches_property_filter, project_rows_checked, resolve_range,
32    scan_node_ids_for_label_groups, value_matches_property_value, variable_length_expand,
33    GroupValueKey,
34};
35
36pub struct ExecutionContext<'a, S: GraphStorage> {
37    pub storage: &'a S,
38    pub params: BTreeMap<String, LoraValue>,
39}
40
41pub struct Executor<'a, S: GraphStorage> {
42    ctx: ExecutionContext<'a, S>,
43    deadline: Option<Instant>,
44}
45
46impl<'a, S: GraphStorage> Executor<'a, S> {
47    pub fn new(ctx: ExecutionContext<'a, S>) -> Self {
48        Self {
49            ctx,
50            deadline: None,
51        }
52    }
53
54    pub fn with_deadline(ctx: ExecutionContext<'a, S>, deadline: Option<Instant>) -> Self {
55        Self { ctx, deadline }
56    }
57
58    #[inline]
59    fn check_deadline(&self) -> ExecResult<()> {
60        if let Some(deadline) = self.deadline {
61            check_deadline_at(deadline)
62        } else {
63            Ok(())
64        }
65    }
66
67    #[inline]
68    fn check_loop_deadline(deadline: Option<Instant>) -> ExecResult<()> {
69        if let Some(deadline) = deadline {
70            check_deadline_at(deadline)
71        } else {
72            Ok(())
73        }
74    }
75}
76
77impl<'a, S: GraphStorage> Executor<'a, S> {
78    pub fn execute(
79        &self,
80        plan: &PhysicalPlan,
81        options: Option<ExecuteOptions>,
82    ) -> ExecResult<QueryResult> {
83        let rows = self.execute_rows(plan)?;
84        Ok(project_rows(rows, options.unwrap_or_default()))
85    }
86
87    pub fn execute_compiled(
88        &self,
89        compiled: &CompiledQuery,
90        options: Option<ExecuteOptions>,
91    ) -> ExecResult<QueryResult> {
92        let rows = self.execute_compiled_rows(compiled)?;
93        Ok(project_rows(rows, options.unwrap_or_default()))
94    }
95
96    pub fn execute_compiled_rows(&self, compiled: &CompiledQuery) -> ExecResult<Vec<Row>> {
97        self.check_deadline()?;
98        if compiled.unions.is_empty() {
99            return self.execute_rows(&compiled.physical);
100        }
101
102        clear_eval_error();
103
104        let mut all_rows = self.execute_rows(&compiled.physical)?;
105        let mut needs_dedup = false;
106
107        for branch in &compiled.unions {
108            self.check_deadline()?;
109            let branch_rows = self.execute_rows(&branch.physical)?;
110            all_rows.extend(branch_rows);
111
112            if !branch.all {
113                needs_dedup = true;
114            }
115        }
116
117        if needs_dedup {
118            all_rows = dedup_rows(all_rows);
119        }
120
121        Ok(all_rows)
122    }
123
124    pub fn execute_rows(&self, plan: &PhysicalPlan) -> ExecResult<Vec<Row>> {
125        self.check_deadline()?;
126        // Clear any error residue that a previous query on this thread may have
127        // left in the thread-local eval-error slot.
128        clear_eval_error();
129
130        let rows = self.execute_node(plan, plan.root)?;
131        Ok(rows
132            .into_iter()
133            .map(|row| self.hydrate_row(row))
134            .collect::<Vec<_>>())
135    }
136
137    fn hydrate_row(&self, row: Row) -> Row {
138        let mut out = Row::new();
139
140        for (var, name, value) in row.into_iter_named() {
141            out.insert_named(var, name, self.hydrate_value(value));
142        }
143
144        out
145    }
146
147    /// Buffered execution of an arbitrary subplan. Public to the
148    /// crate so the pull pipeline can fall back to materialized
149    /// execution for operators that have no streaming source yet.
150    pub(crate) fn execute_subtree(
151        &self,
152        plan: &PhysicalPlan,
153        node_id: PhysicalNodeId,
154    ) -> ExecResult<Vec<Row>> {
155        self.execute_node(plan, node_id)
156    }
157
158    fn execute_node(&self, plan: &PhysicalPlan, node_id: PhysicalNodeId) -> ExecResult<Vec<Row>> {
159        self.check_deadline()?;
160        trace!("read-only execute_node start: node_id={node_id:?}");
161
162        let result = match &plan.nodes[node_id] {
163            PhysicalOp::Argument(op) => self.exec_argument(op),
164            PhysicalOp::NodeScan(op) => self.exec_node_scan(plan, op),
165            PhysicalOp::NodeByLabelScan(op) => self.exec_node_by_label_scan(plan, op),
166            PhysicalOp::NodeByPropertyScan(op) => self.exec_node_by_property_scan(plan, op),
167            PhysicalOp::Expand(op) => self.exec_expand(plan, op),
168            PhysicalOp::Filter(op) => self.exec_filter(plan, op),
169            PhysicalOp::Projection(op) => self.exec_projection(plan, op),
170            PhysicalOp::Unwind(op) => self.exec_unwind(plan, op),
171            PhysicalOp::HashAggregation(op) => self.exec_hash_aggregation(plan, op),
172            PhysicalOp::Sort(op) => self.exec_sort(plan, op),
173            PhysicalOp::Limit(op) => self.exec_limit(plan, op),
174            PhysicalOp::OptionalMatch(op) => self.exec_optional_match(plan, op),
175            PhysicalOp::PathBuild(op) => self.exec_path_build(plan, op),
176            PhysicalOp::Create(_) => Err(ExecutorError::ReadOnlyCreate { node_id }),
177            PhysicalOp::Merge(_) => Err(ExecutorError::ReadOnlyMerge { node_id }),
178            PhysicalOp::Delete(_) => Err(ExecutorError::ReadOnlyDelete { node_id }),
179            PhysicalOp::Set(_) => Err(ExecutorError::ReadOnlySet { node_id }),
180            PhysicalOp::Remove(_) => Err(ExecutorError::ReadOnlyRemove { node_id }),
181        };
182
183        match &result {
184            Ok(rows) => trace!(
185                "read-only execute_node ok: node_id={node_id:?}, rows={}",
186                rows.len()
187            ),
188            Err(err) => error!("read-only execute_node failed: node_id={node_id:?}, error={err}"),
189        }
190
191        result
192    }
193
194    fn exec_argument(&self, _op: &ArgumentExec) -> ExecResult<Vec<Row>> {
195        Ok(vec![Row::new()])
196    }
197
198    fn exec_node_scan(&self, plan: &PhysicalPlan, op: &NodeScanExec) -> ExecResult<Vec<Row>> {
199        let base_rows = match op.input {
200            Some(input) => self.execute_node(plan, input)?,
201            None => vec![Row::new()],
202        };
203
204        let node_ids = self.ctx.storage.all_node_ids();
205        let mut out = Vec::new();
206
207        let deadline = self.deadline;
208        for row in base_rows {
209            Self::check_loop_deadline(deadline)?;
210            if let Some(existing) = row.get(op.var) {
211                match existing {
212                    LoraValue::Node(existing_id) => {
213                        if self.ctx.storage.has_node(*existing_id) {
214                            out.push(row);
215                        }
216                    }
217                    other => {
218                        return Err(ExecutorError::ExpectedNodeForExpand {
219                            var: format!("{:?}", op.var),
220                            found: value_kind(other),
221                        });
222                    }
223                }
224                continue;
225            }
226
227            for &id in &node_ids {
228                Self::check_loop_deadline(deadline)?;
229                let mut new_row = row.clone();
230                new_row.insert(op.var, LoraValue::Node(id));
231                out.push(new_row);
232            }
233        }
234
235        Ok(out)
236    }
237
238    fn exec_node_by_label_scan(
239        &self,
240        plan: &PhysicalPlan,
241        op: &NodeByLabelScanExec,
242    ) -> ExecResult<Vec<Row>> {
243        let base_rows = match op.input {
244            Some(input) => self.execute_node(plan, input)?,
245            None => vec![Row::new()],
246        };
247
248        let candidate_ids = scan_node_ids_for_label_groups(self.ctx.storage, &op.labels);
249        let candidates_prefiltered = label_group_candidates_prefiltered(&op.labels);
250        let mut out = Vec::new();
251
252        match self.deadline {
253            Some(deadline) => {
254                for row in base_rows {
255                    check_deadline_at(deadline)?;
256                    if let Some(existing) = row.get(op.var) {
257                        match existing {
258                            LoraValue::Node(existing_id) => {
259                                let labels_ok = self
260                                    .ctx
261                                    .storage
262                                    .with_node(*existing_id, |n| {
263                                        node_matches_label_groups(&n.labels, &op.labels)
264                                    })
265                                    .unwrap_or(false);
266                                if labels_ok {
267                                    out.push(row);
268                                }
269                            }
270                            other => {
271                                return Err(ExecutorError::ExpectedNodeForExpand {
272                                    var: format!("{:?}", op.var),
273                                    found: value_kind(other),
274                                });
275                            }
276                        }
277                        continue;
278                    }
279
280                    for &id in &candidate_ids {
281                        check_deadline_at(deadline)?;
282                        if !candidates_prefiltered {
283                            let labels_ok = self
284                                .ctx
285                                .storage
286                                .with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
287                                .unwrap_or(false);
288                            if !labels_ok {
289                                continue;
290                            }
291                        }
292                        let mut new_row = row.clone();
293                        new_row.insert(op.var, LoraValue::Node(id));
294                        out.push(new_row);
295                    }
296                }
297            }
298            None => {
299                for row in base_rows {
300                    if let Some(existing) = row.get(op.var) {
301                        match existing {
302                            LoraValue::Node(existing_id) => {
303                                let labels_ok = self
304                                    .ctx
305                                    .storage
306                                    .with_node(*existing_id, |n| {
307                                        node_matches_label_groups(&n.labels, &op.labels)
308                                    })
309                                    .unwrap_or(false);
310                                if labels_ok {
311                                    out.push(row);
312                                }
313                            }
314                            other => {
315                                return Err(ExecutorError::ExpectedNodeForExpand {
316                                    var: format!("{:?}", op.var),
317                                    found: value_kind(other),
318                                });
319                            }
320                        }
321                        continue;
322                    }
323
324                    for &id in &candidate_ids {
325                        if !candidates_prefiltered {
326                            let labels_ok = self
327                                .ctx
328                                .storage
329                                .with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
330                                .unwrap_or(false);
331                            if !labels_ok {
332                                continue;
333                            }
334                        }
335                        let mut new_row = row.clone();
336                        new_row.insert(op.var, LoraValue::Node(id));
337                        out.push(new_row);
338                    }
339                }
340            }
341        }
342
343        Ok(out)
344    }
345
346    fn exec_node_by_property_scan(
347        &self,
348        plan: &PhysicalPlan,
349        op: &NodeByPropertyScanExec,
350    ) -> ExecResult<Vec<Row>> {
351        let base_rows = match op.input {
352            Some(input) => self.execute_node(plan, input)?,
353            None => vec![Row::new()],
354        };
355
356        let eval_ctx = EvalContext {
357            storage: self.ctx.storage,
358            params: &self.ctx.params,
359        };
360        let mut out = Vec::new();
361
362        let deadline = self.deadline;
363        for row in base_rows {
364            Self::check_loop_deadline(deadline)?;
365            let expected = eval_expr(&op.value, &row, &eval_ctx);
366
367            if let Some(existing) = row.get(op.var) {
368                match existing {
369                    LoraValue::Node(existing_id) => {
370                        if node_matches_property_filter(
371                            self.ctx.storage,
372                            *existing_id,
373                            &op.labels,
374                            &op.key,
375                            &expected,
376                        ) {
377                            out.push(row);
378                        }
379                    }
380                    other => {
381                        return Err(ExecutorError::ExpectedNodeForExpand {
382                            var: format!("{:?}", op.var),
383                            found: value_kind(other),
384                        });
385                    }
386                }
387                continue;
388            }
389
390            let candidates =
391                indexed_node_property_candidates(self.ctx.storage, &op.labels, &op.key, &expected);
392            for id in candidates.ids {
393                Self::check_loop_deadline(deadline)?;
394                if !candidates.prefiltered
395                    && !node_matches_property_filter(
396                        self.ctx.storage,
397                        id,
398                        &op.labels,
399                        &op.key,
400                        &expected,
401                    )
402                {
403                    continue;
404                }
405                let mut new_row = row.clone();
406                new_row.insert(op.var, LoraValue::Node(id));
407                out.push(new_row);
408            }
409        }
410
411        Ok(out)
412    }
413
414    fn exec_expand(&self, plan: &PhysicalPlan, op: &ExpandExec) -> ExecResult<Vec<Row>> {
415        // Variable-length expansion: delegate to iterative expander.
416        if let Some(range) = &op.range {
417            return self.exec_expand_var_len(plan, op, range);
418        }
419
420        let input_rows = self.execute_node(plan, op.input)?;
421        let mut out = Vec::new();
422
423        for row in input_rows {
424            let src_node_id = match row.get(op.src) {
425                Some(LoraValue::Node(id)) => *id,
426                Some(other) => {
427                    return Err(ExecutorError::ExpectedNodeForExpand {
428                        var: format!("{:?}", op.src),
429                        found: value_kind(other),
430                    });
431                }
432                None => continue,
433            };
434
435            for (rel_id, dst_id) in
436                self.ctx
437                    .storage
438                    .expand_ids(src_node_id, op.direction, &op.types)
439            {
440                if let Some(expr) = op.rel_properties.as_ref() {
441                    let actual_props = self
442                        .ctx
443                        .storage
444                        .with_relationship(rel_id, |rel| rel.properties.clone());
445                    let matches = match actual_props {
446                        Some(props) => {
447                            self.relationship_matches_properties(&props, Some(expr), &row)?
448                        }
449                        None => false,
450                    };
451                    if !matches {
452                        continue;
453                    }
454                }
455
456                if let Some(existing_dst) = row.get(op.dst) {
457                    match existing_dst {
458                        LoraValue::Node(existing_id) if *existing_id == dst_id => {}
459                        LoraValue::Node(_) => continue,
460                        other => {
461                            return Err(ExecutorError::ExpectedNodeForExpand {
462                                var: format!("{:?}", op.dst),
463                                found: value_kind(other),
464                            });
465                        }
466                    }
467                }
468
469                if let Some(rel_var) = op.rel {
470                    if let Some(existing_rel) = row.get(rel_var) {
471                        match existing_rel {
472                            LoraValue::Relationship(existing_id) if *existing_id == rel_id => {}
473                            LoraValue::Relationship(_) => continue,
474                            other => {
475                                return Err(ExecutorError::ExpectedRelationshipForExpand {
476                                    var: format!("{:?}", rel_var),
477                                    found: value_kind(other),
478                                });
479                            }
480                        }
481                    }
482                }
483
484                let mut new_row = row.clone();
485
486                if !new_row.contains_key(op.dst) {
487                    new_row.insert(op.dst, LoraValue::Node(dst_id));
488                }
489
490                if let Some(rel_var) = op.rel {
491                    if !new_row.contains_key(rel_var) {
492                        new_row.insert(rel_var, LoraValue::Relationship(rel_id));
493                    }
494                }
495
496                out.push(new_row);
497            }
498        }
499
500        Ok(out)
501    }
502
503    fn exec_expand_var_len(
504        &self,
505        plan: &PhysicalPlan,
506        op: &ExpandExec,
507        range: &RangeLiteral,
508    ) -> ExecResult<Vec<Row>> {
509        let input_rows = self.execute_node(plan, op.input)?;
510        let (min_hops, max_hops) = resolve_range(range);
511        let mut out = Vec::new();
512
513        for row in input_rows {
514            let src_node_id = match row.get(op.src) {
515                Some(LoraValue::Node(id)) => *id,
516                Some(other) => {
517                    return Err(ExecutorError::ExpectedNodeForExpand {
518                        var: format!("{:?}", op.src),
519                        found: value_kind(other),
520                    });
521                }
522                None => continue,
523            };
524
525            let expansions = variable_length_expand(
526                self.ctx.storage,
527                src_node_id,
528                op.direction,
529                &op.types,
530                min_hops,
531                max_hops,
532            );
533
534            for result in expansions {
535                let mut new_row = row.clone();
536                new_row.insert(op.dst, LoraValue::Node(result.dst_node_id));
537
538                // For variable-length patterns, bind the relationship variable
539                // to a list of relationship IDs traversed.
540                if let Some(rel_var) = op.rel {
541                    // Consume rel_ids — it's owned and no longer needed after this.
542                    let rel_list = LoraValue::List(
543                        result
544                            .rel_ids
545                            .into_iter()
546                            .map(LoraValue::Relationship)
547                            .collect(),
548                    );
549                    new_row.insert(rel_var, rel_list);
550                }
551
552                out.push(new_row);
553            }
554        }
555
556        Ok(out)
557    }
558
559    fn relationship_matches_properties(
560        &self,
561        actual: &Properties,
562        expected_expr: Option<&ResolvedExpr>,
563        row: &Row,
564    ) -> ExecResult<bool> {
565        let Some(expr) = expected_expr else {
566            return Ok(true);
567        };
568
569        let eval_ctx = EvalContext {
570            storage: self.ctx.storage,
571            params: &self.ctx.params,
572        };
573
574        let expected = eval_expr(expr, row, &eval_ctx);
575
576        let LoraValue::Map(expected_map) = expected else {
577            return Err(ExecutorError::ExpectedPropertyMap {
578                found: value_kind(&expected),
579            });
580        };
581
582        Ok(expected_map.iter().all(|(key, expected_value)| {
583            actual
584                .get(key)
585                .map(|actual_value| value_matches_property_value(expected_value, actual_value))
586                .unwrap_or(false)
587        }))
588    }
589
590    fn exec_filter(&self, plan: &PhysicalPlan, op: &FilterExec) -> ExecResult<Vec<Row>> {
591        let input_rows = self.execute_node(plan, op.input)?;
592        let eval_ctx = EvalContext {
593            storage: self.ctx.storage,
594            params: &self.ctx.params,
595        };
596
597        filter_rows_checked(input_rows, &op.predicate, &eval_ctx)
598    }
599
600    fn exec_projection(&self, plan: &PhysicalPlan, op: &ProjectionExec) -> ExecResult<Vec<Row>> {
601        let input_rows = self.execute_node(plan, op.input)?;
602        let eval_ctx = EvalContext {
603            storage: self.ctx.storage,
604            params: &self.ctx.params,
605        };
606
607        project_rows_checked(input_rows, op, &eval_ctx)
608    }
609
610    fn hydrate_value(&self, value: LoraValue) -> LoraValue {
611        match value {
612            LoraValue::Node(id) => self.hydrate_node(id),
613            LoraValue::Relationship(id) => self.hydrate_relationship(id),
614            LoraValue::List(values) => {
615                LoraValue::List(values.into_iter().map(|v| self.hydrate_value(v)).collect())
616            }
617            LoraValue::Map(map) => LoraValue::Map(
618                map.into_iter()
619                    .map(|(k, v)| (k, self.hydrate_value(v)))
620                    .collect(),
621            ),
622            other => other,
623        }
624    }
625
626    fn hydrate_node(&self, id: u64) -> LoraValue {
627        self.ctx
628            .storage
629            .with_node(id, hydrate_node_record)
630            .unwrap_or(LoraValue::Null)
631    }
632
633    fn hydrate_relationship(&self, id: u64) -> LoraValue {
634        self.ctx
635            .storage
636            .with_relationship(id, hydrate_relationship_record)
637            .unwrap_or(LoraValue::Null)
638    }
639
640    fn exec_unwind(&self, plan: &PhysicalPlan, op: &UnwindExec) -> ExecResult<Vec<Row>> {
641        let input_rows = self.execute_node(plan, op.input)?;
642        let eval_ctx = EvalContext {
643            storage: self.ctx.storage,
644            params: &self.ctx.params,
645        };
646
647        let mut out = Vec::new();
648
649        for row in input_rows {
650            match eval_expr(&op.expr, &row, &eval_ctx) {
651                LoraValue::List(values) => {
652                    for value in values {
653                        let mut new_row = row.clone();
654                        new_row.insert(op.alias, value);
655                        out.push(new_row);
656                    }
657                }
658                LoraValue::Null => {}
659                other => {
660                    let mut new_row = row;
661                    new_row.insert(op.alias, other);
662                    out.push(new_row);
663                }
664            }
665        }
666
667        Ok(out)
668    }
669
670    fn exec_hash_aggregation(
671        &self,
672        plan: &PhysicalPlan,
673        op: &HashAggregationExec,
674    ) -> ExecResult<Vec<Row>> {
675        let input_rows = self.execute_node(plan, op.input)?;
676        let eval_ctx = EvalContext {
677            storage: self.ctx.storage,
678            params: &self.ctx.params,
679        };
680
681        // Streaming fold fast path. When every aggregate is a fold-only
682        // function (count/sum/min/max/avg, no DISTINCT) we never buffer
683        // input rows by group — we fold each row's aggregate value into
684        // running per-group state. Memory drops from O(input_rows) to
685        // O(groups).
686        if let Some(specs) = crate::pull::classify_streamable_aggregates(&op.aggregates) {
687            return self.exec_hash_aggregation_streaming(
688                input_rows,
689                &op.group_by,
690                &op.aggregates,
691                &specs,
692                &eval_ctx,
693            );
694        }
695
696        let mut groups: BTreeMap<Vec<GroupValueKey>, Vec<Row>> = BTreeMap::new();
697
698        if op.group_by.is_empty() {
699            groups.insert(Vec::new(), input_rows);
700        } else {
701            for row in input_rows {
702                let mut key = Vec::with_capacity(op.group_by.len());
703                for proj in &op.group_by {
704                    let value = eval_expr_result(&proj.expr, &row, &eval_ctx)
705                        .map_err(ExecutorError::RuntimeError)?;
706                    key.push(GroupValueKey::from_value(&value));
707                }
708
709                groups.entry(key).or_default().push(row);
710            }
711        }
712
713        let mut out = Vec::new();
714
715        for rows in groups.into_values() {
716            let mut result = Row::new();
717
718            if let Some(first) = rows.first() {
719                for proj in &op.group_by {
720                    let value = eval_expr_result(&proj.expr, first, &eval_ctx)
721                        .map_err(ExecutorError::RuntimeError)?;
722                    let value = self.hydrate_value(value);
723                    result.insert_named(proj.output, proj.name.clone(), value);
724                }
725            }
726
727            for proj in &op.aggregates {
728                let value = compute_aggregate_expr(&proj.expr, &rows, &eval_ctx)?;
729                result.insert_named(proj.output, proj.name.clone(), value);
730            }
731
732            out.push(result);
733        }
734
735        Ok(out)
736    }
737
738    fn exec_hash_aggregation_streaming(
739        &self,
740        input_rows: Vec<Row>,
741        group_by: &[ResolvedProjection],
742        aggregates: &[ResolvedProjection],
743        specs: &[crate::pull::StreamableAggSpec],
744        eval_ctx: &EvalContext<'_, S>,
745    ) -> ExecResult<Vec<Row>> {
746        // No-group-by fast path: a single accumulator, no BTreeMap.
747        if group_by.is_empty() {
748            // `count(*)`-only shortcut. The buffered immutable executor
749            // already materialised every input row into `Vec<Row>`, so the
750            // aggregate is just the row count — no need to fold per row.
751            // This is the v0.6 shape that the streaming refactor lost.
752            if specs
753                .iter()
754                .all(|s| matches!(s.kind, crate::pull::StreamableAggKind::CountAll))
755            {
756                let count = LoraValue::Int(input_rows.len() as i64);
757                let mut result = Row::new();
758                for proj in aggregates {
759                    result.insert_named(proj.output, proj.name.clone(), count.clone());
760                }
761                return Ok(vec![result]);
762            }
763
764            let mut aggs: Vec<crate::pull::AggState> = specs
765                .iter()
766                .map(|s| crate::pull::AggState::seed(s.kind))
767                .collect();
768            for row in &input_rows {
769                for (i, spec) in specs.iter().enumerate() {
770                    let value = match &spec.arg {
771                        Some(arg) => eval_expr_result(arg, row, eval_ctx)
772                            .map_err(ExecutorError::RuntimeError)?,
773                        None => LoraValue::Null,
774                    };
775                    aggs[i].fold(spec.kind, value);
776                }
777            }
778            let mut result = Row::new();
779            for (i, proj) in aggregates.iter().enumerate() {
780                let value =
781                    std::mem::replace(&mut aggs[i], crate::pull::AggState::seed(specs[i].kind))
782                        .finalize(specs[i].kind);
783                result.insert_named(proj.output, proj.name.clone(), value);
784            }
785            return Ok(vec![result]);
786        }
787
788        // Group-by path: per-group running accumulator. The first row in
789        // each group is retained so we can compute the group_by output
790        // expressions later; nothing else from the input is buffered.
791        let mut groups: BTreeMap<Vec<GroupValueKey>, (Row, Vec<crate::pull::AggState>)> =
792            BTreeMap::new();
793
794        for row in input_rows {
795            let mut key = Vec::with_capacity(group_by.len());
796            for proj in group_by {
797                let value = eval_expr_result(&proj.expr, &row, eval_ctx)
798                    .map_err(ExecutorError::RuntimeError)?;
799                key.push(GroupValueKey::from_value(&value));
800            }
801
802            let entry = groups.entry(key).or_insert_with(|| {
803                (
804                    row.clone(),
805                    specs
806                        .iter()
807                        .map(|s| crate::pull::AggState::seed(s.kind))
808                        .collect(),
809                )
810            });
811
812            for (i, spec) in specs.iter().enumerate() {
813                let value = match &spec.arg {
814                    Some(arg) => eval_expr_result(arg, &row, eval_ctx)
815                        .map_err(ExecutorError::RuntimeError)?,
816                    None => LoraValue::Null,
817                };
818                entry.1[i].fold(spec.kind, value);
819            }
820        }
821
822        let mut out = Vec::with_capacity(groups.len());
823        for (_, (first_row, mut aggs)) in groups {
824            let mut result = Row::new();
825            for proj in group_by {
826                let value = eval_expr_result(&proj.expr, &first_row, eval_ctx)
827                    .map_err(ExecutorError::RuntimeError)?;
828                let value = self.hydrate_value(value);
829                result.insert_named(proj.output, proj.name.clone(), value);
830            }
831            for (i, proj) in aggregates.iter().enumerate() {
832                let value =
833                    std::mem::replace(&mut aggs[i], crate::pull::AggState::seed(specs[i].kind))
834                        .finalize(specs[i].kind);
835                result.insert_named(proj.output, proj.name.clone(), value);
836            }
837            out.push(result);
838        }
839        Ok(out)
840    }
841
842    fn exec_sort(&self, plan: &PhysicalPlan, op: &SortExec) -> ExecResult<Vec<Row>> {
843        let mut rows = self.execute_node(plan, op.input)?;
844        let eval_ctx = EvalContext {
845            storage: self.ctx.storage,
846            params: &self.ctx.params,
847        };
848
849        rows.sort_by(|a, b| {
850            for item in &op.items {
851                let ord = compare_sort_item(item, a, b, &eval_ctx);
852                if ord != Ordering::Equal {
853                    return ord;
854                }
855            }
856            Ordering::Equal
857        });
858
859        Ok(rows)
860    }
861
862    fn exec_limit(&self, plan: &PhysicalPlan, op: &LimitExec) -> ExecResult<Vec<Row>> {
863        let mut rows = self.execute_node(plan, op.input)?;
864        let eval_ctx = EvalContext {
865            storage: self.ctx.storage,
866            params: &self.ctx.params,
867        };
868
869        let limit = op
870            .limit
871            .as_ref()
872            .and_then(|e| eval_expr(e, &Row::new(), &eval_ctx).as_i64())
873            .unwrap_or(rows.len() as i64)
874            .max(0) as usize;
875
876        let skip = op
877            .skip
878            .as_ref()
879            .and_then(|e| eval_expr(e, &Row::new(), &eval_ctx).as_i64())
880            .unwrap_or(0)
881            .max(0) as usize;
882
883        if skip >= rows.len() {
884            return Ok(Vec::new());
885        }
886
887        rows.drain(0..skip);
888        rows.truncate(limit);
889        Ok(rows)
890    }
891
892    fn exec_optional_match(
893        &self,
894        plan: &PhysicalPlan,
895        op: &OptionalMatchExec,
896    ) -> ExecResult<Vec<Row>> {
897        let input_rows = self.execute_node(plan, op.input)?;
898
899        // The inner plan is built to start from Argument (an empty row) and is
900        // read-only, so its output does not depend on the upstream input. Execute
901        // it once and reuse the result across every input row, instead of
902        // producing |input_rows| × |inner_rows| allocations.
903        let inner_rows = self.execute_node(plan, op.inner)?;
904
905        let mut out = Vec::new();
906
907        for input_row in input_rows {
908            let mut matched = false;
909
910            for inner_row in &inner_rows {
911                // Each variable already bound in input_row must match.
912                let compatible = input_row
913                    .iter()
914                    .all(|(var, val)| match inner_row.get(*var) {
915                        Some(inner_val) => inner_val == val,
916                        None => true,
917                    });
918                if !compatible {
919                    continue;
920                }
921
922                let mut merged = input_row.clone();
923                for (var, name, val) in inner_row.iter_named() {
924                    if !merged.contains_key(*var) {
925                        merged.insert_named(*var, name.into_owned(), val.clone());
926                    }
927                }
928                out.push(merged);
929                matched = true;
930            }
931
932            if !matched {
933                let mut null_row = input_row;
934                for &var_id in &op.new_vars {
935                    if !null_row.contains_key(var_id) {
936                        null_row.insert(var_id, LoraValue::Null);
937                    }
938                }
939                out.push(null_row);
940            }
941        }
942
943        Ok(out)
944    }
945
946    fn exec_path_build(&self, plan: &PhysicalPlan, op: &PathBuildExec) -> ExecResult<Vec<Row>> {
947        let input_rows = self.execute_node(plan, op.input)?;
948        let mut rows: Vec<Row> = input_rows
949            .into_iter()
950            .map(|mut row| {
951                let path = build_path_value(&row, &op.node_vars, &op.rel_vars, self.ctx.storage);
952                row.insert(op.output, path);
953                row
954            })
955            .collect();
956
957        if let Some(all) = op.shortest_path_all {
958            rows = filter_shortest_paths(rows, op.output, all);
959        }
960        Ok(rows)
961    }
962}