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, LoraPath, 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_row_bound, 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    /// The row a writing `CALL { ... }` body's bottom `Argument` yields:
77    /// the outer row it runs for. `None` outside such a body.
78    argument_seed: Option<Row>,
79    /// When set, existence constraints on created entities are checked
80    /// once the statement finishes rather than at `CREATE`, so a later
81    /// `SET` (or `ON CREATE SET`) in the same statement can supply the
82    /// property. See [`plan_defers_existence`].
83    defer_existence: bool,
84    /// Entities created while `defer_existence` is on, still to check.
85    pending_existence: Vec<EntityTarget>,
86}
87
88impl<'a, S: GraphStorageMut> MutableExecutor<'a, S> {
89    pub fn new(ctx: MutableExecutionContext<'a, S>) -> Self {
90        Self {
91            ctx,
92            deadline: None,
93            argument_seed: None,
94            defer_existence: false,
95            pending_existence: Vec::new(),
96        }
97    }
98
99    pub fn with_deadline(ctx: MutableExecutionContext<'a, S>, deadline: Option<Instant>) -> Self {
100        Self {
101            ctx,
102            deadline,
103            argument_seed: None,
104            defer_existence: false,
105            pending_existence: Vec::new(),
106        }
107    }
108
109    #[inline]
110    fn check_deadline(&self) -> ExecResult<()> {
111        if let Some(deadline) = self.deadline {
112            check_deadline_at(deadline)
113        } else {
114            Ok(())
115        }
116    }
117
118    pub fn execute(
119        &mut self,
120        plan: &PhysicalPlan,
121        options: Option<ExecuteOptions>,
122    ) -> ExecResult<QueryResult> {
123        let _deadline_scope = crate::cancel::DeadlineScope::enter(self.deadline);
124        let rows = self.execute_rows(plan)?;
125        Ok(project_rows(rows, options.unwrap_or_default()))
126    }
127
128    pub fn execute_rows(&mut self, plan: &PhysicalPlan) -> ExecResult<Vec<Row>> {
129        self.defer_existence = plan_defers_existence(plan);
130        let rows = self.execute_plan_rows(plan)?;
131        self.check_pending_existence()?;
132        Ok(rows)
133    }
134
135    /// Defer existence checks on created entities to the end of the
136    /// statement (see [`plan_defers_existence`]); the caller then runs
137    /// [`Self::check_pending_existence`].
138    pub(crate) fn defer_existence_checks(&mut self, defer: bool) {
139        self.defer_existence = defer;
140    }
141
142    /// Check the existence constraints deferred so far, clearing them.
143    pub(crate) fn check_pending_existence(&mut self) -> ExecResult<()> {
144        for target in std::mem::take(&mut self.pending_existence) {
145            let checked = match target {
146                EntityTarget::Node(id) => self.ctx.storage.check_node_existence_constraints(id),
147                EntityTarget::Relationship(id) => self
148                    .ctx
149                    .storage
150                    .check_relationship_existence_constraints(id),
151            };
152            checked.map_err(ExecutorError::ConstraintViolation)?;
153        }
154        Ok(())
155    }
156
157    fn execute_plan_rows(&mut self, plan: &PhysicalPlan) -> ExecResult<Vec<Row>> {
158        let _deadline_scope = crate::cancel::DeadlineScope::enter(self.deadline);
159        self.check_deadline()?;
160        // Clear any error residue that a previous query on this thread may have
161        // left in the thread-local eval-error slot.
162        clear_eval_error();
163
164        let rows = self.execute_node(plan, plan.root)?;
165        if plan_ends_in_write(plan) {
166            return Ok(Vec::new());
167        }
168        if !plan_may_need_hydration(plan) {
169            return Ok(rows);
170        }
171        Ok(rows
172            .into_iter()
173            .map(|row| self.hydrate_row(row))
174            .collect::<Vec<_>>())
175    }
176
177    /// Execute a compiled query that may include UNION branches.
178    pub fn execute_compiled(
179        &mut self,
180        compiled: &CompiledQuery,
181        options: Option<ExecuteOptions>,
182    ) -> ExecResult<QueryResult> {
183        let _deadline_scope = crate::cancel::DeadlineScope::enter(self.deadline);
184        let rows = self.execute_compiled_rows(compiled)?;
185        Ok(project_rows(rows, options.unwrap_or_default()))
186    }
187
188    pub fn execute_compiled_rows(&mut self, compiled: &CompiledQuery) -> ExecResult<Vec<Row>> {
189        let _deadline_scope = crate::cancel::DeadlineScope::enter(self.deadline);
190        self.check_deadline()?;
191        self.defer_existence = plan_defers_existence(&compiled.physical)
192            || !compiled.unions.is_empty()
193                && compiled
194                    .unions
195                    .iter()
196                    .any(|b| plan_defers_existence(&b.physical));
197        if compiled.unions.is_empty() {
198            let rows = self.execute_plan_rows(&compiled.physical)?;
199            self.check_pending_existence()?;
200            return Ok(rows);
201        }
202
203        clear_eval_error();
204
205        // Execute the head branch.
206        let mut all_rows = self.execute_and_hydrate(&compiled.physical)?;
207
208        // Execute each UNION branch and combine.
209        // Track whether any branch uses plain UNION (dedup needed).
210        let mut needs_dedup = false;
211
212        for branch in &compiled.unions {
213            self.check_deadline()?;
214            let branch_rows = self.execute_and_hydrate(&branch.physical)?;
215            all_rows.extend(branch_rows);
216
217            if !branch.all {
218                needs_dedup = true;
219            }
220        }
221
222        if needs_dedup {
223            all_rows = dedup_rows(all_rows);
224        }
225
226        self.check_pending_existence()?;
227        Ok(all_rows)
228    }
229
230    fn execute_and_hydrate(&mut self, plan: &PhysicalPlan) -> ExecResult<Vec<Row>> {
231        self.check_deadline()?;
232        let rows = self.execute_node(plan, plan.root)?;
233        if plan_ends_in_write(plan) {
234            return Ok(Vec::new());
235        }
236        if !plan_may_need_hydration(plan) {
237            return Ok(rows);
238        }
239        Ok(rows.into_iter().map(|row| self.hydrate_row(row)).collect())
240    }
241
242    pub(crate) fn hydrate_row(&self, row: Row) -> Row {
243        let mut out = Row::new();
244
245        for (var, name, value) in row.into_iter_named() {
246            out.insert_named_inline(var, name, self.hydrate_value(value));
247        }
248
249        out
250    }
251
252    fn execute_node(
253        &mut self,
254        plan: &PhysicalPlan,
255        node_id: PhysicalNodeId,
256    ) -> ExecResult<Vec<Row>> {
257        self.check_deadline()?;
258        trace!("mutable execute_node start: node_id={node_id:?}");
259
260        let result = match &plan.nodes[node_id] {
261            PhysicalOp::Argument(op) => self.exec_argument(op),
262            PhysicalOp::NodeScan(op) => self.exec_node_scan(plan, op),
263            PhysicalOp::NodeByLabelScan(op) => self.exec_node_by_label_scan(plan, op),
264            PhysicalOp::NodeByIdSeek(op) => self.exec_node_by_id_seek(plan, op),
265            PhysicalOp::RelByIdSeek(op) => self.exec_rel_by_id_seek(plan, op),
266            PhysicalOp::NodeByPropertyScan(op) => self.exec_node_by_property_scan(plan, op),
267            PhysicalOp::NodeByPropertyRangeScan(op) => {
268                self.exec_node_by_property_range_scan(plan, op)
269            }
270            PhysicalOp::NodeByTextScan(op) => self.exec_node_by_text_scan(plan, op),
271            PhysicalOp::NodeByPointScan(op) => self.exec_node_by_point_scan(plan, op),
272            PhysicalOp::RelByPropertyRangeScan(op) => {
273                self.exec_rel_by_property_range_scan(plan, op)
274            }
275            PhysicalOp::RelByTextScan(op) => self.exec_rel_by_text_scan(plan, op),
276            PhysicalOp::RelByPointScan(op) => self.exec_rel_by_point_scan(plan, op),
277            PhysicalOp::Expand(op) => self.exec_expand(plan, op),
278            PhysicalOp::Filter(op) => self.exec_filter(plan, op),
279            PhysicalOp::Projection(op) => self.exec_projection(plan, op),
280            PhysicalOp::Unwind(op) => self.exec_unwind(plan, op),
281            PhysicalOp::HashAggregation(op) => self.exec_hash_aggregation(plan, op),
282            PhysicalOp::Sort(op) => self.exec_sort(plan, op),
283            PhysicalOp::Limit(op) => self.exec_limit(plan, op),
284            PhysicalOp::Create(op) => self.exec_create(plan, op),
285            PhysicalOp::Merge(op) => self.exec_merge(plan, op),
286            PhysicalOp::Delete(op) => self.exec_delete(plan, op),
287            PhysicalOp::Set(op) => self.exec_set(plan, op),
288            PhysicalOp::Remove(op) => self.exec_remove(plan, op),
289            PhysicalOp::Foreach(op) => self.exec_foreach(plan, op),
290            PhysicalOp::OptionalMatch(op) => self.exec_optional_match(plan, op),
291            PhysicalOp::CallSubquery(op) => self.exec_call_subquery(plan, op),
292            PhysicalOp::PathBuild(op) => self.exec_path_build(plan, op),
293        };
294
295        match &result {
296            Ok(rows) => trace!(
297                "mutable execute_node ok: node_id={node_id:?}, rows={}",
298                rows.len()
299            ),
300            Err(err) => error!("mutable execute_node failed: node_id={node_id:?}, error={err}"),
301        }
302
303        result
304    }
305
306    fn exec_argument(&self, _op: &ArgumentExec) -> ExecResult<Vec<Row>> {
307        Ok(vec![self.argument_seed.clone().unwrap_or_default()])
308    }
309
310    fn exec_node_scan(&mut self, plan: &PhysicalPlan, op: &NodeScanExec) -> ExecResult<Vec<Row>> {
311        let base_rows = match op.input {
312            Some(input) => self.execute_node(plan, input)?,
313            None => vec![Row::new()],
314        };
315
316        node_scan_rows(&*self.ctx.storage, base_rows, op, self.deadline)
317    }
318
319    fn exec_node_by_label_scan(
320        &mut self,
321        plan: &PhysicalPlan,
322        op: &NodeByLabelScanExec,
323    ) -> ExecResult<Vec<Row>> {
324        let base_rows = match op.input {
325            Some(input) => self.execute_node(plan, input)?,
326            None => vec![Row::new()],
327        };
328
329        node_by_label_scan_rows(&*self.ctx.storage, base_rows, op, self.deadline)
330    }
331
332    fn exec_node_by_property_scan(
333        &mut self,
334        plan: &PhysicalPlan,
335        op: &NodeByPropertyScanExec,
336    ) -> ExecResult<Vec<Row>> {
337        let base_rows = match op.input {
338            Some(input) => self.execute_node(plan, input)?,
339            None => vec![Row::new()],
340        };
341
342        node_by_property_scan_rows(
343            &*self.ctx.storage,
344            &self.ctx.params,
345            base_rows,
346            op,
347            self.deadline,
348        )
349    }
350
351    fn exec_node_by_property_range_scan(
352        &mut self,
353        plan: &PhysicalPlan,
354        op: &lora_compiler::NodeByPropertyRangeScanExec,
355    ) -> ExecResult<Vec<Row>> {
356        let base_rows = match op.input {
357            Some(input) => self.execute_node(plan, input)?,
358            None => vec![Row::new()],
359        };
360        super::helpers::node_by_property_range_scan_rows(
361            &*self.ctx.storage,
362            &self.ctx.params,
363            base_rows,
364            op,
365            self.deadline,
366        )
367    }
368
369    fn exec_node_by_text_scan(
370        &mut self,
371        plan: &PhysicalPlan,
372        op: &lora_compiler::NodeByTextScanExec,
373    ) -> ExecResult<Vec<Row>> {
374        let base_rows = match op.input {
375            Some(input) => self.execute_node(plan, input)?,
376            None => vec![Row::new()],
377        };
378        super::helpers::node_by_text_scan_rows(
379            &*self.ctx.storage,
380            &self.ctx.params,
381            base_rows,
382            op,
383            self.deadline,
384        )
385    }
386
387    fn exec_node_by_point_scan(
388        &mut self,
389        plan: &PhysicalPlan,
390        op: &lora_compiler::NodeByPointScanExec,
391    ) -> ExecResult<Vec<Row>> {
392        let base_rows = match op.input {
393            Some(input) => self.execute_node(plan, input)?,
394            None => vec![Row::new()],
395        };
396        super::helpers::node_by_point_scan_rows(
397            &*self.ctx.storage,
398            &self.ctx.params,
399            base_rows,
400            op,
401            self.deadline,
402        )
403    }
404
405    fn exec_rel_by_property_range_scan(
406        &mut self,
407        plan: &PhysicalPlan,
408        op: &lora_compiler::RelByPropertyRangeScanExec,
409    ) -> ExecResult<Vec<Row>> {
410        let base_rows = match op.input {
411            Some(input) => self.execute_node(plan, input)?,
412            None => vec![Row::new()],
413        };
414        super::helpers::rel_by_property_range_scan_rows(
415            &*self.ctx.storage,
416            &self.ctx.params,
417            base_rows,
418            op,
419            self.deadline,
420        )
421    }
422
423    fn exec_rel_by_text_scan(
424        &mut self,
425        plan: &PhysicalPlan,
426        op: &lora_compiler::RelByTextScanExec,
427    ) -> ExecResult<Vec<Row>> {
428        let base_rows = match op.input {
429            Some(input) => self.execute_node(plan, input)?,
430            None => vec![Row::new()],
431        };
432        super::helpers::rel_by_text_scan_rows(
433            &*self.ctx.storage,
434            &self.ctx.params,
435            base_rows,
436            op,
437            self.deadline,
438        )
439    }
440
441    fn exec_node_by_id_seek(
442        &mut self,
443        plan: &PhysicalPlan,
444        op: &lora_compiler::NodeByIdSeekExec,
445    ) -> ExecResult<Vec<Row>> {
446        let base_rows = match op.input {
447            Some(input) => self.execute_node(plan, input)?,
448            None => vec![Row::new()],
449        };
450        super::helpers::node_by_id_seek_rows(
451            &*self.ctx.storage,
452            &self.ctx.params,
453            base_rows,
454            op,
455            self.deadline,
456        )
457    }
458
459    fn exec_rel_by_id_seek(
460        &mut self,
461        plan: &PhysicalPlan,
462        op: &lora_compiler::RelByIdSeekExec,
463    ) -> ExecResult<Vec<Row>> {
464        let base_rows = match op.input {
465            Some(input) => self.execute_node(plan, input)?,
466            None => vec![Row::new()],
467        };
468        super::helpers::rel_by_id_seek_rows(
469            &*self.ctx.storage,
470            &self.ctx.params,
471            base_rows,
472            op,
473            self.deadline,
474        )
475    }
476
477    fn exec_rel_by_point_scan(
478        &mut self,
479        plan: &PhysicalPlan,
480        op: &lora_compiler::RelByPointScanExec,
481    ) -> ExecResult<Vec<Row>> {
482        let base_rows = match op.input {
483            Some(input) => self.execute_node(plan, input)?,
484            None => vec![Row::new()],
485        };
486        super::helpers::rel_by_point_scan_rows(
487            &*self.ctx.storage,
488            &self.ctx.params,
489            base_rows,
490            op,
491            self.deadline,
492        )
493    }
494
495    fn exec_expand(&mut self, plan: &PhysicalPlan, op: &ExpandExec) -> ExecResult<Vec<Row>> {
496        let input_rows = self.execute_node(plan, op.input)?;
497        if let Some(range) = &op.range {
498            expand_var_len_rows(&*self.ctx.storage, input_rows, op, range)
499        } else {
500            expand_rows(&*self.ctx.storage, &self.ctx.params, input_rows, op)
501        }
502    }
503
504    fn exec_filter(&mut self, plan: &PhysicalPlan, op: &FilterExec) -> ExecResult<Vec<Row>> {
505        let input_rows = self.execute_node(plan, op.input)?;
506        let eval_ctx = EvalContext {
507            storage: &*self.ctx.storage,
508            params: &self.ctx.params,
509        };
510
511        filter_rows_checked(input_rows, &op.predicate, &eval_ctx)
512    }
513
514    fn exec_projection(
515        &mut self,
516        plan: &PhysicalPlan,
517        op: &ProjectionExec,
518    ) -> ExecResult<Vec<Row>> {
519        let input_rows = self.execute_node(plan, op.input)?;
520        let eval_ctx = EvalContext {
521            storage: &*self.ctx.storage,
522            params: &self.ctx.params,
523        };
524
525        project_rows_checked(input_rows, op, &eval_ctx)
526    }
527
528    fn hydrate_value(&self, value: LoraValue) -> LoraValue {
529        match value {
530            LoraValue::Node(id) => self.hydrate_node(id),
531            LoraValue::Relationship(id) => self.hydrate_relationship(id),
532            LoraValue::List(values) => {
533                LoraValue::List(values.into_iter().map(|v| self.hydrate_value(v)).collect())
534            }
535            LoraValue::Map(map) => LoraValue::Map(
536                map.into_iter()
537                    .map(|(k, v)| (k, self.hydrate_value(v)))
538                    .collect(),
539            ),
540            other => other,
541        }
542    }
543
544    fn hydrate_node(&self, id: u64) -> LoraValue {
545        self.ctx
546            .storage
547            .with_node(id, hydrate_node_record)
548            .unwrap_or(LoraValue::Null)
549    }
550
551    fn hydrate_relationship(&self, id: u64) -> LoraValue {
552        self.ctx
553            .storage
554            .with_relationship(id, hydrate_relationship_record)
555            .unwrap_or(LoraValue::Null)
556    }
557
558    fn exec_unwind(&mut self, plan: &PhysicalPlan, op: &UnwindExec) -> ExecResult<Vec<Row>> {
559        let input_rows = self.execute_node(plan, op.input)?;
560        let eval_ctx = EvalContext {
561            storage: &*self.ctx.storage,
562            params: &self.ctx.params,
563        };
564
565        unwind_rows(input_rows, op, &eval_ctx)
566    }
567
568    fn exec_hash_aggregation(
569        &mut self,
570        plan: &PhysicalPlan,
571        op: &HashAggregationExec,
572    ) -> ExecResult<Vec<Row>> {
573        if let Some(rows) =
574            super::helpers::count_all_scan_aggregation_rows(&*self.ctx.storage, plan, op)
575        {
576            return Ok(rows);
577        }
578
579        let input_rows = self.execute_node(plan, op.input)?;
580        let eval_ctx = EvalContext {
581            storage: &*self.ctx.storage,
582            params: &self.ctx.params,
583        };
584
585        aggregate_rows(input_rows, &op.group_by, &op.aggregates, &eval_ctx)
586    }
587
588    fn exec_sort(&mut self, plan: &PhysicalPlan, op: &SortExec) -> ExecResult<Vec<Row>> {
589        let mut rows = self.execute_node(plan, op.input)?;
590        let eval_ctx = EvalContext {
591            storage: &*self.ctx.storage,
592            params: &self.ctx.params,
593        };
594
595        let bound = sort_row_bound(op.top_k, op.limit.as_ref(), &eval_ctx);
596        sort_rows_with_top_k(&mut rows, &op.items, &eval_ctx, bound);
597
598        Ok(rows)
599    }
600
601    fn exec_limit(&mut self, plan: &PhysicalPlan, op: &LimitExec) -> ExecResult<Vec<Row>> {
602        let rows = self.execute_node(plan, op.input)?;
603        let eval_ctx = EvalContext {
604            storage: &*self.ctx.storage,
605            params: &self.ctx.params,
606        };
607
608        limit_rows(rows, op, &eval_ctx)
609    }
610
611    fn exec_optional_match(
612        &mut self,
613        plan: &PhysicalPlan,
614        op: &OptionalMatchExec,
615    ) -> ExecResult<Vec<Row>> {
616        let input_rows = self.execute_node(plan, op.input)?;
617
618        if super::optional::optional_can_correlate(plan, op.inner) {
619            let storage_ref: &S = &*self.ctx.storage;
620            return super::optional::correlated_optional_match_rows(
621                storage_ref,
622                &self.ctx.params,
623                plan,
624                op.inner,
625                input_rows,
626                &op.new_vars,
627            );
628        }
629
630        // Fallback: execute the inner plan once, uncorrelated, and join.
631        let inner_rows = self.execute_node(plan, op.inner)?;
632
633        Ok(optional_match_rows(input_rows, &inner_rows, &op.new_vars))
634    }
635
636    fn exec_call_subquery(
637        &mut self,
638        plan: &PhysicalPlan,
639        op: &CallSubqueryExec,
640    ) -> ExecResult<Vec<Row>> {
641        let input_rows = self.execute_node(plan, op.input)?;
642        let mut out = Vec::with_capacity(input_rows.len());
643
644        if crate::pull::subtree_has_write(plan, op.inner) {
645            // A writing body runs on this executor, once per outer row,
646            // with the outer row seeded into its bottom `Argument`. Each
647            // run sees the writes of the runs before it.
648            let unit = op.new_vars.is_empty();
649            for outer_row in input_rows {
650                self.check_deadline()?;
651                let prev = self.argument_seed.replace(outer_row.clone());
652                let inner_rows = self.execute_node(plan, op.inner);
653                self.argument_seed = prev;
654                let inner_rows = inner_rows?;
655                if unit {
656                    // A unit subquery keeps the outer row as it is, once,
657                    // however many rows its body produced.
658                    out.push(outer_row);
659                    continue;
660                }
661                for inner_row in inner_rows {
662                    out.push(crate::executor::merge_optional_rows(&outer_row, &inner_row));
663                }
664            }
665            return Ok(out);
666        }
667
668        let params = std::sync::Arc::new(self.ctx.params.clone());
669        let storage_ref: &S = &*self.ctx.storage;
670        for outer_row in input_rows {
671            let mut inner_source = crate::pull::build_streaming_seeded(
672                plan,
673                op.inner,
674                storage_ref,
675                params.clone(),
676                outer_row.clone(),
677            )?;
678            let inner_rows = crate::pull::drain(inner_source.as_mut())?;
679            for inner_row in inner_rows {
680                out.push(crate::executor::merge_optional_rows(&outer_row, &inner_row));
681            }
682        }
683        Ok(out)
684    }
685
686    fn exec_path_build(&mut self, plan: &PhysicalPlan, op: &PathBuildExec) -> ExecResult<Vec<Row>> {
687        let input_rows = self.execute_node(plan, op.input)?;
688        let mut rows: Vec<Row> = input_rows
689            .into_iter()
690            .map(|mut row| {
691                let path = build_path_value(&row, &op.node_vars, &op.rel_vars, &*self.ctx.storage);
692                row.insert(op.output, path);
693                row
694            })
695            .collect();
696
697        if let Some(all) = op.shortest_path_all {
698            rows = filter_shortest_paths(rows, op.output, all);
699        }
700        Ok(rows)
701    }
702
703    fn exec_create(&mut self, plan: &PhysicalPlan, op: &CreateExec) -> ExecResult<Vec<Row>> {
704        // Fast path: if the input subtree is fully streamable (no
705        // nested writes, no blocking operators), pull rows one at a
706        // time and apply the create pattern per row, instead of
707        // materializing the whole input. The output Vec still
708        // accumulates — auto-commit-side output streaming is M1.b.
709        if crate::pull::subtree_is_fully_streaming(plan, op.input) {
710            return self.exec_create_streaming_input(plan, op);
711        }
712
713        let input_rows = self.execute_node(plan, op.input)?;
714        let mut out = Vec::with_capacity(input_rows.len());
715
716        for mut row in input_rows {
717            self.apply_create_pattern(&mut row, &op.pattern)?;
718            out.push(row);
719        }
720
721        Ok(out)
722    }
723
724    /// Generic streaming-input loop for write operators whose input
725    /// subtree is fully streamable. Opens a pull-based read cursor
726    /// over the input subtree, calls `apply` per row, and accumulates
727    /// the resulting rows.
728    ///
729    /// # Safety
730    ///
731    /// The upstream [`crate::pull::RowSource`] needs `&S` while it
732    /// lives; the per-row `apply` callback needs `&mut S` (via
733    /// `&mut self`). The existing read-side `RowSource` impls
734    /// materialize their iteration state into owned `Vec`s at
735    /// construction time (see `NodeScanSource::cur_ids`,
736    /// `ExpandSource::cur_edges`, etc. in `pull.rs`), so no live
737    /// `&S` borrow into storage persists across `next_row` calls.
738    /// We exploit that by deriving the read borrow from a raw
739    /// pointer — Rust then doesn't see the shared/mutable conflict
740    /// at compile time, and the dynamic access pattern is
741    /// non-aliasing: read-only inside `next_row`, then mutable
742    /// inside `apply`, never both at the same instant.
743    fn streaming_apply<F>(
744        &mut self,
745        plan: &PhysicalPlan,
746        input: PhysicalNodeId,
747        mut apply: F,
748    ) -> ExecResult<Vec<Row>>
749    where
750        F: FnMut(&mut Self, &mut Row) -> ExecResult<()>,
751    {
752        use std::sync::Arc;
753
754        let storage_ptr: *mut S = self.ctx.storage as *mut S;
755        let params = Arc::new(self.ctx.params.clone());
756
757        // SAFETY: see method-level comment.
758        let storage_ref: &S = unsafe { &*storage_ptr };
759        // Inside a writing `CALL { ... }` body the input's bottom
760        // `Argument` yields the outer row.
761        let mut upstream = match self.argument_seed.clone() {
762            Some(seed) => {
763                crate::pull::build_streaming_seeded(plan, input, storage_ref, params, seed)?
764            }
765            None => crate::pull::build_streaming(plan, input, storage_ref, params)?,
766        };
767
768        let mut out = Vec::new();
769        while let Some(mut row) = upstream.next_row()? {
770            apply(self, &mut row)?;
771            out.push(row);
772        }
773
774        Ok(out)
775    }
776
777    /// Streaming-input variant of [`Self::exec_create`]. Delegates
778    /// to [`Self::streaming_apply`].
779    fn exec_create_streaming_input(
780        &mut self,
781        plan: &PhysicalPlan,
782        op: &CreateExec,
783    ) -> ExecResult<Vec<Row>> {
784        self.streaming_apply(plan, op.input, |this, row| {
785            this.apply_create_pattern(row, &op.pattern)
786        })
787    }
788
789    fn apply_remove_item(&mut self, row: &Row, item: &ResolvedRemoveItem) -> ExecResult<()> {
790        match item {
791            ResolvedRemoveItem::Labels { variable, labels } => match row.get(*variable) {
792                Some(LoraValue::Node(node_id)) => {
793                    let node_id = *node_id;
794                    for label in labels {
795                        self.ctx.storage.remove_node_label(node_id, label);
796                    }
797                    Ok(())
798                }
799                Some(other) => Err(ExecutorError::ExpectedNodeForRemoveLabels {
800                    found: value_kind(other),
801                }),
802                None => Err(ExecutorError::UnboundVariableForRemove {
803                    var: format!("{variable:?}"),
804                }),
805            },
806
807            ResolvedRemoveItem::Property { expr } => self.remove_property_from_expr(row, expr),
808        }
809    }
810
811    fn delete_value(&mut self, value: LoraValue, detach: bool) -> ExecResult<()> {
812        match value {
813            LoraValue::Null => Ok(()),
814
815            LoraValue::Node(node_id) => {
816                if detach {
817                    self.ctx.storage.detach_delete_node(node_id);
818                    Ok(())
819                } else {
820                    let ok = self.ctx.storage.delete_node(node_id);
821                    if ok {
822                        Ok(())
823                    } else {
824                        Err(ExecutorError::DeleteNodeWithRelationships { node_id })
825                    }
826                }
827            }
828
829            LoraValue::Relationship(rel_id) => {
830                let ok = self.ctx.storage.delete_relationship(rel_id);
831                if ok {
832                    Ok(())
833                } else {
834                    Err(ExecutorError::DeleteRelationshipFailed { rel_id })
835                }
836            }
837
838            LoraValue::List(values) => {
839                for v in values {
840                    self.delete_value(v, detach)?;
841                }
842                Ok(())
843            }
844
845            other => Err(ExecutorError::InvalidDeleteTarget {
846                found: value_kind(&other),
847            }),
848        }
849    }
850
851    fn collect_delete_targets(
852        &self,
853        value: &LoraValue,
854        targets: &mut BTreeSet<DeleteTarget>,
855    ) -> ExecResult<()> {
856        match value {
857            LoraValue::Null => Ok(()),
858
859            LoraValue::Node(node_id) => {
860                targets.insert(DeleteTarget::Node(*node_id));
861                Ok(())
862            }
863
864            LoraValue::Relationship(rel_id) => {
865                targets.insert(DeleteTarget::Relationship(*rel_id));
866                Ok(())
867            }
868
869            LoraValue::List(values) => {
870                for v in values {
871                    self.collect_delete_targets(v, targets)?;
872                }
873                Ok(())
874            }
875
876            other => Err(ExecutorError::InvalidDeleteTarget {
877                found: value_kind(other),
878            }),
879        }
880    }
881
882    fn validate_delete_targets(
883        &self,
884        targets: &BTreeSet<DeleteTarget>,
885        detach: bool,
886    ) -> ExecResult<()> {
887        for target in targets {
888            match target {
889                DeleteTarget::Relationship(rel_id) => {
890                    if !self.ctx.storage.contains_relationship(*rel_id) {
891                        return Err(ExecutorError::DeleteRelationshipFailed { rel_id: *rel_id });
892                    }
893                }
894                DeleteTarget::Node(node_id) if !detach => {
895                    if !self.ctx.storage.contains_node(*node_id) {
896                        return Err(ExecutorError::DeleteNodeWithRelationships {
897                            node_id: *node_id,
898                        });
899                    }
900                    let has_external_relationship = self
901                        .ctx
902                        .storage
903                        .relationship_ids_of(*node_id, Direction::Undirected)
904                        .into_iter()
905                        .any(|rel_id| !targets.contains(&DeleteTarget::Relationship(rel_id)));
906                    if has_external_relationship {
907                        return Err(ExecutorError::DeleteNodeWithRelationships {
908                            node_id: *node_id,
909                        });
910                    }
911                }
912                DeleteTarget::Node(_) => {}
913            }
914        }
915        Ok(())
916    }
917
918    fn delete_target(&mut self, target: DeleteTarget, detach: bool) -> ExecResult<()> {
919        match target {
920            DeleteTarget::Node(node_id) => {
921                if detach {
922                    self.ctx.storage.detach_delete_node(node_id);
923                    Ok(())
924                } else {
925                    let ok = self.ctx.storage.delete_node(node_id);
926                    if ok {
927                        Ok(())
928                    } else {
929                        Err(ExecutorError::DeleteNodeWithRelationships { node_id })
930                    }
931                }
932            }
933            DeleteTarget::Relationship(rel_id) => {
934                let ok = self.ctx.storage.delete_relationship(rel_id);
935                if ok {
936                    Ok(())
937                } else {
938                    Err(ExecutorError::DeleteRelationshipFailed { rel_id })
939                }
940            }
941        }
942    }
943
944    fn exec_merge(&mut self, plan: &PhysicalPlan, op: &MergeExec) -> ExecResult<Vec<Row>> {
945        // Streaming-input fast path when the input subtree is fully
946        // streamable. Per-row work (probe → optionally create →
947        // ON MATCH / ON CREATE actions) is identical to the
948        // materialized branch below.
949        if crate::pull::subtree_is_fully_streaming(plan, op.input) {
950            return self.streaming_apply(plan, op.input, |this, row| {
951                let matched = this.match_merge_pattern(row, &op.pattern_part)?;
952                if !matched {
953                    this.apply_create_pattern_part(row, &op.pattern_part)?;
954                }
955                for action in &op.actions {
956                    if action.on_match == matched {
957                        for item in &action.set.items {
958                            this.apply_set_item(row, item)?;
959                        }
960                    }
961                }
962                Ok(())
963            });
964        }
965
966        let input_rows = self.execute_node(plan, op.input)?;
967        let mut out = Vec::with_capacity(input_rows.len());
968
969        for mut row in input_rows {
970            let matched = self.match_merge_pattern(&mut row, &op.pattern_part)?;
971
972            if !matched {
973                self.apply_create_pattern_part(&mut row, &op.pattern_part)?;
974            }
975
976            for action in &op.actions {
977                if action.on_match == matched {
978                    for item in &action.set.items {
979                        self.apply_set_item(&row, item)?;
980                    }
981                }
982            }
983
984            out.push(row);
985        }
986
987        Ok(out)
988    }
989
990    /// Whether the MERGE pattern already exists: every variable bound by an
991    /// earlier clause, or a match found in the graph (then bound in the
992    /// row). On a match its path variable (`MERGE p = …`) is bound too; on a
993    /// miss the create binds it.
994    fn match_merge_pattern(&self, row: &mut Row, part: &ResolvedPatternPart) -> ExecResult<bool> {
995        if self.pattern_part_is_bound(row, part)? {
996            if let (Some(var), Some(path)) = (part.binding, bound_pattern_path(row, part)) {
997                row.insert(var, LoraValue::Path(path));
998            }
999            return Ok(true);
1000        }
1001        self.try_match_merge_pattern(row, part)
1002    }
1003
1004    /// Try to find an existing node/pattern in the graph matching the MERGE
1005    /// pattern. If found, bind its variables (and path variable) in the row
1006    /// and return true. On a miss the row is left untouched, so the create
1007    /// path sees only the variables that were bound before the MERGE.
1008    fn try_match_merge_pattern(
1009        &self,
1010        row: &mut Row,
1011        part: &ResolvedPatternPart,
1012    ) -> ExecResult<bool> {
1013        match &part.element {
1014            ResolvedPatternElement::Node {
1015                var,
1016                labels,
1017                properties,
1018            } => {
1019                let expected_props = self.merge_expected_props(properties.as_ref(), row);
1020                let Some(id) = self
1021                    .merge_node_candidates(labels, &expected_props)
1022                    .into_iter()
1023                    .find(|&id| self.merge_node_matches(id, labels, &expected_props))
1024                else {
1025                    return Ok(false);
1026                };
1027                if let Some(var_id) = var {
1028                    row.insert(*var_id, LoraValue::Node(id));
1029                }
1030                if let Some(path_var) = part.binding {
1031                    let path = LoraPath {
1032                        nodes: vec![id],
1033                        rels: Vec::new(),
1034                    };
1035                    row.insert(path_var, LoraValue::Path(path));
1036                }
1037                Ok(true)
1038            }
1039
1040            ResolvedPatternElement::ShortestPath { .. } => {
1041                // ShortestPath is not valid in MERGE context
1042                Ok(false)
1043            }
1044
1045            ResolvedPatternElement::NodeChain { head, chain } => {
1046                // The head is usually bound by an earlier clause; otherwise
1047                // every node matching it is a possible start.
1048                let head_candidates = match head.var.and_then(|v| row.get(v)) {
1049                    Some(LoraValue::Node(id)) => vec![*id],
1050                    _ => {
1051                        let expected = self.merge_expected_props(head.properties.as_ref(), row);
1052                        self.merge_node_candidates(&head.labels, &expected)
1053                            .into_iter()
1054                            .filter(|&id| self.merge_node_matches(id, &head.labels, &expected))
1055                            .collect()
1056                    }
1057                };
1058
1059                for head_id in head_candidates {
1060                    let mut trial = row.clone();
1061                    if let Some(var_id) = head.var {
1062                        trial.insert(var_id, LoraValue::Node(head_id));
1063                    }
1064                    let mut walked = Vec::with_capacity(chain.len());
1065                    if self.match_merge_chain(&mut trial, head_id, chain, &mut walked) {
1066                        if let Some(path_var) = part.binding {
1067                            let path = LoraPath {
1068                                nodes: std::iter::once(head_id)
1069                                    .chain(walked.iter().map(|&(_, node)| node))
1070                                    .collect(),
1071                                rels: walked.iter().map(|&(rel, _)| rel).collect(),
1072                            };
1073                            trial.insert(path_var, LoraValue::Path(path));
1074                        }
1075                        *row = trial;
1076                        return Ok(true);
1077                    }
1078                }
1079                Ok(false)
1080            }
1081        }
1082    }
1083
1084    /// Match `chain` from `current`, backtracking over every candidate
1085    /// edge. A step node or relationship already bound in the row (by an
1086    /// earlier clause or earlier in the chain) must be the one reached;
1087    /// the same relationship is never used twice in one pattern. `walked`
1088    /// holds the `(relationship, node)` steps taken, in order.
1089    fn match_merge_chain(
1090        &self,
1091        row: &mut Row,
1092        current: NodeId,
1093        chain: &[lora_analyzer::ResolvedChain],
1094        walked: &mut Vec<(u64, NodeId)>,
1095    ) -> bool {
1096        let Some((step, rest)) = chain.split_first() else {
1097            return true;
1098        };
1099
1100        let bound_dst = match step.node.var.and_then(|v| row.get(v)) {
1101            Some(LoraValue::Node(id)) => Some(*id),
1102            _ => None,
1103        };
1104        let bound_rel = match step.rel.var.and_then(|v| row.get(v)) {
1105            Some(LoraValue::Relationship(id)) => Some(*id),
1106            _ => None,
1107        };
1108        let expected_node = self.merge_expected_props(step.node.properties.as_ref(), row);
1109        let expected_rel = self.merge_expected_props(step.rel.properties.as_ref(), row);
1110
1111        let edges = self
1112            .ctx
1113            .storage
1114            .expand_ids(current, step.rel.direction, &step.rel.types);
1115        for (rel_id, node_id) in edges {
1116            if bound_dst.is_some_and(|id| id != node_id)
1117                || bound_rel.is_some_and(|id| id != rel_id)
1118                || walked.iter().any(|&(used, _)| used == rel_id)
1119            {
1120                continue;
1121            }
1122            if !self.merge_node_matches(node_id, &step.node.labels, &expected_node) {
1123                continue;
1124            }
1125            if let Some(LoraValue::Map(expected_map)) = &expected_rel {
1126                let rel_ok = self
1127                    .ctx
1128                    .storage
1129                    .with_relationship(rel_id, |rel_rec| {
1130                        expected_map.iter().all(|(key, expected_val)| {
1131                            rel_rec
1132                                .properties
1133                                .get(key.as_str())
1134                                .map(|actual| value_matches_property_value(expected_val, actual))
1135                                .unwrap_or(false)
1136                        })
1137                    })
1138                    .unwrap_or(false);
1139                if !rel_ok {
1140                    continue;
1141                }
1142            }
1143
1144            let mut next = row.clone();
1145            if let Some(rel_var) = step.rel.var {
1146                next.insert(rel_var, LoraValue::Relationship(rel_id));
1147            }
1148            if let Some(node_var) = step.node.var {
1149                next.insert(node_var, LoraValue::Node(node_id));
1150            }
1151            walked.push((rel_id, node_id));
1152            if self.match_merge_chain(&mut next, node_id, rest, walked) {
1153                *row = next;
1154                return true;
1155            }
1156            walked.pop();
1157        }
1158        false
1159    }
1160
1161    fn merge_expected_props(
1162        &self,
1163        properties: Option<&ResolvedExpr>,
1164        row: &Row,
1165    ) -> Option<LoraValue> {
1166        let eval_ctx = EvalContext {
1167            storage: &*self.ctx.storage,
1168            params: &self.ctx.params,
1169        };
1170        properties.map(|e| eval_expr(e, row, &eval_ctx))
1171    }
1172
1173    /// Candidate ids for a MERGE node pattern. `MERGE (n:L {key: $k})`
1174    /// looks the key up in the property index instead of scanning every
1175    /// `:L` node, so an upsert costs the same on a large label as on a
1176    /// small one. Candidates are re-checked by [`Self::merge_node_matches`].
1177    fn merge_node_candidates(
1178        &self,
1179        labels: &[Vec<String>],
1180        expected_props: &Option<LoraValue>,
1181    ) -> Vec<NodeId> {
1182        let indexed = match expected_props {
1183            Some(LoraValue::Map(expected)) => {
1184                merge_candidates_from_index(&*self.ctx.storage, labels, expected)
1185            }
1186            _ => None,
1187        };
1188        match indexed {
1189            Some(ids) => ids,
1190            None if labels.is_empty() => self.ctx.storage.all_node_ids(),
1191            None => scan_node_ids_for_label_groups(&*self.ctx.storage, labels),
1192        }
1193    }
1194
1195    fn merge_node_matches(
1196        &self,
1197        id: NodeId,
1198        labels: &[Vec<String>],
1199        expected_props: &Option<LoraValue>,
1200    ) -> bool {
1201        self.ctx
1202            .storage
1203            .with_node(id, |node| {
1204                if !node_matches_label_groups(&node.labels, labels) {
1205                    return false;
1206                }
1207                if let Some(LoraValue::Map(expected)) = expected_props {
1208                    return expected.iter().all(|(key, expected_value)| {
1209                        node.properties
1210                            .get(key.as_str())
1211                            .map(|actual| value_matches_property_value(expected_value, actual))
1212                            .unwrap_or(false)
1213                    });
1214                }
1215                true
1216            })
1217            .unwrap_or(false)
1218    }
1219
1220    fn exec_delete(&mut self, plan: &PhysicalPlan, op: &DeleteExec) -> ExecResult<Vec<Row>> {
1221        let input_rows = self.execute_node(plan, op.input)?;
1222        let mut targets = BTreeSet::new();
1223
1224        for row in &input_rows {
1225            for expr in &op.expressions {
1226                let value = {
1227                    let eval_ctx = EvalContext {
1228                        storage: &*self.ctx.storage,
1229                        params: &self.ctx.params,
1230                    };
1231                    eval_expr(expr, row, &eval_ctx)
1232                };
1233                self.collect_delete_targets(&value, &mut targets)?;
1234            }
1235        }
1236
1237        self.validate_delete_targets(&targets, op.detach)?;
1238
1239        for target in &targets {
1240            if let DeleteTarget::Relationship(_) = target {
1241                self.delete_target(*target, op.detach)?;
1242            }
1243        }
1244        for target in targets {
1245            if let DeleteTarget::Node(_) = target {
1246                self.delete_target(target, op.detach)?;
1247            }
1248        }
1249
1250        Ok(input_rows)
1251    }
1252
1253    fn exec_set(&mut self, plan: &PhysicalPlan, op: &SetExec) -> ExecResult<Vec<Row>> {
1254        if crate::pull::subtree_is_fully_streaming(plan, op.input) {
1255            return self.streaming_apply(plan, op.input, |this, row| {
1256                for item in &op.items {
1257                    this.apply_set_item(row, item)?;
1258                }
1259                Ok(())
1260            });
1261        }
1262
1263        let input_rows = self.execute_node(plan, op.input)?;
1264
1265        for row in &input_rows {
1266            for item in &op.items {
1267                self.apply_set_item(row, item)?;
1268            }
1269        }
1270
1271        Ok(input_rows)
1272    }
1273
1274    /// `FOREACH (var IN list | body...)` — for each input row, evaluate
1275    /// the list and run the body once per element with `var` bound to
1276    /// that element. Each iteration runs on a fresh clone of the row
1277    /// so any new bindings the body introduces (e.g. anonymous
1278    /// `CREATE` node VarIds) don't leak between iterations or back to
1279    /// the outer scope. Side effects on the graph persist; the outer
1280    /// row is emitted unchanged.
1281    fn exec_foreach(&mut self, plan: &PhysicalPlan, op: &ForeachExec) -> ExecResult<Vec<Row>> {
1282        let input_rows = self.execute_node(plan, op.input)?;
1283        let mut out = Vec::with_capacity(input_rows.len());
1284
1285        for row in input_rows {
1286            let list_value = {
1287                let eval_ctx = EvalContext {
1288                    storage: &*self.ctx.storage,
1289                    params: &self.ctx.params,
1290                };
1291                eval_expr(&op.list, &row, &eval_ctx)
1292            };
1293
1294            let elements: Vec<LoraValue> = match list_value {
1295                LoraValue::List(items) => items,
1296                LoraValue::Null => Vec::new(),
1297                other => {
1298                    return Err(ExecutorError::RuntimeError(format!(
1299                        "FOREACH expects a list, got {}",
1300                        value_kind(&other)
1301                    )));
1302                }
1303            };
1304
1305            for element in elements {
1306                // Fresh row per iteration so body-introduced bindings
1307                // don't reuse VarIds across iterations.
1308                let mut iter_row = row.clone();
1309                iter_row.insert(op.variable, element);
1310                for clause in &op.body {
1311                    self.apply_foreach_body_clause(&mut iter_row, clause)?;
1312                }
1313            }
1314
1315            out.push(row);
1316        }
1317
1318        Ok(out)
1319    }
1320
1321    /// Apply one resolved updating clause to `row` for its side effect
1322    /// inside a `FOREACH` body. Only updating clauses (Create / Merge /
1323    /// Delete / Set / Remove / nested Foreach) are legal here; the
1324    /// analyzer guarantees that.
1325    fn apply_foreach_body_clause(
1326        &mut self,
1327        row: &mut Row,
1328        clause: &lora_analyzer::ResolvedClause,
1329    ) -> ExecResult<()> {
1330        use lora_analyzer::ResolvedClause;
1331        match clause {
1332            ResolvedClause::Create(c) => self.apply_create_pattern(row, &c.pattern),
1333            ResolvedClause::Set(s) => {
1334                for item in &s.items {
1335                    self.apply_set_item(row, item)?;
1336                }
1337                Ok(())
1338            }
1339            ResolvedClause::Remove(r) => {
1340                for item in &r.items {
1341                    self.apply_remove_item(row, item)?;
1342                }
1343                Ok(())
1344            }
1345            ResolvedClause::Delete(d) => {
1346                let detach = d.detach;
1347                for expr in &d.expressions {
1348                    let value = {
1349                        let eval_ctx = EvalContext {
1350                            storage: &*self.ctx.storage,
1351                            params: &self.ctx.params,
1352                        };
1353                        eval_expr(expr, row, &eval_ctx)
1354                    };
1355                    self.delete_value(value, detach)?;
1356                }
1357                Ok(())
1358            }
1359            ResolvedClause::Merge(m) => {
1360                let matched = self.match_merge_pattern(row, &m.pattern_part)?;
1361                if !matched {
1362                    self.apply_create_pattern_part(row, &m.pattern_part)?;
1363                }
1364                for action in &m.actions {
1365                    if action.on_match == matched {
1366                        for item in &action.set.items {
1367                            self.apply_set_item(row, item)?;
1368                        }
1369                    }
1370                }
1371                Ok(())
1372            }
1373            ResolvedClause::Foreach(nested) => {
1374                let list_value = {
1375                    let eval_ctx = EvalContext {
1376                        storage: &*self.ctx.storage,
1377                        params: &self.ctx.params,
1378                    };
1379                    eval_expr(&nested.list, row, &eval_ctx)
1380                };
1381
1382                let elements: Vec<LoraValue> = match list_value {
1383                    LoraValue::List(items) => items,
1384                    LoraValue::Null => Vec::new(),
1385                    other => {
1386                        return Err(ExecutorError::RuntimeError(format!(
1387                            "FOREACH expects a list, got {}",
1388                            value_kind(&other)
1389                        )));
1390                    }
1391                };
1392
1393                for element in elements {
1394                    let mut iter_row = row.clone();
1395                    iter_row.insert(nested.variable, element);
1396                    for inner in &nested.body {
1397                        self.apply_foreach_body_clause(&mut iter_row, inner)?;
1398                    }
1399                }
1400
1401                Ok(())
1402            }
1403            other => Err(ExecutorError::RuntimeError(format!(
1404                "FOREACH body may only contain updating clauses, got {:?}",
1405                std::mem::discriminant(other)
1406            ))),
1407        }
1408    }
1409
1410    fn exec_remove(&mut self, plan: &PhysicalPlan, op: &RemoveExec) -> ExecResult<Vec<Row>> {
1411        if crate::pull::subtree_is_fully_streaming(plan, op.input) {
1412            return self.streaming_apply(plan, op.input, |this, row| {
1413                for item in &op.items {
1414                    this.apply_remove_item(row, item)?;
1415                }
1416                Ok(())
1417            });
1418        }
1419
1420        let input_rows = self.execute_node(plan, op.input)?;
1421
1422        for row in &input_rows {
1423            for item in &op.items {
1424                self.apply_remove_item(row, item)?;
1425            }
1426        }
1427
1428        Ok(input_rows)
1429    }
1430
1431    fn apply_set_item(&mut self, row: &Row, item: &ResolvedSetItem) -> ExecResult<()> {
1432        match item {
1433            ResolvedSetItem::SetProperty { target, value } => {
1434                let new_value = {
1435                    let eval_ctx = EvalContext {
1436                        storage: &*self.ctx.storage,
1437                        params: &self.ctx.params,
1438                    };
1439                    eval_expr(value, row, &eval_ctx)
1440                };
1441
1442                self.set_property_from_expr(row, target, new_value)
1443            }
1444
1445            ResolvedSetItem::SetVariable { variable, value } => {
1446                // Only need the entity's id — peek at the binding by reference.
1447                let entity_ref =
1448                    row.get(*variable)
1449                        .ok_or(ExecutorError::UnboundVariableForSet {
1450                            var: format!("{variable:?}"),
1451                        })?;
1452                let entity_target = entity_target_from_value(entity_ref)?;
1453
1454                let new_value = {
1455                    let eval_ctx = EvalContext {
1456                        storage: &*self.ctx.storage,
1457                        params: &self.ctx.params,
1458                    };
1459                    eval_expr(value, row, &eval_ctx)
1460                };
1461
1462                self.overwrite_entity_target(entity_target, new_value)
1463            }
1464
1465            ResolvedSetItem::MutateVariable { variable, value } => {
1466                let entity_ref =
1467                    row.get(*variable)
1468                        .ok_or(ExecutorError::UnboundVariableForSet {
1469                            var: format!("{variable:?}"),
1470                        })?;
1471                let entity_target = entity_target_from_value(entity_ref)?;
1472
1473                let patch = {
1474                    let eval_ctx = EvalContext {
1475                        storage: &*self.ctx.storage,
1476                        params: &self.ctx.params,
1477                    };
1478                    eval_expr(value, row, &eval_ctx)
1479                };
1480
1481                self.mutate_entity_target(entity_target, patch)
1482            }
1483
1484            ResolvedSetItem::SetLabels { variable, labels } => match row.get(*variable) {
1485                Some(LoraValue::Node(node_id)) => {
1486                    let node_id = *node_id;
1487                    for label in labels {
1488                        // A deferring statement may supply the label's
1489                        // required properties after the label (`SET n:L,
1490                        // n.required = 1`): existence waits for its end.
1491                        let checked = if self.defer_existence {
1492                            self.ctx
1493                                .storage
1494                                .check_node_add_label_deferring_existence(node_id, label)
1495                        } else {
1496                            self.ctx
1497                                .storage
1498                                .check_node_add_label_against_constraints(node_id, label)
1499                        };
1500                        checked.map_err(ExecutorError::ConstraintViolation)?;
1501                        self.ctx.storage.add_node_label(node_id, label);
1502                    }
1503                    if self.defer_existence {
1504                        self.pending_existence.push(EntityTarget::Node(node_id));
1505                    }
1506                    Ok(())
1507                }
1508                Some(other) => Err(ExecutorError::ExpectedNodeForSetLabels {
1509                    found: value_kind(other),
1510                }),
1511                None => Err(ExecutorError::UnboundVariableForSet {
1512                    var: format!("{variable:?}"),
1513                }),
1514            },
1515        }
1516    }
1517
1518    fn set_property_from_expr(
1519        &mut self,
1520        row: &Row,
1521        target_expr: &ResolvedExpr,
1522        new_value: LoraValue,
1523    ) -> ExecResult<()> {
1524        let ResolvedExpr::Property { expr, property } = target_expr else {
1525            return Err(ExecutorError::UnsupportedSetTarget);
1526        };
1527
1528        let owner = {
1529            let eval_ctx = EvalContext {
1530                storage: &*self.ctx.storage,
1531                params: &self.ctx.params,
1532            };
1533            eval_expr(expr, row, &eval_ctx)
1534        };
1535
1536        // `SET n.a = null` removes the property.
1537        if matches!(new_value, LoraValue::Null) {
1538            return match owner {
1539                LoraValue::Node(node_id) => {
1540                    self.remove_entity_property(EntityTarget::Node(node_id), property)
1541                }
1542                LoraValue::Relationship(rel_id) => {
1543                    self.remove_entity_property(EntityTarget::Relationship(rel_id), property)
1544                }
1545                other => Err(ExecutorError::InvalidSetTarget {
1546                    found: value_kind(&other),
1547                }),
1548            };
1549        }
1550
1551        match owner {
1552            LoraValue::Node(node_id) => {
1553                let prop = lora_value_to_property(new_value)
1554                    .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1555                if let Err(msg) = self
1556                    .ctx
1557                    .storage
1558                    .check_node_set_property_against_constraints(node_id, property, &prop)
1559                {
1560                    return Err(ExecutorError::ConstraintViolation(msg));
1561                }
1562                self.ctx
1563                    .storage
1564                    .set_node_property(node_id, property.clone(), prop);
1565                Ok(())
1566            }
1567            LoraValue::Relationship(rel_id) => {
1568                let prop = lora_value_to_property(new_value)
1569                    .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1570                if let Err(msg) = self
1571                    .ctx
1572                    .storage
1573                    .check_relationship_set_property_against_constraints(rel_id, property, &prop)
1574                {
1575                    return Err(ExecutorError::ConstraintViolation(msg));
1576                }
1577                self.ctx
1578                    .storage
1579                    .set_relationship_property(rel_id, property.clone(), prop);
1580                Ok(())
1581            }
1582            other => Err(ExecutorError::InvalidSetTarget {
1583                found: value_kind(&other),
1584            }),
1585        }
1586    }
1587
1588    /// Remove one property, checking constraints first. Removing a
1589    /// property the entity does not have is a no-op. Removal can only
1590    /// break an existence constraint, so a statement that defers those
1591    /// checks it once it finishes (`REMOVE n.a SET n.a = …` keeps it).
1592    fn remove_entity_property(&mut self, target: EntityTarget, property: &str) -> ExecResult<()> {
1593        if self.defer_existence {
1594            match target {
1595                EntityTarget::Node(node_id) => {
1596                    self.ctx.storage.remove_node_property(node_id, property);
1597                }
1598                EntityTarget::Relationship(rel_id) => {
1599                    self.ctx
1600                        .storage
1601                        .remove_relationship_property(rel_id, property);
1602                }
1603            }
1604            self.pending_existence.push(target);
1605            return Ok(());
1606        }
1607        match target {
1608            EntityTarget::Node(node_id) => {
1609                if let Err(msg) = self
1610                    .ctx
1611                    .storage
1612                    .check_node_remove_property_against_constraints(node_id, property)
1613                {
1614                    return Err(ExecutorError::ConstraintViolation(msg));
1615                }
1616                self.ctx.storage.remove_node_property(node_id, property);
1617            }
1618            EntityTarget::Relationship(rel_id) => {
1619                if let Err(msg) = self
1620                    .ctx
1621                    .storage
1622                    .check_relationship_remove_property_against_constraints(rel_id, property)
1623                {
1624                    return Err(ExecutorError::ConstraintViolation(msg));
1625                }
1626                self.ctx
1627                    .storage
1628                    .remove_relationship_property(rel_id, property);
1629            }
1630        }
1631        Ok(())
1632    }
1633
1634    fn remove_property_from_expr(&mut self, row: &Row, expr: &ResolvedExpr) -> ExecResult<()> {
1635        let ResolvedExpr::Property {
1636            expr: owner_expr,
1637            property,
1638        } = expr
1639        else {
1640            return Err(ExecutorError::UnsupportedRemoveTarget);
1641        };
1642
1643        let owner = {
1644            let eval_ctx = EvalContext {
1645                storage: &*self.ctx.storage,
1646                params: &self.ctx.params,
1647            };
1648            eval_expr(owner_expr, row, &eval_ctx)
1649        };
1650
1651        match owner {
1652            LoraValue::Node(node_id) => {
1653                self.remove_entity_property(EntityTarget::Node(node_id), property)
1654            }
1655            LoraValue::Relationship(rel_id) => {
1656                self.remove_entity_property(EntityTarget::Relationship(rel_id), property)
1657            }
1658            other => Err(ExecutorError::InvalidRemoveTarget {
1659                found: value_kind(&other),
1660            }),
1661        }
1662    }
1663
1664    fn overwrite_entity_target(
1665        &mut self,
1666        target: EntityTarget,
1667        new_value: LoraValue,
1668    ) -> ExecResult<()> {
1669        let LoraValue::Map(map) = new_value else {
1670            return Err(ExecutorError::ExpectedPropertyMap {
1671                found: value_kind(&new_value),
1672            });
1673        };
1674
1675        let mut props: Properties = Properties::new();
1676        for (k, v) in map {
1677            // `SET n = {a: null}` leaves `a` absent.
1678            if matches!(v, LoraValue::Null) {
1679                continue;
1680            }
1681            let prop = lora_value_to_property(v)
1682                .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1683            props.insert(lora_store::intern_owned(k), prop);
1684        }
1685
1686        // Replacing drops every property the map leaves out, so a
1687        // deferring statement may restore a required one later (`SET n =
1688        // {…}, n.required = …`): existence waits for its end.
1689        let storage = &*self.ctx.storage;
1690        let checked = match (target, self.defer_existence) {
1691            (EntityTarget::Node(id), true) => {
1692                storage.check_node_replace_properties_deferring_existence(id, &props)
1693            }
1694            (EntityTarget::Node(id), false) => {
1695                storage.check_node_replace_properties_against_constraints(id, &props)
1696            }
1697            (EntityTarget::Relationship(id), true) => {
1698                storage.check_relationship_replace_properties_deferring_existence(id, &props)
1699            }
1700            (EntityTarget::Relationship(id), false) => {
1701                storage.check_relationship_replace_properties_against_constraints(id, &props)
1702            }
1703        };
1704        checked.map_err(ExecutorError::ConstraintViolation)?;
1705        match target {
1706            EntityTarget::Node(node_id) => {
1707                self.ctx.storage.replace_node_properties(node_id, props);
1708            }
1709            EntityTarget::Relationship(rel_id) => {
1710                self.ctx
1711                    .storage
1712                    .replace_relationship_properties(rel_id, props);
1713            }
1714        }
1715        if self.defer_existence {
1716            self.pending_existence.push(target);
1717        }
1718        Ok(())
1719    }
1720
1721    fn mutate_entity_target(
1722        &mut self,
1723        target: EntityTarget,
1724        patch_value: LoraValue,
1725    ) -> ExecResult<()> {
1726        let LoraValue::Map(map) = patch_value else {
1727            return Err(ExecutorError::ExpectedPropertyMap {
1728                found: value_kind(&patch_value),
1729            });
1730        };
1731
1732        match target {
1733            EntityTarget::Node(node_id) => {
1734                for (k, v) in map {
1735                    // `SET n += {a: null}` removes `a`.
1736                    if matches!(v, LoraValue::Null) {
1737                        self.remove_entity_property(target, &k)?;
1738                        continue;
1739                    }
1740                    let prop = lora_value_to_property(v)
1741                        .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1742                    if let Err(msg) = self
1743                        .ctx
1744                        .storage
1745                        .check_node_set_property_against_constraints(node_id, &k, &prop)
1746                    {
1747                        return Err(ExecutorError::ConstraintViolation(msg));
1748                    }
1749                    self.ctx.storage.set_node_property(node_id, k, prop);
1750                }
1751            }
1752            EntityTarget::Relationship(rel_id) => {
1753                for (k, v) in map {
1754                    if matches!(v, LoraValue::Null) {
1755                        self.remove_entity_property(target, &k)?;
1756                        continue;
1757                    }
1758                    let prop = lora_value_to_property(v)
1759                        .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1760                    if let Err(msg) = self
1761                        .ctx
1762                        .storage
1763                        .check_relationship_set_property_against_constraints(rel_id, &k, &prop)
1764                    {
1765                        return Err(ExecutorError::ConstraintViolation(msg));
1766                    }
1767                    self.ctx.storage.set_relationship_property(rel_id, k, prop);
1768                }
1769            }
1770        }
1771        Ok(())
1772    }
1773
1774    pub(crate) fn apply_create_pattern(
1775        &mut self,
1776        row: &mut Row,
1777        pattern: &ResolvedPattern,
1778    ) -> ExecResult<()> {
1779        for part in &pattern.parts {
1780            self.apply_create_pattern_part(row, part)?;
1781        }
1782        Ok(())
1783    }
1784
1785    /// Apply a single per-row write for any of the streamable write
1786    /// operators (Create / Set / Delete / Remove / Merge). Used by
1787    /// the [`crate::pull::StreamingWriteCursor`] auto-commit fast
1788    /// path: the cursor pulls one input row from a read upstream,
1789    /// hands it here for the side effect, and emits the row back.
1790    pub(crate) fn apply_write_op(&mut self, op: &PhysicalOp, row: &mut Row) -> ExecResult<()> {
1791        match op {
1792            PhysicalOp::Create(c) => self.apply_create_pattern(row, &c.pattern),
1793            PhysicalOp::Set(s) => {
1794                for item in &s.items {
1795                    self.apply_set_item(row, item)?;
1796                }
1797                Ok(())
1798            }
1799            PhysicalOp::Delete(d) => {
1800                let detach = d.detach;
1801                for expr in &d.expressions {
1802                    let value = {
1803                        let eval_ctx = EvalContext {
1804                            storage: &*self.ctx.storage,
1805                            params: &self.ctx.params,
1806                        };
1807                        eval_expr(expr, row, &eval_ctx)
1808                    };
1809                    self.delete_value(value, detach)?;
1810                }
1811                Ok(())
1812            }
1813            PhysicalOp::Remove(r) => {
1814                for item in &r.items {
1815                    self.apply_remove_item(row, item)?;
1816                }
1817                Ok(())
1818            }
1819            PhysicalOp::Merge(m) => {
1820                let matched = self.match_merge_pattern(row, &m.pattern_part)?;
1821                if !matched {
1822                    self.apply_create_pattern_part(row, &m.pattern_part)?;
1823                }
1824                for action in &m.actions {
1825                    if action.on_match == matched {
1826                        for item in &action.set.items {
1827                            self.apply_set_item(row, item)?;
1828                        }
1829                    }
1830                }
1831                Ok(())
1832            }
1833            other => Err(ExecutorError::RuntimeError(format!(
1834                "apply_write_op called on non-write op: {other:?}"
1835            ))),
1836        }
1837    }
1838
1839    /// Create a pattern part, binding its path variable (`CREATE p = …`)
1840    /// to the nodes and relationships it created or reused.
1841    fn apply_create_pattern_part(
1842        &mut self,
1843        row: &mut Row,
1844        part: &ResolvedPatternPart,
1845    ) -> ExecResult<()> {
1846        let path = self.apply_create_pattern_element(row, &part.element)?;
1847        if let (Some(var), Some(path)) = (part.binding, path) {
1848            row.insert(var, LoraValue::Path(path));
1849        }
1850        Ok(())
1851    }
1852
1853    fn apply_create_pattern_element(
1854        &mut self,
1855        row: &mut Row,
1856        element: &ResolvedPatternElement,
1857    ) -> ExecResult<Option<LoraPath>> {
1858        match element {
1859            ResolvedPatternElement::Node {
1860                var,
1861                labels,
1862                properties,
1863            } => {
1864                let node_id =
1865                    self.materialize_node_pattern(row, *var, labels, properties.as_ref())?;
1866                Ok(Some(LoraPath {
1867                    nodes: vec![node_id],
1868                    rels: Vec::new(),
1869                }))
1870            }
1871
1872            ResolvedPatternElement::NodeChain { head, chain } => {
1873                let mut current_node_id = self.materialize_node_pattern(
1874                    row,
1875                    head.var,
1876                    &head.labels,
1877                    head.properties.as_ref(),
1878                )?;
1879                let mut path = LoraPath {
1880                    nodes: Vec::with_capacity(chain.len() + 1),
1881                    rels: Vec::with_capacity(chain.len()),
1882                };
1883                path.nodes.push(current_node_id);
1884
1885                for link in chain {
1886                    let next_node_id = self.materialize_node_pattern(
1887                        row,
1888                        link.node.var,
1889                        &link.node.labels,
1890                        link.node.properties.as_ref(),
1891                    )?;
1892
1893                    let rel_id = self.materialize_relationship_pattern(
1894                        row,
1895                        current_node_id,
1896                        next_node_id,
1897                        &link.rel,
1898                    )?;
1899                    path.rels.push(rel_id);
1900                    path.nodes.push(next_node_id);
1901
1902                    current_node_id = next_node_id;
1903                }
1904
1905                Ok(Some(path))
1906            }
1907
1908            ResolvedPatternElement::ShortestPath { .. } => {
1909                // ShortestPath is not valid in CREATE context
1910                Ok(None)
1911            }
1912        }
1913    }
1914
1915    /// Whether every variable of a MERGE pattern part is already bound in
1916    /// `row` (so the MERGE matches trivially). A variable bound to a value
1917    /// of the wrong kind for its position (a map, null or scalar where a
1918    /// node or relationship belongs) is an error rather than "unbound":
1919    /// treating it as unbound would MERGE a fresh, unrelated entity.
1920    fn pattern_part_is_bound(&self, row: &Row, part: &ResolvedPatternPart) -> ExecResult<bool> {
1921        fn node_bound(row: &Row, var: Option<VarId>) -> ExecResult<bool> {
1922            let Some(var) = var else { return Ok(false) };
1923            match row.get(var) {
1924                None => Ok(false),
1925                Some(LoraValue::Node(_)) => Ok(true),
1926                Some(other) => Err(ExecutorError::ExpectedNodeForCreate {
1927                    var: bound_var_name(row, var),
1928                    found: value_kind(other),
1929                }),
1930            }
1931        }
1932        fn rel_bound(row: &Row, var: Option<VarId>) -> ExecResult<bool> {
1933            // For MERGE, anonymous relationships cannot be considered
1934            // "bound" because we have no variable to check. The merge
1935            // must search the graph to see if the relationship exists.
1936            let Some(var) = var else { return Ok(false) };
1937            match row.get(var) {
1938                None => Ok(false),
1939                Some(LoraValue::Relationship(_)) => Ok(true),
1940                Some(other) => Err(ExecutorError::ExpectedRelationshipForCreate {
1941                    var: bound_var_name(row, var),
1942                    found: value_kind(other),
1943                }),
1944            }
1945        }
1946
1947        match &part.element {
1948            ResolvedPatternElement::Node { var, .. } => node_bound(row, *var),
1949
1950            ResolvedPatternElement::ShortestPath { .. } => Ok(false),
1951
1952            ResolvedPatternElement::NodeChain { head, chain } => {
1953                // Check every position (not short-circuiting) so a
1954                // mis-typed binding anywhere in the chain is reported.
1955                let mut all_bound = node_bound(row, head.var)?;
1956                for link in chain {
1957                    let node_ok = node_bound(row, link.node.var)?;
1958                    let rel_ok = rel_bound(row, link.rel.var)?;
1959                    all_bound &= node_ok && rel_ok;
1960                }
1961                Ok(all_bound)
1962            }
1963        }
1964    }
1965
1966    fn materialize_node_pattern(
1967        &mut self,
1968        row: &mut Row,
1969        var: Option<VarId>,
1970        labels: &[Vec<String>],
1971        properties: Option<&ResolvedExpr>,
1972    ) -> ExecResult<u64> {
1973        if let Some(var_id) = var {
1974            match row.get(var_id) {
1975                Some(LoraValue::Node(id)) => return Ok(*id),
1976                // A bound variable in a node position names an existing
1977                // node. Anything else (a map, null, a scalar) is an error,
1978                // never a fresh node: silently creating one would attach
1979                // the pattern to a blank node instead of the intended one.
1980                Some(other) => {
1981                    return Err(ExecutorError::ExpectedNodeForCreate {
1982                        var: bound_var_name(row, var_id),
1983                        found: value_kind(other),
1984                    });
1985                }
1986                None => {}
1987            }
1988        }
1989
1990        let properties = match properties {
1991            Some(expr) => eval_properties_expr(expr, row, &*self.ctx.storage, &self.ctx.params)?,
1992            None => Properties::new(),
1993        };
1994
1995        let flat_labels = flatten_label_groups(labels);
1996        debug!("creating node with labels={flat_labels:?}");
1997        let checked = if self.defer_existence {
1998            self.ctx
1999                .storage
2000                .check_node_create_deferring_existence(&flat_labels, &properties)
2001        } else {
2002            self.ctx
2003                .storage
2004                .check_node_create_against_constraints(&flat_labels, &properties)
2005        };
2006        checked.map_err(ExecutorError::ConstraintViolation)?;
2007        let created = self
2008            .ctx
2009            .storage
2010            .try_create_node(flat_labels, properties)
2011            .ok_or(ExecutorError::NodeCreateFailed)?;
2012        if self.defer_existence {
2013            self.pending_existence.push(EntityTarget::Node(created.id));
2014        }
2015
2016        if let Some(var_id) = var {
2017            row.insert(var_id, LoraValue::Node(created.id));
2018        }
2019
2020        Ok(created.id)
2021    }
2022
2023    fn materialize_relationship_pattern(
2024        &mut self,
2025        row: &mut Row,
2026        left_node_id: u64,
2027        right_node_id: u64,
2028        rel: &lora_analyzer::ResolvedRel,
2029    ) -> ExecResult<u64> {
2030        if let Some(var_id) = rel.var {
2031            if let Some(other) = row
2032                .get(var_id)
2033                .filter(|v| !matches!(v, LoraValue::Relationship(_)))
2034            {
2035                return Err(ExecutorError::ExpectedRelationshipForCreate {
2036                    var: bound_var_name(row, var_id),
2037                    found: value_kind(other),
2038                });
2039            }
2040            if let Some(LoraValue::Relationship(id)) = row.get(var_id) {
2041                let id = *id;
2042                if let Some((src, dst)) = self.ctx.storage.relationship_endpoints(id) {
2043                    let endpoints_match = match rel.direction {
2044                        Direction::Right | Direction::Undirected => {
2045                            src == left_node_id && dst == right_node_id
2046                        }
2047                        Direction::Left => src == right_node_id && dst == left_node_id,
2048                    };
2049
2050                    if endpoints_match {
2051                        return Ok(id);
2052                    }
2053                }
2054            }
2055        }
2056
2057        if rel.range.is_some() {
2058            return Err(ExecutorError::UnsupportedCreateRelationshipRange);
2059        }
2060
2061        let (src, dst) = match rel.direction {
2062            Direction::Right | Direction::Undirected => (left_node_id, right_node_id),
2063            Direction::Left => (right_node_id, left_node_id),
2064        };
2065
2066        let rel_type = rel
2067            .types
2068            .first()
2069            .ok_or(ExecutorError::MissingRelationshipType)?;
2070
2071        if rel_type.is_empty() {
2072            return Err(ExecutorError::MissingRelationshipType);
2073        }
2074
2075        let properties = match rel.properties.as_ref() {
2076            Some(expr) => eval_properties_expr(expr, row, &*self.ctx.storage, &self.ctx.params)?,
2077            None => Properties::new(),
2078        };
2079
2080        debug!("creating relationship: src={src}, dst={dst}, type={rel_type}");
2081
2082        let checked = if self.defer_existence {
2083            self.ctx
2084                .storage
2085                .check_relationship_create_deferring_existence(rel_type, &properties)
2086        } else {
2087            self.ctx
2088                .storage
2089                .check_relationship_create_against_constraints(rel_type, &properties)
2090        };
2091        checked.map_err(ExecutorError::ConstraintViolation)?;
2092
2093        let created = self
2094            .ctx
2095            .storage
2096            .create_relationship(src, dst, rel_type, properties)
2097            .ok_or_else(|| ExecutorError::RelationshipCreateFailed {
2098                src,
2099                dst,
2100                rel_type: rel_type.clone(),
2101            })?;
2102        if self.defer_existence {
2103            self.pending_existence
2104                .push(EntityTarget::Relationship(created.id));
2105        }
2106
2107        if let Some(var_id) = rel.var {
2108            row.insert(var_id, LoraValue::Relationship(created.id));
2109        }
2110
2111        Ok(created.id)
2112    }
2113}
2114
2115/// The path a pattern part names when every node and relationship in it is
2116/// already bound in `row` (see `pattern_part_is_bound`).
2117fn bound_pattern_path(row: &Row, part: &ResolvedPatternPart) -> Option<LoraPath> {
2118    let node = |var: Option<VarId>| match var.and_then(|v| row.get(v)) {
2119        Some(LoraValue::Node(id)) => Some(*id),
2120        _ => None,
2121    };
2122    match &part.element {
2123        ResolvedPatternElement::Node { var, .. } => Some(LoraPath {
2124            nodes: vec![node(*var)?],
2125            rels: Vec::new(),
2126        }),
2127        ResolvedPatternElement::NodeChain { head, chain } => {
2128            let mut path = LoraPath {
2129                nodes: vec![node(head.var)?],
2130                rels: Vec::with_capacity(chain.len()),
2131            };
2132            for link in chain {
2133                match link.rel.var.and_then(|v| row.get(v)) {
2134                    Some(LoraValue::Relationship(id)) => path.rels.push(*id),
2135                    _ => return None,
2136                }
2137                path.nodes.push(node(link.node.var)?);
2138            }
2139            Some(path)
2140        }
2141        ResolvedPatternElement::ShortestPath { .. } => None,
2142    }
2143}
2144
2145/// Whether existence constraints must wait for the end of the statement.
2146/// They can be checked at `CREATE` only when nothing after it can add a
2147/// property or remove the entity: the only write is one `CREATE`. Checking
2148/// early keeps a failing create from mutating anything.
2149pub(crate) fn plan_defers_existence(plan: &PhysicalPlan) -> bool {
2150    let mut creates = 0;
2151    let mut deletes = false;
2152    for op in &plan.nodes {
2153        match op {
2154            PhysicalOp::Create(_) => creates += 1,
2155            // `CREATE (n:L) DELETE n`: a deleted entity needs no property.
2156            PhysicalOp::Delete(_) => deletes = true,
2157            PhysicalOp::Merge(_)
2158            | PhysicalOp::Set(_)
2159            | PhysicalOp::Remove(_)
2160            | PhysicalOp::Foreach(_) => return true,
2161            _ => {}
2162        }
2163    }
2164    creates > 1 || creates == 1 && deletes
2165}
2166
2167/// Whether a plan is a write statement with no `RETURN` (its root is the
2168/// write operator itself). Such a statement produces no result rows, as
2169/// in other Cypher databases; the write operator's pass-through rows
2170/// would otherwise leak as anonymous `_0` columns carrying internal ids.
2171pub(crate) fn plan_ends_in_write(plan: &PhysicalPlan) -> bool {
2172    match &plan.nodes[plan.root] {
2173        PhysicalOp::Create(_)
2174        | PhysicalOp::Merge(_)
2175        | PhysicalOp::Set(_)
2176        | PhysicalOp::Delete(_)
2177        | PhysicalOp::Remove(_)
2178        | PhysicalOp::Foreach(_) => true,
2179        // A query ending in a unit `CALL { ... }` returns no rows.
2180        PhysicalOp::CallSubquery(op) => op.new_vars.is_empty(),
2181        _ => false,
2182    }
2183}
2184
2185/// Candidate nodes for a MERGE node pattern from the property index, or
2186/// `None` to fall back to a label scan. Every candidate is still checked
2187/// against the full pattern, so the only requirement is that no real
2188/// match is missed. MERGE compares `1` and `1.0` as equal while the index
2189/// keys them apart, so numbers look up both images; values without an
2190/// exact index image (lists, maps, NaN, floats beyond 2^53) scan.
2191fn merge_candidates_from_index<S: lora_store::GraphStorage>(
2192    storage: &S,
2193    labels: &[Vec<String>],
2194    expected: &std::collections::BTreeMap<String, LoraValue>,
2195) -> Option<Vec<lora_store::NodeId>> {
2196    use lora_store::PropertyValue;
2197
2198    // A single required label scopes the lookup; otherwise look up
2199    // across labels and let the pattern check filter.
2200    let label = match labels {
2201        [group] if group.len() == 1 => Some(group[0].as_str()),
2202        _ => None,
2203    };
2204    let (key, value) = expected.iter().find(|(_, v)| {
2205        matches!(
2206            v,
2207            LoraValue::String(_) | LoraValue::Bool(_) | LoraValue::Int(_)
2208        ) || matches!(v, LoraValue::Float(f) if f.is_finite() && f.abs() < 9_007_199_254_740_992.0)
2209    })?;
2210    let images: Vec<PropertyValue> = match value {
2211        LoraValue::String(s) => vec![PropertyValue::String(s.clone())],
2212        LoraValue::Bool(b) => vec![PropertyValue::Bool(*b)],
2213        LoraValue::Int(i) => vec![PropertyValue::Int(*i), PropertyValue::Float(*i as f64)],
2214        LoraValue::Float(f) => {
2215            let mut v = vec![PropertyValue::Float(*f)];
2216            if f.fract() == 0.0 {
2217                v.push(PropertyValue::Int(*f as i64));
2218            }
2219            v
2220        }
2221        _ => return None,
2222    };
2223    let mut ids: Vec<lora_store::NodeId> = images
2224        .iter()
2225        .flat_map(|image| storage.find_node_ids_by_property(label, key, image))
2226        .collect();
2227    ids.sort_unstable();
2228    ids.dedup();
2229    Some(ids)
2230}
2231
2232/// The user-facing name of `var` in `row`, for error messages.
2233fn bound_var_name(row: &Row, var: VarId) -> String {
2234    row.iter_named()
2235        .find(|(key, _, _)| **key == var)
2236        .map(|(_, name, _)| name.into_owned())
2237        .unwrap_or_else(|| format!("{var:?}"))
2238}