Skip to main content

lora_executor/executor/
mutable.rs

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