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