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;
28use std::time::Instant;
29use tracing::{debug, error, trace};
30
31use super::aggregate_rows;
32use super::helpers::{
33 build_path_value, check_deadline_at, dedup_rows, eval_properties_expr, expand_rows,
34 expand_var_len_rows, filter_rows_checked, filter_shortest_paths, flatten_label_groups,
35 hydrate_node_record, hydrate_relationship_record, limit_rows, node_by_label_scan_rows,
36 node_by_property_scan_rows, node_matches_label_groups, node_scan_rows, plan_may_need_hydration,
37 project_rows_checked, scan_node_ids_for_label_groups, unwind_rows,
38 value_matches_property_value,
39};
40use super::optional_match_rows;
41use super::sort_rows_with_top_k;
42
43#[derive(Clone, Copy)]
47enum EntityTarget {
48 Node(NodeId),
49 Relationship(u64),
50}
51
52fn entity_target_from_value(value: &LoraValue) -> ExecResult<EntityTarget> {
53 match value {
54 LoraValue::Node(id) => Ok(EntityTarget::Node(*id)),
55 LoraValue::Relationship(id) => Ok(EntityTarget::Relationship(*id)),
56 other => Err(ExecutorError::InvalidSetTarget {
57 found: value_kind(other),
58 }),
59 }
60}
61
62pub struct MutableExecutionContext<'a, S: GraphStorageMut> {
63 pub storage: &'a mut S,
64 pub params: BTreeMap<String, LoraValue>,
65}
66
67pub struct MutableExecutor<'a, S: GraphStorageMut> {
68 ctx: MutableExecutionContext<'a, S>,
69 deadline: Option<Instant>,
70}
71
72impl<'a, S: GraphStorageMut> MutableExecutor<'a, S> {
73 pub fn new(ctx: MutableExecutionContext<'a, S>) -> Self {
74 Self {
75 ctx,
76 deadline: None,
77 }
78 }
79
80 pub fn with_deadline(ctx: MutableExecutionContext<'a, S>, deadline: Option<Instant>) -> Self {
81 Self { ctx, deadline }
82 }
83
84 #[inline]
85 fn check_deadline(&self) -> ExecResult<()> {
86 if let Some(deadline) = self.deadline {
87 check_deadline_at(deadline)
88 } else {
89 Ok(())
90 }
91 }
92
93 pub fn execute(
94 &mut self,
95 plan: &PhysicalPlan,
96 options: Option<ExecuteOptions>,
97 ) -> ExecResult<QueryResult> {
98 let rows = self.execute_rows(plan)?;
99 Ok(project_rows(rows, options.unwrap_or_default()))
100 }
101
102 pub fn execute_rows(&mut self, plan: &PhysicalPlan) -> ExecResult<Vec<Row>> {
103 self.check_deadline()?;
104 clear_eval_error();
107
108 let rows = self.execute_node(plan, plan.root)?;
109 if !plan_may_need_hydration(plan) {
110 return Ok(rows);
111 }
112 Ok(rows
113 .into_iter()
114 .map(|row| self.hydrate_row(row))
115 .collect::<Vec<_>>())
116 }
117
118 pub fn execute_compiled(
120 &mut self,
121 compiled: &CompiledQuery,
122 options: Option<ExecuteOptions>,
123 ) -> ExecResult<QueryResult> {
124 let rows = self.execute_compiled_rows(compiled)?;
125 Ok(project_rows(rows, options.unwrap_or_default()))
126 }
127
128 pub fn execute_compiled_rows(&mut self, compiled: &CompiledQuery) -> ExecResult<Vec<Row>> {
129 self.check_deadline()?;
130 if compiled.unions.is_empty() {
131 return self.execute_rows(&compiled.physical);
132 }
133
134 clear_eval_error();
135
136 let mut all_rows = self.execute_and_hydrate(&compiled.physical)?;
138
139 let mut needs_dedup = false;
142
143 for branch in &compiled.unions {
144 self.check_deadline()?;
145 let branch_rows = self.execute_and_hydrate(&branch.physical)?;
146 all_rows.extend(branch_rows);
147
148 if !branch.all {
149 needs_dedup = true;
150 }
151 }
152
153 if needs_dedup {
154 all_rows = dedup_rows(all_rows);
155 }
156
157 Ok(all_rows)
158 }
159
160 fn execute_and_hydrate(&mut self, plan: &PhysicalPlan) -> ExecResult<Vec<Row>> {
161 self.check_deadline()?;
162 let rows = self.execute_node(plan, plan.root)?;
163 if !plan_may_need_hydration(plan) {
164 return Ok(rows);
165 }
166 Ok(rows.into_iter().map(|row| self.hydrate_row(row)).collect())
167 }
168
169 pub(crate) fn hydrate_row(&self, row: Row) -> Row {
170 let mut out = Row::new();
171
172 for (var, name, value) in row.into_iter_named() {
173 out.insert_named(var, name, self.hydrate_value(value));
174 }
175
176 out
177 }
178
179 fn execute_node(
180 &mut self,
181 plan: &PhysicalPlan,
182 node_id: PhysicalNodeId,
183 ) -> ExecResult<Vec<Row>> {
184 self.check_deadline()?;
185 trace!("mutable execute_node start: node_id={node_id:?}");
186
187 let result = match &plan.nodes[node_id] {
188 PhysicalOp::Argument(op) => self.exec_argument(op),
189 PhysicalOp::NodeScan(op) => self.exec_node_scan(plan, op),
190 PhysicalOp::NodeByLabelScan(op) => self.exec_node_by_label_scan(plan, op),
191 PhysicalOp::NodeByPropertyScan(op) => self.exec_node_by_property_scan(plan, op),
192 PhysicalOp::NodeByPropertyRangeScan(op) => {
193 self.exec_node_by_property_range_scan(plan, op)
194 }
195 PhysicalOp::NodeByTextScan(op) => self.exec_node_by_text_scan(plan, op),
196 PhysicalOp::NodeByPointScan(op) => self.exec_node_by_point_scan(plan, op),
197 PhysicalOp::RelByPropertyRangeScan(op) => {
198 self.exec_rel_by_property_range_scan(plan, op)
199 }
200 PhysicalOp::RelByTextScan(op) => self.exec_rel_by_text_scan(plan, op),
201 PhysicalOp::RelByPointScan(op) => self.exec_rel_by_point_scan(plan, op),
202 PhysicalOp::Expand(op) => self.exec_expand(plan, op),
203 PhysicalOp::Filter(op) => self.exec_filter(plan, op),
204 PhysicalOp::Projection(op) => self.exec_projection(plan, op),
205 PhysicalOp::Unwind(op) => self.exec_unwind(plan, op),
206 PhysicalOp::HashAggregation(op) => self.exec_hash_aggregation(plan, op),
207 PhysicalOp::Sort(op) => self.exec_sort(plan, op),
208 PhysicalOp::Limit(op) => self.exec_limit(plan, op),
209 PhysicalOp::Create(op) => self.exec_create(plan, op),
210 PhysicalOp::Merge(op) => self.exec_merge(plan, op),
211 PhysicalOp::Delete(op) => self.exec_delete(plan, op),
212 PhysicalOp::Set(op) => self.exec_set(plan, op),
213 PhysicalOp::Remove(op) => self.exec_remove(plan, op),
214 PhysicalOp::OptionalMatch(op) => self.exec_optional_match(plan, op),
215 PhysicalOp::PathBuild(op) => self.exec_path_build(plan, op),
216 };
217
218 match &result {
219 Ok(rows) => trace!(
220 "mutable execute_node ok: node_id={node_id:?}, rows={}",
221 rows.len()
222 ),
223 Err(err) => error!("mutable execute_node failed: node_id={node_id:?}, error={err}"),
224 }
225
226 result
227 }
228
229 fn exec_argument(&self, _op: &ArgumentExec) -> ExecResult<Vec<Row>> {
230 Ok(vec![Row::new()])
231 }
232
233 fn exec_node_scan(&mut self, plan: &PhysicalPlan, op: &NodeScanExec) -> ExecResult<Vec<Row>> {
234 let base_rows = match op.input {
235 Some(input) => self.execute_node(plan, input)?,
236 None => vec![Row::new()],
237 };
238
239 node_scan_rows(&*self.ctx.storage, base_rows, op, self.deadline)
240 }
241
242 fn exec_node_by_label_scan(
243 &mut self,
244 plan: &PhysicalPlan,
245 op: &NodeByLabelScanExec,
246 ) -> ExecResult<Vec<Row>> {
247 let base_rows = match op.input {
248 Some(input) => self.execute_node(plan, input)?,
249 None => vec![Row::new()],
250 };
251
252 node_by_label_scan_rows(&*self.ctx.storage, base_rows, op, self.deadline)
253 }
254
255 fn exec_node_by_property_scan(
256 &mut self,
257 plan: &PhysicalPlan,
258 op: &NodeByPropertyScanExec,
259 ) -> ExecResult<Vec<Row>> {
260 let base_rows = match op.input {
261 Some(input) => self.execute_node(plan, input)?,
262 None => vec![Row::new()],
263 };
264
265 node_by_property_scan_rows(
266 &*self.ctx.storage,
267 &self.ctx.params,
268 base_rows,
269 op,
270 self.deadline,
271 )
272 }
273
274 fn exec_node_by_property_range_scan(
275 &mut self,
276 plan: &PhysicalPlan,
277 op: &lora_compiler::NodeByPropertyRangeScanExec,
278 ) -> ExecResult<Vec<Row>> {
279 let base_rows = match op.input {
280 Some(input) => self.execute_node(plan, input)?,
281 None => vec![Row::new()],
282 };
283 super::helpers::node_by_property_range_scan_rows(
284 &*self.ctx.storage,
285 &self.ctx.params,
286 base_rows,
287 op,
288 self.deadline,
289 )
290 }
291
292 fn exec_node_by_text_scan(
293 &mut self,
294 plan: &PhysicalPlan,
295 op: &lora_compiler::NodeByTextScanExec,
296 ) -> ExecResult<Vec<Row>> {
297 let base_rows = match op.input {
298 Some(input) => self.execute_node(plan, input)?,
299 None => vec![Row::new()],
300 };
301 super::helpers::node_by_text_scan_rows(
302 &*self.ctx.storage,
303 &self.ctx.params,
304 base_rows,
305 op,
306 self.deadline,
307 )
308 }
309
310 fn exec_node_by_point_scan(
311 &mut self,
312 plan: &PhysicalPlan,
313 op: &lora_compiler::NodeByPointScanExec,
314 ) -> ExecResult<Vec<Row>> {
315 let base_rows = match op.input {
316 Some(input) => self.execute_node(plan, input)?,
317 None => vec![Row::new()],
318 };
319 super::helpers::node_by_point_scan_rows(
320 &*self.ctx.storage,
321 &self.ctx.params,
322 base_rows,
323 op,
324 self.deadline,
325 )
326 }
327
328 fn exec_rel_by_property_range_scan(
329 &mut self,
330 plan: &PhysicalPlan,
331 op: &lora_compiler::RelByPropertyRangeScanExec,
332 ) -> ExecResult<Vec<Row>> {
333 let base_rows = match op.input {
334 Some(input) => self.execute_node(plan, input)?,
335 None => vec![Row::new()],
336 };
337 super::helpers::rel_by_property_range_scan_rows(
338 &*self.ctx.storage,
339 &self.ctx.params,
340 base_rows,
341 op,
342 self.deadline,
343 )
344 }
345
346 fn exec_rel_by_text_scan(
347 &mut self,
348 plan: &PhysicalPlan,
349 op: &lora_compiler::RelByTextScanExec,
350 ) -> ExecResult<Vec<Row>> {
351 let base_rows = match op.input {
352 Some(input) => self.execute_node(plan, input)?,
353 None => vec![Row::new()],
354 };
355 super::helpers::rel_by_text_scan_rows(
356 &*self.ctx.storage,
357 &self.ctx.params,
358 base_rows,
359 op,
360 self.deadline,
361 )
362 }
363
364 fn exec_rel_by_point_scan(
365 &mut self,
366 plan: &PhysicalPlan,
367 op: &lora_compiler::RelByPointScanExec,
368 ) -> ExecResult<Vec<Row>> {
369 let base_rows = match op.input {
370 Some(input) => self.execute_node(plan, input)?,
371 None => vec![Row::new()],
372 };
373 super::helpers::rel_by_point_scan_rows(
374 &*self.ctx.storage,
375 &self.ctx.params,
376 base_rows,
377 op,
378 self.deadline,
379 )
380 }
381
382 fn exec_expand(&mut self, plan: &PhysicalPlan, op: &ExpandExec) -> ExecResult<Vec<Row>> {
383 let input_rows = self.execute_node(plan, op.input)?;
384 if let Some(range) = &op.range {
385 expand_var_len_rows(&*self.ctx.storage, input_rows, op, range)
386 } else {
387 expand_rows(&*self.ctx.storage, &self.ctx.params, input_rows, op)
388 }
389 }
390
391 fn exec_filter(&mut self, plan: &PhysicalPlan, op: &FilterExec) -> ExecResult<Vec<Row>> {
392 let input_rows = self.execute_node(plan, op.input)?;
393 let eval_ctx = EvalContext {
394 storage: &*self.ctx.storage,
395 params: &self.ctx.params,
396 };
397
398 filter_rows_checked(input_rows, &op.predicate, &eval_ctx)
399 }
400
401 fn exec_projection(
402 &mut self,
403 plan: &PhysicalPlan,
404 op: &ProjectionExec,
405 ) -> ExecResult<Vec<Row>> {
406 let input_rows = self.execute_node(plan, op.input)?;
407 let eval_ctx = EvalContext {
408 storage: &*self.ctx.storage,
409 params: &self.ctx.params,
410 };
411
412 project_rows_checked(input_rows, op, &eval_ctx)
413 }
414
415 fn hydrate_value(&self, value: LoraValue) -> LoraValue {
416 match value {
417 LoraValue::Node(id) => self.hydrate_node(id),
418 LoraValue::Relationship(id) => self.hydrate_relationship(id),
419 LoraValue::List(values) => {
420 LoraValue::List(values.into_iter().map(|v| self.hydrate_value(v)).collect())
421 }
422 LoraValue::Map(map) => LoraValue::Map(
423 map.into_iter()
424 .map(|(k, v)| (k, self.hydrate_value(v)))
425 .collect(),
426 ),
427 other => other,
428 }
429 }
430
431 fn hydrate_node(&self, id: u64) -> LoraValue {
432 self.ctx
433 .storage
434 .with_node(id, hydrate_node_record)
435 .unwrap_or(LoraValue::Null)
436 }
437
438 fn hydrate_relationship(&self, id: u64) -> LoraValue {
439 self.ctx
440 .storage
441 .with_relationship(id, hydrate_relationship_record)
442 .unwrap_or(LoraValue::Null)
443 }
444
445 fn exec_unwind(&mut self, plan: &PhysicalPlan, op: &UnwindExec) -> ExecResult<Vec<Row>> {
446 let input_rows = self.execute_node(plan, op.input)?;
447 let eval_ctx = EvalContext {
448 storage: &*self.ctx.storage,
449 params: &self.ctx.params,
450 };
451
452 Ok(unwind_rows(input_rows, op, &eval_ctx))
453 }
454
455 fn exec_hash_aggregation(
456 &mut self,
457 plan: &PhysicalPlan,
458 op: &HashAggregationExec,
459 ) -> ExecResult<Vec<Row>> {
460 if let Some(rows) =
461 super::helpers::count_all_scan_aggregation_rows(&*self.ctx.storage, plan, op)
462 {
463 return Ok(rows);
464 }
465
466 let input_rows = self.execute_node(plan, op.input)?;
467 let eval_ctx = EvalContext {
468 storage: &*self.ctx.storage,
469 params: &self.ctx.params,
470 };
471
472 aggregate_rows(
473 input_rows,
474 &op.group_by,
475 &op.aggregates,
476 &eval_ctx,
477 |value| self.hydrate_value(value),
478 )
479 }
480
481 fn exec_sort(&mut self, plan: &PhysicalPlan, op: &SortExec) -> ExecResult<Vec<Row>> {
482 let mut rows = self.execute_node(plan, op.input)?;
483 let eval_ctx = EvalContext {
484 storage: &*self.ctx.storage,
485 params: &self.ctx.params,
486 };
487
488 sort_rows_with_top_k(&mut rows, &op.items, &eval_ctx, op.top_k);
489
490 Ok(rows)
491 }
492
493 fn exec_limit(&mut self, plan: &PhysicalPlan, op: &LimitExec) -> ExecResult<Vec<Row>> {
494 let rows = self.execute_node(plan, op.input)?;
495 let eval_ctx = EvalContext {
496 storage: &*self.ctx.storage,
497 params: &self.ctx.params,
498 };
499
500 Ok(limit_rows(rows, op, &eval_ctx))
501 }
502
503 fn exec_optional_match(
504 &mut self,
505 plan: &PhysicalPlan,
506 op: &OptionalMatchExec,
507 ) -> ExecResult<Vec<Row>> {
508 let input_rows = self.execute_node(plan, op.input)?;
509
510 let inner_rows = self.execute_node(plan, op.inner)?;
512
513 Ok(optional_match_rows(input_rows, &inner_rows, &op.new_vars))
514 }
515
516 fn exec_path_build(&mut self, plan: &PhysicalPlan, op: &PathBuildExec) -> ExecResult<Vec<Row>> {
517 let input_rows = self.execute_node(plan, op.input)?;
518 let mut rows: Vec<Row> = input_rows
519 .into_iter()
520 .map(|mut row| {
521 let path = build_path_value(&row, &op.node_vars, &op.rel_vars, &*self.ctx.storage);
522 row.insert(op.output, path);
523 row
524 })
525 .collect();
526
527 if let Some(all) = op.shortest_path_all {
528 rows = filter_shortest_paths(rows, op.output, all);
529 }
530 Ok(rows)
531 }
532
533 fn exec_create(&mut self, plan: &PhysicalPlan, op: &CreateExec) -> ExecResult<Vec<Row>> {
534 if crate::pull::subtree_is_fully_streaming(plan, op.input) {
540 return self.exec_create_streaming_input(plan, op);
541 }
542
543 let input_rows = self.execute_node(plan, op.input)?;
544 let mut out = Vec::with_capacity(input_rows.len());
545
546 for mut row in input_rows {
547 self.apply_create_pattern(&mut row, &op.pattern)?;
548 out.push(row);
549 }
550
551 Ok(out)
552 }
553
554 fn streaming_apply<F>(
574 &mut self,
575 plan: &PhysicalPlan,
576 input: PhysicalNodeId,
577 mut apply: F,
578 ) -> ExecResult<Vec<Row>>
579 where
580 F: FnMut(&mut Self, &mut Row) -> ExecResult<()>,
581 {
582 use std::sync::Arc;
583
584 let storage_ptr: *mut S = self.ctx.storage as *mut S;
585 let params = Arc::new(self.ctx.params.clone());
586
587 let storage_ref: &S = unsafe { &*storage_ptr };
589 let mut upstream = crate::pull::build_streaming(plan, input, storage_ref, params)?;
590
591 let mut out = Vec::new();
592 while let Some(mut row) = upstream.next_row()? {
593 apply(self, &mut row)?;
594 out.push(row);
595 }
596
597 Ok(out)
598 }
599
600 fn exec_create_streaming_input(
603 &mut self,
604 plan: &PhysicalPlan,
605 op: &CreateExec,
606 ) -> ExecResult<Vec<Row>> {
607 self.streaming_apply(plan, op.input, |this, row| {
608 this.apply_create_pattern(row, &op.pattern)
609 })
610 }
611
612 fn apply_remove_item(&mut self, row: &Row, item: &ResolvedRemoveItem) -> ExecResult<()> {
613 match item {
614 ResolvedRemoveItem::Labels { variable, labels } => match row.get(*variable) {
615 Some(LoraValue::Node(node_id)) => {
616 let node_id = *node_id;
617 for label in labels {
618 self.ctx.storage.remove_node_label(node_id, label);
619 }
620 Ok(())
621 }
622 Some(other) => Err(ExecutorError::ExpectedNodeForRemoveLabels {
623 found: value_kind(other),
624 }),
625 None => Err(ExecutorError::UnboundVariableForRemove {
626 var: format!("{variable:?}"),
627 }),
628 },
629
630 ResolvedRemoveItem::Property { expr } => self.remove_property_from_expr(row, expr),
631 }
632 }
633
634 fn delete_value(&mut self, value: LoraValue, detach: bool) -> ExecResult<()> {
635 match value {
636 LoraValue::Null => Ok(()),
637
638 LoraValue::Node(node_id) => {
639 if detach {
640 self.ctx.storage.detach_delete_node(node_id);
641 Ok(())
642 } else {
643 let ok = self.ctx.storage.delete_node(node_id);
644 if ok {
645 Ok(())
646 } else {
647 Err(ExecutorError::DeleteNodeWithRelationships { node_id })
648 }
649 }
650 }
651
652 LoraValue::Relationship(rel_id) => {
653 let ok = self.ctx.storage.delete_relationship(rel_id);
654 if ok {
655 Ok(())
656 } else {
657 Err(ExecutorError::DeleteRelationshipFailed { rel_id })
658 }
659 }
660
661 LoraValue::List(values) => {
662 for v in values {
663 self.delete_value(v, detach)?;
664 }
665 Ok(())
666 }
667
668 other => Err(ExecutorError::InvalidDeleteTarget {
669 found: value_kind(&other),
670 }),
671 }
672 }
673
674 fn exec_merge(&mut self, plan: &PhysicalPlan, op: &MergeExec) -> ExecResult<Vec<Row>> {
675 if crate::pull::subtree_is_fully_streaming(plan, op.input) {
680 return self.streaming_apply(plan, op.input, |this, row| {
681 let already_bound = this.pattern_part_is_bound(row, &op.pattern_part);
682 let matched = if already_bound {
683 true
684 } else {
685 this.try_match_merge_pattern(row, &op.pattern_part)?
686 };
687 if !matched {
688 this.apply_create_pattern_part(row, &op.pattern_part)?;
689 }
690 for action in &op.actions {
691 if action.on_match == matched {
692 for item in &action.set.items {
693 this.apply_set_item(row, item)?;
694 }
695 }
696 }
697 Ok(())
698 });
699 }
700
701 let input_rows = self.execute_node(plan, op.input)?;
702 let mut out = Vec::with_capacity(input_rows.len());
703
704 for mut row in input_rows {
705 let already_bound = self.pattern_part_is_bound(&row, &op.pattern_part);
707
708 let matched = if already_bound {
709 true
710 } else {
711 self.try_match_merge_pattern(&mut row, &op.pattern_part)?
713 };
714
715 if !matched {
716 self.apply_create_pattern_part(&mut row, &op.pattern_part)?;
717 }
718
719 for action in &op.actions {
720 if action.on_match == matched {
721 for item in &action.set.items {
722 self.apply_set_item(&row, item)?;
723 }
724 }
725 }
726
727 out.push(row);
728 }
729
730 Ok(out)
731 }
732
733 fn try_match_merge_pattern(
736 &self,
737 row: &mut Row,
738 part: &ResolvedPatternPart,
739 ) -> ExecResult<bool> {
740 match &part.element {
741 ResolvedPatternElement::Node {
742 var,
743 labels,
744 properties,
745 } => {
746 let candidate_ids = if labels.is_empty() {
749 self.ctx.storage.all_node_ids()
750 } else {
751 scan_node_ids_for_label_groups(&*self.ctx.storage, labels)
752 };
753
754 let eval_ctx = EvalContext {
756 storage: &*self.ctx.storage,
757 params: &self.ctx.params,
758 };
759 let expected_props = properties.as_ref().map(|e| eval_expr(e, row, &eval_ctx));
760
761 for id in candidate_ids {
762 let matched = self
763 .ctx
764 .storage
765 .with_node(id, |node| {
766 if !node_matches_label_groups(&node.labels, labels) {
767 return false;
768 }
769 if let Some(LoraValue::Map(expected)) = &expected_props {
770 let all_match = expected.iter().all(|(key, expected_value)| {
771 node.properties
772 .get(key)
773 .map(|actual| {
774 value_matches_property_value(expected_value, actual)
775 })
776 .unwrap_or(false)
777 });
778 if !all_match {
779 return false;
780 }
781 }
782 true
783 })
784 .unwrap_or(false);
785
786 if !matched {
787 continue;
788 }
789
790 if let Some(var_id) = var {
792 row.insert(*var_id, LoraValue::Node(id));
793 }
794 return Ok(true);
795 }
796
797 Ok(false)
798 }
799
800 ResolvedPatternElement::ShortestPath { .. } => {
801 Ok(false)
803 }
804
805 ResolvedPatternElement::NodeChain { head, chain } => {
806 let head_node_id = if let Some(var_id) = head.var {
808 if let Some(LoraValue::Node(id)) = row.get(var_id) {
809 *id
810 } else {
811 let node_matched = self.try_match_merge_pattern(
813 row,
814 &ResolvedPatternPart {
815 binding: None,
816 element: ResolvedPatternElement::Node {
817 var: head.var,
818 labels: head.labels.clone(),
819 properties: head.properties.clone(),
820 },
821 },
822 )?;
823 if !node_matched {
824 return Ok(false);
825 }
826 match row.get(var_id) {
827 Some(LoraValue::Node(id)) => *id,
828 _ => return Ok(false),
829 }
830 }
831 } else {
832 return Ok(false);
833 };
834
835 let mut current_node_id = head_node_id;
836
837 for step in chain {
838 let eval_ctx = EvalContext {
839 storage: &*self.ctx.storage,
840 params: &self.ctx.params,
841 };
842
843 let direction = step.rel.direction;
844
845 let mut found = false;
848 let _ = self.ctx.storage.try_for_each_expand_id(
849 current_node_id,
850 direction,
851 &step.rel.types,
852 |rel_id, node_id| {
853 let node_ok = self
855 .ctx
856 .storage
857 .with_node(node_id, |node_rec| {
858 if !node_matches_label_groups(
859 &node_rec.labels,
860 &step.node.labels,
861 ) {
862 return false;
863 }
864 if let Some(props_expr) = &step.node.properties {
865 let expected = eval_expr(props_expr, row, &eval_ctx);
866 if let LoraValue::Map(expected_map) = &expected {
867 let all_match =
868 expected_map.iter().all(|(key, expected_val)| {
869 node_rec
870 .properties
871 .get(key)
872 .map(|actual| {
873 value_matches_property_value(
874 expected_val,
875 actual,
876 )
877 })
878 .unwrap_or(false)
879 });
880 if !all_match {
881 return false;
882 }
883 }
884 }
885 true
886 })
887 .unwrap_or(false);
888 if !node_ok {
889 return Ok::<(), ()>(());
890 }
891
892 let rel_ok = self
894 .ctx
895 .storage
896 .with_relationship(rel_id, |rel_rec| {
897 if let Some(rel_props_expr) = &step.rel.properties {
898 let expected = eval_expr(rel_props_expr, row, &eval_ctx);
899 if let LoraValue::Map(expected_map) = &expected {
900 let all_match =
901 expected_map.iter().all(|(key, expected_val)| {
902 rel_rec
903 .properties
904 .get(key)
905 .map(|actual| {
906 value_matches_property_value(
907 expected_val,
908 actual,
909 )
910 })
911 .unwrap_or(false)
912 });
913 if !all_match {
914 return false;
915 }
916 }
917 }
918 true
919 })
920 .unwrap_or(false);
921 if !rel_ok {
922 return Ok(());
923 }
924
925 if let Some(rel_var) = step.rel.var {
927 row.insert(rel_var, LoraValue::Relationship(rel_id));
928 }
929 if let Some(node_var) = step.node.var {
930 row.insert(node_var, LoraValue::Node(node_id));
931 }
932 current_node_id = node_id;
933 found = true;
934 Err(())
935 },
936 );
937
938 if !found {
939 return Ok(false);
940 }
941 }
942
943 Ok(true)
944 }
945 }
946 }
947
948 fn exec_delete(&mut self, plan: &PhysicalPlan, op: &DeleteExec) -> ExecResult<Vec<Row>> {
949 if crate::pull::subtree_is_fully_streaming(plan, op.input) {
950 let detach = op.detach;
951 return self.streaming_apply(plan, op.input, |this, row| {
952 for expr in &op.expressions {
953 let value = {
954 let eval_ctx = EvalContext {
955 storage: &*this.ctx.storage,
956 params: &this.ctx.params,
957 };
958 eval_expr(expr, row, &eval_ctx)
959 };
960 this.delete_value(value, detach)?;
961 }
962 Ok(())
963 });
964 }
965
966 let input_rows = self.execute_node(plan, op.input)?;
967
968 for row in &input_rows {
969 for expr in &op.expressions {
970 let value = {
971 let eval_ctx = EvalContext {
972 storage: &*self.ctx.storage,
973 params: &self.ctx.params,
974 };
975 eval_expr(expr, row, &eval_ctx)
976 };
977
978 self.delete_value(value, op.detach)?;
979 }
980 }
981
982 Ok(input_rows)
983 }
984
985 fn exec_set(&mut self, plan: &PhysicalPlan, op: &SetExec) -> ExecResult<Vec<Row>> {
986 if crate::pull::subtree_is_fully_streaming(plan, op.input) {
987 return self.streaming_apply(plan, op.input, |this, row| {
988 for item in &op.items {
989 this.apply_set_item(row, item)?;
990 }
991 Ok(())
992 });
993 }
994
995 let input_rows = self.execute_node(plan, op.input)?;
996
997 for row in &input_rows {
998 for item in &op.items {
999 self.apply_set_item(row, item)?;
1000 }
1001 }
1002
1003 Ok(input_rows)
1004 }
1005
1006 fn exec_remove(&mut self, plan: &PhysicalPlan, op: &RemoveExec) -> ExecResult<Vec<Row>> {
1007 if crate::pull::subtree_is_fully_streaming(plan, op.input) {
1008 return self.streaming_apply(plan, op.input, |this, row| {
1009 for item in &op.items {
1010 this.apply_remove_item(row, item)?;
1011 }
1012 Ok(())
1013 });
1014 }
1015
1016 let input_rows = self.execute_node(plan, op.input)?;
1017
1018 for row in &input_rows {
1019 for item in &op.items {
1020 self.apply_remove_item(row, item)?;
1021 }
1022 }
1023
1024 Ok(input_rows)
1025 }
1026
1027 fn apply_set_item(&mut self, row: &Row, item: &ResolvedSetItem) -> ExecResult<()> {
1028 match item {
1029 ResolvedSetItem::SetProperty { target, value } => {
1030 let new_value = {
1031 let eval_ctx = EvalContext {
1032 storage: &*self.ctx.storage,
1033 params: &self.ctx.params,
1034 };
1035 eval_expr(value, row, &eval_ctx)
1036 };
1037
1038 self.set_property_from_expr(row, target, new_value)
1039 }
1040
1041 ResolvedSetItem::SetVariable { variable, value } => {
1042 let entity_ref =
1044 row.get(*variable)
1045 .ok_or(ExecutorError::UnboundVariableForSet {
1046 var: format!("{variable:?}"),
1047 })?;
1048 let entity_target = entity_target_from_value(entity_ref)?;
1049
1050 let new_value = {
1051 let eval_ctx = EvalContext {
1052 storage: &*self.ctx.storage,
1053 params: &self.ctx.params,
1054 };
1055 eval_expr(value, row, &eval_ctx)
1056 };
1057
1058 self.overwrite_entity_target(entity_target, new_value)
1059 }
1060
1061 ResolvedSetItem::MutateVariable { variable, value } => {
1062 let entity_ref =
1063 row.get(*variable)
1064 .ok_or(ExecutorError::UnboundVariableForSet {
1065 var: format!("{variable:?}"),
1066 })?;
1067 let entity_target = entity_target_from_value(entity_ref)?;
1068
1069 let patch = {
1070 let eval_ctx = EvalContext {
1071 storage: &*self.ctx.storage,
1072 params: &self.ctx.params,
1073 };
1074 eval_expr(value, row, &eval_ctx)
1075 };
1076
1077 self.mutate_entity_target(entity_target, patch)
1078 }
1079
1080 ResolvedSetItem::SetLabels { variable, labels } => match row.get(*variable) {
1081 Some(LoraValue::Node(node_id)) => {
1082 let node_id = *node_id;
1083 for label in labels {
1084 if let Err(msg) = self
1085 .ctx
1086 .storage
1087 .check_node_add_label_against_constraints(node_id, label)
1088 {
1089 return Err(ExecutorError::ConstraintViolation(msg));
1090 }
1091 self.ctx.storage.add_node_label(node_id, label);
1092 }
1093 Ok(())
1094 }
1095 Some(other) => Err(ExecutorError::ExpectedNodeForSetLabels {
1096 found: value_kind(other),
1097 }),
1098 None => Err(ExecutorError::UnboundVariableForSet {
1099 var: format!("{variable:?}"),
1100 }),
1101 },
1102 }
1103 }
1104
1105 fn set_property_from_expr(
1106 &mut self,
1107 row: &Row,
1108 target_expr: &ResolvedExpr,
1109 new_value: LoraValue,
1110 ) -> ExecResult<()> {
1111 let ResolvedExpr::Property { expr, property } = target_expr else {
1112 return Err(ExecutorError::UnsupportedSetTarget);
1113 };
1114
1115 let owner = {
1116 let eval_ctx = EvalContext {
1117 storage: &*self.ctx.storage,
1118 params: &self.ctx.params,
1119 };
1120 eval_expr(expr, row, &eval_ctx)
1121 };
1122
1123 match owner {
1124 LoraValue::Node(node_id) => {
1125 let prop = lora_value_to_property(new_value)
1126 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1127 if let Err(msg) = self
1128 .ctx
1129 .storage
1130 .check_node_set_property_against_constraints(node_id, property, &prop)
1131 {
1132 return Err(ExecutorError::ConstraintViolation(msg));
1133 }
1134 self.ctx
1135 .storage
1136 .set_node_property(node_id, property.clone(), prop);
1137 Ok(())
1138 }
1139 LoraValue::Relationship(rel_id) => {
1140 let prop = lora_value_to_property(new_value)
1141 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1142 if let Err(msg) = self
1143 .ctx
1144 .storage
1145 .check_relationship_set_property_against_constraints(rel_id, property, &prop)
1146 {
1147 return Err(ExecutorError::ConstraintViolation(msg));
1148 }
1149 self.ctx
1150 .storage
1151 .set_relationship_property(rel_id, property.clone(), prop);
1152 Ok(())
1153 }
1154 other => Err(ExecutorError::InvalidSetTarget {
1155 found: value_kind(&other),
1156 }),
1157 }
1158 }
1159
1160 fn remove_property_from_expr(&mut self, row: &Row, expr: &ResolvedExpr) -> ExecResult<()> {
1161 let ResolvedExpr::Property {
1162 expr: owner_expr,
1163 property,
1164 } = expr
1165 else {
1166 return Err(ExecutorError::UnsupportedRemoveTarget);
1167 };
1168
1169 let owner = {
1170 let eval_ctx = EvalContext {
1171 storage: &*self.ctx.storage,
1172 params: &self.ctx.params,
1173 };
1174 eval_expr(owner_expr, row, &eval_ctx)
1175 };
1176
1177 match owner {
1178 LoraValue::Node(node_id) => {
1179 if let Err(msg) = self
1180 .ctx
1181 .storage
1182 .check_node_remove_property_against_constraints(node_id, property)
1183 {
1184 return Err(ExecutorError::ConstraintViolation(msg));
1185 }
1186 self.ctx.storage.remove_node_property(node_id, property);
1187 Ok(())
1188 }
1189 LoraValue::Relationship(rel_id) => {
1190 if let Err(msg) = self
1191 .ctx
1192 .storage
1193 .check_relationship_remove_property_against_constraints(rel_id, property)
1194 {
1195 return Err(ExecutorError::ConstraintViolation(msg));
1196 }
1197 self.ctx
1198 .storage
1199 .remove_relationship_property(rel_id, property);
1200 Ok(())
1201 }
1202 other => Err(ExecutorError::InvalidRemoveTarget {
1203 found: value_kind(&other),
1204 }),
1205 }
1206 }
1207
1208 fn overwrite_entity_target(
1209 &mut self,
1210 target: EntityTarget,
1211 new_value: LoraValue,
1212 ) -> ExecResult<()> {
1213 let LoraValue::Map(map) = new_value else {
1214 return Err(ExecutorError::ExpectedPropertyMap {
1215 found: value_kind(&new_value),
1216 });
1217 };
1218
1219 let mut props: Properties = Properties::new();
1220 for (k, v) in map {
1221 let prop = lora_value_to_property(v)
1222 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1223 props.insert(k, prop);
1224 }
1225
1226 match target {
1227 EntityTarget::Node(node_id) => {
1228 if let Err(msg) = self
1229 .ctx
1230 .storage
1231 .check_node_replace_properties_against_constraints(node_id, &props)
1232 {
1233 return Err(ExecutorError::ConstraintViolation(msg));
1234 }
1235 self.ctx.storage.replace_node_properties(node_id, props);
1236 }
1237 EntityTarget::Relationship(rel_id) => {
1238 if let Err(msg) = self
1239 .ctx
1240 .storage
1241 .check_relationship_replace_properties_against_constraints(rel_id, &props)
1242 {
1243 return Err(ExecutorError::ConstraintViolation(msg));
1244 }
1245 self.ctx
1246 .storage
1247 .replace_relationship_properties(rel_id, props);
1248 }
1249 }
1250 Ok(())
1251 }
1252
1253 fn mutate_entity_target(
1254 &mut self,
1255 target: EntityTarget,
1256 patch_value: LoraValue,
1257 ) -> ExecResult<()> {
1258 let LoraValue::Map(map) = patch_value else {
1259 return Err(ExecutorError::ExpectedPropertyMap {
1260 found: value_kind(&patch_value),
1261 });
1262 };
1263
1264 match target {
1265 EntityTarget::Node(node_id) => {
1266 for (k, v) in map {
1267 let prop = lora_value_to_property(v)
1268 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1269 if let Err(msg) = self
1270 .ctx
1271 .storage
1272 .check_node_set_property_against_constraints(node_id, &k, &prop)
1273 {
1274 return Err(ExecutorError::ConstraintViolation(msg));
1275 }
1276 self.ctx.storage.set_node_property(node_id, k, prop);
1277 }
1278 }
1279 EntityTarget::Relationship(rel_id) => {
1280 for (k, v) in map {
1281 let prop = lora_value_to_property(v)
1282 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1283 if let Err(msg) = self
1284 .ctx
1285 .storage
1286 .check_relationship_set_property_against_constraints(rel_id, &k, &prop)
1287 {
1288 return Err(ExecutorError::ConstraintViolation(msg));
1289 }
1290 self.ctx.storage.set_relationship_property(rel_id, k, prop);
1291 }
1292 }
1293 }
1294 Ok(())
1295 }
1296
1297 pub(crate) fn apply_create_pattern(
1298 &mut self,
1299 row: &mut Row,
1300 pattern: &ResolvedPattern,
1301 ) -> ExecResult<()> {
1302 for part in &pattern.parts {
1303 self.apply_create_pattern_part(row, part)?;
1304 }
1305 Ok(())
1306 }
1307
1308 pub(crate) fn apply_write_op(&mut self, op: &PhysicalOp, row: &mut Row) -> ExecResult<()> {
1314 match op {
1315 PhysicalOp::Create(c) => self.apply_create_pattern(row, &c.pattern),
1316 PhysicalOp::Set(s) => {
1317 for item in &s.items {
1318 self.apply_set_item(row, item)?;
1319 }
1320 Ok(())
1321 }
1322 PhysicalOp::Delete(d) => {
1323 let detach = d.detach;
1324 for expr in &d.expressions {
1325 let value = {
1326 let eval_ctx = EvalContext {
1327 storage: &*self.ctx.storage,
1328 params: &self.ctx.params,
1329 };
1330 eval_expr(expr, row, &eval_ctx)
1331 };
1332 self.delete_value(value, detach)?;
1333 }
1334 Ok(())
1335 }
1336 PhysicalOp::Remove(r) => {
1337 for item in &r.items {
1338 self.apply_remove_item(row, item)?;
1339 }
1340 Ok(())
1341 }
1342 PhysicalOp::Merge(m) => {
1343 let already_bound = self.pattern_part_is_bound(row, &m.pattern_part);
1344 let matched = if already_bound {
1345 true
1346 } else {
1347 self.try_match_merge_pattern(row, &m.pattern_part)?
1348 };
1349 if !matched {
1350 self.apply_create_pattern_part(row, &m.pattern_part)?;
1351 }
1352 for action in &m.actions {
1353 if action.on_match == matched {
1354 for item in &action.set.items {
1355 self.apply_set_item(row, item)?;
1356 }
1357 }
1358 }
1359 Ok(())
1360 }
1361 other => Err(ExecutorError::RuntimeError(format!(
1362 "apply_write_op called on non-write op: {other:?}"
1363 ))),
1364 }
1365 }
1366
1367 fn apply_create_pattern_part(
1368 &mut self,
1369 row: &mut Row,
1370 part: &ResolvedPatternPart,
1371 ) -> ExecResult<()> {
1372 if part.binding.is_some() {
1373 trace!("create pattern part has path binding; path materialization not implemented");
1374 }
1375
1376 let _ = self.apply_create_pattern_element(row, &part.element)?;
1377 Ok(())
1378 }
1379
1380 fn apply_create_pattern_element(
1381 &mut self,
1382 row: &mut Row,
1383 element: &ResolvedPatternElement,
1384 ) -> ExecResult<Option<LoraValue>> {
1385 match element {
1386 ResolvedPatternElement::Node {
1387 var,
1388 labels,
1389 properties,
1390 } => {
1391 let node_id =
1392 self.materialize_node_pattern(row, *var, labels, properties.as_ref())?;
1393 Ok(Some(LoraValue::Node(node_id)))
1394 }
1395
1396 ResolvedPatternElement::NodeChain { head, chain } => {
1397 let mut current_node_id = self.materialize_node_pattern(
1398 row,
1399 head.var,
1400 &head.labels,
1401 head.properties.as_ref(),
1402 )?;
1403
1404 for link in chain {
1405 let next_node_id = self.materialize_node_pattern(
1406 row,
1407 link.node.var,
1408 &link.node.labels,
1409 link.node.properties.as_ref(),
1410 )?;
1411
1412 let _ = self.materialize_relationship_pattern(
1413 row,
1414 current_node_id,
1415 next_node_id,
1416 &link.rel,
1417 )?;
1418
1419 current_node_id = next_node_id;
1420 }
1421
1422 Ok(Some(LoraValue::Node(current_node_id)))
1423 }
1424
1425 ResolvedPatternElement::ShortestPath { .. } => {
1426 Ok(None)
1428 }
1429 }
1430 }
1431
1432 fn pattern_part_is_bound(&self, row: &Row, part: &ResolvedPatternPart) -> bool {
1433 match &part.element {
1434 ResolvedPatternElement::Node { var, .. } => var.and_then(|v| row.get(v)).is_some(),
1435
1436 ResolvedPatternElement::ShortestPath { .. } => false,
1437
1438 ResolvedPatternElement::NodeChain { head, chain } => {
1439 let head_ok = head.var.and_then(|v| row.get(v)).is_some();
1440
1441 let chain_ok = chain.iter().all(|link| {
1442 let node_ok = link.node.var.and_then(|v| row.get(v)).is_some();
1443 let rel_ok = match link.rel.var {
1447 Some(v) => row.get(v).is_some(),
1448 None => false,
1449 };
1450 node_ok && rel_ok
1451 });
1452
1453 head_ok && chain_ok
1454 }
1455 }
1456 }
1457
1458 fn materialize_node_pattern(
1459 &mut self,
1460 row: &mut Row,
1461 var: Option<VarId>,
1462 labels: &[Vec<String>],
1463 properties: Option<&ResolvedExpr>,
1464 ) -> ExecResult<u64> {
1465 if let Some(var_id) = var {
1466 if let Some(LoraValue::Node(id)) = row.get(var_id) {
1467 return Ok(*id);
1468 }
1469 }
1470
1471 let properties = match properties {
1472 Some(expr) => eval_properties_expr(expr, row, &*self.ctx.storage, &self.ctx.params)?,
1473 None => Properties::new(),
1474 };
1475
1476 let flat_labels = flatten_label_groups(labels);
1477 debug!("creating node with labels={flat_labels:?}");
1478 if let Err(msg) = self
1479 .ctx
1480 .storage
1481 .check_node_create_against_constraints(&flat_labels, &properties)
1482 {
1483 return Err(ExecutorError::ConstraintViolation(msg));
1484 }
1485 let created = self.ctx.storage.create_node(flat_labels, properties);
1486
1487 if let Some(var_id) = var {
1488 row.insert(var_id, LoraValue::Node(created.id));
1489 }
1490
1491 Ok(created.id)
1492 }
1493
1494 fn materialize_relationship_pattern(
1495 &mut self,
1496 row: &mut Row,
1497 left_node_id: u64,
1498 right_node_id: u64,
1499 rel: &lora_analyzer::ResolvedRel,
1500 ) -> ExecResult<u64> {
1501 if let Some(var_id) = rel.var {
1502 if let Some(LoraValue::Relationship(id)) = row.get(var_id) {
1503 let id = *id;
1504 if let Some((src, dst)) = self.ctx.storage.relationship_endpoints(id) {
1505 let endpoints_match = match rel.direction {
1506 Direction::Right | Direction::Undirected => {
1507 src == left_node_id && dst == right_node_id
1508 }
1509 Direction::Left => src == right_node_id && dst == left_node_id,
1510 };
1511
1512 if endpoints_match {
1513 return Ok(id);
1514 }
1515 }
1516 }
1517 }
1518
1519 if rel.range.is_some() {
1520 return Err(ExecutorError::UnsupportedCreateRelationshipRange);
1521 }
1522
1523 let (src, dst) = match rel.direction {
1524 Direction::Right | Direction::Undirected => (left_node_id, right_node_id),
1525 Direction::Left => (right_node_id, left_node_id),
1526 };
1527
1528 let rel_type = rel
1529 .types
1530 .first()
1531 .ok_or(ExecutorError::MissingRelationshipType)?;
1532
1533 if rel_type.is_empty() {
1534 return Err(ExecutorError::MissingRelationshipType);
1535 }
1536
1537 let properties = match rel.properties.as_ref() {
1538 Some(expr) => eval_properties_expr(expr, row, &*self.ctx.storage, &self.ctx.params)?,
1539 None => Properties::new(),
1540 };
1541
1542 debug!("creating relationship: src={src}, dst={dst}, type={rel_type}");
1543
1544 if let Err(msg) = self
1545 .ctx
1546 .storage
1547 .check_relationship_create_against_constraints(rel_type, &properties)
1548 {
1549 return Err(ExecutorError::ConstraintViolation(msg));
1550 }
1551
1552 let created = self
1553 .ctx
1554 .storage
1555 .create_relationship(src, dst, rel_type, properties)
1556 .ok_or_else(|| ExecutorError::RelationshipCreateFailed {
1557 src,
1558 dst,
1559 rel_type: rel_type.clone(),
1560 })?;
1561
1562 if let Some(var_id) = rel.var {
1563 row.insert(var_id, LoraValue::Relationship(created.id));
1564 }
1565
1566 Ok(created.id)
1567 }
1568}