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