Skip to main content

lora_executor/executor/
mutable.rs

1//! Mutable buffered executor: applies CREATE / MERGE / DELETE / SET /
2//! REMOVE on top of the read-side operator set.
3//!
4//! [`MutableExecutor`] mirrors the read-only [`super::immutable::Executor`]
5//! for all read operators (so a write op above any read subtree
6//! materializes the same way) and adds the per-row write
7//! implementations. The streaming pull pipeline in `crate::pull` runs
8//! `MutableExecutor::apply_write_op` row-by-row through the
9//! `StreamingWriteCursor` fast path; the buffered `exec_*` methods
10//! here handle the fallback when a write op's input subtree is not
11//! fully streamable.
12
13use crate::errors::{value_kind, ExecResult, ExecutorError};
14use crate::eval::{clear_eval_error, eval_expr, EvalContext};
15use crate::value::{lora_value_to_property, LoraValue, Row};
16use crate::{project_rows, ExecuteOptions, QueryResult};
17
18use lora_analyzer::{
19    symbols::VarId, ResolvedExpr, ResolvedPattern, ResolvedPatternElement, ResolvedPatternPart,
20    ResolvedRemoveItem, ResolvedSetItem,
21};
22use lora_ast::Direction;
23use lora_compiler::physical::*;
24use lora_compiler::CompiledQuery;
25use lora_store::{GraphStorageMut, NodeId, Properties};
26
27use std::collections::{BTreeMap, BTreeSet};
28use tracing::{debug, error, trace};
29use web_time::Instant;
30
31use super::aggregate_rows;
32use super::helpers::{
33    build_path_value, check_deadline_at, dedup_rows, eval_properties_expr, expand_rows,
34    expand_var_len_rows, filter_rows_checked, filter_shortest_paths, flatten_label_groups,
35    hydrate_node_record, hydrate_relationship_record, limit_rows, node_by_label_scan_rows,
36    node_by_property_scan_rows, node_matches_label_groups, node_scan_rows, plan_may_need_hydration,
37    project_rows_checked, scan_node_ids_for_label_groups, unwind_rows,
38    value_matches_property_value,
39};
40use super::optional_match_rows;
41use super::{sort_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(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(
548            input_rows,
549            &op.group_by,
550            &op.aggregates,
551            &eval_ctx,
552            |value| self.hydrate_value(value),
553        )
554    }
555
556    fn exec_sort(&mut self, plan: &PhysicalPlan, op: &SortExec) -> ExecResult<Vec<Row>> {
557        let mut rows = self.execute_node(plan, op.input)?;
558        let eval_ctx = EvalContext {
559            storage: &*self.ctx.storage,
560            params: &self.ctx.params,
561        };
562
563        let bound = sort_row_bound(op.top_k, op.limit.as_ref(), &eval_ctx);
564        sort_rows_with_top_k(&mut rows, &op.items, &eval_ctx, bound);
565
566        Ok(rows)
567    }
568
569    fn exec_limit(&mut self, plan: &PhysicalPlan, op: &LimitExec) -> ExecResult<Vec<Row>> {
570        let rows = self.execute_node(plan, op.input)?;
571        let eval_ctx = EvalContext {
572            storage: &*self.ctx.storage,
573            params: &self.ctx.params,
574        };
575
576        limit_rows(rows, op, &eval_ctx)
577    }
578
579    fn exec_optional_match(
580        &mut self,
581        plan: &PhysicalPlan,
582        op: &OptionalMatchExec,
583    ) -> ExecResult<Vec<Row>> {
584        let input_rows = self.execute_node(plan, op.input)?;
585
586        if super::optional::optional_can_correlate(plan, op.inner) {
587            let storage_ref: &S = &*self.ctx.storage;
588            return super::optional::correlated_optional_match_rows(
589                storage_ref,
590                &self.ctx.params,
591                plan,
592                op.inner,
593                input_rows,
594                &op.new_vars,
595            );
596        }
597
598        // Fallback: execute the inner plan once, uncorrelated, and join.
599        let inner_rows = self.execute_node(plan, op.inner)?;
600
601        Ok(optional_match_rows(input_rows, &inner_rows, &op.new_vars))
602    }
603
604    fn exec_call_subquery(
605        &mut self,
606        plan: &PhysicalPlan,
607        op: &CallSubqueryExec,
608    ) -> ExecResult<Vec<Row>> {
609        let input_rows = self.execute_node(plan, op.input)?;
610        let mut out = Vec::with_capacity(input_rows.len());
611
612        if crate::pull::subtree_has_write(plan, op.inner) {
613            // A writing body runs on this executor, once per outer row,
614            // with the outer row seeded into its bottom `Argument`. Each
615            // run sees the writes of the runs before it.
616            let unit = op.new_vars.is_empty();
617            for outer_row in input_rows {
618                self.check_deadline()?;
619                let prev = self.argument_seed.replace(outer_row.clone());
620                let inner_rows = self.execute_node(plan, op.inner);
621                self.argument_seed = prev;
622                let inner_rows = inner_rows?;
623                if unit {
624                    // A unit subquery keeps the outer row as it is, once,
625                    // however many rows its body produced.
626                    out.push(outer_row);
627                    continue;
628                }
629                for inner_row in inner_rows {
630                    out.push(crate::executor::merge_optional_rows(&outer_row, &inner_row));
631                }
632            }
633            return Ok(out);
634        }
635
636        let params = std::sync::Arc::new(self.ctx.params.clone());
637        let storage_ref: &S = &*self.ctx.storage;
638        for outer_row in input_rows {
639            let mut inner_source = crate::pull::build_streaming_seeded(
640                plan,
641                op.inner,
642                storage_ref,
643                params.clone(),
644                outer_row.clone(),
645            )?;
646            let inner_rows = crate::pull::drain(inner_source.as_mut())?;
647            for inner_row in inner_rows {
648                out.push(crate::executor::merge_optional_rows(&outer_row, &inner_row));
649            }
650        }
651        Ok(out)
652    }
653
654    fn exec_path_build(&mut self, plan: &PhysicalPlan, op: &PathBuildExec) -> ExecResult<Vec<Row>> {
655        let input_rows = self.execute_node(plan, op.input)?;
656        let mut rows: Vec<Row> = input_rows
657            .into_iter()
658            .map(|mut row| {
659                let path = build_path_value(&row, &op.node_vars, &op.rel_vars, &*self.ctx.storage);
660                row.insert(op.output, path);
661                row
662            })
663            .collect();
664
665        if let Some(all) = op.shortest_path_all {
666            rows = filter_shortest_paths(rows, op.output, all);
667        }
668        Ok(rows)
669    }
670
671    fn exec_create(&mut self, plan: &PhysicalPlan, op: &CreateExec) -> ExecResult<Vec<Row>> {
672        // Fast path: if the input subtree is fully streamable (no
673        // nested writes, no blocking operators), pull rows one at a
674        // time and apply the create pattern per row, instead of
675        // materializing the whole input. The output Vec still
676        // accumulates — auto-commit-side output streaming is M1.b.
677        if crate::pull::subtree_is_fully_streaming(plan, op.input) {
678            return self.exec_create_streaming_input(plan, op);
679        }
680
681        let input_rows = self.execute_node(plan, op.input)?;
682        let mut out = Vec::with_capacity(input_rows.len());
683
684        for mut row in input_rows {
685            self.apply_create_pattern(&mut row, &op.pattern)?;
686            out.push(row);
687        }
688
689        Ok(out)
690    }
691
692    /// Generic streaming-input loop for write operators whose input
693    /// subtree is fully streamable. Opens a pull-based read cursor
694    /// over the input subtree, calls `apply` per row, and accumulates
695    /// the resulting rows.
696    ///
697    /// # Safety
698    ///
699    /// The upstream [`crate::pull::RowSource`] needs `&S` while it
700    /// lives; the per-row `apply` callback needs `&mut S` (via
701    /// `&mut self`). The existing read-side `RowSource` impls
702    /// materialize their iteration state into owned `Vec`s at
703    /// construction time (see `NodeScanSource::cur_ids`,
704    /// `ExpandSource::cur_edges`, etc. in `pull.rs`), so no live
705    /// `&S` borrow into storage persists across `next_row` calls.
706    /// We exploit that by deriving the read borrow from a raw
707    /// pointer — Rust then doesn't see the shared/mutable conflict
708    /// at compile time, and the dynamic access pattern is
709    /// non-aliasing: read-only inside `next_row`, then mutable
710    /// inside `apply`, never both at the same instant.
711    fn streaming_apply<F>(
712        &mut self,
713        plan: &PhysicalPlan,
714        input: PhysicalNodeId,
715        mut apply: F,
716    ) -> ExecResult<Vec<Row>>
717    where
718        F: FnMut(&mut Self, &mut Row) -> ExecResult<()>,
719    {
720        use std::sync::Arc;
721
722        let storage_ptr: *mut S = self.ctx.storage as *mut S;
723        let params = Arc::new(self.ctx.params.clone());
724
725        // SAFETY: see method-level comment.
726        let storage_ref: &S = unsafe { &*storage_ptr };
727        // Inside a writing `CALL { ... }` body the input's bottom
728        // `Argument` yields the outer row.
729        let mut upstream = match self.argument_seed.clone() {
730            Some(seed) => {
731                crate::pull::build_streaming_seeded(plan, input, storage_ref, params, seed)?
732            }
733            None => crate::pull::build_streaming(plan, input, storage_ref, params)?,
734        };
735
736        let mut out = Vec::new();
737        while let Some(mut row) = upstream.next_row()? {
738            apply(self, &mut row)?;
739            out.push(row);
740        }
741
742        Ok(out)
743    }
744
745    /// Streaming-input variant of [`Self::exec_create`]. Delegates
746    /// to [`Self::streaming_apply`].
747    fn exec_create_streaming_input(
748        &mut self,
749        plan: &PhysicalPlan,
750        op: &CreateExec,
751    ) -> ExecResult<Vec<Row>> {
752        self.streaming_apply(plan, op.input, |this, row| {
753            this.apply_create_pattern(row, &op.pattern)
754        })
755    }
756
757    fn apply_remove_item(&mut self, row: &Row, item: &ResolvedRemoveItem) -> ExecResult<()> {
758        match item {
759            ResolvedRemoveItem::Labels { variable, labels } => match row.get(*variable) {
760                Some(LoraValue::Node(node_id)) => {
761                    let node_id = *node_id;
762                    for label in labels {
763                        self.ctx.storage.remove_node_label(node_id, label);
764                    }
765                    Ok(())
766                }
767                Some(other) => Err(ExecutorError::ExpectedNodeForRemoveLabels {
768                    found: value_kind(other),
769                }),
770                None => Err(ExecutorError::UnboundVariableForRemove {
771                    var: format!("{variable:?}"),
772                }),
773            },
774
775            ResolvedRemoveItem::Property { expr } => self.remove_property_from_expr(row, expr),
776        }
777    }
778
779    fn delete_value(&mut self, value: LoraValue, detach: bool) -> ExecResult<()> {
780        match value {
781            LoraValue::Null => Ok(()),
782
783            LoraValue::Node(node_id) => {
784                if detach {
785                    self.ctx.storage.detach_delete_node(node_id);
786                    Ok(())
787                } else {
788                    let ok = self.ctx.storage.delete_node(node_id);
789                    if ok {
790                        Ok(())
791                    } else {
792                        Err(ExecutorError::DeleteNodeWithRelationships { node_id })
793                    }
794                }
795            }
796
797            LoraValue::Relationship(rel_id) => {
798                let ok = self.ctx.storage.delete_relationship(rel_id);
799                if ok {
800                    Ok(())
801                } else {
802                    Err(ExecutorError::DeleteRelationshipFailed { rel_id })
803                }
804            }
805
806            LoraValue::List(values) => {
807                for v in values {
808                    self.delete_value(v, detach)?;
809                }
810                Ok(())
811            }
812
813            other => Err(ExecutorError::InvalidDeleteTarget {
814                found: value_kind(&other),
815            }),
816        }
817    }
818
819    fn collect_delete_targets(
820        &self,
821        value: &LoraValue,
822        targets: &mut BTreeSet<DeleteTarget>,
823    ) -> ExecResult<()> {
824        match value {
825            LoraValue::Null => Ok(()),
826
827            LoraValue::Node(node_id) => {
828                targets.insert(DeleteTarget::Node(*node_id));
829                Ok(())
830            }
831
832            LoraValue::Relationship(rel_id) => {
833                targets.insert(DeleteTarget::Relationship(*rel_id));
834                Ok(())
835            }
836
837            LoraValue::List(values) => {
838                for v in values {
839                    self.collect_delete_targets(v, targets)?;
840                }
841                Ok(())
842            }
843
844            other => Err(ExecutorError::InvalidDeleteTarget {
845                found: value_kind(other),
846            }),
847        }
848    }
849
850    fn validate_delete_targets(
851        &self,
852        targets: &BTreeSet<DeleteTarget>,
853        detach: bool,
854    ) -> ExecResult<()> {
855        for target in targets {
856            match target {
857                DeleteTarget::Relationship(rel_id) => {
858                    if !self.ctx.storage.contains_relationship(*rel_id) {
859                        return Err(ExecutorError::DeleteRelationshipFailed { rel_id: *rel_id });
860                    }
861                }
862                DeleteTarget::Node(node_id) if !detach => {
863                    if !self.ctx.storage.contains_node(*node_id) {
864                        return Err(ExecutorError::DeleteNodeWithRelationships {
865                            node_id: *node_id,
866                        });
867                    }
868                    let has_external_relationship = self
869                        .ctx
870                        .storage
871                        .relationship_ids_of(*node_id, Direction::Undirected)
872                        .into_iter()
873                        .any(|rel_id| !targets.contains(&DeleteTarget::Relationship(rel_id)));
874                    if has_external_relationship {
875                        return Err(ExecutorError::DeleteNodeWithRelationships {
876                            node_id: *node_id,
877                        });
878                    }
879                }
880                DeleteTarget::Node(_) => {}
881            }
882        }
883        Ok(())
884    }
885
886    fn delete_target(&mut self, target: DeleteTarget, detach: bool) -> ExecResult<()> {
887        match target {
888            DeleteTarget::Node(node_id) => {
889                if detach {
890                    self.ctx.storage.detach_delete_node(node_id);
891                    Ok(())
892                } else {
893                    let ok = self.ctx.storage.delete_node(node_id);
894                    if ok {
895                        Ok(())
896                    } else {
897                        Err(ExecutorError::DeleteNodeWithRelationships { node_id })
898                    }
899                }
900            }
901            DeleteTarget::Relationship(rel_id) => {
902                let ok = self.ctx.storage.delete_relationship(rel_id);
903                if ok {
904                    Ok(())
905                } else {
906                    Err(ExecutorError::DeleteRelationshipFailed { rel_id })
907                }
908            }
909        }
910    }
911
912    fn exec_merge(&mut self, plan: &PhysicalPlan, op: &MergeExec) -> ExecResult<Vec<Row>> {
913        // Streaming-input fast path when the input subtree is fully
914        // streamable. Per-row work (probe → optionally create →
915        // ON MATCH / ON CREATE actions) is identical to the
916        // materialized branch below.
917        if crate::pull::subtree_is_fully_streaming(plan, op.input) {
918            return self.streaming_apply(plan, op.input, |this, row| {
919                let already_bound = this.pattern_part_is_bound(row, &op.pattern_part);
920                let matched = if already_bound {
921                    true
922                } else {
923                    this.try_match_merge_pattern(row, &op.pattern_part)?
924                };
925                if !matched {
926                    this.apply_create_pattern_part(row, &op.pattern_part)?;
927                }
928                for action in &op.actions {
929                    if action.on_match == matched {
930                        for item in &action.set.items {
931                            this.apply_set_item(row, item)?;
932                        }
933                    }
934                }
935                Ok(())
936            });
937        }
938
939        let input_rows = self.execute_node(plan, op.input)?;
940        let mut out = Vec::with_capacity(input_rows.len());
941
942        for mut row in input_rows {
943            // First check if the pattern variable is already bound in the row.
944            let already_bound = self.pattern_part_is_bound(&row, &op.pattern_part);
945
946            let matched = if already_bound {
947                true
948            } else {
949                // Try to find an existing match in the graph.
950                self.try_match_merge_pattern(&mut row, &op.pattern_part)?
951            };
952
953            if !matched {
954                self.apply_create_pattern_part(&mut row, &op.pattern_part)?;
955            }
956
957            for action in &op.actions {
958                if action.on_match == matched {
959                    for item in &action.set.items {
960                        self.apply_set_item(&row, item)?;
961                    }
962                }
963            }
964
965            out.push(row);
966        }
967
968        Ok(out)
969    }
970
971    /// Try to find an existing node/pattern in the graph matching the MERGE
972    /// pattern. If found, bind its variables in the row and return true.
973    /// On a miss the row is left untouched, so the create path sees only
974    /// the variables that were bound before the MERGE.
975    fn try_match_merge_pattern(
976        &self,
977        row: &mut Row,
978        part: &ResolvedPatternPart,
979    ) -> ExecResult<bool> {
980        match &part.element {
981            ResolvedPatternElement::Node {
982                var,
983                labels,
984                properties,
985            } => {
986                let expected_props = self.merge_expected_props(properties.as_ref(), row);
987                let Some(id) = self
988                    .merge_node_candidates(labels, &expected_props)
989                    .into_iter()
990                    .find(|&id| self.merge_node_matches(id, labels, &expected_props))
991                else {
992                    return Ok(false);
993                };
994                if let Some(var_id) = var {
995                    row.insert(*var_id, LoraValue::Node(id));
996                }
997                Ok(true)
998            }
999
1000            ResolvedPatternElement::ShortestPath { .. } => {
1001                // ShortestPath is not valid in MERGE context
1002                Ok(false)
1003            }
1004
1005            ResolvedPatternElement::NodeChain { head, chain } => {
1006                // The head is usually bound by an earlier clause; otherwise
1007                // every node matching it is a possible start.
1008                let head_candidates = match head.var.and_then(|v| row.get(v)) {
1009                    Some(LoraValue::Node(id)) => vec![*id],
1010                    _ => {
1011                        let expected = self.merge_expected_props(head.properties.as_ref(), row);
1012                        self.merge_node_candidates(&head.labels, &expected)
1013                            .into_iter()
1014                            .filter(|&id| self.merge_node_matches(id, &head.labels, &expected))
1015                            .collect()
1016                    }
1017                };
1018
1019                for head_id in head_candidates {
1020                    let mut trial = row.clone();
1021                    if let Some(var_id) = head.var {
1022                        trial.insert(var_id, LoraValue::Node(head_id));
1023                    }
1024                    let mut used_rels = Vec::with_capacity(chain.len());
1025                    if self.match_merge_chain(&mut trial, head_id, chain, &mut used_rels) {
1026                        *row = trial;
1027                        return Ok(true);
1028                    }
1029                }
1030                Ok(false)
1031            }
1032        }
1033    }
1034
1035    /// Match `chain` from `current`, backtracking over every candidate
1036    /// edge. A step node or relationship already bound in the row (by an
1037    /// earlier clause or earlier in the chain) must be the one reached;
1038    /// the same relationship is never used twice in one pattern.
1039    fn match_merge_chain(
1040        &self,
1041        row: &mut Row,
1042        current: NodeId,
1043        chain: &[lora_analyzer::ResolvedChain],
1044        used_rels: &mut Vec<u64>,
1045    ) -> bool {
1046        let Some((step, rest)) = chain.split_first() else {
1047            return true;
1048        };
1049
1050        let bound_dst = match step.node.var.and_then(|v| row.get(v)) {
1051            Some(LoraValue::Node(id)) => Some(*id),
1052            _ => None,
1053        };
1054        let bound_rel = match step.rel.var.and_then(|v| row.get(v)) {
1055            Some(LoraValue::Relationship(id)) => Some(*id),
1056            _ => None,
1057        };
1058        let expected_node = self.merge_expected_props(step.node.properties.as_ref(), row);
1059        let expected_rel = self.merge_expected_props(step.rel.properties.as_ref(), row);
1060
1061        let edges = self
1062            .ctx
1063            .storage
1064            .expand_ids(current, step.rel.direction, &step.rel.types);
1065        for (rel_id, node_id) in edges {
1066            if bound_dst.is_some_and(|id| id != node_id)
1067                || bound_rel.is_some_and(|id| id != rel_id)
1068                || used_rels.contains(&rel_id)
1069            {
1070                continue;
1071            }
1072            if !self.merge_node_matches(node_id, &step.node.labels, &expected_node) {
1073                continue;
1074            }
1075            if let Some(LoraValue::Map(expected_map)) = &expected_rel {
1076                let rel_ok = self
1077                    .ctx
1078                    .storage
1079                    .with_relationship(rel_id, |rel_rec| {
1080                        expected_map.iter().all(|(key, expected_val)| {
1081                            rel_rec
1082                                .properties
1083                                .get(key.as_str())
1084                                .map(|actual| value_matches_property_value(expected_val, actual))
1085                                .unwrap_or(false)
1086                        })
1087                    })
1088                    .unwrap_or(false);
1089                if !rel_ok {
1090                    continue;
1091                }
1092            }
1093
1094            let mut next = row.clone();
1095            if let Some(rel_var) = step.rel.var {
1096                next.insert(rel_var, LoraValue::Relationship(rel_id));
1097            }
1098            if let Some(node_var) = step.node.var {
1099                next.insert(node_var, LoraValue::Node(node_id));
1100            }
1101            used_rels.push(rel_id);
1102            if self.match_merge_chain(&mut next, node_id, rest, used_rels) {
1103                *row = next;
1104                return true;
1105            }
1106            used_rels.pop();
1107        }
1108        false
1109    }
1110
1111    fn merge_expected_props(
1112        &self,
1113        properties: Option<&ResolvedExpr>,
1114        row: &Row,
1115    ) -> Option<LoraValue> {
1116        let eval_ctx = EvalContext {
1117            storage: &*self.ctx.storage,
1118            params: &self.ctx.params,
1119        };
1120        properties.map(|e| eval_expr(e, row, &eval_ctx))
1121    }
1122
1123    /// Candidate ids for a MERGE node pattern. `MERGE (n:L {key: $k})`
1124    /// looks the key up in the property index instead of scanning every
1125    /// `:L` node, so an upsert costs the same on a large label as on a
1126    /// small one. Candidates are re-checked by [`Self::merge_node_matches`].
1127    fn merge_node_candidates(
1128        &self,
1129        labels: &[Vec<String>],
1130        expected_props: &Option<LoraValue>,
1131    ) -> Vec<NodeId> {
1132        let indexed = match expected_props {
1133            Some(LoraValue::Map(expected)) => {
1134                merge_candidates_from_index(&*self.ctx.storage, labels, expected)
1135            }
1136            _ => None,
1137        };
1138        match indexed {
1139            Some(ids) => ids,
1140            None if labels.is_empty() => self.ctx.storage.all_node_ids(),
1141            None => scan_node_ids_for_label_groups(&*self.ctx.storage, labels),
1142        }
1143    }
1144
1145    fn merge_node_matches(
1146        &self,
1147        id: NodeId,
1148        labels: &[Vec<String>],
1149        expected_props: &Option<LoraValue>,
1150    ) -> bool {
1151        self.ctx
1152            .storage
1153            .with_node(id, |node| {
1154                if !node_matches_label_groups(&node.labels, labels) {
1155                    return false;
1156                }
1157                if let Some(LoraValue::Map(expected)) = expected_props {
1158                    return expected.iter().all(|(key, expected_value)| {
1159                        node.properties
1160                            .get(key.as_str())
1161                            .map(|actual| value_matches_property_value(expected_value, actual))
1162                            .unwrap_or(false)
1163                    });
1164                }
1165                true
1166            })
1167            .unwrap_or(false)
1168    }
1169
1170    fn exec_delete(&mut self, plan: &PhysicalPlan, op: &DeleteExec) -> ExecResult<Vec<Row>> {
1171        let input_rows = self.execute_node(plan, op.input)?;
1172        let mut targets = BTreeSet::new();
1173
1174        for row in &input_rows {
1175            for expr in &op.expressions {
1176                let value = {
1177                    let eval_ctx = EvalContext {
1178                        storage: &*self.ctx.storage,
1179                        params: &self.ctx.params,
1180                    };
1181                    eval_expr(expr, row, &eval_ctx)
1182                };
1183                self.collect_delete_targets(&value, &mut targets)?;
1184            }
1185        }
1186
1187        self.validate_delete_targets(&targets, op.detach)?;
1188
1189        for target in &targets {
1190            if let DeleteTarget::Relationship(_) = target {
1191                self.delete_target(*target, op.detach)?;
1192            }
1193        }
1194        for target in targets {
1195            if let DeleteTarget::Node(_) = target {
1196                self.delete_target(target, op.detach)?;
1197            }
1198        }
1199
1200        Ok(input_rows)
1201    }
1202
1203    fn exec_set(&mut self, plan: &PhysicalPlan, op: &SetExec) -> ExecResult<Vec<Row>> {
1204        if crate::pull::subtree_is_fully_streaming(plan, op.input) {
1205            return self.streaming_apply(plan, op.input, |this, row| {
1206                for item in &op.items {
1207                    this.apply_set_item(row, item)?;
1208                }
1209                Ok(())
1210            });
1211        }
1212
1213        let input_rows = self.execute_node(plan, op.input)?;
1214
1215        for row in &input_rows {
1216            for item in &op.items {
1217                self.apply_set_item(row, item)?;
1218            }
1219        }
1220
1221        Ok(input_rows)
1222    }
1223
1224    /// `FOREACH (var IN list | body...)` — for each input row, evaluate
1225    /// the list and run the body once per element with `var` bound to
1226    /// that element. Each iteration runs on a fresh clone of the row
1227    /// so any new bindings the body introduces (e.g. anonymous
1228    /// `CREATE` node VarIds) don't leak between iterations or back to
1229    /// the outer scope. Side effects on the graph persist; the outer
1230    /// row is emitted unchanged.
1231    fn exec_foreach(&mut self, plan: &PhysicalPlan, op: &ForeachExec) -> ExecResult<Vec<Row>> {
1232        let input_rows = self.execute_node(plan, op.input)?;
1233        let mut out = Vec::with_capacity(input_rows.len());
1234
1235        for row in input_rows {
1236            let list_value = {
1237                let eval_ctx = EvalContext {
1238                    storage: &*self.ctx.storage,
1239                    params: &self.ctx.params,
1240                };
1241                eval_expr(&op.list, &row, &eval_ctx)
1242            };
1243
1244            let elements: Vec<LoraValue> = match list_value {
1245                LoraValue::List(items) => items,
1246                LoraValue::Null => Vec::new(),
1247                other => {
1248                    return Err(ExecutorError::RuntimeError(format!(
1249                        "FOREACH expects a list, got {}",
1250                        value_kind(&other)
1251                    )));
1252                }
1253            };
1254
1255            for element in elements {
1256                // Fresh row per iteration so body-introduced bindings
1257                // don't reuse VarIds across iterations.
1258                let mut iter_row = row.clone();
1259                iter_row.insert(op.variable, element);
1260                for clause in &op.body {
1261                    self.apply_foreach_body_clause(&mut iter_row, clause)?;
1262                }
1263            }
1264
1265            out.push(row);
1266        }
1267
1268        Ok(out)
1269    }
1270
1271    /// Apply one resolved updating clause to `row` for its side effect
1272    /// inside a `FOREACH` body. Only updating clauses (Create / Merge /
1273    /// Delete / Set / Remove / nested Foreach) are legal here; the
1274    /// analyzer guarantees that.
1275    fn apply_foreach_body_clause(
1276        &mut self,
1277        row: &mut Row,
1278        clause: &lora_analyzer::ResolvedClause,
1279    ) -> ExecResult<()> {
1280        use lora_analyzer::ResolvedClause;
1281        match clause {
1282            ResolvedClause::Create(c) => self.apply_create_pattern(row, &c.pattern),
1283            ResolvedClause::Set(s) => {
1284                for item in &s.items {
1285                    self.apply_set_item(row, item)?;
1286                }
1287                Ok(())
1288            }
1289            ResolvedClause::Remove(r) => {
1290                for item in &r.items {
1291                    self.apply_remove_item(row, item)?;
1292                }
1293                Ok(())
1294            }
1295            ResolvedClause::Delete(d) => {
1296                let detach = d.detach;
1297                for expr in &d.expressions {
1298                    let value = {
1299                        let eval_ctx = EvalContext {
1300                            storage: &*self.ctx.storage,
1301                            params: &self.ctx.params,
1302                        };
1303                        eval_expr(expr, row, &eval_ctx)
1304                    };
1305                    self.delete_value(value, detach)?;
1306                }
1307                Ok(())
1308            }
1309            ResolvedClause::Merge(m) => {
1310                let already_bound = self.pattern_part_is_bound(row, &m.pattern_part);
1311                let matched = if already_bound {
1312                    true
1313                } else {
1314                    self.try_match_merge_pattern(row, &m.pattern_part)?
1315                };
1316                if !matched {
1317                    self.apply_create_pattern_part(row, &m.pattern_part)?;
1318                }
1319                for action in &m.actions {
1320                    if action.on_match == matched {
1321                        for item in &action.set.items {
1322                            self.apply_set_item(row, item)?;
1323                        }
1324                    }
1325                }
1326                Ok(())
1327            }
1328            ResolvedClause::Foreach(nested) => {
1329                let list_value = {
1330                    let eval_ctx = EvalContext {
1331                        storage: &*self.ctx.storage,
1332                        params: &self.ctx.params,
1333                    };
1334                    eval_expr(&nested.list, row, &eval_ctx)
1335                };
1336
1337                let elements: Vec<LoraValue> = match list_value {
1338                    LoraValue::List(items) => items,
1339                    LoraValue::Null => Vec::new(),
1340                    other => {
1341                        return Err(ExecutorError::RuntimeError(format!(
1342                            "FOREACH expects a list, got {}",
1343                            value_kind(&other)
1344                        )));
1345                    }
1346                };
1347
1348                for element in elements {
1349                    let mut iter_row = row.clone();
1350                    iter_row.insert(nested.variable, element);
1351                    for inner in &nested.body {
1352                        self.apply_foreach_body_clause(&mut iter_row, inner)?;
1353                    }
1354                }
1355
1356                Ok(())
1357            }
1358            other => Err(ExecutorError::RuntimeError(format!(
1359                "FOREACH body may only contain updating clauses, got {:?}",
1360                std::mem::discriminant(other)
1361            ))),
1362        }
1363    }
1364
1365    fn exec_remove(&mut self, plan: &PhysicalPlan, op: &RemoveExec) -> ExecResult<Vec<Row>> {
1366        if crate::pull::subtree_is_fully_streaming(plan, op.input) {
1367            return self.streaming_apply(plan, op.input, |this, row| {
1368                for item in &op.items {
1369                    this.apply_remove_item(row, item)?;
1370                }
1371                Ok(())
1372            });
1373        }
1374
1375        let input_rows = self.execute_node(plan, op.input)?;
1376
1377        for row in &input_rows {
1378            for item in &op.items {
1379                self.apply_remove_item(row, item)?;
1380            }
1381        }
1382
1383        Ok(input_rows)
1384    }
1385
1386    fn apply_set_item(&mut self, row: &Row, item: &ResolvedSetItem) -> ExecResult<()> {
1387        match item {
1388            ResolvedSetItem::SetProperty { target, value } => {
1389                let new_value = {
1390                    let eval_ctx = EvalContext {
1391                        storage: &*self.ctx.storage,
1392                        params: &self.ctx.params,
1393                    };
1394                    eval_expr(value, row, &eval_ctx)
1395                };
1396
1397                self.set_property_from_expr(row, target, new_value)
1398            }
1399
1400            ResolvedSetItem::SetVariable { variable, value } => {
1401                // Only need the entity's id — peek at the binding by reference.
1402                let entity_ref =
1403                    row.get(*variable)
1404                        .ok_or(ExecutorError::UnboundVariableForSet {
1405                            var: format!("{variable:?}"),
1406                        })?;
1407                let entity_target = entity_target_from_value(entity_ref)?;
1408
1409                let new_value = {
1410                    let eval_ctx = EvalContext {
1411                        storage: &*self.ctx.storage,
1412                        params: &self.ctx.params,
1413                    };
1414                    eval_expr(value, row, &eval_ctx)
1415                };
1416
1417                self.overwrite_entity_target(entity_target, new_value)
1418            }
1419
1420            ResolvedSetItem::MutateVariable { variable, value } => {
1421                let entity_ref =
1422                    row.get(*variable)
1423                        .ok_or(ExecutorError::UnboundVariableForSet {
1424                            var: format!("{variable:?}"),
1425                        })?;
1426                let entity_target = entity_target_from_value(entity_ref)?;
1427
1428                let patch = {
1429                    let eval_ctx = EvalContext {
1430                        storage: &*self.ctx.storage,
1431                        params: &self.ctx.params,
1432                    };
1433                    eval_expr(value, row, &eval_ctx)
1434                };
1435
1436                self.mutate_entity_target(entity_target, patch)
1437            }
1438
1439            ResolvedSetItem::SetLabels { variable, labels } => match row.get(*variable) {
1440                Some(LoraValue::Node(node_id)) => {
1441                    let node_id = *node_id;
1442                    for label in labels {
1443                        if let Err(msg) = self
1444                            .ctx
1445                            .storage
1446                            .check_node_add_label_against_constraints(node_id, label)
1447                        {
1448                            return Err(ExecutorError::ConstraintViolation(msg));
1449                        }
1450                        self.ctx.storage.add_node_label(node_id, label);
1451                    }
1452                    Ok(())
1453                }
1454                Some(other) => Err(ExecutorError::ExpectedNodeForSetLabels {
1455                    found: value_kind(other),
1456                }),
1457                None => Err(ExecutorError::UnboundVariableForSet {
1458                    var: format!("{variable:?}"),
1459                }),
1460            },
1461        }
1462    }
1463
1464    fn set_property_from_expr(
1465        &mut self,
1466        row: &Row,
1467        target_expr: &ResolvedExpr,
1468        new_value: LoraValue,
1469    ) -> ExecResult<()> {
1470        let ResolvedExpr::Property { expr, property } = target_expr else {
1471            return Err(ExecutorError::UnsupportedSetTarget);
1472        };
1473
1474        let owner = {
1475            let eval_ctx = EvalContext {
1476                storage: &*self.ctx.storage,
1477                params: &self.ctx.params,
1478            };
1479            eval_expr(expr, row, &eval_ctx)
1480        };
1481
1482        // `SET n.a = null` removes the property.
1483        if matches!(new_value, LoraValue::Null) {
1484            return match owner {
1485                LoraValue::Node(node_id) => {
1486                    self.remove_entity_property(EntityTarget::Node(node_id), property)
1487                }
1488                LoraValue::Relationship(rel_id) => {
1489                    self.remove_entity_property(EntityTarget::Relationship(rel_id), property)
1490                }
1491                other => Err(ExecutorError::InvalidSetTarget {
1492                    found: value_kind(&other),
1493                }),
1494            };
1495        }
1496
1497        match owner {
1498            LoraValue::Node(node_id) => {
1499                let prop = lora_value_to_property(new_value)
1500                    .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1501                if let Err(msg) = self
1502                    .ctx
1503                    .storage
1504                    .check_node_set_property_against_constraints(node_id, property, &prop)
1505                {
1506                    return Err(ExecutorError::ConstraintViolation(msg));
1507                }
1508                self.ctx
1509                    .storage
1510                    .set_node_property(node_id, property.clone(), prop);
1511                Ok(())
1512            }
1513            LoraValue::Relationship(rel_id) => {
1514                let prop = lora_value_to_property(new_value)
1515                    .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1516                if let Err(msg) = self
1517                    .ctx
1518                    .storage
1519                    .check_relationship_set_property_against_constraints(rel_id, property, &prop)
1520                {
1521                    return Err(ExecutorError::ConstraintViolation(msg));
1522                }
1523                self.ctx
1524                    .storage
1525                    .set_relationship_property(rel_id, property.clone(), prop);
1526                Ok(())
1527            }
1528            other => Err(ExecutorError::InvalidSetTarget {
1529                found: value_kind(&other),
1530            }),
1531        }
1532    }
1533
1534    /// Remove one property, checking constraints first. Removing a
1535    /// property the entity does not have is a no-op.
1536    fn remove_entity_property(&mut self, target: EntityTarget, property: &str) -> ExecResult<()> {
1537        match target {
1538            EntityTarget::Node(node_id) => {
1539                if let Err(msg) = self
1540                    .ctx
1541                    .storage
1542                    .check_node_remove_property_against_constraints(node_id, property)
1543                {
1544                    return Err(ExecutorError::ConstraintViolation(msg));
1545                }
1546                self.ctx.storage.remove_node_property(node_id, property);
1547            }
1548            EntityTarget::Relationship(rel_id) => {
1549                if let Err(msg) = self
1550                    .ctx
1551                    .storage
1552                    .check_relationship_remove_property_against_constraints(rel_id, property)
1553                {
1554                    return Err(ExecutorError::ConstraintViolation(msg));
1555                }
1556                self.ctx
1557                    .storage
1558                    .remove_relationship_property(rel_id, property);
1559            }
1560        }
1561        Ok(())
1562    }
1563
1564    fn remove_property_from_expr(&mut self, row: &Row, expr: &ResolvedExpr) -> ExecResult<()> {
1565        let ResolvedExpr::Property {
1566            expr: owner_expr,
1567            property,
1568        } = expr
1569        else {
1570            return Err(ExecutorError::UnsupportedRemoveTarget);
1571        };
1572
1573        let owner = {
1574            let eval_ctx = EvalContext {
1575                storage: &*self.ctx.storage,
1576                params: &self.ctx.params,
1577            };
1578            eval_expr(owner_expr, row, &eval_ctx)
1579        };
1580
1581        match owner {
1582            LoraValue::Node(node_id) => {
1583                if let Err(msg) = self
1584                    .ctx
1585                    .storage
1586                    .check_node_remove_property_against_constraints(node_id, property)
1587                {
1588                    return Err(ExecutorError::ConstraintViolation(msg));
1589                }
1590                self.ctx.storage.remove_node_property(node_id, property);
1591                Ok(())
1592            }
1593            LoraValue::Relationship(rel_id) => {
1594                if let Err(msg) = self
1595                    .ctx
1596                    .storage
1597                    .check_relationship_remove_property_against_constraints(rel_id, property)
1598                {
1599                    return Err(ExecutorError::ConstraintViolation(msg));
1600                }
1601                self.ctx
1602                    .storage
1603                    .remove_relationship_property(rel_id, property);
1604                Ok(())
1605            }
1606            other => Err(ExecutorError::InvalidRemoveTarget {
1607                found: value_kind(&other),
1608            }),
1609        }
1610    }
1611
1612    fn overwrite_entity_target(
1613        &mut self,
1614        target: EntityTarget,
1615        new_value: LoraValue,
1616    ) -> ExecResult<()> {
1617        let LoraValue::Map(map) = new_value else {
1618            return Err(ExecutorError::ExpectedPropertyMap {
1619                found: value_kind(&new_value),
1620            });
1621        };
1622
1623        let mut props: Properties = Properties::new();
1624        for (k, v) in map {
1625            // `SET n = {a: null}` leaves `a` absent.
1626            if matches!(v, LoraValue::Null) {
1627                continue;
1628            }
1629            let prop = lora_value_to_property(v)
1630                .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1631            props.insert(lora_store::intern_owned(k), prop);
1632        }
1633
1634        match target {
1635            EntityTarget::Node(node_id) => {
1636                if let Err(msg) = self
1637                    .ctx
1638                    .storage
1639                    .check_node_replace_properties_against_constraints(node_id, &props)
1640                {
1641                    return Err(ExecutorError::ConstraintViolation(msg));
1642                }
1643                self.ctx.storage.replace_node_properties(node_id, props);
1644            }
1645            EntityTarget::Relationship(rel_id) => {
1646                if let Err(msg) = self
1647                    .ctx
1648                    .storage
1649                    .check_relationship_replace_properties_against_constraints(rel_id, &props)
1650                {
1651                    return Err(ExecutorError::ConstraintViolation(msg));
1652                }
1653                self.ctx
1654                    .storage
1655                    .replace_relationship_properties(rel_id, props);
1656            }
1657        }
1658        Ok(())
1659    }
1660
1661    fn mutate_entity_target(
1662        &mut self,
1663        target: EntityTarget,
1664        patch_value: LoraValue,
1665    ) -> ExecResult<()> {
1666        let LoraValue::Map(map) = patch_value else {
1667            return Err(ExecutorError::ExpectedPropertyMap {
1668                found: value_kind(&patch_value),
1669            });
1670        };
1671
1672        match target {
1673            EntityTarget::Node(node_id) => {
1674                for (k, v) in map {
1675                    // `SET n += {a: null}` removes `a`.
1676                    if matches!(v, LoraValue::Null) {
1677                        self.remove_entity_property(target, &k)?;
1678                        continue;
1679                    }
1680                    let prop = lora_value_to_property(v)
1681                        .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1682                    if let Err(msg) = self
1683                        .ctx
1684                        .storage
1685                        .check_node_set_property_against_constraints(node_id, &k, &prop)
1686                    {
1687                        return Err(ExecutorError::ConstraintViolation(msg));
1688                    }
1689                    self.ctx.storage.set_node_property(node_id, k, prop);
1690                }
1691            }
1692            EntityTarget::Relationship(rel_id) => {
1693                for (k, v) in map {
1694                    if matches!(v, LoraValue::Null) {
1695                        self.remove_entity_property(target, &k)?;
1696                        continue;
1697                    }
1698                    let prop = lora_value_to_property(v)
1699                        .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1700                    if let Err(msg) = self
1701                        .ctx
1702                        .storage
1703                        .check_relationship_set_property_against_constraints(rel_id, &k, &prop)
1704                    {
1705                        return Err(ExecutorError::ConstraintViolation(msg));
1706                    }
1707                    self.ctx.storage.set_relationship_property(rel_id, k, prop);
1708                }
1709            }
1710        }
1711        Ok(())
1712    }
1713
1714    pub(crate) fn apply_create_pattern(
1715        &mut self,
1716        row: &mut Row,
1717        pattern: &ResolvedPattern,
1718    ) -> ExecResult<()> {
1719        for part in &pattern.parts {
1720            self.apply_create_pattern_part(row, part)?;
1721        }
1722        Ok(())
1723    }
1724
1725    /// Apply a single per-row write for any of the streamable write
1726    /// operators (Create / Set / Delete / Remove / Merge). Used by
1727    /// the [`crate::pull::StreamingWriteCursor`] auto-commit fast
1728    /// path: the cursor pulls one input row from a read upstream,
1729    /// hands it here for the side effect, and emits the row back.
1730    pub(crate) fn apply_write_op(&mut self, op: &PhysicalOp, row: &mut Row) -> ExecResult<()> {
1731        match op {
1732            PhysicalOp::Create(c) => self.apply_create_pattern(row, &c.pattern),
1733            PhysicalOp::Set(s) => {
1734                for item in &s.items {
1735                    self.apply_set_item(row, item)?;
1736                }
1737                Ok(())
1738            }
1739            PhysicalOp::Delete(d) => {
1740                let detach = d.detach;
1741                for expr in &d.expressions {
1742                    let value = {
1743                        let eval_ctx = EvalContext {
1744                            storage: &*self.ctx.storage,
1745                            params: &self.ctx.params,
1746                        };
1747                        eval_expr(expr, row, &eval_ctx)
1748                    };
1749                    self.delete_value(value, detach)?;
1750                }
1751                Ok(())
1752            }
1753            PhysicalOp::Remove(r) => {
1754                for item in &r.items {
1755                    self.apply_remove_item(row, item)?;
1756                }
1757                Ok(())
1758            }
1759            PhysicalOp::Merge(m) => {
1760                let already_bound = self.pattern_part_is_bound(row, &m.pattern_part);
1761                let matched = if already_bound {
1762                    true
1763                } else {
1764                    self.try_match_merge_pattern(row, &m.pattern_part)?
1765                };
1766                if !matched {
1767                    self.apply_create_pattern_part(row, &m.pattern_part)?;
1768                }
1769                for action in &m.actions {
1770                    if action.on_match == matched {
1771                        for item in &action.set.items {
1772                            self.apply_set_item(row, item)?;
1773                        }
1774                    }
1775                }
1776                Ok(())
1777            }
1778            other => Err(ExecutorError::RuntimeError(format!(
1779                "apply_write_op called on non-write op: {other:?}"
1780            ))),
1781        }
1782    }
1783
1784    fn apply_create_pattern_part(
1785        &mut self,
1786        row: &mut Row,
1787        part: &ResolvedPatternPart,
1788    ) -> ExecResult<()> {
1789        if part.binding.is_some() {
1790            trace!("create pattern part has path binding; path materialization not implemented");
1791        }
1792
1793        let _ = self.apply_create_pattern_element(row, &part.element)?;
1794        Ok(())
1795    }
1796
1797    fn apply_create_pattern_element(
1798        &mut self,
1799        row: &mut Row,
1800        element: &ResolvedPatternElement,
1801    ) -> ExecResult<Option<LoraValue>> {
1802        match element {
1803            ResolvedPatternElement::Node {
1804                var,
1805                labels,
1806                properties,
1807            } => {
1808                let node_id =
1809                    self.materialize_node_pattern(row, *var, labels, properties.as_ref())?;
1810                Ok(Some(LoraValue::Node(node_id)))
1811            }
1812
1813            ResolvedPatternElement::NodeChain { head, chain } => {
1814                let mut current_node_id = self.materialize_node_pattern(
1815                    row,
1816                    head.var,
1817                    &head.labels,
1818                    head.properties.as_ref(),
1819                )?;
1820
1821                for link in chain {
1822                    let next_node_id = self.materialize_node_pattern(
1823                        row,
1824                        link.node.var,
1825                        &link.node.labels,
1826                        link.node.properties.as_ref(),
1827                    )?;
1828
1829                    let _ = self.materialize_relationship_pattern(
1830                        row,
1831                        current_node_id,
1832                        next_node_id,
1833                        &link.rel,
1834                    )?;
1835
1836                    current_node_id = next_node_id;
1837                }
1838
1839                Ok(Some(LoraValue::Node(current_node_id)))
1840            }
1841
1842            ResolvedPatternElement::ShortestPath { .. } => {
1843                // ShortestPath is not valid in CREATE context
1844                Ok(None)
1845            }
1846        }
1847    }
1848
1849    fn pattern_part_is_bound(&self, row: &Row, part: &ResolvedPatternPart) -> bool {
1850        match &part.element {
1851            ResolvedPatternElement::Node { var, .. } => var.and_then(|v| row.get(v)).is_some(),
1852
1853            ResolvedPatternElement::ShortestPath { .. } => false,
1854
1855            ResolvedPatternElement::NodeChain { head, chain } => {
1856                let head_ok = head.var.and_then(|v| row.get(v)).is_some();
1857
1858                let chain_ok = chain.iter().all(|link| {
1859                    let node_ok = link.node.var.and_then(|v| row.get(v)).is_some();
1860                    // For MERGE, anonymous relationships cannot be considered
1861                    // "bound" because we have no variable to check. The merge
1862                    // must search the graph to see if the relationship exists.
1863                    let rel_ok = match link.rel.var {
1864                        Some(v) => row.get(v).is_some(),
1865                        None => false,
1866                    };
1867                    node_ok && rel_ok
1868                });
1869
1870                head_ok && chain_ok
1871            }
1872        }
1873    }
1874
1875    fn materialize_node_pattern(
1876        &mut self,
1877        row: &mut Row,
1878        var: Option<VarId>,
1879        labels: &[Vec<String>],
1880        properties: Option<&ResolvedExpr>,
1881    ) -> ExecResult<u64> {
1882        if let Some(var_id) = var {
1883            if let Some(LoraValue::Node(id)) = row.get(var_id) {
1884                return Ok(*id);
1885            }
1886        }
1887
1888        let properties = match properties {
1889            Some(expr) => eval_properties_expr(expr, row, &*self.ctx.storage, &self.ctx.params)?,
1890            None => Properties::new(),
1891        };
1892
1893        let flat_labels = flatten_label_groups(labels);
1894        debug!("creating node with labels={flat_labels:?}");
1895        let checked = if self.defer_existence {
1896            self.ctx
1897                .storage
1898                .check_node_create_deferring_existence(&flat_labels, &properties)
1899        } else {
1900            self.ctx
1901                .storage
1902                .check_node_create_against_constraints(&flat_labels, &properties)
1903        };
1904        checked.map_err(ExecutorError::ConstraintViolation)?;
1905        let created = self
1906            .ctx
1907            .storage
1908            .try_create_node(flat_labels, properties)
1909            .ok_or(ExecutorError::NodeCreateFailed)?;
1910        if self.defer_existence {
1911            self.pending_existence.push(EntityTarget::Node(created.id));
1912        }
1913
1914        if let Some(var_id) = var {
1915            row.insert(var_id, LoraValue::Node(created.id));
1916        }
1917
1918        Ok(created.id)
1919    }
1920
1921    fn materialize_relationship_pattern(
1922        &mut self,
1923        row: &mut Row,
1924        left_node_id: u64,
1925        right_node_id: u64,
1926        rel: &lora_analyzer::ResolvedRel,
1927    ) -> ExecResult<u64> {
1928        if let Some(var_id) = rel.var {
1929            if let Some(LoraValue::Relationship(id)) = row.get(var_id) {
1930                let id = *id;
1931                if let Some((src, dst)) = self.ctx.storage.relationship_endpoints(id) {
1932                    let endpoints_match = match rel.direction {
1933                        Direction::Right | Direction::Undirected => {
1934                            src == left_node_id && dst == right_node_id
1935                        }
1936                        Direction::Left => src == right_node_id && dst == left_node_id,
1937                    };
1938
1939                    if endpoints_match {
1940                        return Ok(id);
1941                    }
1942                }
1943            }
1944        }
1945
1946        if rel.range.is_some() {
1947            return Err(ExecutorError::UnsupportedCreateRelationshipRange);
1948        }
1949
1950        let (src, dst) = match rel.direction {
1951            Direction::Right | Direction::Undirected => (left_node_id, right_node_id),
1952            Direction::Left => (right_node_id, left_node_id),
1953        };
1954
1955        let rel_type = rel
1956            .types
1957            .first()
1958            .ok_or(ExecutorError::MissingRelationshipType)?;
1959
1960        if rel_type.is_empty() {
1961            return Err(ExecutorError::MissingRelationshipType);
1962        }
1963
1964        let properties = match rel.properties.as_ref() {
1965            Some(expr) => eval_properties_expr(expr, row, &*self.ctx.storage, &self.ctx.params)?,
1966            None => Properties::new(),
1967        };
1968
1969        debug!("creating relationship: src={src}, dst={dst}, type={rel_type}");
1970
1971        let checked = if self.defer_existence {
1972            self.ctx
1973                .storage
1974                .check_relationship_create_deferring_existence(rel_type, &properties)
1975        } else {
1976            self.ctx
1977                .storage
1978                .check_relationship_create_against_constraints(rel_type, &properties)
1979        };
1980        checked.map_err(ExecutorError::ConstraintViolation)?;
1981
1982        let created = self
1983            .ctx
1984            .storage
1985            .create_relationship(src, dst, rel_type, properties)
1986            .ok_or_else(|| ExecutorError::RelationshipCreateFailed {
1987                src,
1988                dst,
1989                rel_type: rel_type.clone(),
1990            })?;
1991        if self.defer_existence {
1992            self.pending_existence
1993                .push(EntityTarget::Relationship(created.id));
1994        }
1995
1996        if let Some(var_id) = rel.var {
1997            row.insert(var_id, LoraValue::Relationship(created.id));
1998        }
1999
2000        Ok(created.id)
2001    }
2002}
2003
2004/// Whether existence constraints on entities a plan creates must wait
2005/// for the end of the statement. They can be checked at `CREATE` only
2006/// when nothing after it can add a property: every write is a `CREATE`
2007/// or a `DELETE`, with at most one `CREATE`. Checking early keeps a
2008/// failing create from mutating anything, which the in-place write path
2009/// relies on.
2010pub(crate) fn plan_defers_existence(plan: &PhysicalPlan) -> bool {
2011    let mut creates = 0;
2012    for op in &plan.nodes {
2013        match op {
2014            PhysicalOp::Create(_) => creates += 1,
2015            PhysicalOp::Delete(_) => {}
2016            PhysicalOp::Merge(_)
2017            | PhysicalOp::Set(_)
2018            | PhysicalOp::Remove(_)
2019            | PhysicalOp::Foreach(_) => return true,
2020            _ => {}
2021        }
2022    }
2023    creates > 1
2024}
2025
2026/// Whether a plan is a write statement with no `RETURN` (its root is the
2027/// write operator itself). Such a statement produces no result rows, as
2028/// in other Cypher databases; the write operator's pass-through rows
2029/// would otherwise leak as anonymous `_0` columns carrying internal ids.
2030pub(crate) fn plan_ends_in_write(plan: &PhysicalPlan) -> bool {
2031    match &plan.nodes[plan.root] {
2032        PhysicalOp::Create(_)
2033        | PhysicalOp::Merge(_)
2034        | PhysicalOp::Set(_)
2035        | PhysicalOp::Delete(_)
2036        | PhysicalOp::Remove(_)
2037        | PhysicalOp::Foreach(_) => true,
2038        // A query ending in a unit `CALL { ... }` returns no rows.
2039        PhysicalOp::CallSubquery(op) => op.new_vars.is_empty(),
2040        _ => false,
2041    }
2042}
2043
2044/// Candidate nodes for a MERGE node pattern from the property index, or
2045/// `None` to fall back to a label scan. Every candidate is still checked
2046/// against the full pattern, so the only requirement is that no real
2047/// match is missed. MERGE compares `1` and `1.0` as equal while the index
2048/// keys them apart, so numbers look up both images; values without an
2049/// exact index image (lists, maps, NaN, floats beyond 2^53) scan.
2050fn merge_candidates_from_index<S: lora_store::GraphStorage>(
2051    storage: &S,
2052    labels: &[Vec<String>],
2053    expected: &std::collections::BTreeMap<String, LoraValue>,
2054) -> Option<Vec<lora_store::NodeId>> {
2055    use lora_store::PropertyValue;
2056
2057    // A single required label scopes the lookup; otherwise look up
2058    // across labels and let the pattern check filter.
2059    let label = match labels {
2060        [group] if group.len() == 1 => Some(group[0].as_str()),
2061        _ => None,
2062    };
2063    let (key, value) = expected.iter().find(|(_, v)| {
2064        matches!(
2065            v,
2066            LoraValue::String(_) | LoraValue::Bool(_) | LoraValue::Int(_)
2067        ) || matches!(v, LoraValue::Float(f) if f.is_finite() && f.abs() < 9_007_199_254_740_992.0)
2068    })?;
2069    let images: Vec<PropertyValue> = match value {
2070        LoraValue::String(s) => vec![PropertyValue::String(s.clone())],
2071        LoraValue::Bool(b) => vec![PropertyValue::Bool(*b)],
2072        LoraValue::Int(i) => vec![PropertyValue::Int(*i), PropertyValue::Float(*i as f64)],
2073        LoraValue::Float(f) => {
2074            let mut v = vec![PropertyValue::Float(*f)];
2075            if f.fract() == 0.0 {
2076                v.push(PropertyValue::Int(*f as i64));
2077            }
2078            v
2079        }
2080        _ => return None,
2081    };
2082    let mut ids: Vec<lora_store::NodeId> = images
2083        .iter()
2084        .flat_map(|image| storage.find_node_ids_by_property(label, key, image))
2085        .collect();
2086    ids.sort_unstable();
2087    ids.dedup();
2088    Some(ids)
2089}