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