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