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