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