Skip to main content

lora_executor/executor/
mutable.rs

1//! Mutable buffered executor: applies CREATE / MERGE / DELETE / SET /
2//! REMOVE on top of the read-side operator set.
3//!
4//! [`MutableExecutor`] mirrors the read-only [`super::immutable::Executor`]
5//! for all read operators (so a write op above any read subtree
6//! materializes the same way) and adds the per-row write
7//! implementations. The streaming pull pipeline in `crate::pull` runs
8//! `MutableExecutor::apply_write_op` row-by-row through the
9//! `StreamingWriteCursor` fast path; the buffered `exec_*` methods
10//! here handle the fallback when a write op's input subtree is not
11//! fully streamable.
12
13use crate::errors::{value_kind, ExecResult, ExecutorError};
14use crate::eval::{clear_eval_error, eval_expr, eval_expr_result, EvalContext};
15use crate::value::{lora_value_to_property, LoraValue, Row};
16use crate::{project_rows, ExecuteOptions, QueryResult};
17
18use lora_analyzer::{
19    symbols::VarId, ResolvedExpr, ResolvedPattern, ResolvedPatternElement, ResolvedPatternPart,
20    ResolvedProjection, ResolvedRemoveItem, ResolvedSetItem,
21};
22use lora_ast::{Direction, RangeLiteral};
23use lora_compiler::physical::*;
24use lora_compiler::CompiledQuery;
25use lora_store::{GraphStorageMut, NodeId, Properties};
26
27use std::cmp::Ordering;
28use std::collections::BTreeMap;
29use std::time::Instant;
30use tracing::{debug, error, trace};
31
32use super::helpers::{
33    build_path_value, check_deadline_at, compare_sort_item, compute_aggregate_expr, dedup_rows,
34    eval_properties_expr, filter_rows_checked, filter_shortest_paths, flatten_label_groups,
35    hydrate_node_record, hydrate_relationship_record, indexed_node_property_candidates,
36    label_group_candidates_prefiltered, node_matches_label_groups, node_matches_property_filter,
37    project_rows_checked, resolve_range, scan_node_ids_for_label_groups,
38    value_matches_property_value, variable_length_expand, GroupValueKey,
39};
40
41/// Lightweight target for SET property-mutation paths. Lets the SET logic
42/// borrow the row entry (just pulling out the id) instead of cloning the
43/// whole `LoraValue`.
44#[derive(Clone, Copy)]
45enum EntityTarget {
46    Node(NodeId),
47    Relationship(u64),
48}
49
50fn entity_target_from_value(value: &LoraValue) -> ExecResult<EntityTarget> {
51    match value {
52        LoraValue::Node(id) => Ok(EntityTarget::Node(*id)),
53        LoraValue::Relationship(id) => Ok(EntityTarget::Relationship(*id)),
54        other => Err(ExecutorError::InvalidSetTarget {
55            found: value_kind(other),
56        }),
57    }
58}
59
60pub struct MutableExecutionContext<'a, S: GraphStorageMut> {
61    pub storage: &'a mut S,
62    pub params: BTreeMap<String, LoraValue>,
63}
64
65pub struct MutableExecutor<'a, S: GraphStorageMut> {
66    ctx: MutableExecutionContext<'a, S>,
67    deadline: Option<Instant>,
68}
69
70impl<'a, S: GraphStorageMut> MutableExecutor<'a, S> {
71    pub fn new(ctx: MutableExecutionContext<'a, S>) -> Self {
72        Self {
73            ctx,
74            deadline: None,
75        }
76    }
77
78    pub fn with_deadline(ctx: MutableExecutionContext<'a, S>, deadline: Option<Instant>) -> Self {
79        Self { ctx, deadline }
80    }
81
82    #[inline]
83    fn check_deadline(&self) -> ExecResult<()> {
84        if let Some(deadline) = self.deadline {
85            check_deadline_at(deadline)
86        } else {
87            Ok(())
88        }
89    }
90
91    #[inline]
92    fn check_loop_deadline(deadline: Option<Instant>) -> ExecResult<()> {
93        if let Some(deadline) = deadline {
94            check_deadline_at(deadline)
95        } else {
96            Ok(())
97        }
98    }
99
100    pub fn execute(
101        &mut self,
102        plan: &PhysicalPlan,
103        options: Option<ExecuteOptions>,
104    ) -> ExecResult<QueryResult> {
105        let rows = self.execute_rows(plan)?;
106        Ok(project_rows(rows, options.unwrap_or_default()))
107    }
108
109    pub fn execute_rows(&mut self, plan: &PhysicalPlan) -> ExecResult<Vec<Row>> {
110        self.check_deadline()?;
111        // Clear any error residue that a previous query on this thread may have
112        // left in the thread-local eval-error slot.
113        clear_eval_error();
114
115        let rows = self.execute_node(plan, plan.root)?;
116        Ok(rows
117            .into_iter()
118            .map(|row| self.hydrate_row(row))
119            .collect::<Vec<_>>())
120    }
121
122    /// Execute a compiled query that may include UNION branches.
123    pub fn execute_compiled(
124        &mut self,
125        compiled: &CompiledQuery,
126        options: Option<ExecuteOptions>,
127    ) -> ExecResult<QueryResult> {
128        let rows = self.execute_compiled_rows(compiled)?;
129        Ok(project_rows(rows, options.unwrap_or_default()))
130    }
131
132    pub fn execute_compiled_rows(&mut self, compiled: &CompiledQuery) -> ExecResult<Vec<Row>> {
133        self.check_deadline()?;
134        if compiled.unions.is_empty() {
135            return self.execute_rows(&compiled.physical);
136        }
137
138        clear_eval_error();
139
140        // Execute the head branch.
141        let mut all_rows = self.execute_and_hydrate(&compiled.physical)?;
142
143        // Execute each UNION branch and combine.
144        // Track whether any branch uses plain UNION (dedup needed).
145        let mut needs_dedup = false;
146
147        for branch in &compiled.unions {
148            self.check_deadline()?;
149            let branch_rows = self.execute_and_hydrate(&branch.physical)?;
150            all_rows.extend(branch_rows);
151
152            if !branch.all {
153                needs_dedup = true;
154            }
155        }
156
157        if needs_dedup {
158            all_rows = dedup_rows(all_rows);
159        }
160
161        Ok(all_rows)
162    }
163
164    fn execute_and_hydrate(&mut self, plan: &PhysicalPlan) -> ExecResult<Vec<Row>> {
165        self.check_deadline()?;
166        let rows = self.execute_node(plan, plan.root)?;
167        Ok(rows.into_iter().map(|row| self.hydrate_row(row)).collect())
168    }
169
170    pub(crate) fn hydrate_row(&self, row: Row) -> Row {
171        let mut out = Row::new();
172
173        for (var, name, value) in row.into_iter_named() {
174            out.insert_named(var, name, self.hydrate_value(value));
175        }
176
177        out
178    }
179
180    fn execute_node(
181        &mut self,
182        plan: &PhysicalPlan,
183        node_id: PhysicalNodeId,
184    ) -> ExecResult<Vec<Row>> {
185        self.check_deadline()?;
186        trace!("mutable execute_node start: node_id={node_id:?}");
187
188        let result = match &plan.nodes[node_id] {
189            PhysicalOp::Argument(op) => self.exec_argument(op),
190            PhysicalOp::NodeScan(op) => self.exec_node_scan(plan, op),
191            PhysicalOp::NodeByLabelScan(op) => self.exec_node_by_label_scan(plan, op),
192            PhysicalOp::NodeByPropertyScan(op) => self.exec_node_by_property_scan(plan, op),
193            PhysicalOp::Expand(op) => self.exec_expand(plan, op),
194            PhysicalOp::Filter(op) => self.exec_filter(plan, op),
195            PhysicalOp::Projection(op) => self.exec_projection(plan, op),
196            PhysicalOp::Unwind(op) => self.exec_unwind(plan, op),
197            PhysicalOp::HashAggregation(op) => self.exec_hash_aggregation(plan, op),
198            PhysicalOp::Sort(op) => self.exec_sort(plan, op),
199            PhysicalOp::Limit(op) => self.exec_limit(plan, op),
200            PhysicalOp::Create(op) => self.exec_create(plan, op),
201            PhysicalOp::Merge(op) => self.exec_merge(plan, op),
202            PhysicalOp::Delete(op) => self.exec_delete(plan, op),
203            PhysicalOp::Set(op) => self.exec_set(plan, op),
204            PhysicalOp::Remove(op) => self.exec_remove(plan, op),
205            PhysicalOp::OptionalMatch(op) => self.exec_optional_match(plan, op),
206            PhysicalOp::PathBuild(op) => self.exec_path_build(plan, op),
207        };
208
209        match &result {
210            Ok(rows) => trace!(
211                "mutable execute_node ok: node_id={node_id:?}, rows={}",
212                rows.len()
213            ),
214            Err(err) => error!("mutable execute_node failed: node_id={node_id:?}, error={err}"),
215        }
216
217        result
218    }
219
220    fn exec_argument(&self, _op: &ArgumentExec) -> ExecResult<Vec<Row>> {
221        Ok(vec![Row::new()])
222    }
223
224    fn exec_node_scan(&mut self, plan: &PhysicalPlan, op: &NodeScanExec) -> ExecResult<Vec<Row>> {
225        let base_rows = match op.input {
226            Some(input) => self.execute_node(plan, input)?,
227            None => vec![Row::new()],
228        };
229
230        let node_ids = self.ctx.storage.all_node_ids();
231        let mut out = Vec::new();
232
233        let deadline = self.deadline;
234        for row in base_rows {
235            Self::check_loop_deadline(deadline)?;
236            if let Some(existing) = row.get(op.var) {
237                match existing {
238                    LoraValue::Node(existing_id) => {
239                        if self.ctx.storage.has_node(*existing_id) {
240                            out.push(row);
241                        }
242                    }
243                    other => {
244                        return Err(ExecutorError::ExpectedNodeForExpand {
245                            var: format!("{:?}", op.var),
246                            found: value_kind(other),
247                        });
248                    }
249                }
250                continue;
251            }
252
253            for &id in &node_ids {
254                Self::check_loop_deadline(deadline)?;
255                let mut new_row = row.clone();
256                new_row.insert(op.var, LoraValue::Node(id));
257                out.push(new_row);
258            }
259        }
260
261        Ok(out)
262    }
263
264    fn exec_node_by_label_scan(
265        &mut self,
266        plan: &PhysicalPlan,
267        op: &NodeByLabelScanExec,
268    ) -> ExecResult<Vec<Row>> {
269        let base_rows = match op.input {
270            Some(input) => self.execute_node(plan, input)?,
271            None => vec![Row::new()],
272        };
273
274        let candidate_ids = scan_node_ids_for_label_groups(&*self.ctx.storage, &op.labels);
275        let candidates_prefiltered = label_group_candidates_prefiltered(&op.labels);
276        let mut out = Vec::new();
277
278        let deadline = self.deadline;
279        for row in base_rows {
280            Self::check_loop_deadline(deadline)?;
281            if let Some(existing) = row.get(op.var) {
282                match existing {
283                    LoraValue::Node(existing_id) => {
284                        let labels_ok = self
285                            .ctx
286                            .storage
287                            .with_node(*existing_id, |n| {
288                                node_matches_label_groups(&n.labels, &op.labels)
289                            })
290                            .unwrap_or(false);
291                        if labels_ok {
292                            out.push(row);
293                        }
294                    }
295                    other => {
296                        return Err(ExecutorError::ExpectedNodeForExpand {
297                            var: format!("{:?}", op.var),
298                            found: value_kind(other),
299                        });
300                    }
301                }
302                continue;
303            }
304
305            for &id in &candidate_ids {
306                Self::check_loop_deadline(deadline)?;
307                if !candidates_prefiltered {
308                    let labels_ok = self
309                        .ctx
310                        .storage
311                        .with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
312                        .unwrap_or(false);
313                    if !labels_ok {
314                        continue;
315                    }
316                }
317                let mut new_row = row.clone();
318                new_row.insert(op.var, LoraValue::Node(id));
319                out.push(new_row);
320            }
321        }
322
323        Ok(out)
324    }
325
326    fn exec_node_by_property_scan(
327        &mut self,
328        plan: &PhysicalPlan,
329        op: &NodeByPropertyScanExec,
330    ) -> ExecResult<Vec<Row>> {
331        let base_rows = match op.input {
332            Some(input) => self.execute_node(plan, input)?,
333            None => vec![Row::new()],
334        };
335
336        let mut out = Vec::new();
337
338        let deadline = self.deadline;
339        for row in base_rows {
340            Self::check_loop_deadline(deadline)?;
341            let expected = {
342                let eval_ctx = EvalContext {
343                    storage: &*self.ctx.storage,
344                    params: &self.ctx.params,
345                };
346                eval_expr(&op.value, &row, &eval_ctx)
347            };
348
349            if let Some(existing) = row.get(op.var) {
350                match existing {
351                    LoraValue::Node(existing_id) => {
352                        if node_matches_property_filter(
353                            &*self.ctx.storage,
354                            *existing_id,
355                            &op.labels,
356                            &op.key,
357                            &expected,
358                        ) {
359                            out.push(row);
360                        }
361                    }
362                    other => {
363                        return Err(ExecutorError::ExpectedNodeForExpand {
364                            var: format!("{:?}", op.var),
365                            found: value_kind(other),
366                        });
367                    }
368                }
369                continue;
370            }
371
372            let candidates = indexed_node_property_candidates(
373                &*self.ctx.storage,
374                &op.labels,
375                &op.key,
376                &expected,
377            );
378            for id in candidates.ids {
379                Self::check_loop_deadline(deadline)?;
380                if !candidates.prefiltered
381                    && !node_matches_property_filter(
382                        &*self.ctx.storage,
383                        id,
384                        &op.labels,
385                        &op.key,
386                        &expected,
387                    )
388                {
389                    continue;
390                }
391                let mut new_row = row.clone();
392                new_row.insert(op.var, LoraValue::Node(id));
393                out.push(new_row);
394            }
395        }
396
397        Ok(out)
398    }
399
400    fn exec_expand(&mut self, plan: &PhysicalPlan, op: &ExpandExec) -> ExecResult<Vec<Row>> {
401        // Variable-length expansion: delegate to iterative expander.
402        if let Some(range) = &op.range {
403            return self.exec_expand_var_len(plan, op, range);
404        }
405
406        let input_rows = self.execute_node(plan, op.input)?;
407        let mut out = Vec::new();
408
409        for row in input_rows {
410            let src_node_id = match row.get(op.src) {
411                Some(LoraValue::Node(id)) => *id,
412                Some(other) => {
413                    return Err(ExecutorError::ExpectedNodeForExpand {
414                        var: format!("{:?}", op.src),
415                        found: value_kind(other),
416                    });
417                }
418                None => continue,
419            };
420
421            for (rel_id, dst_id) in
422                self.ctx
423                    .storage
424                    .expand_ids(src_node_id, op.direction, &op.types)
425            {
426                if let Some(expr) = op.rel_properties.as_ref() {
427                    let actual_props = self
428                        .ctx
429                        .storage
430                        .with_relationship(rel_id, |rel| rel.properties.clone());
431                    let matches = match actual_props {
432                        Some(props) => {
433                            self.relationship_matches_properties(&props, Some(expr), &row)?
434                        }
435                        None => false,
436                    };
437                    if !matches {
438                        continue;
439                    }
440                }
441
442                if let Some(existing_dst) = row.get(op.dst) {
443                    match existing_dst {
444                        LoraValue::Node(existing_id) if *existing_id == dst_id => {}
445                        LoraValue::Node(_) => continue,
446                        other => {
447                            return Err(ExecutorError::ExpectedNodeForExpand {
448                                var: format!("{:?}", op.dst),
449                                found: value_kind(other),
450                            });
451                        }
452                    }
453                }
454
455                if let Some(rel_var) = op.rel {
456                    if let Some(existing_rel) = row.get(rel_var) {
457                        match existing_rel {
458                            LoraValue::Relationship(existing_id) if *existing_id == rel_id => {}
459                            LoraValue::Relationship(_) => continue,
460                            other => {
461                                return Err(ExecutorError::ExpectedRelationshipForExpand {
462                                    var: format!("{:?}", rel_var),
463                                    found: value_kind(other),
464                                });
465                            }
466                        }
467                    }
468                }
469
470                let mut new_row = row.clone();
471
472                if !new_row.contains_key(op.dst) {
473                    new_row.insert(op.dst, LoraValue::Node(dst_id));
474                }
475
476                if let Some(rel_var) = op.rel {
477                    if !new_row.contains_key(rel_var) {
478                        new_row.insert(rel_var, LoraValue::Relationship(rel_id));
479                    }
480                }
481
482                out.push(new_row);
483            }
484        }
485
486        Ok(out)
487    }
488
489    fn exec_expand_var_len(
490        &mut self,
491        plan: &PhysicalPlan,
492        op: &ExpandExec,
493        range: &RangeLiteral,
494    ) -> ExecResult<Vec<Row>> {
495        let input_rows = self.execute_node(plan, op.input)?;
496        let (min_hops, max_hops) = resolve_range(range);
497        let mut out = Vec::new();
498
499        for row in input_rows {
500            let src_node_id = match row.get(op.src) {
501                Some(LoraValue::Node(id)) => *id,
502                Some(other) => {
503                    return Err(ExecutorError::ExpectedNodeForExpand {
504                        var: format!("{:?}", op.src),
505                        found: value_kind(other),
506                    });
507                }
508                None => continue,
509            };
510
511            let expansions = variable_length_expand(
512                &*self.ctx.storage,
513                src_node_id,
514                op.direction,
515                &op.types,
516                min_hops,
517                max_hops,
518            );
519
520            for result in expansions {
521                let mut new_row = row.clone();
522                new_row.insert(op.dst, LoraValue::Node(result.dst_node_id));
523
524                if let Some(rel_var) = op.rel {
525                    // Consume rel_ids — it's owned and no longer needed after this.
526                    let rel_list = LoraValue::List(
527                        result
528                            .rel_ids
529                            .into_iter()
530                            .map(LoraValue::Relationship)
531                            .collect(),
532                    );
533                    new_row.insert(rel_var, rel_list);
534                }
535
536                out.push(new_row);
537            }
538        }
539
540        Ok(out)
541    }
542
543    fn relationship_matches_properties(
544        &self,
545        actual: &Properties,
546        expected_expr: Option<&ResolvedExpr>,
547        row: &Row,
548    ) -> ExecResult<bool> {
549        let Some(expr) = expected_expr else {
550            return Ok(true);
551        };
552
553        let eval_ctx = EvalContext {
554            storage: &*self.ctx.storage,
555            params: &self.ctx.params,
556        };
557
558        let expected = eval_expr(expr, row, &eval_ctx);
559
560        let LoraValue::Map(expected_map) = expected else {
561            return Err(ExecutorError::ExpectedPropertyMap {
562                found: value_kind(&expected),
563            });
564        };
565
566        Ok(expected_map.iter().all(|(key, expected_value)| {
567            actual
568                .get(key)
569                .map(|actual_value| value_matches_property_value(expected_value, actual_value))
570                .unwrap_or(false)
571        }))
572    }
573
574    fn exec_filter(&mut self, plan: &PhysicalPlan, op: &FilterExec) -> ExecResult<Vec<Row>> {
575        let input_rows = self.execute_node(plan, op.input)?;
576        let eval_ctx = EvalContext {
577            storage: &*self.ctx.storage,
578            params: &self.ctx.params,
579        };
580
581        filter_rows_checked(input_rows, &op.predicate, &eval_ctx)
582    }
583
584    fn exec_projection(
585        &mut self,
586        plan: &PhysicalPlan,
587        op: &ProjectionExec,
588    ) -> ExecResult<Vec<Row>> {
589        let input_rows = self.execute_node(plan, op.input)?;
590        let eval_ctx = EvalContext {
591            storage: &*self.ctx.storage,
592            params: &self.ctx.params,
593        };
594
595        project_rows_checked(input_rows, op, &eval_ctx)
596    }
597
598    fn hydrate_value(&self, value: LoraValue) -> LoraValue {
599        match value {
600            LoraValue::Node(id) => self.hydrate_node(id),
601            LoraValue::Relationship(id) => self.hydrate_relationship(id),
602            LoraValue::List(values) => {
603                LoraValue::List(values.into_iter().map(|v| self.hydrate_value(v)).collect())
604            }
605            LoraValue::Map(map) => LoraValue::Map(
606                map.into_iter()
607                    .map(|(k, v)| (k, self.hydrate_value(v)))
608                    .collect(),
609            ),
610            other => other,
611        }
612    }
613
614    fn hydrate_node(&self, id: u64) -> LoraValue {
615        self.ctx
616            .storage
617            .with_node(id, hydrate_node_record)
618            .unwrap_or(LoraValue::Null)
619    }
620
621    fn hydrate_relationship(&self, id: u64) -> LoraValue {
622        self.ctx
623            .storage
624            .with_relationship(id, hydrate_relationship_record)
625            .unwrap_or(LoraValue::Null)
626    }
627
628    fn exec_unwind(&mut self, plan: &PhysicalPlan, op: &UnwindExec) -> ExecResult<Vec<Row>> {
629        let input_rows = self.execute_node(plan, op.input)?;
630        let eval_ctx = EvalContext {
631            storage: &*self.ctx.storage,
632            params: &self.ctx.params,
633        };
634
635        let mut out = Vec::new();
636
637        for row in input_rows {
638            match eval_expr(&op.expr, &row, &eval_ctx) {
639                LoraValue::List(values) => {
640                    for value in values {
641                        let mut new_row = row.clone();
642                        new_row.insert(op.alias, value);
643                        out.push(new_row);
644                    }
645                }
646                LoraValue::Null => {}
647                other => {
648                    let mut new_row = row;
649                    new_row.insert(op.alias, other);
650                    out.push(new_row);
651                }
652            }
653        }
654
655        Ok(out)
656    }
657
658    fn exec_hash_aggregation(
659        &mut self,
660        plan: &PhysicalPlan,
661        op: &HashAggregationExec,
662    ) -> ExecResult<Vec<Row>> {
663        let input_rows = self.execute_node(plan, op.input)?;
664        let eval_ctx = EvalContext {
665            storage: &*self.ctx.storage,
666            params: &self.ctx.params,
667        };
668
669        // Streaming fold fast path — same logic as the read-side
670        // `Executor::exec_hash_aggregation`. See that method for the full
671        // rationale.
672        if let Some(specs) = crate::pull::classify_streamable_aggregates(&op.aggregates) {
673            return self.exec_hash_aggregation_streaming(
674                input_rows,
675                &op.group_by,
676                &op.aggregates,
677                &specs,
678                &eval_ctx,
679            );
680        }
681
682        let mut groups: BTreeMap<Vec<GroupValueKey>, Vec<Row>> = BTreeMap::new();
683
684        if op.group_by.is_empty() {
685            groups.insert(Vec::new(), input_rows);
686        } else {
687            for row in input_rows {
688                let mut key = Vec::with_capacity(op.group_by.len());
689                for proj in &op.group_by {
690                    let value = eval_expr_result(&proj.expr, &row, &eval_ctx)
691                        .map_err(ExecutorError::RuntimeError)?;
692                    key.push(GroupValueKey::from_value(&value));
693                }
694
695                groups.entry(key).or_default().push(row);
696            }
697        }
698
699        let mut out = Vec::new();
700
701        for rows in groups.into_values() {
702            let mut result = Row::new();
703
704            if let Some(first) = rows.first() {
705                for proj in &op.group_by {
706                    let value = eval_expr_result(&proj.expr, first, &eval_ctx)
707                        .map_err(ExecutorError::RuntimeError)?;
708                    let value = self.hydrate_value(value);
709                    result.insert_named(proj.output, proj.name.clone(), value);
710                }
711            }
712
713            for proj in &op.aggregates {
714                let value = compute_aggregate_expr(&proj.expr, &rows, &eval_ctx)?;
715                result.insert_named(proj.output, proj.name.clone(), value);
716            }
717
718            out.push(result);
719        }
720
721        Ok(out)
722    }
723
724    fn exec_hash_aggregation_streaming(
725        &self,
726        input_rows: Vec<Row>,
727        group_by: &[ResolvedProjection],
728        aggregates: &[ResolvedProjection],
729        specs: &[crate::pull::StreamableAggSpec],
730        eval_ctx: &EvalContext<'_, S>,
731    ) -> ExecResult<Vec<Row>> {
732        if group_by.is_empty() {
733            let mut aggs: Vec<crate::pull::AggState> = specs
734                .iter()
735                .map(|s| crate::pull::AggState::seed(s.kind))
736                .collect();
737            for row in &input_rows {
738                for (i, spec) in specs.iter().enumerate() {
739                    let value = match &spec.arg {
740                        Some(arg) => eval_expr_result(arg, row, eval_ctx)
741                            .map_err(ExecutorError::RuntimeError)?,
742                        None => LoraValue::Null,
743                    };
744                    aggs[i].fold(spec.kind, value);
745                }
746            }
747            let mut result = Row::new();
748            for (i, proj) in aggregates.iter().enumerate() {
749                let value =
750                    std::mem::replace(&mut aggs[i], crate::pull::AggState::seed(specs[i].kind))
751                        .finalize(specs[i].kind);
752                result.insert_named(proj.output, proj.name.clone(), value);
753            }
754            return Ok(vec![result]);
755        }
756
757        let mut groups: BTreeMap<Vec<GroupValueKey>, (Row, Vec<crate::pull::AggState>)> =
758            BTreeMap::new();
759
760        for row in input_rows {
761            let mut key = Vec::with_capacity(group_by.len());
762            for proj in group_by {
763                let value = eval_expr_result(&proj.expr, &row, eval_ctx)
764                    .map_err(ExecutorError::RuntimeError)?;
765                key.push(GroupValueKey::from_value(&value));
766            }
767
768            let entry = groups.entry(key).or_insert_with(|| {
769                (
770                    row.clone(),
771                    specs
772                        .iter()
773                        .map(|s| crate::pull::AggState::seed(s.kind))
774                        .collect(),
775                )
776            });
777
778            for (i, spec) in specs.iter().enumerate() {
779                let value = match &spec.arg {
780                    Some(arg) => eval_expr_result(arg, &row, eval_ctx)
781                        .map_err(ExecutorError::RuntimeError)?,
782                    None => LoraValue::Null,
783                };
784                entry.1[i].fold(spec.kind, value);
785            }
786        }
787
788        let mut out = Vec::with_capacity(groups.len());
789        for (_, (first_row, mut aggs)) in groups {
790            let mut result = Row::new();
791            for proj in group_by {
792                let value = eval_expr_result(&proj.expr, &first_row, eval_ctx)
793                    .map_err(ExecutorError::RuntimeError)?;
794                let value = self.hydrate_value(value);
795                result.insert_named(proj.output, proj.name.clone(), value);
796            }
797            for (i, proj) in aggregates.iter().enumerate() {
798                let value =
799                    std::mem::replace(&mut aggs[i], crate::pull::AggState::seed(specs[i].kind))
800                        .finalize(specs[i].kind);
801                result.insert_named(proj.output, proj.name.clone(), value);
802            }
803            out.push(result);
804        }
805        Ok(out)
806    }
807
808    fn exec_sort(&mut self, plan: &PhysicalPlan, op: &SortExec) -> ExecResult<Vec<Row>> {
809        let mut rows = self.execute_node(plan, op.input)?;
810        let eval_ctx = EvalContext {
811            storage: &*self.ctx.storage,
812            params: &self.ctx.params,
813        };
814
815        rows.sort_by(|a, b| {
816            for item in &op.items {
817                let ord = compare_sort_item(item, a, b, &eval_ctx);
818                if ord != Ordering::Equal {
819                    return ord;
820                }
821            }
822            Ordering::Equal
823        });
824
825        Ok(rows)
826    }
827
828    fn exec_limit(&mut self, plan: &PhysicalPlan, op: &LimitExec) -> ExecResult<Vec<Row>> {
829        let mut rows = self.execute_node(plan, op.input)?;
830        let eval_ctx = EvalContext {
831            storage: &*self.ctx.storage,
832            params: &self.ctx.params,
833        };
834
835        let limit = op
836            .limit
837            .as_ref()
838            .and_then(|e| eval_expr(e, &Row::new(), &eval_ctx).as_i64())
839            .unwrap_or(rows.len() as i64)
840            .max(0) as usize;
841
842        let skip = op
843            .skip
844            .as_ref()
845            .and_then(|e| eval_expr(e, &Row::new(), &eval_ctx).as_i64())
846            .unwrap_or(0)
847            .max(0) as usize;
848
849        if skip >= rows.len() {
850            return Ok(Vec::new());
851        }
852
853        rows.drain(0..skip);
854        rows.truncate(limit);
855        Ok(rows)
856    }
857
858    fn exec_optional_match(
859        &mut self,
860        plan: &PhysicalPlan,
861        op: &OptionalMatchExec,
862    ) -> ExecResult<Vec<Row>> {
863        let input_rows = self.execute_node(plan, op.input)?;
864
865        // Inner plan is read-only and input-independent; execute once and reuse.
866        let inner_rows = self.execute_node(plan, op.inner)?;
867
868        let mut out = Vec::new();
869
870        for input_row in input_rows {
871            let mut matched = false;
872
873            for inner_row in &inner_rows {
874                let compatible = input_row
875                    .iter()
876                    .all(|(var, val)| match inner_row.get(*var) {
877                        Some(inner_val) => inner_val == val,
878                        None => true,
879                    });
880                if !compatible {
881                    continue;
882                }
883
884                let mut merged = input_row.clone();
885                for (var, name, val) in inner_row.iter_named() {
886                    if !merged.contains_key(*var) {
887                        merged.insert_named(*var, name.into_owned(), val.clone());
888                    }
889                }
890                out.push(merged);
891                matched = true;
892            }
893
894            if !matched {
895                let mut null_row = input_row;
896                for &var_id in &op.new_vars {
897                    if !null_row.contains_key(var_id) {
898                        null_row.insert(var_id, LoraValue::Null);
899                    }
900                }
901                out.push(null_row);
902            }
903        }
904
905        Ok(out)
906    }
907
908    fn exec_path_build(&mut self, plan: &PhysicalPlan, op: &PathBuildExec) -> ExecResult<Vec<Row>> {
909        let input_rows = self.execute_node(plan, op.input)?;
910        let mut rows: Vec<Row> = input_rows
911            .into_iter()
912            .map(|mut row| {
913                let path = build_path_value(&row, &op.node_vars, &op.rel_vars, &*self.ctx.storage);
914                row.insert(op.output, path);
915                row
916            })
917            .collect();
918
919        if let Some(all) = op.shortest_path_all {
920            rows = filter_shortest_paths(rows, op.output, all);
921        }
922        Ok(rows)
923    }
924
925    fn exec_create(&mut self, plan: &PhysicalPlan, op: &CreateExec) -> ExecResult<Vec<Row>> {
926        // Fast path: if the input subtree is fully streamable (no
927        // nested writes, no blocking operators), pull rows one at a
928        // time and apply the create pattern per row, instead of
929        // materializing the whole input. The output Vec still
930        // accumulates — auto-commit-side output streaming is M1.b.
931        if crate::pull::subtree_is_fully_streaming(plan, op.input) {
932            return self.exec_create_streaming_input(plan, op);
933        }
934
935        let input_rows = self.execute_node(plan, op.input)?;
936        let mut out = Vec::with_capacity(input_rows.len());
937
938        for mut row in input_rows {
939            self.apply_create_pattern(&mut row, &op.pattern)?;
940            out.push(row);
941        }
942
943        Ok(out)
944    }
945
946    /// Generic streaming-input loop for write operators whose input
947    /// subtree is fully streamable. Opens a pull-based read cursor
948    /// over the input subtree, calls `apply` per row, and accumulates
949    /// the resulting rows.
950    ///
951    /// # Safety
952    ///
953    /// The upstream [`crate::pull::RowSource`] needs `&S` while it
954    /// lives; the per-row `apply` callback needs `&mut S` (via
955    /// `&mut self`). The existing read-side `RowSource` impls
956    /// materialize their iteration state into owned `Vec`s at
957    /// construction time (see `NodeScanSource::cur_ids`,
958    /// `ExpandSource::cur_edges`, etc. in `pull.rs`), so no live
959    /// `&S` borrow into storage persists across `next_row` calls.
960    /// We exploit that by deriving the read borrow from a raw
961    /// pointer — Rust then doesn't see the shared/mutable conflict
962    /// at compile time, and the dynamic access pattern is
963    /// non-aliasing: read-only inside `next_row`, then mutable
964    /// inside `apply`, never both at the same instant.
965    fn streaming_apply<F>(
966        &mut self,
967        plan: &PhysicalPlan,
968        input: PhysicalNodeId,
969        mut apply: F,
970    ) -> ExecResult<Vec<Row>>
971    where
972        F: FnMut(&mut Self, &mut Row) -> ExecResult<()>,
973    {
974        use std::sync::Arc;
975
976        let storage_ptr: *mut S = self.ctx.storage as *mut S;
977        let params = Arc::new(self.ctx.params.clone());
978
979        // SAFETY: see method-level comment.
980        let storage_ref: &S = unsafe { &*storage_ptr };
981        let mut upstream = crate::pull::build_streaming(plan, input, storage_ref, params)?;
982
983        let mut out = Vec::new();
984        while let Some(mut row) = upstream.next_row()? {
985            apply(self, &mut row)?;
986            out.push(row);
987        }
988
989        Ok(out)
990    }
991
992    /// Streaming-input variant of [`Self::exec_create`]. Delegates
993    /// to [`Self::streaming_apply`].
994    fn exec_create_streaming_input(
995        &mut self,
996        plan: &PhysicalPlan,
997        op: &CreateExec,
998    ) -> ExecResult<Vec<Row>> {
999        self.streaming_apply(plan, op.input, |this, row| {
1000            this.apply_create_pattern(row, &op.pattern)
1001        })
1002    }
1003
1004    fn apply_remove_item(&mut self, row: &Row, item: &ResolvedRemoveItem) -> ExecResult<()> {
1005        match item {
1006            ResolvedRemoveItem::Labels { variable, labels } => match row.get(*variable) {
1007                Some(LoraValue::Node(node_id)) => {
1008                    let node_id = *node_id;
1009                    for label in labels {
1010                        self.ctx.storage.remove_node_label(node_id, label);
1011                    }
1012                    Ok(())
1013                }
1014                Some(other) => Err(ExecutorError::ExpectedNodeForRemoveLabels {
1015                    found: value_kind(other),
1016                }),
1017                None => Err(ExecutorError::UnboundVariableForRemove {
1018                    var: format!("{variable:?}"),
1019                }),
1020            },
1021
1022            ResolvedRemoveItem::Property { expr } => self.remove_property_from_expr(row, expr),
1023        }
1024    }
1025
1026    fn delete_value(&mut self, value: LoraValue, detach: bool) -> ExecResult<()> {
1027        match value {
1028            LoraValue::Null => Ok(()),
1029
1030            LoraValue::Node(node_id) => {
1031                if detach {
1032                    self.ctx.storage.detach_delete_node(node_id);
1033                    Ok(())
1034                } else {
1035                    let ok = self.ctx.storage.delete_node(node_id);
1036                    if ok {
1037                        Ok(())
1038                    } else {
1039                        Err(ExecutorError::DeleteNodeWithRelationships { node_id })
1040                    }
1041                }
1042            }
1043
1044            LoraValue::Relationship(rel_id) => {
1045                let ok = self.ctx.storage.delete_relationship(rel_id);
1046                if ok {
1047                    Ok(())
1048                } else {
1049                    Err(ExecutorError::DeleteRelationshipFailed { rel_id })
1050                }
1051            }
1052
1053            LoraValue::List(values) => {
1054                for v in values {
1055                    self.delete_value(v, detach)?;
1056                }
1057                Ok(())
1058            }
1059
1060            other => Err(ExecutorError::InvalidDeleteTarget {
1061                found: value_kind(&other),
1062            }),
1063        }
1064    }
1065
1066    fn exec_merge(&mut self, plan: &PhysicalPlan, op: &MergeExec) -> ExecResult<Vec<Row>> {
1067        // Streaming-input fast path when the input subtree is fully
1068        // streamable. Per-row work (probe → optionally create →
1069        // ON MATCH / ON CREATE actions) is identical to the
1070        // materialized branch below.
1071        if crate::pull::subtree_is_fully_streaming(plan, op.input) {
1072            return self.streaming_apply(plan, op.input, |this, row| {
1073                let already_bound = this.pattern_part_is_bound(row, &op.pattern_part);
1074                let matched = if already_bound {
1075                    true
1076                } else {
1077                    this.try_match_merge_pattern(row, &op.pattern_part)?
1078                };
1079                if !matched {
1080                    this.apply_create_pattern_part(row, &op.pattern_part)?;
1081                }
1082                for action in &op.actions {
1083                    if action.on_match == matched {
1084                        for item in &action.set.items {
1085                            this.apply_set_item(row, item)?;
1086                        }
1087                    }
1088                }
1089                Ok(())
1090            });
1091        }
1092
1093        let input_rows = self.execute_node(plan, op.input)?;
1094        let mut out = Vec::with_capacity(input_rows.len());
1095
1096        for mut row in input_rows {
1097            // First check if the pattern variable is already bound in the row.
1098            let already_bound = self.pattern_part_is_bound(&row, &op.pattern_part);
1099
1100            let matched = if already_bound {
1101                true
1102            } else {
1103                // Try to find an existing match in the graph.
1104                self.try_match_merge_pattern(&mut row, &op.pattern_part)?
1105            };
1106
1107            if !matched {
1108                self.apply_create_pattern_part(&mut row, &op.pattern_part)?;
1109            }
1110
1111            for action in &op.actions {
1112                if action.on_match == matched {
1113                    for item in &action.set.items {
1114                        self.apply_set_item(&row, item)?;
1115                    }
1116                }
1117            }
1118
1119            out.push(row);
1120        }
1121
1122        Ok(out)
1123    }
1124
1125    /// Try to find an existing node/pattern in the graph matching the MERGE
1126    /// pattern. If found, bind the variable in the row and return true.
1127    fn try_match_merge_pattern(
1128        &self,
1129        row: &mut Row,
1130        part: &ResolvedPatternPart,
1131    ) -> ExecResult<bool> {
1132        match &part.element {
1133            ResolvedPatternElement::Node {
1134                var,
1135                labels,
1136                properties,
1137            } => {
1138                // ID-only candidate discovery; borrow the record during
1139                // label/property filtering to avoid cloning non-matches.
1140                let candidate_ids = if labels.is_empty() {
1141                    self.ctx.storage.all_node_ids()
1142                } else {
1143                    scan_node_ids_for_label_groups(&*self.ctx.storage, labels)
1144                };
1145
1146                // Filter by properties if specified
1147                let eval_ctx = EvalContext {
1148                    storage: &*self.ctx.storage,
1149                    params: &self.ctx.params,
1150                };
1151                let expected_props = properties.as_ref().map(|e| eval_expr(e, row, &eval_ctx));
1152
1153                for id in candidate_ids {
1154                    let matched = self
1155                        .ctx
1156                        .storage
1157                        .with_node(id, |node| {
1158                            if !node_matches_label_groups(&node.labels, labels) {
1159                                return false;
1160                            }
1161                            if let Some(LoraValue::Map(expected)) = &expected_props {
1162                                let all_match = expected.iter().all(|(key, expected_value)| {
1163                                    node.properties
1164                                        .get(key)
1165                                        .map(|actual| {
1166                                            value_matches_property_value(expected_value, actual)
1167                                        })
1168                                        .unwrap_or(false)
1169                                });
1170                                if !all_match {
1171                                    return false;
1172                                }
1173                            }
1174                            true
1175                        })
1176                        .unwrap_or(false);
1177
1178                    if !matched {
1179                        continue;
1180                    }
1181
1182                    // Found a match — bind the variable
1183                    if let Some(var_id) = var {
1184                        row.insert(*var_id, LoraValue::Node(id));
1185                    }
1186                    return Ok(true);
1187                }
1188
1189                Ok(false)
1190            }
1191
1192            ResolvedPatternElement::ShortestPath { .. } => {
1193                // ShortestPath is not valid in MERGE context
1194                Ok(false)
1195            }
1196
1197            ResolvedPatternElement::NodeChain { head, chain } => {
1198                // Resolve the head node — it should be already bound in the row.
1199                let head_node_id = if let Some(var_id) = head.var {
1200                    if let Some(LoraValue::Node(id)) = row.get(var_id) {
1201                        *id
1202                    } else {
1203                        // Try to match head node as a standalone node pattern.
1204                        let node_matched = self.try_match_merge_pattern(
1205                            row,
1206                            &ResolvedPatternPart {
1207                                binding: None,
1208                                element: ResolvedPatternElement::Node {
1209                                    var: head.var,
1210                                    labels: head.labels.clone(),
1211                                    properties: head.properties.clone(),
1212                                },
1213                            },
1214                        )?;
1215                        if !node_matched {
1216                            return Ok(false);
1217                        }
1218                        match row.get(var_id) {
1219                            Some(LoraValue::Node(id)) => *id,
1220                            _ => return Ok(false),
1221                        }
1222                    }
1223                } else {
1224                    return Ok(false);
1225                };
1226
1227                let mut current_node_id = head_node_id;
1228
1229                for step in chain {
1230                    let eval_ctx = EvalContext {
1231                        storage: &*self.ctx.storage,
1232                        params: &self.ctx.params,
1233                    };
1234
1235                    let _ = step.rel.types.first();
1236                    let direction = step.rel.direction;
1237
1238                    // ID-only traversal; look up records by reference only for
1239                    // candidates that pass the label/property filters.
1240                    let edges =
1241                        self.ctx
1242                            .storage
1243                            .expand_ids(current_node_id, direction, &step.rel.types);
1244
1245                    // Try to find a matching edge + target node
1246                    let mut found = false;
1247                    for (rel_id, node_id) in edges {
1248                        // Check target node labels and (optional) properties.
1249                        let node_ok = self
1250                            .ctx
1251                            .storage
1252                            .with_node(node_id, |node_rec| {
1253                                if !node_matches_label_groups(&node_rec.labels, &step.node.labels) {
1254                                    return false;
1255                                }
1256                                if let Some(props_expr) = &step.node.properties {
1257                                    let expected = eval_expr(props_expr, row, &eval_ctx);
1258                                    if let LoraValue::Map(expected_map) = &expected {
1259                                        let all_match =
1260                                            expected_map.iter().all(|(key, expected_val)| {
1261                                                node_rec
1262                                                    .properties
1263                                                    .get(key)
1264                                                    .map(|actual| {
1265                                                        value_matches_property_value(
1266                                                            expected_val,
1267                                                            actual,
1268                                                        )
1269                                                    })
1270                                                    .unwrap_or(false)
1271                                            });
1272                                        if !all_match {
1273                                            return false;
1274                                        }
1275                                    }
1276                                }
1277                                true
1278                            })
1279                            .unwrap_or(false);
1280                        if !node_ok {
1281                            continue;
1282                        }
1283
1284                        // Check relationship properties.
1285                        let rel_ok = self
1286                            .ctx
1287                            .storage
1288                            .with_relationship(rel_id, |rel_rec| {
1289                                if let Some(rel_props_expr) = &step.rel.properties {
1290                                    let expected = eval_expr(rel_props_expr, row, &eval_ctx);
1291                                    if let LoraValue::Map(expected_map) = &expected {
1292                                        let all_match =
1293                                            expected_map.iter().all(|(key, expected_val)| {
1294                                                rel_rec
1295                                                    .properties
1296                                                    .get(key)
1297                                                    .map(|actual| {
1298                                                        value_matches_property_value(
1299                                                            expected_val,
1300                                                            actual,
1301                                                        )
1302                                                    })
1303                                                    .unwrap_or(false)
1304                                            });
1305                                        if !all_match {
1306                                            return false;
1307                                        }
1308                                    }
1309                                }
1310                                true
1311                            })
1312                            .unwrap_or(false);
1313                        if !rel_ok {
1314                            continue;
1315                        }
1316
1317                        // Match found — bind variables
1318                        if let Some(rel_var) = step.rel.var {
1319                            row.insert(rel_var, LoraValue::Relationship(rel_id));
1320                        }
1321                        if let Some(node_var) = step.node.var {
1322                            row.insert(node_var, LoraValue::Node(node_id));
1323                        }
1324                        current_node_id = node_id;
1325                        found = true;
1326                        break;
1327                    }
1328
1329                    if !found {
1330                        return Ok(false);
1331                    }
1332                }
1333
1334                Ok(true)
1335            }
1336        }
1337    }
1338
1339    fn exec_delete(&mut self, plan: &PhysicalPlan, op: &DeleteExec) -> ExecResult<Vec<Row>> {
1340        if crate::pull::subtree_is_fully_streaming(plan, op.input) {
1341            let detach = op.detach;
1342            return self.streaming_apply(plan, op.input, |this, row| {
1343                for expr in &op.expressions {
1344                    let value = {
1345                        let eval_ctx = EvalContext {
1346                            storage: &*this.ctx.storage,
1347                            params: &this.ctx.params,
1348                        };
1349                        eval_expr(expr, row, &eval_ctx)
1350                    };
1351                    this.delete_value(value, detach)?;
1352                }
1353                Ok(())
1354            });
1355        }
1356
1357        let input_rows = self.execute_node(plan, op.input)?;
1358
1359        for row in &input_rows {
1360            for expr in &op.expressions {
1361                let value = {
1362                    let eval_ctx = EvalContext {
1363                        storage: &*self.ctx.storage,
1364                        params: &self.ctx.params,
1365                    };
1366                    eval_expr(expr, row, &eval_ctx)
1367                };
1368
1369                self.delete_value(value, op.detach)?;
1370            }
1371        }
1372
1373        Ok(input_rows)
1374    }
1375
1376    fn exec_set(&mut self, plan: &PhysicalPlan, op: &SetExec) -> ExecResult<Vec<Row>> {
1377        if crate::pull::subtree_is_fully_streaming(plan, op.input) {
1378            return self.streaming_apply(plan, op.input, |this, row| {
1379                for item in &op.items {
1380                    this.apply_set_item(row, item)?;
1381                }
1382                Ok(())
1383            });
1384        }
1385
1386        let input_rows = self.execute_node(plan, op.input)?;
1387
1388        for row in &input_rows {
1389            for item in &op.items {
1390                self.apply_set_item(row, item)?;
1391            }
1392        }
1393
1394        Ok(input_rows)
1395    }
1396
1397    fn exec_remove(&mut self, plan: &PhysicalPlan, op: &RemoveExec) -> ExecResult<Vec<Row>> {
1398        if crate::pull::subtree_is_fully_streaming(plan, op.input) {
1399            return self.streaming_apply(plan, op.input, |this, row| {
1400                for item in &op.items {
1401                    this.apply_remove_item(row, item)?;
1402                }
1403                Ok(())
1404            });
1405        }
1406
1407        let input_rows = self.execute_node(plan, op.input)?;
1408
1409        for row in &input_rows {
1410            for item in &op.items {
1411                self.apply_remove_item(row, item)?;
1412            }
1413        }
1414
1415        Ok(input_rows)
1416    }
1417
1418    fn apply_set_item(&mut self, row: &Row, item: &ResolvedSetItem) -> ExecResult<()> {
1419        match item {
1420            ResolvedSetItem::SetProperty { target, value } => {
1421                let new_value = {
1422                    let eval_ctx = EvalContext {
1423                        storage: &*self.ctx.storage,
1424                        params: &self.ctx.params,
1425                    };
1426                    eval_expr(value, row, &eval_ctx)
1427                };
1428
1429                self.set_property_from_expr(row, target, new_value)
1430            }
1431
1432            ResolvedSetItem::SetVariable { variable, value } => {
1433                // Only need the entity's id — peek at the binding by reference.
1434                let entity_ref =
1435                    row.get(*variable)
1436                        .ok_or(ExecutorError::UnboundVariableForSet {
1437                            var: format!("{variable:?}"),
1438                        })?;
1439                let entity_target = entity_target_from_value(entity_ref)?;
1440
1441                let new_value = {
1442                    let eval_ctx = EvalContext {
1443                        storage: &*self.ctx.storage,
1444                        params: &self.ctx.params,
1445                    };
1446                    eval_expr(value, row, &eval_ctx)
1447                };
1448
1449                self.overwrite_entity_target(entity_target, new_value)
1450            }
1451
1452            ResolvedSetItem::MutateVariable { variable, value } => {
1453                let entity_ref =
1454                    row.get(*variable)
1455                        .ok_or(ExecutorError::UnboundVariableForSet {
1456                            var: format!("{variable:?}"),
1457                        })?;
1458                let entity_target = entity_target_from_value(entity_ref)?;
1459
1460                let patch = {
1461                    let eval_ctx = EvalContext {
1462                        storage: &*self.ctx.storage,
1463                        params: &self.ctx.params,
1464                    };
1465                    eval_expr(value, row, &eval_ctx)
1466                };
1467
1468                self.mutate_entity_target(entity_target, patch)
1469            }
1470
1471            ResolvedSetItem::SetLabels { variable, labels } => match row.get(*variable) {
1472                Some(LoraValue::Node(node_id)) => {
1473                    let node_id = *node_id;
1474                    for label in labels {
1475                        self.ctx.storage.add_node_label(node_id, label);
1476                    }
1477                    Ok(())
1478                }
1479                Some(other) => Err(ExecutorError::ExpectedNodeForSetLabels {
1480                    found: value_kind(other),
1481                }),
1482                None => Err(ExecutorError::UnboundVariableForSet {
1483                    var: format!("{variable:?}"),
1484                }),
1485            },
1486        }
1487    }
1488
1489    fn set_property_from_expr(
1490        &mut self,
1491        row: &Row,
1492        target_expr: &ResolvedExpr,
1493        new_value: LoraValue,
1494    ) -> ExecResult<()> {
1495        let ResolvedExpr::Property { expr, property } = target_expr else {
1496            return Err(ExecutorError::UnsupportedSetTarget);
1497        };
1498
1499        let owner = {
1500            let eval_ctx = EvalContext {
1501                storage: &*self.ctx.storage,
1502                params: &self.ctx.params,
1503            };
1504            eval_expr(expr, row, &eval_ctx)
1505        };
1506
1507        match owner {
1508            LoraValue::Node(node_id) => {
1509                let prop = lora_value_to_property(new_value)
1510                    .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1511                self.ctx
1512                    .storage
1513                    .set_node_property(node_id, property.clone(), prop);
1514                Ok(())
1515            }
1516            LoraValue::Relationship(rel_id) => {
1517                let prop = lora_value_to_property(new_value)
1518                    .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1519                self.ctx
1520                    .storage
1521                    .set_relationship_property(rel_id, property.clone(), prop);
1522                Ok(())
1523            }
1524            other => Err(ExecutorError::InvalidSetTarget {
1525                found: value_kind(&other),
1526            }),
1527        }
1528    }
1529
1530    fn remove_property_from_expr(&mut self, row: &Row, expr: &ResolvedExpr) -> ExecResult<()> {
1531        let ResolvedExpr::Property {
1532            expr: owner_expr,
1533            property,
1534        } = expr
1535        else {
1536            return Err(ExecutorError::UnsupportedRemoveTarget);
1537        };
1538
1539        let owner = {
1540            let eval_ctx = EvalContext {
1541                storage: &*self.ctx.storage,
1542                params: &self.ctx.params,
1543            };
1544            eval_expr(owner_expr, row, &eval_ctx)
1545        };
1546
1547        match owner {
1548            LoraValue::Node(node_id) => {
1549                self.ctx.storage.remove_node_property(node_id, property);
1550                Ok(())
1551            }
1552            LoraValue::Relationship(rel_id) => {
1553                self.ctx
1554                    .storage
1555                    .remove_relationship_property(rel_id, property);
1556                Ok(())
1557            }
1558            other => Err(ExecutorError::InvalidRemoveTarget {
1559                found: value_kind(&other),
1560            }),
1561        }
1562    }
1563
1564    fn overwrite_entity_target(
1565        &mut self,
1566        target: EntityTarget,
1567        new_value: LoraValue,
1568    ) -> ExecResult<()> {
1569        let LoraValue::Map(map) = new_value else {
1570            return Err(ExecutorError::ExpectedPropertyMap {
1571                found: value_kind(&new_value),
1572            });
1573        };
1574
1575        let mut props: Properties = Properties::new();
1576        for (k, v) in map {
1577            let prop = lora_value_to_property(v)
1578                .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1579            props.insert(k, prop);
1580        }
1581
1582        match target {
1583            EntityTarget::Node(node_id) => {
1584                self.ctx.storage.replace_node_properties(node_id, props);
1585            }
1586            EntityTarget::Relationship(rel_id) => {
1587                self.ctx
1588                    .storage
1589                    .replace_relationship_properties(rel_id, props);
1590            }
1591        }
1592        Ok(())
1593    }
1594
1595    fn mutate_entity_target(
1596        &mut self,
1597        target: EntityTarget,
1598        patch_value: LoraValue,
1599    ) -> ExecResult<()> {
1600        let LoraValue::Map(map) = patch_value else {
1601            return Err(ExecutorError::ExpectedPropertyMap {
1602                found: value_kind(&patch_value),
1603            });
1604        };
1605
1606        match target {
1607            EntityTarget::Node(node_id) => {
1608                for (k, v) in map {
1609                    let prop = lora_value_to_property(v)
1610                        .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1611                    self.ctx.storage.set_node_property(node_id, k, prop);
1612                }
1613            }
1614            EntityTarget::Relationship(rel_id) => {
1615                for (k, v) in map {
1616                    let prop = lora_value_to_property(v)
1617                        .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1618                    self.ctx.storage.set_relationship_property(rel_id, k, prop);
1619                }
1620            }
1621        }
1622        Ok(())
1623    }
1624
1625    pub(crate) fn apply_create_pattern(
1626        &mut self,
1627        row: &mut Row,
1628        pattern: &ResolvedPattern,
1629    ) -> ExecResult<()> {
1630        for part in &pattern.parts {
1631            self.apply_create_pattern_part(row, part)?;
1632        }
1633        Ok(())
1634    }
1635
1636    /// Apply a single per-row write for any of the streamable write
1637    /// operators (Create / Set / Delete / Remove / Merge). Used by
1638    /// the [`crate::pull::StreamingWriteCursor`] auto-commit fast
1639    /// path: the cursor pulls one input row from a read upstream,
1640    /// hands it here for the side effect, and emits the row back.
1641    pub(crate) fn apply_write_op(&mut self, op: &PhysicalOp, row: &mut Row) -> ExecResult<()> {
1642        match op {
1643            PhysicalOp::Create(c) => self.apply_create_pattern(row, &c.pattern),
1644            PhysicalOp::Set(s) => {
1645                for item in &s.items {
1646                    self.apply_set_item(row, item)?;
1647                }
1648                Ok(())
1649            }
1650            PhysicalOp::Delete(d) => {
1651                let detach = d.detach;
1652                for expr in &d.expressions {
1653                    let value = {
1654                        let eval_ctx = EvalContext {
1655                            storage: &*self.ctx.storage,
1656                            params: &self.ctx.params,
1657                        };
1658                        eval_expr(expr, row, &eval_ctx)
1659                    };
1660                    self.delete_value(value, detach)?;
1661                }
1662                Ok(())
1663            }
1664            PhysicalOp::Remove(r) => {
1665                for item in &r.items {
1666                    self.apply_remove_item(row, item)?;
1667                }
1668                Ok(())
1669            }
1670            PhysicalOp::Merge(m) => {
1671                let already_bound = self.pattern_part_is_bound(row, &m.pattern_part);
1672                let matched = if already_bound {
1673                    true
1674                } else {
1675                    self.try_match_merge_pattern(row, &m.pattern_part)?
1676                };
1677                if !matched {
1678                    self.apply_create_pattern_part(row, &m.pattern_part)?;
1679                }
1680                for action in &m.actions {
1681                    if action.on_match == matched {
1682                        for item in &action.set.items {
1683                            self.apply_set_item(row, item)?;
1684                        }
1685                    }
1686                }
1687                Ok(())
1688            }
1689            other => Err(ExecutorError::RuntimeError(format!(
1690                "apply_write_op called on non-write op: {other:?}"
1691            ))),
1692        }
1693    }
1694
1695    fn apply_create_pattern_part(
1696        &mut self,
1697        row: &mut Row,
1698        part: &ResolvedPatternPart,
1699    ) -> ExecResult<()> {
1700        if part.binding.is_some() {
1701            trace!("create pattern part has path binding; path materialization not implemented");
1702        }
1703
1704        let _ = self.apply_create_pattern_element(row, &part.element)?;
1705        Ok(())
1706    }
1707
1708    fn apply_create_pattern_element(
1709        &mut self,
1710        row: &mut Row,
1711        element: &ResolvedPatternElement,
1712    ) -> ExecResult<Option<LoraValue>> {
1713        match element {
1714            ResolvedPatternElement::Node {
1715                var,
1716                labels,
1717                properties,
1718            } => {
1719                let node_id =
1720                    self.materialize_node_pattern(row, *var, labels, properties.as_ref())?;
1721                Ok(Some(LoraValue::Node(node_id)))
1722            }
1723
1724            ResolvedPatternElement::NodeChain { head, chain } => {
1725                let mut current_node_id = self.materialize_node_pattern(
1726                    row,
1727                    head.var,
1728                    &head.labels,
1729                    head.properties.as_ref(),
1730                )?;
1731
1732                for link in chain {
1733                    let next_node_id = self.materialize_node_pattern(
1734                        row,
1735                        link.node.var,
1736                        &link.node.labels,
1737                        link.node.properties.as_ref(),
1738                    )?;
1739
1740                    let _ = self.materialize_relationship_pattern(
1741                        row,
1742                        current_node_id,
1743                        next_node_id,
1744                        &link.rel,
1745                    )?;
1746
1747                    current_node_id = next_node_id;
1748                }
1749
1750                Ok(Some(LoraValue::Node(current_node_id)))
1751            }
1752
1753            ResolvedPatternElement::ShortestPath { .. } => {
1754                // ShortestPath is not valid in CREATE context
1755                Ok(None)
1756            }
1757        }
1758    }
1759
1760    fn pattern_part_is_bound(&self, row: &Row, part: &ResolvedPatternPart) -> bool {
1761        match &part.element {
1762            ResolvedPatternElement::Node { var, .. } => var.and_then(|v| row.get(v)).is_some(),
1763
1764            ResolvedPatternElement::ShortestPath { .. } => false,
1765
1766            ResolvedPatternElement::NodeChain { head, chain } => {
1767                let head_ok = head.var.and_then(|v| row.get(v)).is_some();
1768
1769                let chain_ok = chain.iter().all(|link| {
1770                    let node_ok = link.node.var.and_then(|v| row.get(v)).is_some();
1771                    // For MERGE, anonymous relationships cannot be considered
1772                    // "bound" because we have no variable to check. The merge
1773                    // must search the graph to see if the relationship exists.
1774                    let rel_ok = match link.rel.var {
1775                        Some(v) => row.get(v).is_some(),
1776                        None => false,
1777                    };
1778                    node_ok && rel_ok
1779                });
1780
1781                head_ok && chain_ok
1782            }
1783        }
1784    }
1785
1786    fn materialize_node_pattern(
1787        &mut self,
1788        row: &mut Row,
1789        var: Option<VarId>,
1790        labels: &[Vec<String>],
1791        properties: Option<&ResolvedExpr>,
1792    ) -> ExecResult<u64> {
1793        if let Some(var_id) = var {
1794            if let Some(LoraValue::Node(id)) = row.get(var_id) {
1795                return Ok(*id);
1796            }
1797        }
1798
1799        let properties = match properties {
1800            Some(expr) => eval_properties_expr(expr, row, &*self.ctx.storage, &self.ctx.params)?,
1801            None => Properties::new(),
1802        };
1803
1804        let flat_labels = flatten_label_groups(labels);
1805        debug!("creating node with labels={flat_labels:?}");
1806        let created = self.ctx.storage.create_node(flat_labels, properties);
1807
1808        if let Some(var_id) = var {
1809            row.insert(var_id, LoraValue::Node(created.id));
1810        }
1811
1812        Ok(created.id)
1813    }
1814
1815    fn materialize_relationship_pattern(
1816        &mut self,
1817        row: &mut Row,
1818        left_node_id: u64,
1819        right_node_id: u64,
1820        rel: &lora_analyzer::ResolvedRel,
1821    ) -> ExecResult<u64> {
1822        if let Some(var_id) = rel.var {
1823            if let Some(LoraValue::Relationship(id)) = row.get(var_id) {
1824                let id = *id;
1825                if let Some((src, dst)) = self.ctx.storage.relationship_endpoints(id) {
1826                    let endpoints_match = match rel.direction {
1827                        Direction::Right | Direction::Undirected => {
1828                            src == left_node_id && dst == right_node_id
1829                        }
1830                        Direction::Left => src == right_node_id && dst == left_node_id,
1831                    };
1832
1833                    if endpoints_match {
1834                        return Ok(id);
1835                    }
1836                }
1837            }
1838        }
1839
1840        if rel.range.is_some() {
1841            return Err(ExecutorError::UnsupportedCreateRelationshipRange);
1842        }
1843
1844        let (src, dst) = match rel.direction {
1845            Direction::Right | Direction::Undirected => (left_node_id, right_node_id),
1846            Direction::Left => (right_node_id, left_node_id),
1847        };
1848
1849        let rel_type = rel
1850            .types
1851            .first()
1852            .ok_or(ExecutorError::MissingRelationshipType)?;
1853
1854        if rel_type.is_empty() {
1855            return Err(ExecutorError::MissingRelationshipType);
1856        }
1857
1858        let properties = match rel.properties.as_ref() {
1859            Some(expr) => eval_properties_expr(expr, row, &*self.ctx.storage, &self.ctx.params)?,
1860            None => Properties::new(),
1861        };
1862
1863        debug!("creating relationship: src={src}, dst={dst}, type={rel_type}");
1864
1865        let created = self
1866            .ctx
1867            .storage
1868            .create_relationship(src, dst, rel_type, properties)
1869            .ok_or_else(|| ExecutorError::RelationshipCreateFailed {
1870                src,
1871                dst,
1872                rel_type: rel_type.clone(),
1873            })?;
1874
1875        if let Some(var_id) = rel.var {
1876            row.insert(var_id, LoraValue::Relationship(created.id));
1877        }
1878
1879        Ok(created.id)
1880    }
1881}