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}
77
78impl<'a, S: GraphStorageMut> MutableExecutor<'a, S> {
79 pub fn new(ctx: MutableExecutionContext<'a, S>) -> Self {
80 Self {
81 ctx,
82 deadline: None,
83 }
84 }
85
86 pub fn with_deadline(ctx: MutableExecutionContext<'a, S>, deadline: Option<Instant>) -> Self {
87 Self { ctx, deadline }
88 }
89
90 #[inline]
91 fn check_deadline(&self) -> ExecResult<()> {
92 if let Some(deadline) = self.deadline {
93 check_deadline_at(deadline)
94 } else {
95 Ok(())
96 }
97 }
98
99 pub fn execute(
100 &mut self,
101 plan: &PhysicalPlan,
102 options: Option<ExecuteOptions>,
103 ) -> ExecResult<QueryResult> {
104 let rows = self.execute_rows(plan)?;
105 Ok(project_rows(rows, options.unwrap_or_default()))
106 }
107
108 pub fn execute_rows(&mut self, plan: &PhysicalPlan) -> ExecResult<Vec<Row>> {
109 self.check_deadline()?;
110 clear_eval_error();
113
114 let rows = self.execute_node(plan, plan.root)?;
115 if !plan_may_need_hydration(plan) {
116 return Ok(rows);
117 }
118 Ok(rows
119 .into_iter()
120 .map(|row| self.hydrate_row(row))
121 .collect::<Vec<_>>())
122 }
123
124 pub fn execute_compiled(
126 &mut self,
127 compiled: &CompiledQuery,
128 options: Option<ExecuteOptions>,
129 ) -> ExecResult<QueryResult> {
130 let rows = self.execute_compiled_rows(compiled)?;
131 Ok(project_rows(rows, options.unwrap_or_default()))
132 }
133
134 pub fn execute_compiled_rows(&mut self, compiled: &CompiledQuery) -> ExecResult<Vec<Row>> {
135 self.check_deadline()?;
136 if compiled.unions.is_empty() {
137 return self.execute_rows(&compiled.physical);
138 }
139
140 clear_eval_error();
141
142 let mut all_rows = self.execute_and_hydrate(&compiled.physical)?;
144
145 let mut needs_dedup = false;
148
149 for branch in &compiled.unions {
150 self.check_deadline()?;
151 let branch_rows = self.execute_and_hydrate(&branch.physical)?;
152 all_rows.extend(branch_rows);
153
154 if !branch.all {
155 needs_dedup = true;
156 }
157 }
158
159 if needs_dedup {
160 all_rows = dedup_rows(all_rows);
161 }
162
163 Ok(all_rows)
164 }
165
166 fn execute_and_hydrate(&mut self, plan: &PhysicalPlan) -> ExecResult<Vec<Row>> {
167 self.check_deadline()?;
168 let rows = self.execute_node(plan, plan.root)?;
169 if !plan_may_need_hydration(plan) {
170 return Ok(rows);
171 }
172 Ok(rows.into_iter().map(|row| self.hydrate_row(row)).collect())
173 }
174
175 pub(crate) fn hydrate_row(&self, row: Row) -> Row {
176 let mut out = Row::new();
177
178 for (var, name, value) in row.into_iter_named() {
179 out.insert_named(var, name, self.hydrate_value(value));
180 }
181
182 out
183 }
184
185 fn execute_node(
186 &mut self,
187 plan: &PhysicalPlan,
188 node_id: PhysicalNodeId,
189 ) -> ExecResult<Vec<Row>> {
190 self.check_deadline()?;
191 trace!("mutable execute_node start: node_id={node_id:?}");
192
193 let result = match &plan.nodes[node_id] {
194 PhysicalOp::Argument(op) => self.exec_argument(op),
195 PhysicalOp::NodeScan(op) => self.exec_node_scan(plan, op),
196 PhysicalOp::NodeByLabelScan(op) => self.exec_node_by_label_scan(plan, op),
197 PhysicalOp::NodeByPropertyScan(op) => self.exec_node_by_property_scan(plan, op),
198 PhysicalOp::NodeByPropertyRangeScan(op) => {
199 self.exec_node_by_property_range_scan(plan, op)
200 }
201 PhysicalOp::NodeByTextScan(op) => self.exec_node_by_text_scan(plan, op),
202 PhysicalOp::NodeByPointScan(op) => self.exec_node_by_point_scan(plan, op),
203 PhysicalOp::RelByPropertyRangeScan(op) => {
204 self.exec_rel_by_property_range_scan(plan, op)
205 }
206 PhysicalOp::RelByTextScan(op) => self.exec_rel_by_text_scan(plan, op),
207 PhysicalOp::RelByPointScan(op) => self.exec_rel_by_point_scan(plan, op),
208 PhysicalOp::Expand(op) => self.exec_expand(plan, op),
209 PhysicalOp::Filter(op) => self.exec_filter(plan, op),
210 PhysicalOp::Projection(op) => self.exec_projection(plan, op),
211 PhysicalOp::Unwind(op) => self.exec_unwind(plan, op),
212 PhysicalOp::HashAggregation(op) => self.exec_hash_aggregation(plan, op),
213 PhysicalOp::Sort(op) => self.exec_sort(plan, op),
214 PhysicalOp::Limit(op) => self.exec_limit(plan, op),
215 PhysicalOp::Create(op) => self.exec_create(plan, op),
216 PhysicalOp::Merge(op) => self.exec_merge(plan, op),
217 PhysicalOp::Delete(op) => self.exec_delete(plan, op),
218 PhysicalOp::Set(op) => self.exec_set(plan, op),
219 PhysicalOp::Remove(op) => self.exec_remove(plan, op),
220 PhysicalOp::Foreach(op) => self.exec_foreach(plan, op),
221 PhysicalOp::OptionalMatch(op) => self.exec_optional_match(plan, op),
222 PhysicalOp::CallSubquery(op) => self.exec_call_subquery(plan, op),
223 PhysicalOp::PathBuild(op) => self.exec_path_build(plan, op),
224 };
225
226 match &result {
227 Ok(rows) => trace!(
228 "mutable execute_node ok: node_id={node_id:?}, rows={}",
229 rows.len()
230 ),
231 Err(err) => error!("mutable execute_node failed: node_id={node_id:?}, error={err}"),
232 }
233
234 result
235 }
236
237 fn exec_argument(&self, _op: &ArgumentExec) -> ExecResult<Vec<Row>> {
238 Ok(vec![Row::new()])
239 }
240
241 fn exec_node_scan(&mut self, plan: &PhysicalPlan, op: &NodeScanExec) -> ExecResult<Vec<Row>> {
242 let base_rows = match op.input {
243 Some(input) => self.execute_node(plan, input)?,
244 None => vec![Row::new()],
245 };
246
247 node_scan_rows(&*self.ctx.storage, base_rows, op, self.deadline)
248 }
249
250 fn exec_node_by_label_scan(
251 &mut self,
252 plan: &PhysicalPlan,
253 op: &NodeByLabelScanExec,
254 ) -> ExecResult<Vec<Row>> {
255 let base_rows = match op.input {
256 Some(input) => self.execute_node(plan, input)?,
257 None => vec![Row::new()],
258 };
259
260 node_by_label_scan_rows(&*self.ctx.storage, base_rows, op, self.deadline)
261 }
262
263 fn exec_node_by_property_scan(
264 &mut self,
265 plan: &PhysicalPlan,
266 op: &NodeByPropertyScanExec,
267 ) -> ExecResult<Vec<Row>> {
268 let base_rows = match op.input {
269 Some(input) => self.execute_node(plan, input)?,
270 None => vec![Row::new()],
271 };
272
273 node_by_property_scan_rows(
274 &*self.ctx.storage,
275 &self.ctx.params,
276 base_rows,
277 op,
278 self.deadline,
279 )
280 }
281
282 fn exec_node_by_property_range_scan(
283 &mut self,
284 plan: &PhysicalPlan,
285 op: &lora_compiler::NodeByPropertyRangeScanExec,
286 ) -> ExecResult<Vec<Row>> {
287 let base_rows = match op.input {
288 Some(input) => self.execute_node(plan, input)?,
289 None => vec![Row::new()],
290 };
291 super::helpers::node_by_property_range_scan_rows(
292 &*self.ctx.storage,
293 &self.ctx.params,
294 base_rows,
295 op,
296 self.deadline,
297 )
298 }
299
300 fn exec_node_by_text_scan(
301 &mut self,
302 plan: &PhysicalPlan,
303 op: &lora_compiler::NodeByTextScanExec,
304 ) -> ExecResult<Vec<Row>> {
305 let base_rows = match op.input {
306 Some(input) => self.execute_node(plan, input)?,
307 None => vec![Row::new()],
308 };
309 super::helpers::node_by_text_scan_rows(
310 &*self.ctx.storage,
311 &self.ctx.params,
312 base_rows,
313 op,
314 self.deadline,
315 )
316 }
317
318 fn exec_node_by_point_scan(
319 &mut self,
320 plan: &PhysicalPlan,
321 op: &lora_compiler::NodeByPointScanExec,
322 ) -> ExecResult<Vec<Row>> {
323 let base_rows = match op.input {
324 Some(input) => self.execute_node(plan, input)?,
325 None => vec![Row::new()],
326 };
327 super::helpers::node_by_point_scan_rows(
328 &*self.ctx.storage,
329 &self.ctx.params,
330 base_rows,
331 op,
332 self.deadline,
333 )
334 }
335
336 fn exec_rel_by_property_range_scan(
337 &mut self,
338 plan: &PhysicalPlan,
339 op: &lora_compiler::RelByPropertyRangeScanExec,
340 ) -> ExecResult<Vec<Row>> {
341 let base_rows = match op.input {
342 Some(input) => self.execute_node(plan, input)?,
343 None => vec![Row::new()],
344 };
345 super::helpers::rel_by_property_range_scan_rows(
346 &*self.ctx.storage,
347 &self.ctx.params,
348 base_rows,
349 op,
350 self.deadline,
351 )
352 }
353
354 fn exec_rel_by_text_scan(
355 &mut self,
356 plan: &PhysicalPlan,
357 op: &lora_compiler::RelByTextScanExec,
358 ) -> ExecResult<Vec<Row>> {
359 let base_rows = match op.input {
360 Some(input) => self.execute_node(plan, input)?,
361 None => vec![Row::new()],
362 };
363 super::helpers::rel_by_text_scan_rows(
364 &*self.ctx.storage,
365 &self.ctx.params,
366 base_rows,
367 op,
368 self.deadline,
369 )
370 }
371
372 fn exec_rel_by_point_scan(
373 &mut self,
374 plan: &PhysicalPlan,
375 op: &lora_compiler::RelByPointScanExec,
376 ) -> ExecResult<Vec<Row>> {
377 let base_rows = match op.input {
378 Some(input) => self.execute_node(plan, input)?,
379 None => vec![Row::new()],
380 };
381 super::helpers::rel_by_point_scan_rows(
382 &*self.ctx.storage,
383 &self.ctx.params,
384 base_rows,
385 op,
386 self.deadline,
387 )
388 }
389
390 fn exec_expand(&mut self, plan: &PhysicalPlan, op: &ExpandExec) -> ExecResult<Vec<Row>> {
391 let input_rows = self.execute_node(plan, op.input)?;
392 if let Some(range) = &op.range {
393 expand_var_len_rows(&*self.ctx.storage, input_rows, op, range)
394 } else {
395 expand_rows(&*self.ctx.storage, &self.ctx.params, input_rows, op)
396 }
397 }
398
399 fn exec_filter(&mut self, plan: &PhysicalPlan, op: &FilterExec) -> ExecResult<Vec<Row>> {
400 let input_rows = self.execute_node(plan, op.input)?;
401 let eval_ctx = EvalContext {
402 storage: &*self.ctx.storage,
403 params: &self.ctx.params,
404 };
405
406 filter_rows_checked(input_rows, &op.predicate, &eval_ctx)
407 }
408
409 fn exec_projection(
410 &mut self,
411 plan: &PhysicalPlan,
412 op: &ProjectionExec,
413 ) -> ExecResult<Vec<Row>> {
414 let input_rows = self.execute_node(plan, op.input)?;
415 let eval_ctx = EvalContext {
416 storage: &*self.ctx.storage,
417 params: &self.ctx.params,
418 };
419
420 project_rows_checked(input_rows, op, &eval_ctx)
421 }
422
423 fn hydrate_value(&self, value: LoraValue) -> LoraValue {
424 match value {
425 LoraValue::Node(id) => self.hydrate_node(id),
426 LoraValue::Relationship(id) => self.hydrate_relationship(id),
427 LoraValue::List(values) => {
428 LoraValue::List(values.into_iter().map(|v| self.hydrate_value(v)).collect())
429 }
430 LoraValue::Map(map) => LoraValue::Map(
431 map.into_iter()
432 .map(|(k, v)| (k, self.hydrate_value(v)))
433 .collect(),
434 ),
435 other => other,
436 }
437 }
438
439 fn hydrate_node(&self, id: u64) -> LoraValue {
440 self.ctx
441 .storage
442 .with_node(id, hydrate_node_record)
443 .unwrap_or(LoraValue::Null)
444 }
445
446 fn hydrate_relationship(&self, id: u64) -> LoraValue {
447 self.ctx
448 .storage
449 .with_relationship(id, hydrate_relationship_record)
450 .unwrap_or(LoraValue::Null)
451 }
452
453 fn exec_unwind(&mut self, plan: &PhysicalPlan, op: &UnwindExec) -> ExecResult<Vec<Row>> {
454 let input_rows = self.execute_node(plan, op.input)?;
455 let eval_ctx = EvalContext {
456 storage: &*self.ctx.storage,
457 params: &self.ctx.params,
458 };
459
460 Ok(unwind_rows(input_rows, op, &eval_ctx))
461 }
462
463 fn exec_hash_aggregation(
464 &mut self,
465 plan: &PhysicalPlan,
466 op: &HashAggregationExec,
467 ) -> ExecResult<Vec<Row>> {
468 if let Some(rows) =
469 super::helpers::count_all_scan_aggregation_rows(&*self.ctx.storage, plan, op)
470 {
471 return Ok(rows);
472 }
473
474 let input_rows = self.execute_node(plan, op.input)?;
475 let eval_ctx = EvalContext {
476 storage: &*self.ctx.storage,
477 params: &self.ctx.params,
478 };
479
480 aggregate_rows(
481 input_rows,
482 &op.group_by,
483 &op.aggregates,
484 &eval_ctx,
485 |value| self.hydrate_value(value),
486 )
487 }
488
489 fn exec_sort(&mut self, plan: &PhysicalPlan, op: &SortExec) -> ExecResult<Vec<Row>> {
490 let mut rows = self.execute_node(plan, op.input)?;
491 let eval_ctx = EvalContext {
492 storage: &*self.ctx.storage,
493 params: &self.ctx.params,
494 };
495
496 sort_rows_with_top_k(&mut rows, &op.items, &eval_ctx, op.top_k);
497
498 Ok(rows)
499 }
500
501 fn exec_limit(&mut self, plan: &PhysicalPlan, op: &LimitExec) -> ExecResult<Vec<Row>> {
502 let rows = self.execute_node(plan, op.input)?;
503 let eval_ctx = EvalContext {
504 storage: &*self.ctx.storage,
505 params: &self.ctx.params,
506 };
507
508 Ok(limit_rows(rows, op, &eval_ctx))
509 }
510
511 fn exec_optional_match(
512 &mut self,
513 plan: &PhysicalPlan,
514 op: &OptionalMatchExec,
515 ) -> ExecResult<Vec<Row>> {
516 let input_rows = self.execute_node(plan, op.input)?;
517
518 let inner_rows = self.execute_node(plan, op.inner)?;
520
521 Ok(optional_match_rows(input_rows, &inner_rows, &op.new_vars))
522 }
523
524 fn exec_call_subquery(
525 &mut self,
526 plan: &PhysicalPlan,
527 op: &CallSubqueryExec,
528 ) -> ExecResult<Vec<Row>> {
529 let input_rows = self.execute_node(plan, op.input)?;
530 let mut out = Vec::with_capacity(input_rows.len());
531 let params = std::sync::Arc::new(self.ctx.params.clone());
532 let storage_ref: &S = &*self.ctx.storage;
533 for outer_row in input_rows {
534 let mut inner_source = crate::pull::build_streaming_seeded(
535 plan,
536 op.inner,
537 storage_ref,
538 params.clone(),
539 outer_row.clone(),
540 )?;
541 let inner_rows = crate::pull::drain(inner_source.as_mut())?;
542 for inner_row in inner_rows {
543 out.push(crate::executor::merge_optional_rows(&outer_row, &inner_row));
544 }
545 }
546 Ok(out)
547 }
548
549 fn exec_path_build(&mut self, plan: &PhysicalPlan, op: &PathBuildExec) -> ExecResult<Vec<Row>> {
550 let input_rows = self.execute_node(plan, op.input)?;
551 let mut rows: Vec<Row> = input_rows
552 .into_iter()
553 .map(|mut row| {
554 let path = build_path_value(&row, &op.node_vars, &op.rel_vars, &*self.ctx.storage);
555 row.insert(op.output, path);
556 row
557 })
558 .collect();
559
560 if let Some(all) = op.shortest_path_all {
561 rows = filter_shortest_paths(rows, op.output, all);
562 }
563 Ok(rows)
564 }
565
566 fn exec_create(&mut self, plan: &PhysicalPlan, op: &CreateExec) -> ExecResult<Vec<Row>> {
567 if crate::pull::subtree_is_fully_streaming(plan, op.input) {
573 return self.exec_create_streaming_input(plan, op);
574 }
575
576 let input_rows = self.execute_node(plan, op.input)?;
577 let mut out = Vec::with_capacity(input_rows.len());
578
579 for mut row in input_rows {
580 self.apply_create_pattern(&mut row, &op.pattern)?;
581 out.push(row);
582 }
583
584 Ok(out)
585 }
586
587 fn streaming_apply<F>(
607 &mut self,
608 plan: &PhysicalPlan,
609 input: PhysicalNodeId,
610 mut apply: F,
611 ) -> ExecResult<Vec<Row>>
612 where
613 F: FnMut(&mut Self, &mut Row) -> ExecResult<()>,
614 {
615 use std::sync::Arc;
616
617 let storage_ptr: *mut S = self.ctx.storage as *mut S;
618 let params = Arc::new(self.ctx.params.clone());
619
620 let storage_ref: &S = unsafe { &*storage_ptr };
622 let mut upstream = crate::pull::build_streaming(plan, input, storage_ref, params)?;
623
624 let mut out = Vec::new();
625 while let Some(mut row) = upstream.next_row()? {
626 apply(self, &mut row)?;
627 out.push(row);
628 }
629
630 Ok(out)
631 }
632
633 fn exec_create_streaming_input(
636 &mut self,
637 plan: &PhysicalPlan,
638 op: &CreateExec,
639 ) -> ExecResult<Vec<Row>> {
640 self.streaming_apply(plan, op.input, |this, row| {
641 this.apply_create_pattern(row, &op.pattern)
642 })
643 }
644
645 fn apply_remove_item(&mut self, row: &Row, item: &ResolvedRemoveItem) -> ExecResult<()> {
646 match item {
647 ResolvedRemoveItem::Labels { variable, labels } => match row.get(*variable) {
648 Some(LoraValue::Node(node_id)) => {
649 let node_id = *node_id;
650 for label in labels {
651 self.ctx.storage.remove_node_label(node_id, label);
652 }
653 Ok(())
654 }
655 Some(other) => Err(ExecutorError::ExpectedNodeForRemoveLabels {
656 found: value_kind(other),
657 }),
658 None => Err(ExecutorError::UnboundVariableForRemove {
659 var: format!("{variable:?}"),
660 }),
661 },
662
663 ResolvedRemoveItem::Property { expr } => self.remove_property_from_expr(row, expr),
664 }
665 }
666
667 fn delete_value(&mut self, value: LoraValue, detach: bool) -> ExecResult<()> {
668 match value {
669 LoraValue::Null => Ok(()),
670
671 LoraValue::Node(node_id) => {
672 if detach {
673 self.ctx.storage.detach_delete_node(node_id);
674 Ok(())
675 } else {
676 let ok = self.ctx.storage.delete_node(node_id);
677 if ok {
678 Ok(())
679 } else {
680 Err(ExecutorError::DeleteNodeWithRelationships { node_id })
681 }
682 }
683 }
684
685 LoraValue::Relationship(rel_id) => {
686 let ok = self.ctx.storage.delete_relationship(rel_id);
687 if ok {
688 Ok(())
689 } else {
690 Err(ExecutorError::DeleteRelationshipFailed { rel_id })
691 }
692 }
693
694 LoraValue::List(values) => {
695 for v in values {
696 self.delete_value(v, detach)?;
697 }
698 Ok(())
699 }
700
701 other => Err(ExecutorError::InvalidDeleteTarget {
702 found: value_kind(&other),
703 }),
704 }
705 }
706
707 fn collect_delete_targets(
708 &self,
709 value: &LoraValue,
710 targets: &mut BTreeSet<DeleteTarget>,
711 ) -> ExecResult<()> {
712 match value {
713 LoraValue::Null => Ok(()),
714
715 LoraValue::Node(node_id) => {
716 targets.insert(DeleteTarget::Node(*node_id));
717 Ok(())
718 }
719
720 LoraValue::Relationship(rel_id) => {
721 targets.insert(DeleteTarget::Relationship(*rel_id));
722 Ok(())
723 }
724
725 LoraValue::List(values) => {
726 for v in values {
727 self.collect_delete_targets(v, targets)?;
728 }
729 Ok(())
730 }
731
732 other => Err(ExecutorError::InvalidDeleteTarget {
733 found: value_kind(other),
734 }),
735 }
736 }
737
738 fn validate_delete_targets(
739 &self,
740 targets: &BTreeSet<DeleteTarget>,
741 detach: bool,
742 ) -> ExecResult<()> {
743 for target in targets {
744 match target {
745 DeleteTarget::Relationship(rel_id) => {
746 if !self.ctx.storage.contains_relationship(*rel_id) {
747 return Err(ExecutorError::DeleteRelationshipFailed { rel_id: *rel_id });
748 }
749 }
750 DeleteTarget::Node(node_id) if !detach => {
751 if !self.ctx.storage.contains_node(*node_id) {
752 return Err(ExecutorError::DeleteNodeWithRelationships {
753 node_id: *node_id,
754 });
755 }
756 let has_external_relationship = self
757 .ctx
758 .storage
759 .relationship_ids_of(*node_id, Direction::Undirected)
760 .into_iter()
761 .any(|rel_id| !targets.contains(&DeleteTarget::Relationship(rel_id)));
762 if has_external_relationship {
763 return Err(ExecutorError::DeleteNodeWithRelationships {
764 node_id: *node_id,
765 });
766 }
767 }
768 DeleteTarget::Node(_) => {}
769 }
770 }
771 Ok(())
772 }
773
774 fn delete_target(&mut self, target: DeleteTarget, detach: bool) -> ExecResult<()> {
775 match target {
776 DeleteTarget::Node(node_id) => {
777 if detach {
778 self.ctx.storage.detach_delete_node(node_id);
779 Ok(())
780 } else {
781 let ok = self.ctx.storage.delete_node(node_id);
782 if ok {
783 Ok(())
784 } else {
785 Err(ExecutorError::DeleteNodeWithRelationships { node_id })
786 }
787 }
788 }
789 DeleteTarget::Relationship(rel_id) => {
790 let ok = self.ctx.storage.delete_relationship(rel_id);
791 if ok {
792 Ok(())
793 } else {
794 Err(ExecutorError::DeleteRelationshipFailed { rel_id })
795 }
796 }
797 }
798 }
799
800 fn exec_merge(&mut self, plan: &PhysicalPlan, op: &MergeExec) -> ExecResult<Vec<Row>> {
801 if crate::pull::subtree_is_fully_streaming(plan, op.input) {
806 return self.streaming_apply(plan, op.input, |this, row| {
807 let already_bound = this.pattern_part_is_bound(row, &op.pattern_part);
808 let matched = if already_bound {
809 true
810 } else {
811 this.try_match_merge_pattern(row, &op.pattern_part)?
812 };
813 if !matched {
814 this.apply_create_pattern_part(row, &op.pattern_part)?;
815 }
816 for action in &op.actions {
817 if action.on_match == matched {
818 for item in &action.set.items {
819 this.apply_set_item(row, item)?;
820 }
821 }
822 }
823 Ok(())
824 });
825 }
826
827 let input_rows = self.execute_node(plan, op.input)?;
828 let mut out = Vec::with_capacity(input_rows.len());
829
830 for mut row in input_rows {
831 let already_bound = self.pattern_part_is_bound(&row, &op.pattern_part);
833
834 let matched = if already_bound {
835 true
836 } else {
837 self.try_match_merge_pattern(&mut row, &op.pattern_part)?
839 };
840
841 if !matched {
842 self.apply_create_pattern_part(&mut row, &op.pattern_part)?;
843 }
844
845 for action in &op.actions {
846 if action.on_match == matched {
847 for item in &action.set.items {
848 self.apply_set_item(&row, item)?;
849 }
850 }
851 }
852
853 out.push(row);
854 }
855
856 Ok(out)
857 }
858
859 fn try_match_merge_pattern(
862 &self,
863 row: &mut Row,
864 part: &ResolvedPatternPart,
865 ) -> ExecResult<bool> {
866 match &part.element {
867 ResolvedPatternElement::Node {
868 var,
869 labels,
870 properties,
871 } => {
872 let candidate_ids = if labels.is_empty() {
875 self.ctx.storage.all_node_ids()
876 } else {
877 scan_node_ids_for_label_groups(&*self.ctx.storage, labels)
878 };
879
880 let eval_ctx = EvalContext {
882 storage: &*self.ctx.storage,
883 params: &self.ctx.params,
884 };
885 let expected_props = properties.as_ref().map(|e| eval_expr(e, row, &eval_ctx));
886
887 for id in candidate_ids {
888 let matched = self
889 .ctx
890 .storage
891 .with_node(id, |node| {
892 if !node_matches_label_groups(&node.labels, labels) {
893 return false;
894 }
895 if let Some(LoraValue::Map(expected)) = &expected_props {
896 let all_match = expected.iter().all(|(key, expected_value)| {
897 node.properties
898 .get(key.as_str())
899 .map(|actual| {
900 value_matches_property_value(expected_value, actual)
901 })
902 .unwrap_or(false)
903 });
904 if !all_match {
905 return false;
906 }
907 }
908 true
909 })
910 .unwrap_or(false);
911
912 if !matched {
913 continue;
914 }
915
916 if let Some(var_id) = var {
918 row.insert(*var_id, LoraValue::Node(id));
919 }
920 return Ok(true);
921 }
922
923 Ok(false)
924 }
925
926 ResolvedPatternElement::ShortestPath { .. } => {
927 Ok(false)
929 }
930
931 ResolvedPatternElement::NodeChain { head, chain } => {
932 let head_node_id = if let Some(var_id) = head.var {
934 if let Some(LoraValue::Node(id)) = row.get(var_id) {
935 *id
936 } else {
937 let node_matched = self.try_match_merge_pattern(
939 row,
940 &ResolvedPatternPart {
941 binding: None,
942 element: ResolvedPatternElement::Node {
943 var: head.var,
944 labels: head.labels.clone(),
945 properties: head.properties.clone(),
946 },
947 },
948 )?;
949 if !node_matched {
950 return Ok(false);
951 }
952 match row.get(var_id) {
953 Some(LoraValue::Node(id)) => *id,
954 _ => return Ok(false),
955 }
956 }
957 } else {
958 return Ok(false);
959 };
960
961 let mut current_node_id = head_node_id;
962
963 for step in chain {
964 let eval_ctx = EvalContext {
965 storage: &*self.ctx.storage,
966 params: &self.ctx.params,
967 };
968
969 let direction = step.rel.direction;
970
971 let mut found = false;
974 let _ = self.ctx.storage.try_for_each_expand_id(
975 current_node_id,
976 direction,
977 &step.rel.types,
978 |rel_id, node_id| {
979 let node_ok = self
981 .ctx
982 .storage
983 .with_node(node_id, |node_rec| {
984 if !node_matches_label_groups(
985 &node_rec.labels,
986 &step.node.labels,
987 ) {
988 return false;
989 }
990 if let Some(props_expr) = &step.node.properties {
991 let expected = eval_expr(props_expr, row, &eval_ctx);
992 if let LoraValue::Map(expected_map) = &expected {
993 let all_match =
994 expected_map.iter().all(|(key, expected_val)| {
995 node_rec
996 .properties
997 .get(key.as_str())
998 .map(|actual| {
999 value_matches_property_value(
1000 expected_val,
1001 actual,
1002 )
1003 })
1004 .unwrap_or(false)
1005 });
1006 if !all_match {
1007 return false;
1008 }
1009 }
1010 }
1011 true
1012 })
1013 .unwrap_or(false);
1014 if !node_ok {
1015 return Ok::<(), ()>(());
1016 }
1017
1018 let rel_ok = self
1020 .ctx
1021 .storage
1022 .with_relationship(rel_id, |rel_rec| {
1023 if let Some(rel_props_expr) = &step.rel.properties {
1024 let expected = eval_expr(rel_props_expr, row, &eval_ctx);
1025 if let LoraValue::Map(expected_map) = &expected {
1026 let all_match =
1027 expected_map.iter().all(|(key, expected_val)| {
1028 rel_rec
1029 .properties
1030 .get(key.as_str())
1031 .map(|actual| {
1032 value_matches_property_value(
1033 expected_val,
1034 actual,
1035 )
1036 })
1037 .unwrap_or(false)
1038 });
1039 if !all_match {
1040 return false;
1041 }
1042 }
1043 }
1044 true
1045 })
1046 .unwrap_or(false);
1047 if !rel_ok {
1048 return Ok(());
1049 }
1050
1051 if let Some(rel_var) = step.rel.var {
1053 row.insert(rel_var, LoraValue::Relationship(rel_id));
1054 }
1055 if let Some(node_var) = step.node.var {
1056 row.insert(node_var, LoraValue::Node(node_id));
1057 }
1058 current_node_id = node_id;
1059 found = true;
1060 Err(())
1061 },
1062 );
1063
1064 if !found {
1065 return Ok(false);
1066 }
1067 }
1068
1069 Ok(true)
1070 }
1071 }
1072 }
1073
1074 fn exec_delete(&mut self, plan: &PhysicalPlan, op: &DeleteExec) -> ExecResult<Vec<Row>> {
1075 let input_rows = self.execute_node(plan, op.input)?;
1076 let mut targets = BTreeSet::new();
1077
1078 for row in &input_rows {
1079 for expr in &op.expressions {
1080 let value = {
1081 let eval_ctx = EvalContext {
1082 storage: &*self.ctx.storage,
1083 params: &self.ctx.params,
1084 };
1085 eval_expr(expr, row, &eval_ctx)
1086 };
1087 self.collect_delete_targets(&value, &mut targets)?;
1088 }
1089 }
1090
1091 self.validate_delete_targets(&targets, op.detach)?;
1092
1093 for target in &targets {
1094 if let DeleteTarget::Relationship(_) = target {
1095 self.delete_target(*target, op.detach)?;
1096 }
1097 }
1098 for target in targets {
1099 if let DeleteTarget::Node(_) = target {
1100 self.delete_target(target, op.detach)?;
1101 }
1102 }
1103
1104 Ok(input_rows)
1105 }
1106
1107 fn exec_set(&mut self, plan: &PhysicalPlan, op: &SetExec) -> ExecResult<Vec<Row>> {
1108 if crate::pull::subtree_is_fully_streaming(plan, op.input) {
1109 return self.streaming_apply(plan, op.input, |this, row| {
1110 for item in &op.items {
1111 this.apply_set_item(row, item)?;
1112 }
1113 Ok(())
1114 });
1115 }
1116
1117 let input_rows = self.execute_node(plan, op.input)?;
1118
1119 for row in &input_rows {
1120 for item in &op.items {
1121 self.apply_set_item(row, item)?;
1122 }
1123 }
1124
1125 Ok(input_rows)
1126 }
1127
1128 fn exec_foreach(&mut self, plan: &PhysicalPlan, op: &ForeachExec) -> ExecResult<Vec<Row>> {
1136 let input_rows = self.execute_node(plan, op.input)?;
1137 let mut out = Vec::with_capacity(input_rows.len());
1138
1139 for row in input_rows {
1140 let list_value = {
1141 let eval_ctx = EvalContext {
1142 storage: &*self.ctx.storage,
1143 params: &self.ctx.params,
1144 };
1145 eval_expr(&op.list, &row, &eval_ctx)
1146 };
1147
1148 let elements: Vec<LoraValue> = match list_value {
1149 LoraValue::List(items) => items,
1150 LoraValue::Null => Vec::new(),
1151 other => {
1152 return Err(ExecutorError::RuntimeError(format!(
1153 "FOREACH expects a list, got {}",
1154 value_kind(&other)
1155 )));
1156 }
1157 };
1158
1159 for element in elements {
1160 let mut iter_row = row.clone();
1163 iter_row.insert(op.variable, element);
1164 for clause in &op.body {
1165 self.apply_foreach_body_clause(&mut iter_row, clause)?;
1166 }
1167 }
1168
1169 out.push(row);
1170 }
1171
1172 Ok(out)
1173 }
1174
1175 fn apply_foreach_body_clause(
1180 &mut self,
1181 row: &mut Row,
1182 clause: &lora_analyzer::ResolvedClause,
1183 ) -> ExecResult<()> {
1184 use lora_analyzer::ResolvedClause;
1185 match clause {
1186 ResolvedClause::Create(c) => self.apply_create_pattern(row, &c.pattern),
1187 ResolvedClause::Set(s) => {
1188 for item in &s.items {
1189 self.apply_set_item(row, item)?;
1190 }
1191 Ok(())
1192 }
1193 ResolvedClause::Remove(r) => {
1194 for item in &r.items {
1195 self.apply_remove_item(row, item)?;
1196 }
1197 Ok(())
1198 }
1199 ResolvedClause::Delete(d) => {
1200 let detach = d.detach;
1201 for expr in &d.expressions {
1202 let value = {
1203 let eval_ctx = EvalContext {
1204 storage: &*self.ctx.storage,
1205 params: &self.ctx.params,
1206 };
1207 eval_expr(expr, row, &eval_ctx)
1208 };
1209 self.delete_value(value, detach)?;
1210 }
1211 Ok(())
1212 }
1213 ResolvedClause::Merge(m) => {
1214 let already_bound = self.pattern_part_is_bound(row, &m.pattern_part);
1215 let matched = if already_bound {
1216 true
1217 } else {
1218 self.try_match_merge_pattern(row, &m.pattern_part)?
1219 };
1220 if !matched {
1221 self.apply_create_pattern_part(row, &m.pattern_part)?;
1222 }
1223 for action in &m.actions {
1224 if action.on_match == matched {
1225 for item in &action.set.items {
1226 self.apply_set_item(row, item)?;
1227 }
1228 }
1229 }
1230 Ok(())
1231 }
1232 ResolvedClause::Foreach(nested) => {
1233 let list_value = {
1234 let eval_ctx = EvalContext {
1235 storage: &*self.ctx.storage,
1236 params: &self.ctx.params,
1237 };
1238 eval_expr(&nested.list, row, &eval_ctx)
1239 };
1240
1241 let elements: Vec<LoraValue> = match list_value {
1242 LoraValue::List(items) => items,
1243 LoraValue::Null => Vec::new(),
1244 other => {
1245 return Err(ExecutorError::RuntimeError(format!(
1246 "FOREACH expects a list, got {}",
1247 value_kind(&other)
1248 )));
1249 }
1250 };
1251
1252 for element in elements {
1253 let mut iter_row = row.clone();
1254 iter_row.insert(nested.variable, element);
1255 for inner in &nested.body {
1256 self.apply_foreach_body_clause(&mut iter_row, inner)?;
1257 }
1258 }
1259
1260 Ok(())
1261 }
1262 other => Err(ExecutorError::RuntimeError(format!(
1263 "FOREACH body may only contain updating clauses, got {:?}",
1264 std::mem::discriminant(other)
1265 ))),
1266 }
1267 }
1268
1269 fn exec_remove(&mut self, plan: &PhysicalPlan, op: &RemoveExec) -> ExecResult<Vec<Row>> {
1270 if crate::pull::subtree_is_fully_streaming(plan, op.input) {
1271 return self.streaming_apply(plan, op.input, |this, row| {
1272 for item in &op.items {
1273 this.apply_remove_item(row, item)?;
1274 }
1275 Ok(())
1276 });
1277 }
1278
1279 let input_rows = self.execute_node(plan, op.input)?;
1280
1281 for row in &input_rows {
1282 for item in &op.items {
1283 self.apply_remove_item(row, item)?;
1284 }
1285 }
1286
1287 Ok(input_rows)
1288 }
1289
1290 fn apply_set_item(&mut self, row: &Row, item: &ResolvedSetItem) -> ExecResult<()> {
1291 match item {
1292 ResolvedSetItem::SetProperty { target, value } => {
1293 let new_value = {
1294 let eval_ctx = EvalContext {
1295 storage: &*self.ctx.storage,
1296 params: &self.ctx.params,
1297 };
1298 eval_expr(value, row, &eval_ctx)
1299 };
1300
1301 self.set_property_from_expr(row, target, new_value)
1302 }
1303
1304 ResolvedSetItem::SetVariable { variable, value } => {
1305 let entity_ref =
1307 row.get(*variable)
1308 .ok_or(ExecutorError::UnboundVariableForSet {
1309 var: format!("{variable:?}"),
1310 })?;
1311 let entity_target = entity_target_from_value(entity_ref)?;
1312
1313 let new_value = {
1314 let eval_ctx = EvalContext {
1315 storage: &*self.ctx.storage,
1316 params: &self.ctx.params,
1317 };
1318 eval_expr(value, row, &eval_ctx)
1319 };
1320
1321 self.overwrite_entity_target(entity_target, new_value)
1322 }
1323
1324 ResolvedSetItem::MutateVariable { variable, value } => {
1325 let entity_ref =
1326 row.get(*variable)
1327 .ok_or(ExecutorError::UnboundVariableForSet {
1328 var: format!("{variable:?}"),
1329 })?;
1330 let entity_target = entity_target_from_value(entity_ref)?;
1331
1332 let patch = {
1333 let eval_ctx = EvalContext {
1334 storage: &*self.ctx.storage,
1335 params: &self.ctx.params,
1336 };
1337 eval_expr(value, row, &eval_ctx)
1338 };
1339
1340 self.mutate_entity_target(entity_target, patch)
1341 }
1342
1343 ResolvedSetItem::SetLabels { variable, labels } => match row.get(*variable) {
1344 Some(LoraValue::Node(node_id)) => {
1345 let node_id = *node_id;
1346 for label in labels {
1347 if let Err(msg) = self
1348 .ctx
1349 .storage
1350 .check_node_add_label_against_constraints(node_id, label)
1351 {
1352 return Err(ExecutorError::ConstraintViolation(msg));
1353 }
1354 self.ctx.storage.add_node_label(node_id, label);
1355 }
1356 Ok(())
1357 }
1358 Some(other) => Err(ExecutorError::ExpectedNodeForSetLabels {
1359 found: value_kind(other),
1360 }),
1361 None => Err(ExecutorError::UnboundVariableForSet {
1362 var: format!("{variable:?}"),
1363 }),
1364 },
1365 }
1366 }
1367
1368 fn set_property_from_expr(
1369 &mut self,
1370 row: &Row,
1371 target_expr: &ResolvedExpr,
1372 new_value: LoraValue,
1373 ) -> ExecResult<()> {
1374 let ResolvedExpr::Property { expr, property } = target_expr else {
1375 return Err(ExecutorError::UnsupportedSetTarget);
1376 };
1377
1378 let owner = {
1379 let eval_ctx = EvalContext {
1380 storage: &*self.ctx.storage,
1381 params: &self.ctx.params,
1382 };
1383 eval_expr(expr, row, &eval_ctx)
1384 };
1385
1386 match owner {
1387 LoraValue::Node(node_id) => {
1388 let prop = lora_value_to_property(new_value)
1389 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1390 if let Err(msg) = self
1391 .ctx
1392 .storage
1393 .check_node_set_property_against_constraints(node_id, property, &prop)
1394 {
1395 return Err(ExecutorError::ConstraintViolation(msg));
1396 }
1397 self.ctx
1398 .storage
1399 .set_node_property(node_id, property.clone(), prop);
1400 Ok(())
1401 }
1402 LoraValue::Relationship(rel_id) => {
1403 let prop = lora_value_to_property(new_value)
1404 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1405 if let Err(msg) = self
1406 .ctx
1407 .storage
1408 .check_relationship_set_property_against_constraints(rel_id, property, &prop)
1409 {
1410 return Err(ExecutorError::ConstraintViolation(msg));
1411 }
1412 self.ctx
1413 .storage
1414 .set_relationship_property(rel_id, property.clone(), prop);
1415 Ok(())
1416 }
1417 other => Err(ExecutorError::InvalidSetTarget {
1418 found: value_kind(&other),
1419 }),
1420 }
1421 }
1422
1423 fn remove_property_from_expr(&mut self, row: &Row, expr: &ResolvedExpr) -> ExecResult<()> {
1424 let ResolvedExpr::Property {
1425 expr: owner_expr,
1426 property,
1427 } = expr
1428 else {
1429 return Err(ExecutorError::UnsupportedRemoveTarget);
1430 };
1431
1432 let owner = {
1433 let eval_ctx = EvalContext {
1434 storage: &*self.ctx.storage,
1435 params: &self.ctx.params,
1436 };
1437 eval_expr(owner_expr, row, &eval_ctx)
1438 };
1439
1440 match owner {
1441 LoraValue::Node(node_id) => {
1442 if let Err(msg) = self
1443 .ctx
1444 .storage
1445 .check_node_remove_property_against_constraints(node_id, property)
1446 {
1447 return Err(ExecutorError::ConstraintViolation(msg));
1448 }
1449 self.ctx.storage.remove_node_property(node_id, property);
1450 Ok(())
1451 }
1452 LoraValue::Relationship(rel_id) => {
1453 if let Err(msg) = self
1454 .ctx
1455 .storage
1456 .check_relationship_remove_property_against_constraints(rel_id, property)
1457 {
1458 return Err(ExecutorError::ConstraintViolation(msg));
1459 }
1460 self.ctx
1461 .storage
1462 .remove_relationship_property(rel_id, property);
1463 Ok(())
1464 }
1465 other => Err(ExecutorError::InvalidRemoveTarget {
1466 found: value_kind(&other),
1467 }),
1468 }
1469 }
1470
1471 fn overwrite_entity_target(
1472 &mut self,
1473 target: EntityTarget,
1474 new_value: LoraValue,
1475 ) -> ExecResult<()> {
1476 let LoraValue::Map(map) = new_value else {
1477 return Err(ExecutorError::ExpectedPropertyMap {
1478 found: value_kind(&new_value),
1479 });
1480 };
1481
1482 let mut props: Properties = Properties::new();
1483 for (k, v) in map {
1484 let prop = lora_value_to_property(v)
1485 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1486 props.insert(lora_store::intern_owned(k), prop);
1487 }
1488
1489 match target {
1490 EntityTarget::Node(node_id) => {
1491 if let Err(msg) = self
1492 .ctx
1493 .storage
1494 .check_node_replace_properties_against_constraints(node_id, &props)
1495 {
1496 return Err(ExecutorError::ConstraintViolation(msg));
1497 }
1498 self.ctx.storage.replace_node_properties(node_id, props);
1499 }
1500 EntityTarget::Relationship(rel_id) => {
1501 if let Err(msg) = self
1502 .ctx
1503 .storage
1504 .check_relationship_replace_properties_against_constraints(rel_id, &props)
1505 {
1506 return Err(ExecutorError::ConstraintViolation(msg));
1507 }
1508 self.ctx
1509 .storage
1510 .replace_relationship_properties(rel_id, props);
1511 }
1512 }
1513 Ok(())
1514 }
1515
1516 fn mutate_entity_target(
1517 &mut self,
1518 target: EntityTarget,
1519 patch_value: LoraValue,
1520 ) -> ExecResult<()> {
1521 let LoraValue::Map(map) = patch_value else {
1522 return Err(ExecutorError::ExpectedPropertyMap {
1523 found: value_kind(&patch_value),
1524 });
1525 };
1526
1527 match target {
1528 EntityTarget::Node(node_id) => {
1529 for (k, v) in map {
1530 let prop = lora_value_to_property(v)
1531 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1532 if let Err(msg) = self
1533 .ctx
1534 .storage
1535 .check_node_set_property_against_constraints(node_id, &k, &prop)
1536 {
1537 return Err(ExecutorError::ConstraintViolation(msg));
1538 }
1539 self.ctx.storage.set_node_property(node_id, k, prop);
1540 }
1541 }
1542 EntityTarget::Relationship(rel_id) => {
1543 for (k, v) in map {
1544 let prop = lora_value_to_property(v)
1545 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1546 if let Err(msg) = self
1547 .ctx
1548 .storage
1549 .check_relationship_set_property_against_constraints(rel_id, &k, &prop)
1550 {
1551 return Err(ExecutorError::ConstraintViolation(msg));
1552 }
1553 self.ctx.storage.set_relationship_property(rel_id, k, prop);
1554 }
1555 }
1556 }
1557 Ok(())
1558 }
1559
1560 pub(crate) fn apply_create_pattern(
1561 &mut self,
1562 row: &mut Row,
1563 pattern: &ResolvedPattern,
1564 ) -> ExecResult<()> {
1565 for part in &pattern.parts {
1566 self.apply_create_pattern_part(row, part)?;
1567 }
1568 Ok(())
1569 }
1570
1571 pub(crate) fn apply_write_op(&mut self, op: &PhysicalOp, row: &mut Row) -> ExecResult<()> {
1577 match op {
1578 PhysicalOp::Create(c) => self.apply_create_pattern(row, &c.pattern),
1579 PhysicalOp::Set(s) => {
1580 for item in &s.items {
1581 self.apply_set_item(row, item)?;
1582 }
1583 Ok(())
1584 }
1585 PhysicalOp::Delete(d) => {
1586 let detach = d.detach;
1587 for expr in &d.expressions {
1588 let value = {
1589 let eval_ctx = EvalContext {
1590 storage: &*self.ctx.storage,
1591 params: &self.ctx.params,
1592 };
1593 eval_expr(expr, row, &eval_ctx)
1594 };
1595 self.delete_value(value, detach)?;
1596 }
1597 Ok(())
1598 }
1599 PhysicalOp::Remove(r) => {
1600 for item in &r.items {
1601 self.apply_remove_item(row, item)?;
1602 }
1603 Ok(())
1604 }
1605 PhysicalOp::Merge(m) => {
1606 let already_bound = self.pattern_part_is_bound(row, &m.pattern_part);
1607 let matched = if already_bound {
1608 true
1609 } else {
1610 self.try_match_merge_pattern(row, &m.pattern_part)?
1611 };
1612 if !matched {
1613 self.apply_create_pattern_part(row, &m.pattern_part)?;
1614 }
1615 for action in &m.actions {
1616 if action.on_match == matched {
1617 for item in &action.set.items {
1618 self.apply_set_item(row, item)?;
1619 }
1620 }
1621 }
1622 Ok(())
1623 }
1624 other => Err(ExecutorError::RuntimeError(format!(
1625 "apply_write_op called on non-write op: {other:?}"
1626 ))),
1627 }
1628 }
1629
1630 fn apply_create_pattern_part(
1631 &mut self,
1632 row: &mut Row,
1633 part: &ResolvedPatternPart,
1634 ) -> ExecResult<()> {
1635 if part.binding.is_some() {
1636 trace!("create pattern part has path binding; path materialization not implemented");
1637 }
1638
1639 let _ = self.apply_create_pattern_element(row, &part.element)?;
1640 Ok(())
1641 }
1642
1643 fn apply_create_pattern_element(
1644 &mut self,
1645 row: &mut Row,
1646 element: &ResolvedPatternElement,
1647 ) -> ExecResult<Option<LoraValue>> {
1648 match element {
1649 ResolvedPatternElement::Node {
1650 var,
1651 labels,
1652 properties,
1653 } => {
1654 let node_id =
1655 self.materialize_node_pattern(row, *var, labels, properties.as_ref())?;
1656 Ok(Some(LoraValue::Node(node_id)))
1657 }
1658
1659 ResolvedPatternElement::NodeChain { head, chain } => {
1660 let mut current_node_id = self.materialize_node_pattern(
1661 row,
1662 head.var,
1663 &head.labels,
1664 head.properties.as_ref(),
1665 )?;
1666
1667 for link in chain {
1668 let next_node_id = self.materialize_node_pattern(
1669 row,
1670 link.node.var,
1671 &link.node.labels,
1672 link.node.properties.as_ref(),
1673 )?;
1674
1675 let _ = self.materialize_relationship_pattern(
1676 row,
1677 current_node_id,
1678 next_node_id,
1679 &link.rel,
1680 )?;
1681
1682 current_node_id = next_node_id;
1683 }
1684
1685 Ok(Some(LoraValue::Node(current_node_id)))
1686 }
1687
1688 ResolvedPatternElement::ShortestPath { .. } => {
1689 Ok(None)
1691 }
1692 }
1693 }
1694
1695 fn pattern_part_is_bound(&self, row: &Row, part: &ResolvedPatternPart) -> bool {
1696 match &part.element {
1697 ResolvedPatternElement::Node { var, .. } => var.and_then(|v| row.get(v)).is_some(),
1698
1699 ResolvedPatternElement::ShortestPath { .. } => false,
1700
1701 ResolvedPatternElement::NodeChain { head, chain } => {
1702 let head_ok = head.var.and_then(|v| row.get(v)).is_some();
1703
1704 let chain_ok = chain.iter().all(|link| {
1705 let node_ok = link.node.var.and_then(|v| row.get(v)).is_some();
1706 let rel_ok = match link.rel.var {
1710 Some(v) => row.get(v).is_some(),
1711 None => false,
1712 };
1713 node_ok && rel_ok
1714 });
1715
1716 head_ok && chain_ok
1717 }
1718 }
1719 }
1720
1721 fn materialize_node_pattern(
1722 &mut self,
1723 row: &mut Row,
1724 var: Option<VarId>,
1725 labels: &[Vec<String>],
1726 properties: Option<&ResolvedExpr>,
1727 ) -> ExecResult<u64> {
1728 if let Some(var_id) = var {
1729 if let Some(LoraValue::Node(id)) = row.get(var_id) {
1730 return Ok(*id);
1731 }
1732 }
1733
1734 let properties = match properties {
1735 Some(expr) => eval_properties_expr(expr, row, &*self.ctx.storage, &self.ctx.params)?,
1736 None => Properties::new(),
1737 };
1738
1739 let flat_labels = flatten_label_groups(labels);
1740 debug!("creating node with labels={flat_labels:?}");
1741 if let Err(msg) = self
1742 .ctx
1743 .storage
1744 .check_node_create_against_constraints(&flat_labels, &properties)
1745 {
1746 return Err(ExecutorError::ConstraintViolation(msg));
1747 }
1748 let created = self
1749 .ctx
1750 .storage
1751 .try_create_node(flat_labels, properties)
1752 .ok_or(ExecutorError::NodeCreateFailed)?;
1753
1754 if let Some(var_id) = var {
1755 row.insert(var_id, LoraValue::Node(created.id));
1756 }
1757
1758 Ok(created.id)
1759 }
1760
1761 fn materialize_relationship_pattern(
1762 &mut self,
1763 row: &mut Row,
1764 left_node_id: u64,
1765 right_node_id: u64,
1766 rel: &lora_analyzer::ResolvedRel,
1767 ) -> ExecResult<u64> {
1768 if let Some(var_id) = rel.var {
1769 if let Some(LoraValue::Relationship(id)) = row.get(var_id) {
1770 let id = *id;
1771 if let Some((src, dst)) = self.ctx.storage.relationship_endpoints(id) {
1772 let endpoints_match = match rel.direction {
1773 Direction::Right | Direction::Undirected => {
1774 src == left_node_id && dst == right_node_id
1775 }
1776 Direction::Left => src == right_node_id && dst == left_node_id,
1777 };
1778
1779 if endpoints_match {
1780 return Ok(id);
1781 }
1782 }
1783 }
1784 }
1785
1786 if rel.range.is_some() {
1787 return Err(ExecutorError::UnsupportedCreateRelationshipRange);
1788 }
1789
1790 let (src, dst) = match rel.direction {
1791 Direction::Right | Direction::Undirected => (left_node_id, right_node_id),
1792 Direction::Left => (right_node_id, left_node_id),
1793 };
1794
1795 let rel_type = rel
1796 .types
1797 .first()
1798 .ok_or(ExecutorError::MissingRelationshipType)?;
1799
1800 if rel_type.is_empty() {
1801 return Err(ExecutorError::MissingRelationshipType);
1802 }
1803
1804 let properties = match rel.properties.as_ref() {
1805 Some(expr) => eval_properties_expr(expr, row, &*self.ctx.storage, &self.ctx.params)?,
1806 None => Properties::new(),
1807 };
1808
1809 debug!("creating relationship: src={src}, dst={dst}, type={rel_type}");
1810
1811 if let Err(msg) = self
1812 .ctx
1813 .storage
1814 .check_relationship_create_against_constraints(rel_type, &properties)
1815 {
1816 return Err(ExecutorError::ConstraintViolation(msg));
1817 }
1818
1819 let created = self
1820 .ctx
1821 .storage
1822 .create_relationship(src, dst, rel_type, properties)
1823 .ok_or_else(|| ExecutorError::RelationshipCreateFailed {
1824 src,
1825 dst,
1826 rel_type: rel_type.clone(),
1827 })?;
1828
1829 if let Some(var_id) = rel.var {
1830 row.insert(var_id, LoraValue::Relationship(created.id));
1831 }
1832
1833 Ok(created.id)
1834 }
1835}