1use crate::errors::{value_kind, ExecResult, ExecutorError};
14use crate::eval::{clear_eval_error, eval_expr, EvalContext};
15use crate::value::{lora_value_to_property, LoraPath, 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 matched = this.match_merge_pattern(row, &op.pattern_part)?;
914 if !matched {
915 this.apply_create_pattern_part(row, &op.pattern_part)?;
916 }
917 for action in &op.actions {
918 if action.on_match == matched {
919 for item in &action.set.items {
920 this.apply_set_item(row, item)?;
921 }
922 }
923 }
924 Ok(())
925 });
926 }
927
928 let input_rows = self.execute_node(plan, op.input)?;
929 let mut out = Vec::with_capacity(input_rows.len());
930
931 for mut row in input_rows {
932 let matched = self.match_merge_pattern(&mut row, &op.pattern_part)?;
933
934 if !matched {
935 self.apply_create_pattern_part(&mut row, &op.pattern_part)?;
936 }
937
938 for action in &op.actions {
939 if action.on_match == matched {
940 for item in &action.set.items {
941 self.apply_set_item(&row, item)?;
942 }
943 }
944 }
945
946 out.push(row);
947 }
948
949 Ok(out)
950 }
951
952 fn match_merge_pattern(&self, row: &mut Row, part: &ResolvedPatternPart) -> ExecResult<bool> {
957 if self.pattern_part_is_bound(row, part)? {
958 if let (Some(var), Some(path)) = (part.binding, bound_pattern_path(row, part)) {
959 row.insert(var, LoraValue::Path(path));
960 }
961 return Ok(true);
962 }
963 self.try_match_merge_pattern(row, part)
964 }
965
966 fn try_match_merge_pattern(
971 &self,
972 row: &mut Row,
973 part: &ResolvedPatternPart,
974 ) -> ExecResult<bool> {
975 match &part.element {
976 ResolvedPatternElement::Node {
977 var,
978 labels,
979 properties,
980 } => {
981 let expected_props = self.merge_expected_props(properties.as_ref(), row);
982 let Some(id) = self
983 .merge_node_candidates(labels, &expected_props)
984 .into_iter()
985 .find(|&id| self.merge_node_matches(id, labels, &expected_props))
986 else {
987 return Ok(false);
988 };
989 if let Some(var_id) = var {
990 row.insert(*var_id, LoraValue::Node(id));
991 }
992 if let Some(path_var) = part.binding {
993 let path = LoraPath {
994 nodes: vec![id],
995 rels: Vec::new(),
996 };
997 row.insert(path_var, LoraValue::Path(path));
998 }
999 Ok(true)
1000 }
1001
1002 ResolvedPatternElement::ShortestPath { .. } => {
1003 Ok(false)
1005 }
1006
1007 ResolvedPatternElement::NodeChain { head, chain } => {
1008 let head_candidates = match head.var.and_then(|v| row.get(v)) {
1011 Some(LoraValue::Node(id)) => vec![*id],
1012 _ => {
1013 let expected = self.merge_expected_props(head.properties.as_ref(), row);
1014 self.merge_node_candidates(&head.labels, &expected)
1015 .into_iter()
1016 .filter(|&id| self.merge_node_matches(id, &head.labels, &expected))
1017 .collect()
1018 }
1019 };
1020
1021 for head_id in head_candidates {
1022 let mut trial = row.clone();
1023 if let Some(var_id) = head.var {
1024 trial.insert(var_id, LoraValue::Node(head_id));
1025 }
1026 let mut walked = Vec::with_capacity(chain.len());
1027 if self.match_merge_chain(&mut trial, head_id, chain, &mut walked) {
1028 if let Some(path_var) = part.binding {
1029 let path = LoraPath {
1030 nodes: std::iter::once(head_id)
1031 .chain(walked.iter().map(|&(_, node)| node))
1032 .collect(),
1033 rels: walked.iter().map(|&(rel, _)| rel).collect(),
1034 };
1035 trial.insert(path_var, LoraValue::Path(path));
1036 }
1037 *row = trial;
1038 return Ok(true);
1039 }
1040 }
1041 Ok(false)
1042 }
1043 }
1044 }
1045
1046 fn match_merge_chain(
1052 &self,
1053 row: &mut Row,
1054 current: NodeId,
1055 chain: &[lora_analyzer::ResolvedChain],
1056 walked: &mut Vec<(u64, NodeId)>,
1057 ) -> bool {
1058 let Some((step, rest)) = chain.split_first() else {
1059 return true;
1060 };
1061
1062 let bound_dst = match step.node.var.and_then(|v| row.get(v)) {
1063 Some(LoraValue::Node(id)) => Some(*id),
1064 _ => None,
1065 };
1066 let bound_rel = match step.rel.var.and_then(|v| row.get(v)) {
1067 Some(LoraValue::Relationship(id)) => Some(*id),
1068 _ => None,
1069 };
1070 let expected_node = self.merge_expected_props(step.node.properties.as_ref(), row);
1071 let expected_rel = self.merge_expected_props(step.rel.properties.as_ref(), row);
1072
1073 let edges = self
1074 .ctx
1075 .storage
1076 .expand_ids(current, step.rel.direction, &step.rel.types);
1077 for (rel_id, node_id) in edges {
1078 if bound_dst.is_some_and(|id| id != node_id)
1079 || bound_rel.is_some_and(|id| id != rel_id)
1080 || walked.iter().any(|&(used, _)| used == rel_id)
1081 {
1082 continue;
1083 }
1084 if !self.merge_node_matches(node_id, &step.node.labels, &expected_node) {
1085 continue;
1086 }
1087 if let Some(LoraValue::Map(expected_map)) = &expected_rel {
1088 let rel_ok = self
1089 .ctx
1090 .storage
1091 .with_relationship(rel_id, |rel_rec| {
1092 expected_map.iter().all(|(key, expected_val)| {
1093 rel_rec
1094 .properties
1095 .get(key.as_str())
1096 .map(|actual| value_matches_property_value(expected_val, actual))
1097 .unwrap_or(false)
1098 })
1099 })
1100 .unwrap_or(false);
1101 if !rel_ok {
1102 continue;
1103 }
1104 }
1105
1106 let mut next = row.clone();
1107 if let Some(rel_var) = step.rel.var {
1108 next.insert(rel_var, LoraValue::Relationship(rel_id));
1109 }
1110 if let Some(node_var) = step.node.var {
1111 next.insert(node_var, LoraValue::Node(node_id));
1112 }
1113 walked.push((rel_id, node_id));
1114 if self.match_merge_chain(&mut next, node_id, rest, walked) {
1115 *row = next;
1116 return true;
1117 }
1118 walked.pop();
1119 }
1120 false
1121 }
1122
1123 fn merge_expected_props(
1124 &self,
1125 properties: Option<&ResolvedExpr>,
1126 row: &Row,
1127 ) -> Option<LoraValue> {
1128 let eval_ctx = EvalContext {
1129 storage: &*self.ctx.storage,
1130 params: &self.ctx.params,
1131 };
1132 properties.map(|e| eval_expr(e, row, &eval_ctx))
1133 }
1134
1135 fn merge_node_candidates(
1140 &self,
1141 labels: &[Vec<String>],
1142 expected_props: &Option<LoraValue>,
1143 ) -> Vec<NodeId> {
1144 let indexed = match expected_props {
1145 Some(LoraValue::Map(expected)) => {
1146 merge_candidates_from_index(&*self.ctx.storage, labels, expected)
1147 }
1148 _ => None,
1149 };
1150 match indexed {
1151 Some(ids) => ids,
1152 None if labels.is_empty() => self.ctx.storage.all_node_ids(),
1153 None => scan_node_ids_for_label_groups(&*self.ctx.storage, labels),
1154 }
1155 }
1156
1157 fn merge_node_matches(
1158 &self,
1159 id: NodeId,
1160 labels: &[Vec<String>],
1161 expected_props: &Option<LoraValue>,
1162 ) -> bool {
1163 self.ctx
1164 .storage
1165 .with_node(id, |node| {
1166 if !node_matches_label_groups(&node.labels, labels) {
1167 return false;
1168 }
1169 if let Some(LoraValue::Map(expected)) = expected_props {
1170 return expected.iter().all(|(key, expected_value)| {
1171 node.properties
1172 .get(key.as_str())
1173 .map(|actual| value_matches_property_value(expected_value, actual))
1174 .unwrap_or(false)
1175 });
1176 }
1177 true
1178 })
1179 .unwrap_or(false)
1180 }
1181
1182 fn exec_delete(&mut self, plan: &PhysicalPlan, op: &DeleteExec) -> ExecResult<Vec<Row>> {
1183 let input_rows = self.execute_node(plan, op.input)?;
1184 let mut targets = BTreeSet::new();
1185
1186 for row in &input_rows {
1187 for expr in &op.expressions {
1188 let value = {
1189 let eval_ctx = EvalContext {
1190 storage: &*self.ctx.storage,
1191 params: &self.ctx.params,
1192 };
1193 eval_expr(expr, row, &eval_ctx)
1194 };
1195 self.collect_delete_targets(&value, &mut targets)?;
1196 }
1197 }
1198
1199 self.validate_delete_targets(&targets, op.detach)?;
1200
1201 for target in &targets {
1202 if let DeleteTarget::Relationship(_) = target {
1203 self.delete_target(*target, op.detach)?;
1204 }
1205 }
1206 for target in targets {
1207 if let DeleteTarget::Node(_) = target {
1208 self.delete_target(target, op.detach)?;
1209 }
1210 }
1211
1212 Ok(input_rows)
1213 }
1214
1215 fn exec_set(&mut self, plan: &PhysicalPlan, op: &SetExec) -> ExecResult<Vec<Row>> {
1216 if crate::pull::subtree_is_fully_streaming(plan, op.input) {
1217 return self.streaming_apply(plan, op.input, |this, row| {
1218 for item in &op.items {
1219 this.apply_set_item(row, item)?;
1220 }
1221 Ok(())
1222 });
1223 }
1224
1225 let input_rows = self.execute_node(plan, op.input)?;
1226
1227 for row in &input_rows {
1228 for item in &op.items {
1229 self.apply_set_item(row, item)?;
1230 }
1231 }
1232
1233 Ok(input_rows)
1234 }
1235
1236 fn exec_foreach(&mut self, plan: &PhysicalPlan, op: &ForeachExec) -> ExecResult<Vec<Row>> {
1244 let input_rows = self.execute_node(plan, op.input)?;
1245 let mut out = Vec::with_capacity(input_rows.len());
1246
1247 for row in input_rows {
1248 let list_value = {
1249 let eval_ctx = EvalContext {
1250 storage: &*self.ctx.storage,
1251 params: &self.ctx.params,
1252 };
1253 eval_expr(&op.list, &row, &eval_ctx)
1254 };
1255
1256 let elements: Vec<LoraValue> = match list_value {
1257 LoraValue::List(items) => items,
1258 LoraValue::Null => Vec::new(),
1259 other => {
1260 return Err(ExecutorError::RuntimeError(format!(
1261 "FOREACH expects a list, got {}",
1262 value_kind(&other)
1263 )));
1264 }
1265 };
1266
1267 for element in elements {
1268 let mut iter_row = row.clone();
1271 iter_row.insert(op.variable, element);
1272 for clause in &op.body {
1273 self.apply_foreach_body_clause(&mut iter_row, clause)?;
1274 }
1275 }
1276
1277 out.push(row);
1278 }
1279
1280 Ok(out)
1281 }
1282
1283 fn apply_foreach_body_clause(
1288 &mut self,
1289 row: &mut Row,
1290 clause: &lora_analyzer::ResolvedClause,
1291 ) -> ExecResult<()> {
1292 use lora_analyzer::ResolvedClause;
1293 match clause {
1294 ResolvedClause::Create(c) => self.apply_create_pattern(row, &c.pattern),
1295 ResolvedClause::Set(s) => {
1296 for item in &s.items {
1297 self.apply_set_item(row, item)?;
1298 }
1299 Ok(())
1300 }
1301 ResolvedClause::Remove(r) => {
1302 for item in &r.items {
1303 self.apply_remove_item(row, item)?;
1304 }
1305 Ok(())
1306 }
1307 ResolvedClause::Delete(d) => {
1308 let detach = d.detach;
1309 for expr in &d.expressions {
1310 let value = {
1311 let eval_ctx = EvalContext {
1312 storage: &*self.ctx.storage,
1313 params: &self.ctx.params,
1314 };
1315 eval_expr(expr, row, &eval_ctx)
1316 };
1317 self.delete_value(value, detach)?;
1318 }
1319 Ok(())
1320 }
1321 ResolvedClause::Merge(m) => {
1322 let matched = self.match_merge_pattern(row, &m.pattern_part)?;
1323 if !matched {
1324 self.apply_create_pattern_part(row, &m.pattern_part)?;
1325 }
1326 for action in &m.actions {
1327 if action.on_match == matched {
1328 for item in &action.set.items {
1329 self.apply_set_item(row, item)?;
1330 }
1331 }
1332 }
1333 Ok(())
1334 }
1335 ResolvedClause::Foreach(nested) => {
1336 let list_value = {
1337 let eval_ctx = EvalContext {
1338 storage: &*self.ctx.storage,
1339 params: &self.ctx.params,
1340 };
1341 eval_expr(&nested.list, row, &eval_ctx)
1342 };
1343
1344 let elements: Vec<LoraValue> = match list_value {
1345 LoraValue::List(items) => items,
1346 LoraValue::Null => Vec::new(),
1347 other => {
1348 return Err(ExecutorError::RuntimeError(format!(
1349 "FOREACH expects a list, got {}",
1350 value_kind(&other)
1351 )));
1352 }
1353 };
1354
1355 for element in elements {
1356 let mut iter_row = row.clone();
1357 iter_row.insert(nested.variable, element);
1358 for inner in &nested.body {
1359 self.apply_foreach_body_clause(&mut iter_row, inner)?;
1360 }
1361 }
1362
1363 Ok(())
1364 }
1365 other => Err(ExecutorError::RuntimeError(format!(
1366 "FOREACH body may only contain updating clauses, got {:?}",
1367 std::mem::discriminant(other)
1368 ))),
1369 }
1370 }
1371
1372 fn exec_remove(&mut self, plan: &PhysicalPlan, op: &RemoveExec) -> ExecResult<Vec<Row>> {
1373 if crate::pull::subtree_is_fully_streaming(plan, op.input) {
1374 return self.streaming_apply(plan, op.input, |this, row| {
1375 for item in &op.items {
1376 this.apply_remove_item(row, item)?;
1377 }
1378 Ok(())
1379 });
1380 }
1381
1382 let input_rows = self.execute_node(plan, op.input)?;
1383
1384 for row in &input_rows {
1385 for item in &op.items {
1386 self.apply_remove_item(row, item)?;
1387 }
1388 }
1389
1390 Ok(input_rows)
1391 }
1392
1393 fn apply_set_item(&mut self, row: &Row, item: &ResolvedSetItem) -> ExecResult<()> {
1394 match item {
1395 ResolvedSetItem::SetProperty { target, value } => {
1396 let new_value = {
1397 let eval_ctx = EvalContext {
1398 storage: &*self.ctx.storage,
1399 params: &self.ctx.params,
1400 };
1401 eval_expr(value, row, &eval_ctx)
1402 };
1403
1404 self.set_property_from_expr(row, target, new_value)
1405 }
1406
1407 ResolvedSetItem::SetVariable { variable, value } => {
1408 let entity_ref =
1410 row.get(*variable)
1411 .ok_or(ExecutorError::UnboundVariableForSet {
1412 var: format!("{variable:?}"),
1413 })?;
1414 let entity_target = entity_target_from_value(entity_ref)?;
1415
1416 let new_value = {
1417 let eval_ctx = EvalContext {
1418 storage: &*self.ctx.storage,
1419 params: &self.ctx.params,
1420 };
1421 eval_expr(value, row, &eval_ctx)
1422 };
1423
1424 self.overwrite_entity_target(entity_target, new_value)
1425 }
1426
1427 ResolvedSetItem::MutateVariable { variable, value } => {
1428 let entity_ref =
1429 row.get(*variable)
1430 .ok_or(ExecutorError::UnboundVariableForSet {
1431 var: format!("{variable:?}"),
1432 })?;
1433 let entity_target = entity_target_from_value(entity_ref)?;
1434
1435 let patch = {
1436 let eval_ctx = EvalContext {
1437 storage: &*self.ctx.storage,
1438 params: &self.ctx.params,
1439 };
1440 eval_expr(value, row, &eval_ctx)
1441 };
1442
1443 self.mutate_entity_target(entity_target, patch)
1444 }
1445
1446 ResolvedSetItem::SetLabels { variable, labels } => match row.get(*variable) {
1447 Some(LoraValue::Node(node_id)) => {
1448 let node_id = *node_id;
1449 for label in labels {
1450 let checked = if self.defer_existence {
1454 self.ctx
1455 .storage
1456 .check_node_add_label_deferring_existence(node_id, label)
1457 } else {
1458 self.ctx
1459 .storage
1460 .check_node_add_label_against_constraints(node_id, label)
1461 };
1462 checked.map_err(ExecutorError::ConstraintViolation)?;
1463 self.ctx.storage.add_node_label(node_id, label);
1464 }
1465 if self.defer_existence {
1466 self.pending_existence.push(EntityTarget::Node(node_id));
1467 }
1468 Ok(())
1469 }
1470 Some(other) => Err(ExecutorError::ExpectedNodeForSetLabels {
1471 found: value_kind(other),
1472 }),
1473 None => Err(ExecutorError::UnboundVariableForSet {
1474 var: format!("{variable:?}"),
1475 }),
1476 },
1477 }
1478 }
1479
1480 fn set_property_from_expr(
1481 &mut self,
1482 row: &Row,
1483 target_expr: &ResolvedExpr,
1484 new_value: LoraValue,
1485 ) -> ExecResult<()> {
1486 let ResolvedExpr::Property { expr, property } = target_expr else {
1487 return Err(ExecutorError::UnsupportedSetTarget);
1488 };
1489
1490 let owner = {
1491 let eval_ctx = EvalContext {
1492 storage: &*self.ctx.storage,
1493 params: &self.ctx.params,
1494 };
1495 eval_expr(expr, row, &eval_ctx)
1496 };
1497
1498 if matches!(new_value, LoraValue::Null) {
1500 return match owner {
1501 LoraValue::Node(node_id) => {
1502 self.remove_entity_property(EntityTarget::Node(node_id), property)
1503 }
1504 LoraValue::Relationship(rel_id) => {
1505 self.remove_entity_property(EntityTarget::Relationship(rel_id), property)
1506 }
1507 other => Err(ExecutorError::InvalidSetTarget {
1508 found: value_kind(&other),
1509 }),
1510 };
1511 }
1512
1513 match owner {
1514 LoraValue::Node(node_id) => {
1515 let prop = lora_value_to_property(new_value)
1516 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1517 if let Err(msg) = self
1518 .ctx
1519 .storage
1520 .check_node_set_property_against_constraints(node_id, property, &prop)
1521 {
1522 return Err(ExecutorError::ConstraintViolation(msg));
1523 }
1524 self.ctx
1525 .storage
1526 .set_node_property(node_id, property.clone(), prop);
1527 Ok(())
1528 }
1529 LoraValue::Relationship(rel_id) => {
1530 let prop = lora_value_to_property(new_value)
1531 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1532 if let Err(msg) = self
1533 .ctx
1534 .storage
1535 .check_relationship_set_property_against_constraints(rel_id, property, &prop)
1536 {
1537 return Err(ExecutorError::ConstraintViolation(msg));
1538 }
1539 self.ctx
1540 .storage
1541 .set_relationship_property(rel_id, property.clone(), prop);
1542 Ok(())
1543 }
1544 other => Err(ExecutorError::InvalidSetTarget {
1545 found: value_kind(&other),
1546 }),
1547 }
1548 }
1549
1550 fn remove_entity_property(&mut self, target: EntityTarget, property: &str) -> ExecResult<()> {
1555 if self.defer_existence {
1556 match target {
1557 EntityTarget::Node(node_id) => {
1558 self.ctx.storage.remove_node_property(node_id, property);
1559 }
1560 EntityTarget::Relationship(rel_id) => {
1561 self.ctx
1562 .storage
1563 .remove_relationship_property(rel_id, property);
1564 }
1565 }
1566 self.pending_existence.push(target);
1567 return Ok(());
1568 }
1569 match target {
1570 EntityTarget::Node(node_id) => {
1571 if let Err(msg) = self
1572 .ctx
1573 .storage
1574 .check_node_remove_property_against_constraints(node_id, property)
1575 {
1576 return Err(ExecutorError::ConstraintViolation(msg));
1577 }
1578 self.ctx.storage.remove_node_property(node_id, property);
1579 }
1580 EntityTarget::Relationship(rel_id) => {
1581 if let Err(msg) = self
1582 .ctx
1583 .storage
1584 .check_relationship_remove_property_against_constraints(rel_id, property)
1585 {
1586 return Err(ExecutorError::ConstraintViolation(msg));
1587 }
1588 self.ctx
1589 .storage
1590 .remove_relationship_property(rel_id, property);
1591 }
1592 }
1593 Ok(())
1594 }
1595
1596 fn remove_property_from_expr(&mut self, row: &Row, expr: &ResolvedExpr) -> ExecResult<()> {
1597 let ResolvedExpr::Property {
1598 expr: owner_expr,
1599 property,
1600 } = expr
1601 else {
1602 return Err(ExecutorError::UnsupportedRemoveTarget);
1603 };
1604
1605 let owner = {
1606 let eval_ctx = EvalContext {
1607 storage: &*self.ctx.storage,
1608 params: &self.ctx.params,
1609 };
1610 eval_expr(owner_expr, row, &eval_ctx)
1611 };
1612
1613 match owner {
1614 LoraValue::Node(node_id) => {
1615 self.remove_entity_property(EntityTarget::Node(node_id), property)
1616 }
1617 LoraValue::Relationship(rel_id) => {
1618 self.remove_entity_property(EntityTarget::Relationship(rel_id), property)
1619 }
1620 other => Err(ExecutorError::InvalidRemoveTarget {
1621 found: value_kind(&other),
1622 }),
1623 }
1624 }
1625
1626 fn overwrite_entity_target(
1627 &mut self,
1628 target: EntityTarget,
1629 new_value: LoraValue,
1630 ) -> ExecResult<()> {
1631 let LoraValue::Map(map) = new_value else {
1632 return Err(ExecutorError::ExpectedPropertyMap {
1633 found: value_kind(&new_value),
1634 });
1635 };
1636
1637 let mut props: Properties = Properties::new();
1638 for (k, v) in map {
1639 if matches!(v, LoraValue::Null) {
1641 continue;
1642 }
1643 let prop = lora_value_to_property(v)
1644 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1645 props.insert(lora_store::intern_owned(k), prop);
1646 }
1647
1648 let storage = &*self.ctx.storage;
1652 let checked = match (target, self.defer_existence) {
1653 (EntityTarget::Node(id), true) => {
1654 storage.check_node_replace_properties_deferring_existence(id, &props)
1655 }
1656 (EntityTarget::Node(id), false) => {
1657 storage.check_node_replace_properties_against_constraints(id, &props)
1658 }
1659 (EntityTarget::Relationship(id), true) => {
1660 storage.check_relationship_replace_properties_deferring_existence(id, &props)
1661 }
1662 (EntityTarget::Relationship(id), false) => {
1663 storage.check_relationship_replace_properties_against_constraints(id, &props)
1664 }
1665 };
1666 checked.map_err(ExecutorError::ConstraintViolation)?;
1667 match target {
1668 EntityTarget::Node(node_id) => {
1669 self.ctx.storage.replace_node_properties(node_id, props);
1670 }
1671 EntityTarget::Relationship(rel_id) => {
1672 self.ctx
1673 .storage
1674 .replace_relationship_properties(rel_id, props);
1675 }
1676 }
1677 if self.defer_existence {
1678 self.pending_existence.push(target);
1679 }
1680 Ok(())
1681 }
1682
1683 fn mutate_entity_target(
1684 &mut self,
1685 target: EntityTarget,
1686 patch_value: LoraValue,
1687 ) -> ExecResult<()> {
1688 let LoraValue::Map(map) = patch_value else {
1689 return Err(ExecutorError::ExpectedPropertyMap {
1690 found: value_kind(&patch_value),
1691 });
1692 };
1693
1694 match target {
1695 EntityTarget::Node(node_id) => {
1696 for (k, v) in map {
1697 if matches!(v, LoraValue::Null) {
1699 self.remove_entity_property(target, &k)?;
1700 continue;
1701 }
1702 let prop = lora_value_to_property(v)
1703 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1704 if let Err(msg) = self
1705 .ctx
1706 .storage
1707 .check_node_set_property_against_constraints(node_id, &k, &prop)
1708 {
1709 return Err(ExecutorError::ConstraintViolation(msg));
1710 }
1711 self.ctx.storage.set_node_property(node_id, k, prop);
1712 }
1713 }
1714 EntityTarget::Relationship(rel_id) => {
1715 for (k, v) in map {
1716 if matches!(v, LoraValue::Null) {
1717 self.remove_entity_property(target, &k)?;
1718 continue;
1719 }
1720 let prop = lora_value_to_property(v)
1721 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1722 if let Err(msg) = self
1723 .ctx
1724 .storage
1725 .check_relationship_set_property_against_constraints(rel_id, &k, &prop)
1726 {
1727 return Err(ExecutorError::ConstraintViolation(msg));
1728 }
1729 self.ctx.storage.set_relationship_property(rel_id, k, prop);
1730 }
1731 }
1732 }
1733 Ok(())
1734 }
1735
1736 pub(crate) fn apply_create_pattern(
1737 &mut self,
1738 row: &mut Row,
1739 pattern: &ResolvedPattern,
1740 ) -> ExecResult<()> {
1741 for part in &pattern.parts {
1742 self.apply_create_pattern_part(row, part)?;
1743 }
1744 Ok(())
1745 }
1746
1747 pub(crate) fn apply_write_op(&mut self, op: &PhysicalOp, row: &mut Row) -> ExecResult<()> {
1753 match op {
1754 PhysicalOp::Create(c) => self.apply_create_pattern(row, &c.pattern),
1755 PhysicalOp::Set(s) => {
1756 for item in &s.items {
1757 self.apply_set_item(row, item)?;
1758 }
1759 Ok(())
1760 }
1761 PhysicalOp::Delete(d) => {
1762 let detach = d.detach;
1763 for expr in &d.expressions {
1764 let value = {
1765 let eval_ctx = EvalContext {
1766 storage: &*self.ctx.storage,
1767 params: &self.ctx.params,
1768 };
1769 eval_expr(expr, row, &eval_ctx)
1770 };
1771 self.delete_value(value, detach)?;
1772 }
1773 Ok(())
1774 }
1775 PhysicalOp::Remove(r) => {
1776 for item in &r.items {
1777 self.apply_remove_item(row, item)?;
1778 }
1779 Ok(())
1780 }
1781 PhysicalOp::Merge(m) => {
1782 let matched = self.match_merge_pattern(row, &m.pattern_part)?;
1783 if !matched {
1784 self.apply_create_pattern_part(row, &m.pattern_part)?;
1785 }
1786 for action in &m.actions {
1787 if action.on_match == matched {
1788 for item in &action.set.items {
1789 self.apply_set_item(row, item)?;
1790 }
1791 }
1792 }
1793 Ok(())
1794 }
1795 other => Err(ExecutorError::RuntimeError(format!(
1796 "apply_write_op called on non-write op: {other:?}"
1797 ))),
1798 }
1799 }
1800
1801 fn apply_create_pattern_part(
1804 &mut self,
1805 row: &mut Row,
1806 part: &ResolvedPatternPart,
1807 ) -> ExecResult<()> {
1808 let path = self.apply_create_pattern_element(row, &part.element)?;
1809 if let (Some(var), Some(path)) = (part.binding, path) {
1810 row.insert(var, LoraValue::Path(path));
1811 }
1812 Ok(())
1813 }
1814
1815 fn apply_create_pattern_element(
1816 &mut self,
1817 row: &mut Row,
1818 element: &ResolvedPatternElement,
1819 ) -> ExecResult<Option<LoraPath>> {
1820 match element {
1821 ResolvedPatternElement::Node {
1822 var,
1823 labels,
1824 properties,
1825 } => {
1826 let node_id =
1827 self.materialize_node_pattern(row, *var, labels, properties.as_ref())?;
1828 Ok(Some(LoraPath {
1829 nodes: vec![node_id],
1830 rels: Vec::new(),
1831 }))
1832 }
1833
1834 ResolvedPatternElement::NodeChain { head, chain } => {
1835 let mut current_node_id = self.materialize_node_pattern(
1836 row,
1837 head.var,
1838 &head.labels,
1839 head.properties.as_ref(),
1840 )?;
1841 let mut path = LoraPath {
1842 nodes: Vec::with_capacity(chain.len() + 1),
1843 rels: Vec::with_capacity(chain.len()),
1844 };
1845 path.nodes.push(current_node_id);
1846
1847 for link in chain {
1848 let next_node_id = self.materialize_node_pattern(
1849 row,
1850 link.node.var,
1851 &link.node.labels,
1852 link.node.properties.as_ref(),
1853 )?;
1854
1855 let rel_id = self.materialize_relationship_pattern(
1856 row,
1857 current_node_id,
1858 next_node_id,
1859 &link.rel,
1860 )?;
1861 path.rels.push(rel_id);
1862 path.nodes.push(next_node_id);
1863
1864 current_node_id = next_node_id;
1865 }
1866
1867 Ok(Some(path))
1868 }
1869
1870 ResolvedPatternElement::ShortestPath { .. } => {
1871 Ok(None)
1873 }
1874 }
1875 }
1876
1877 fn pattern_part_is_bound(&self, row: &Row, part: &ResolvedPatternPart) -> ExecResult<bool> {
1883 fn node_bound(row: &Row, var: Option<VarId>) -> ExecResult<bool> {
1884 let Some(var) = var else { return Ok(false) };
1885 match row.get(var) {
1886 None => Ok(false),
1887 Some(LoraValue::Node(_)) => Ok(true),
1888 Some(other) => Err(ExecutorError::ExpectedNodeForCreate {
1889 var: bound_var_name(row, var),
1890 found: value_kind(other),
1891 }),
1892 }
1893 }
1894 fn rel_bound(row: &Row, var: Option<VarId>) -> ExecResult<bool> {
1895 let Some(var) = var else { return Ok(false) };
1899 match row.get(var) {
1900 None => Ok(false),
1901 Some(LoraValue::Relationship(_)) => Ok(true),
1902 Some(other) => Err(ExecutorError::ExpectedRelationshipForCreate {
1903 var: bound_var_name(row, var),
1904 found: value_kind(other),
1905 }),
1906 }
1907 }
1908
1909 match &part.element {
1910 ResolvedPatternElement::Node { var, .. } => node_bound(row, *var),
1911
1912 ResolvedPatternElement::ShortestPath { .. } => Ok(false),
1913
1914 ResolvedPatternElement::NodeChain { head, chain } => {
1915 let mut all_bound = node_bound(row, head.var)?;
1918 for link in chain {
1919 let node_ok = node_bound(row, link.node.var)?;
1920 let rel_ok = rel_bound(row, link.rel.var)?;
1921 all_bound &= node_ok && rel_ok;
1922 }
1923 Ok(all_bound)
1924 }
1925 }
1926 }
1927
1928 fn materialize_node_pattern(
1929 &mut self,
1930 row: &mut Row,
1931 var: Option<VarId>,
1932 labels: &[Vec<String>],
1933 properties: Option<&ResolvedExpr>,
1934 ) -> ExecResult<u64> {
1935 if let Some(var_id) = var {
1936 match row.get(var_id) {
1937 Some(LoraValue::Node(id)) => return Ok(*id),
1938 Some(other) => {
1943 return Err(ExecutorError::ExpectedNodeForCreate {
1944 var: bound_var_name(row, var_id),
1945 found: value_kind(other),
1946 });
1947 }
1948 None => {}
1949 }
1950 }
1951
1952 let properties = match properties {
1953 Some(expr) => eval_properties_expr(expr, row, &*self.ctx.storage, &self.ctx.params)?,
1954 None => Properties::new(),
1955 };
1956
1957 let flat_labels = flatten_label_groups(labels);
1958 debug!("creating node with labels={flat_labels:?}");
1959 let checked = if self.defer_existence {
1960 self.ctx
1961 .storage
1962 .check_node_create_deferring_existence(&flat_labels, &properties)
1963 } else {
1964 self.ctx
1965 .storage
1966 .check_node_create_against_constraints(&flat_labels, &properties)
1967 };
1968 checked.map_err(ExecutorError::ConstraintViolation)?;
1969 let created = self
1970 .ctx
1971 .storage
1972 .try_create_node(flat_labels, properties)
1973 .ok_or(ExecutorError::NodeCreateFailed)?;
1974 if self.defer_existence {
1975 self.pending_existence.push(EntityTarget::Node(created.id));
1976 }
1977
1978 if let Some(var_id) = var {
1979 row.insert(var_id, LoraValue::Node(created.id));
1980 }
1981
1982 Ok(created.id)
1983 }
1984
1985 fn materialize_relationship_pattern(
1986 &mut self,
1987 row: &mut Row,
1988 left_node_id: u64,
1989 right_node_id: u64,
1990 rel: &lora_analyzer::ResolvedRel,
1991 ) -> ExecResult<u64> {
1992 if let Some(var_id) = rel.var {
1993 if let Some(other) = row
1994 .get(var_id)
1995 .filter(|v| !matches!(v, LoraValue::Relationship(_)))
1996 {
1997 return Err(ExecutorError::ExpectedRelationshipForCreate {
1998 var: bound_var_name(row, var_id),
1999 found: value_kind(other),
2000 });
2001 }
2002 if let Some(LoraValue::Relationship(id)) = row.get(var_id) {
2003 let id = *id;
2004 if let Some((src, dst)) = self.ctx.storage.relationship_endpoints(id) {
2005 let endpoints_match = match rel.direction {
2006 Direction::Right | Direction::Undirected => {
2007 src == left_node_id && dst == right_node_id
2008 }
2009 Direction::Left => src == right_node_id && dst == left_node_id,
2010 };
2011
2012 if endpoints_match {
2013 return Ok(id);
2014 }
2015 }
2016 }
2017 }
2018
2019 if rel.range.is_some() {
2020 return Err(ExecutorError::UnsupportedCreateRelationshipRange);
2021 }
2022
2023 let (src, dst) = match rel.direction {
2024 Direction::Right | Direction::Undirected => (left_node_id, right_node_id),
2025 Direction::Left => (right_node_id, left_node_id),
2026 };
2027
2028 let rel_type = rel
2029 .types
2030 .first()
2031 .ok_or(ExecutorError::MissingRelationshipType)?;
2032
2033 if rel_type.is_empty() {
2034 return Err(ExecutorError::MissingRelationshipType);
2035 }
2036
2037 let properties = match rel.properties.as_ref() {
2038 Some(expr) => eval_properties_expr(expr, row, &*self.ctx.storage, &self.ctx.params)?,
2039 None => Properties::new(),
2040 };
2041
2042 debug!("creating relationship: src={src}, dst={dst}, type={rel_type}");
2043
2044 let checked = if self.defer_existence {
2045 self.ctx
2046 .storage
2047 .check_relationship_create_deferring_existence(rel_type, &properties)
2048 } else {
2049 self.ctx
2050 .storage
2051 .check_relationship_create_against_constraints(rel_type, &properties)
2052 };
2053 checked.map_err(ExecutorError::ConstraintViolation)?;
2054
2055 let created = self
2056 .ctx
2057 .storage
2058 .create_relationship(src, dst, rel_type, properties)
2059 .ok_or_else(|| ExecutorError::RelationshipCreateFailed {
2060 src,
2061 dst,
2062 rel_type: rel_type.clone(),
2063 })?;
2064 if self.defer_existence {
2065 self.pending_existence
2066 .push(EntityTarget::Relationship(created.id));
2067 }
2068
2069 if let Some(var_id) = rel.var {
2070 row.insert(var_id, LoraValue::Relationship(created.id));
2071 }
2072
2073 Ok(created.id)
2074 }
2075}
2076
2077fn bound_pattern_path(row: &Row, part: &ResolvedPatternPart) -> Option<LoraPath> {
2080 let node = |var: Option<VarId>| match var.and_then(|v| row.get(v)) {
2081 Some(LoraValue::Node(id)) => Some(*id),
2082 _ => None,
2083 };
2084 match &part.element {
2085 ResolvedPatternElement::Node { var, .. } => Some(LoraPath {
2086 nodes: vec![node(*var)?],
2087 rels: Vec::new(),
2088 }),
2089 ResolvedPatternElement::NodeChain { head, chain } => {
2090 let mut path = LoraPath {
2091 nodes: vec![node(head.var)?],
2092 rels: Vec::with_capacity(chain.len()),
2093 };
2094 for link in chain {
2095 match link.rel.var.and_then(|v| row.get(v)) {
2096 Some(LoraValue::Relationship(id)) => path.rels.push(*id),
2097 _ => return None,
2098 }
2099 path.nodes.push(node(link.node.var)?);
2100 }
2101 Some(path)
2102 }
2103 ResolvedPatternElement::ShortestPath { .. } => None,
2104 }
2105}
2106
2107pub(crate) fn plan_defers_existence(plan: &PhysicalPlan) -> bool {
2112 let mut creates = 0;
2113 let mut deletes = false;
2114 for op in &plan.nodes {
2115 match op {
2116 PhysicalOp::Create(_) => creates += 1,
2117 PhysicalOp::Delete(_) => deletes = true,
2119 PhysicalOp::Merge(_)
2120 | PhysicalOp::Set(_)
2121 | PhysicalOp::Remove(_)
2122 | PhysicalOp::Foreach(_) => return true,
2123 _ => {}
2124 }
2125 }
2126 creates > 1 || creates == 1 && deletes
2127}
2128
2129pub(crate) fn plan_ends_in_write(plan: &PhysicalPlan) -> bool {
2134 match &plan.nodes[plan.root] {
2135 PhysicalOp::Create(_)
2136 | PhysicalOp::Merge(_)
2137 | PhysicalOp::Set(_)
2138 | PhysicalOp::Delete(_)
2139 | PhysicalOp::Remove(_)
2140 | PhysicalOp::Foreach(_) => true,
2141 PhysicalOp::CallSubquery(op) => op.new_vars.is_empty(),
2143 _ => false,
2144 }
2145}
2146
2147fn merge_candidates_from_index<S: lora_store::GraphStorage>(
2154 storage: &S,
2155 labels: &[Vec<String>],
2156 expected: &std::collections::BTreeMap<String, LoraValue>,
2157) -> Option<Vec<lora_store::NodeId>> {
2158 use lora_store::PropertyValue;
2159
2160 let label = match labels {
2163 [group] if group.len() == 1 => Some(group[0].as_str()),
2164 _ => None,
2165 };
2166 let (key, value) = expected.iter().find(|(_, v)| {
2167 matches!(
2168 v,
2169 LoraValue::String(_) | LoraValue::Bool(_) | LoraValue::Int(_)
2170 ) || matches!(v, LoraValue::Float(f) if f.is_finite() && f.abs() < 9_007_199_254_740_992.0)
2171 })?;
2172 let images: Vec<PropertyValue> = match value {
2173 LoraValue::String(s) => vec![PropertyValue::String(s.clone())],
2174 LoraValue::Bool(b) => vec![PropertyValue::Bool(*b)],
2175 LoraValue::Int(i) => vec![PropertyValue::Int(*i), PropertyValue::Float(*i as f64)],
2176 LoraValue::Float(f) => {
2177 let mut v = vec![PropertyValue::Float(*f)];
2178 if f.fract() == 0.0 {
2179 v.push(PropertyValue::Int(*f as i64));
2180 }
2181 v
2182 }
2183 _ => return None,
2184 };
2185 let mut ids: Vec<lora_store::NodeId> = images
2186 .iter()
2187 .flat_map(|image| storage.find_node_ids_by_property(label, key, image))
2188 .collect();
2189 ids.sort_unstable();
2190 ids.dedup();
2191 Some(ids)
2192}
2193
2194fn bound_var_name(row: &Row, var: VarId) -> String {
2196 row.iter_named()
2197 .find(|(key, _, _)| **key == var)
2198 .map(|(_, name, _)| name.into_owned())
2199 .unwrap_or_else(|| format!("{var:?}"))
2200}