1use 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#[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 argument_seed: Option<Row>,
79 defer_existence: bool,
84 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 pub(crate) fn defer_existence_checks(&mut self, defer: bool) {
139 self.defer_existence = defer;
140 }
141
142 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_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 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 let mut all_rows = self.execute_and_hydrate(&compiled.physical)?;
207
208 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 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 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 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 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 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 let storage_ref: &S = unsafe { &*storage_ptr };
726 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 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 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 let already_bound = self.pattern_part_is_bound(&row, &op.pattern_part);
944
945 let matched = if already_bound {
946 true
947 } else {
948 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 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 Ok(false)
1002 }
1003
1004 ResolvedPatternElement::NodeChain { head, chain } => {
1005 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 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 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 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 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 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 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 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 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 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 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 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 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 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
2003pub(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
2025pub(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 PhysicalOp::CallSubquery(op) => op.new_vars.is_empty(),
2039 _ => false,
2040 }
2041}
2042
2043fn 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 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}