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