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::NodeByIdSeek(op) => self.exec_node_by_id_seek(plan, op),
265 PhysicalOp::RelByIdSeek(op) => self.exec_rel_by_id_seek(plan, op),
266 PhysicalOp::NodeByPropertyScan(op) => self.exec_node_by_property_scan(plan, op),
267 PhysicalOp::NodeByPropertyRangeScan(op) => {
268 self.exec_node_by_property_range_scan(plan, op)
269 }
270 PhysicalOp::NodeByTextScan(op) => self.exec_node_by_text_scan(plan, op),
271 PhysicalOp::NodeByPointScan(op) => self.exec_node_by_point_scan(plan, op),
272 PhysicalOp::RelByPropertyRangeScan(op) => {
273 self.exec_rel_by_property_range_scan(plan, op)
274 }
275 PhysicalOp::RelByTextScan(op) => self.exec_rel_by_text_scan(plan, op),
276 PhysicalOp::RelByPointScan(op) => self.exec_rel_by_point_scan(plan, op),
277 PhysicalOp::Expand(op) => self.exec_expand(plan, op),
278 PhysicalOp::Filter(op) => self.exec_filter(plan, op),
279 PhysicalOp::Projection(op) => self.exec_projection(plan, op),
280 PhysicalOp::Unwind(op) => self.exec_unwind(plan, op),
281 PhysicalOp::HashAggregation(op) => self.exec_hash_aggregation(plan, op),
282 PhysicalOp::Sort(op) => self.exec_sort(plan, op),
283 PhysicalOp::Limit(op) => self.exec_limit(plan, op),
284 PhysicalOp::Create(op) => self.exec_create(plan, op),
285 PhysicalOp::Merge(op) => self.exec_merge(plan, op),
286 PhysicalOp::Delete(op) => self.exec_delete(plan, op),
287 PhysicalOp::Set(op) => self.exec_set(plan, op),
288 PhysicalOp::Remove(op) => self.exec_remove(plan, op),
289 PhysicalOp::Foreach(op) => self.exec_foreach(plan, op),
290 PhysicalOp::OptionalMatch(op) => self.exec_optional_match(plan, op),
291 PhysicalOp::CallSubquery(op) => self.exec_call_subquery(plan, op),
292 PhysicalOp::PathBuild(op) => self.exec_path_build(plan, op),
293 };
294
295 match &result {
296 Ok(rows) => trace!(
297 "mutable execute_node ok: node_id={node_id:?}, rows={}",
298 rows.len()
299 ),
300 Err(err) => error!("mutable execute_node failed: node_id={node_id:?}, error={err}"),
301 }
302
303 result
304 }
305
306 fn exec_argument(&self, _op: &ArgumentExec) -> ExecResult<Vec<Row>> {
307 Ok(vec![self.argument_seed.clone().unwrap_or_default()])
308 }
309
310 fn exec_node_scan(&mut self, plan: &PhysicalPlan, op: &NodeScanExec) -> ExecResult<Vec<Row>> {
311 let base_rows = match op.input {
312 Some(input) => self.execute_node(plan, input)?,
313 None => vec![Row::new()],
314 };
315
316 node_scan_rows(&*self.ctx.storage, base_rows, op, self.deadline)
317 }
318
319 fn exec_node_by_label_scan(
320 &mut self,
321 plan: &PhysicalPlan,
322 op: &NodeByLabelScanExec,
323 ) -> ExecResult<Vec<Row>> {
324 let base_rows = match op.input {
325 Some(input) => self.execute_node(plan, input)?,
326 None => vec![Row::new()],
327 };
328
329 node_by_label_scan_rows(&*self.ctx.storage, base_rows, op, self.deadline)
330 }
331
332 fn exec_node_by_property_scan(
333 &mut self,
334 plan: &PhysicalPlan,
335 op: &NodeByPropertyScanExec,
336 ) -> ExecResult<Vec<Row>> {
337 let base_rows = match op.input {
338 Some(input) => self.execute_node(plan, input)?,
339 None => vec![Row::new()],
340 };
341
342 node_by_property_scan_rows(
343 &*self.ctx.storage,
344 &self.ctx.params,
345 base_rows,
346 op,
347 self.deadline,
348 )
349 }
350
351 fn exec_node_by_property_range_scan(
352 &mut self,
353 plan: &PhysicalPlan,
354 op: &lora_compiler::NodeByPropertyRangeScanExec,
355 ) -> ExecResult<Vec<Row>> {
356 let base_rows = match op.input {
357 Some(input) => self.execute_node(plan, input)?,
358 None => vec![Row::new()],
359 };
360 super::helpers::node_by_property_range_scan_rows(
361 &*self.ctx.storage,
362 &self.ctx.params,
363 base_rows,
364 op,
365 self.deadline,
366 )
367 }
368
369 fn exec_node_by_text_scan(
370 &mut self,
371 plan: &PhysicalPlan,
372 op: &lora_compiler::NodeByTextScanExec,
373 ) -> ExecResult<Vec<Row>> {
374 let base_rows = match op.input {
375 Some(input) => self.execute_node(plan, input)?,
376 None => vec![Row::new()],
377 };
378 super::helpers::node_by_text_scan_rows(
379 &*self.ctx.storage,
380 &self.ctx.params,
381 base_rows,
382 op,
383 self.deadline,
384 )
385 }
386
387 fn exec_node_by_point_scan(
388 &mut self,
389 plan: &PhysicalPlan,
390 op: &lora_compiler::NodeByPointScanExec,
391 ) -> ExecResult<Vec<Row>> {
392 let base_rows = match op.input {
393 Some(input) => self.execute_node(plan, input)?,
394 None => vec![Row::new()],
395 };
396 super::helpers::node_by_point_scan_rows(
397 &*self.ctx.storage,
398 &self.ctx.params,
399 base_rows,
400 op,
401 self.deadline,
402 )
403 }
404
405 fn exec_rel_by_property_range_scan(
406 &mut self,
407 plan: &PhysicalPlan,
408 op: &lora_compiler::RelByPropertyRangeScanExec,
409 ) -> ExecResult<Vec<Row>> {
410 let base_rows = match op.input {
411 Some(input) => self.execute_node(plan, input)?,
412 None => vec![Row::new()],
413 };
414 super::helpers::rel_by_property_range_scan_rows(
415 &*self.ctx.storage,
416 &self.ctx.params,
417 base_rows,
418 op,
419 self.deadline,
420 )
421 }
422
423 fn exec_rel_by_text_scan(
424 &mut self,
425 plan: &PhysicalPlan,
426 op: &lora_compiler::RelByTextScanExec,
427 ) -> ExecResult<Vec<Row>> {
428 let base_rows = match op.input {
429 Some(input) => self.execute_node(plan, input)?,
430 None => vec![Row::new()],
431 };
432 super::helpers::rel_by_text_scan_rows(
433 &*self.ctx.storage,
434 &self.ctx.params,
435 base_rows,
436 op,
437 self.deadline,
438 )
439 }
440
441 fn exec_node_by_id_seek(
442 &mut self,
443 plan: &PhysicalPlan,
444 op: &lora_compiler::NodeByIdSeekExec,
445 ) -> ExecResult<Vec<Row>> {
446 let base_rows = match op.input {
447 Some(input) => self.execute_node(plan, input)?,
448 None => vec![Row::new()],
449 };
450 super::helpers::node_by_id_seek_rows(
451 &*self.ctx.storage,
452 &self.ctx.params,
453 base_rows,
454 op,
455 self.deadline,
456 )
457 }
458
459 fn exec_rel_by_id_seek(
460 &mut self,
461 plan: &PhysicalPlan,
462 op: &lora_compiler::RelByIdSeekExec,
463 ) -> ExecResult<Vec<Row>> {
464 let base_rows = match op.input {
465 Some(input) => self.execute_node(plan, input)?,
466 None => vec![Row::new()],
467 };
468 super::helpers::rel_by_id_seek_rows(
469 &*self.ctx.storage,
470 &self.ctx.params,
471 base_rows,
472 op,
473 self.deadline,
474 )
475 }
476
477 fn exec_rel_by_point_scan(
478 &mut self,
479 plan: &PhysicalPlan,
480 op: &lora_compiler::RelByPointScanExec,
481 ) -> ExecResult<Vec<Row>> {
482 let base_rows = match op.input {
483 Some(input) => self.execute_node(plan, input)?,
484 None => vec![Row::new()],
485 };
486 super::helpers::rel_by_point_scan_rows(
487 &*self.ctx.storage,
488 &self.ctx.params,
489 base_rows,
490 op,
491 self.deadline,
492 )
493 }
494
495 fn exec_expand(&mut self, plan: &PhysicalPlan, op: &ExpandExec) -> ExecResult<Vec<Row>> {
496 let input_rows = self.execute_node(plan, op.input)?;
497 if let Some(range) = &op.range {
498 expand_var_len_rows(&*self.ctx.storage, input_rows, op, range)
499 } else {
500 expand_rows(&*self.ctx.storage, &self.ctx.params, input_rows, op)
501 }
502 }
503
504 fn exec_filter(&mut self, plan: &PhysicalPlan, op: &FilterExec) -> ExecResult<Vec<Row>> {
505 let input_rows = self.execute_node(plan, op.input)?;
506 let eval_ctx = EvalContext {
507 storage: &*self.ctx.storage,
508 params: &self.ctx.params,
509 };
510
511 filter_rows_checked(input_rows, &op.predicate, &eval_ctx)
512 }
513
514 fn exec_projection(
515 &mut self,
516 plan: &PhysicalPlan,
517 op: &ProjectionExec,
518 ) -> ExecResult<Vec<Row>> {
519 let input_rows = self.execute_node(plan, op.input)?;
520 let eval_ctx = EvalContext {
521 storage: &*self.ctx.storage,
522 params: &self.ctx.params,
523 };
524
525 project_rows_checked(input_rows, op, &eval_ctx)
526 }
527
528 fn hydrate_value(&self, value: LoraValue) -> LoraValue {
529 match value {
530 LoraValue::Node(id) => self.hydrate_node(id),
531 LoraValue::Relationship(id) => self.hydrate_relationship(id),
532 LoraValue::List(values) => {
533 LoraValue::List(values.into_iter().map(|v| self.hydrate_value(v)).collect())
534 }
535 LoraValue::Map(map) => LoraValue::Map(
536 map.into_iter()
537 .map(|(k, v)| (k, self.hydrate_value(v)))
538 .collect(),
539 ),
540 other => other,
541 }
542 }
543
544 fn hydrate_node(&self, id: u64) -> LoraValue {
545 self.ctx
546 .storage
547 .with_node(id, hydrate_node_record)
548 .unwrap_or(LoraValue::Null)
549 }
550
551 fn hydrate_relationship(&self, id: u64) -> LoraValue {
552 self.ctx
553 .storage
554 .with_relationship(id, hydrate_relationship_record)
555 .unwrap_or(LoraValue::Null)
556 }
557
558 fn exec_unwind(&mut self, plan: &PhysicalPlan, op: &UnwindExec) -> ExecResult<Vec<Row>> {
559 let input_rows = self.execute_node(plan, op.input)?;
560 let eval_ctx = EvalContext {
561 storage: &*self.ctx.storage,
562 params: &self.ctx.params,
563 };
564
565 unwind_rows(input_rows, op, &eval_ctx)
566 }
567
568 fn exec_hash_aggregation(
569 &mut self,
570 plan: &PhysicalPlan,
571 op: &HashAggregationExec,
572 ) -> ExecResult<Vec<Row>> {
573 if let Some(rows) =
574 super::helpers::count_all_scan_aggregation_rows(&*self.ctx.storage, plan, op)
575 {
576 return Ok(rows);
577 }
578
579 let input_rows = self.execute_node(plan, op.input)?;
580 let eval_ctx = EvalContext {
581 storage: &*self.ctx.storage,
582 params: &self.ctx.params,
583 };
584
585 aggregate_rows(input_rows, &op.group_by, &op.aggregates, &eval_ctx)
586 }
587
588 fn exec_sort(&mut self, plan: &PhysicalPlan, op: &SortExec) -> ExecResult<Vec<Row>> {
589 let mut rows = self.execute_node(plan, op.input)?;
590 let eval_ctx = EvalContext {
591 storage: &*self.ctx.storage,
592 params: &self.ctx.params,
593 };
594
595 let bound = sort_row_bound(op.top_k, op.limit.as_ref(), &eval_ctx);
596 sort_rows_with_top_k(&mut rows, &op.items, &eval_ctx, bound);
597
598 Ok(rows)
599 }
600
601 fn exec_limit(&mut self, plan: &PhysicalPlan, op: &LimitExec) -> ExecResult<Vec<Row>> {
602 let rows = self.execute_node(plan, op.input)?;
603 let eval_ctx = EvalContext {
604 storage: &*self.ctx.storage,
605 params: &self.ctx.params,
606 };
607
608 limit_rows(rows, op, &eval_ctx)
609 }
610
611 fn exec_optional_match(
612 &mut self,
613 plan: &PhysicalPlan,
614 op: &OptionalMatchExec,
615 ) -> ExecResult<Vec<Row>> {
616 let input_rows = self.execute_node(plan, op.input)?;
617
618 if super::optional::optional_can_correlate(plan, op.inner) {
619 let storage_ref: &S = &*self.ctx.storage;
620 return super::optional::correlated_optional_match_rows(
621 storage_ref,
622 &self.ctx.params,
623 plan,
624 op.inner,
625 input_rows,
626 &op.new_vars,
627 );
628 }
629
630 let inner_rows = self.execute_node(plan, op.inner)?;
632
633 Ok(optional_match_rows(input_rows, &inner_rows, &op.new_vars))
634 }
635
636 fn exec_call_subquery(
637 &mut self,
638 plan: &PhysicalPlan,
639 op: &CallSubqueryExec,
640 ) -> ExecResult<Vec<Row>> {
641 let input_rows = self.execute_node(plan, op.input)?;
642 let mut out = Vec::with_capacity(input_rows.len());
643
644 if crate::pull::subtree_has_write(plan, op.inner) {
645 let unit = op.new_vars.is_empty();
649 for outer_row in input_rows {
650 self.check_deadline()?;
651 let prev = self.argument_seed.replace(outer_row.clone());
652 let inner_rows = self.execute_node(plan, op.inner);
653 self.argument_seed = prev;
654 let inner_rows = inner_rows?;
655 if unit {
656 out.push(outer_row);
659 continue;
660 }
661 for inner_row in inner_rows {
662 out.push(crate::executor::merge_optional_rows(&outer_row, &inner_row));
663 }
664 }
665 return Ok(out);
666 }
667
668 let params = std::sync::Arc::new(self.ctx.params.clone());
669 let storage_ref: &S = &*self.ctx.storage;
670 for outer_row in input_rows {
671 let mut inner_source = crate::pull::build_streaming_seeded(
672 plan,
673 op.inner,
674 storage_ref,
675 params.clone(),
676 outer_row.clone(),
677 )?;
678 let inner_rows = crate::pull::drain(inner_source.as_mut())?;
679 for inner_row in inner_rows {
680 out.push(crate::executor::merge_optional_rows(&outer_row, &inner_row));
681 }
682 }
683 Ok(out)
684 }
685
686 fn exec_path_build(&mut self, plan: &PhysicalPlan, op: &PathBuildExec) -> ExecResult<Vec<Row>> {
687 let input_rows = self.execute_node(plan, op.input)?;
688 let mut rows: Vec<Row> = input_rows
689 .into_iter()
690 .map(|mut row| {
691 let path = build_path_value(&row, &op.node_vars, &op.rel_vars, &*self.ctx.storage);
692 row.insert(op.output, path);
693 row
694 })
695 .collect();
696
697 if let Some(all) = op.shortest_path_all {
698 rows = filter_shortest_paths(rows, op.output, all);
699 }
700 Ok(rows)
701 }
702
703 fn exec_create(&mut self, plan: &PhysicalPlan, op: &CreateExec) -> ExecResult<Vec<Row>> {
704 if crate::pull::subtree_is_fully_streaming(plan, op.input) {
710 return self.exec_create_streaming_input(plan, op);
711 }
712
713 let input_rows = self.execute_node(plan, op.input)?;
714 let mut out = Vec::with_capacity(input_rows.len());
715
716 for mut row in input_rows {
717 self.apply_create_pattern(&mut row, &op.pattern)?;
718 out.push(row);
719 }
720
721 Ok(out)
722 }
723
724 fn streaming_apply<F>(
744 &mut self,
745 plan: &PhysicalPlan,
746 input: PhysicalNodeId,
747 mut apply: F,
748 ) -> ExecResult<Vec<Row>>
749 where
750 F: FnMut(&mut Self, &mut Row) -> ExecResult<()>,
751 {
752 use std::sync::Arc;
753
754 let storage_ptr: *mut S = self.ctx.storage as *mut S;
755 let params = Arc::new(self.ctx.params.clone());
756
757 let storage_ref: &S = unsafe { &*storage_ptr };
759 let mut upstream = match self.argument_seed.clone() {
762 Some(seed) => {
763 crate::pull::build_streaming_seeded(plan, input, storage_ref, params, seed)?
764 }
765 None => crate::pull::build_streaming(plan, input, storage_ref, params)?,
766 };
767
768 let mut out = Vec::new();
769 while let Some(mut row) = upstream.next_row()? {
770 apply(self, &mut row)?;
771 out.push(row);
772 }
773
774 Ok(out)
775 }
776
777 fn exec_create_streaming_input(
780 &mut self,
781 plan: &PhysicalPlan,
782 op: &CreateExec,
783 ) -> ExecResult<Vec<Row>> {
784 self.streaming_apply(plan, op.input, |this, row| {
785 this.apply_create_pattern(row, &op.pattern)
786 })
787 }
788
789 fn apply_remove_item(&mut self, row: &Row, item: &ResolvedRemoveItem) -> ExecResult<()> {
790 match item {
791 ResolvedRemoveItem::Labels { variable, labels } => match row.get(*variable) {
792 Some(LoraValue::Node(node_id)) => {
793 let node_id = *node_id;
794 for label in labels {
795 self.ctx.storage.remove_node_label(node_id, label);
796 }
797 Ok(())
798 }
799 Some(other) => Err(ExecutorError::ExpectedNodeForRemoveLabels {
800 found: value_kind(other),
801 }),
802 None => Err(ExecutorError::UnboundVariableForRemove {
803 var: format!("{variable:?}"),
804 }),
805 },
806
807 ResolvedRemoveItem::Property { expr } => self.remove_property_from_expr(row, expr),
808 }
809 }
810
811 fn delete_value(&mut self, value: LoraValue, detach: bool) -> ExecResult<()> {
812 match value {
813 LoraValue::Null => Ok(()),
814
815 LoraValue::Node(node_id) => {
816 if detach {
817 self.ctx.storage.detach_delete_node(node_id);
818 Ok(())
819 } else {
820 let ok = self.ctx.storage.delete_node(node_id);
821 if ok {
822 Ok(())
823 } else {
824 Err(ExecutorError::DeleteNodeWithRelationships { node_id })
825 }
826 }
827 }
828
829 LoraValue::Relationship(rel_id) => {
830 let ok = self.ctx.storage.delete_relationship(rel_id);
831 if ok {
832 Ok(())
833 } else {
834 Err(ExecutorError::DeleteRelationshipFailed { rel_id })
835 }
836 }
837
838 LoraValue::List(values) => {
839 for v in values {
840 self.delete_value(v, detach)?;
841 }
842 Ok(())
843 }
844
845 other => Err(ExecutorError::InvalidDeleteTarget {
846 found: value_kind(&other),
847 }),
848 }
849 }
850
851 fn collect_delete_targets(
852 &self,
853 value: &LoraValue,
854 targets: &mut BTreeSet<DeleteTarget>,
855 ) -> ExecResult<()> {
856 match value {
857 LoraValue::Null => Ok(()),
858
859 LoraValue::Node(node_id) => {
860 targets.insert(DeleteTarget::Node(*node_id));
861 Ok(())
862 }
863
864 LoraValue::Relationship(rel_id) => {
865 targets.insert(DeleteTarget::Relationship(*rel_id));
866 Ok(())
867 }
868
869 LoraValue::List(values) => {
870 for v in values {
871 self.collect_delete_targets(v, targets)?;
872 }
873 Ok(())
874 }
875
876 other => Err(ExecutorError::InvalidDeleteTarget {
877 found: value_kind(other),
878 }),
879 }
880 }
881
882 fn validate_delete_targets(
883 &self,
884 targets: &BTreeSet<DeleteTarget>,
885 detach: bool,
886 ) -> ExecResult<()> {
887 for target in targets {
888 match target {
889 DeleteTarget::Relationship(rel_id) => {
890 if !self.ctx.storage.contains_relationship(*rel_id) {
891 return Err(ExecutorError::DeleteRelationshipFailed { rel_id: *rel_id });
892 }
893 }
894 DeleteTarget::Node(node_id) if !detach => {
895 if !self.ctx.storage.contains_node(*node_id) {
896 return Err(ExecutorError::DeleteNodeWithRelationships {
897 node_id: *node_id,
898 });
899 }
900 let has_external_relationship = self
901 .ctx
902 .storage
903 .relationship_ids_of(*node_id, Direction::Undirected)
904 .into_iter()
905 .any(|rel_id| !targets.contains(&DeleteTarget::Relationship(rel_id)));
906 if has_external_relationship {
907 return Err(ExecutorError::DeleteNodeWithRelationships {
908 node_id: *node_id,
909 });
910 }
911 }
912 DeleteTarget::Node(_) => {}
913 }
914 }
915 Ok(())
916 }
917
918 fn delete_target(&mut self, target: DeleteTarget, detach: bool) -> ExecResult<()> {
919 match target {
920 DeleteTarget::Node(node_id) => {
921 if detach {
922 self.ctx.storage.detach_delete_node(node_id);
923 Ok(())
924 } else {
925 let ok = self.ctx.storage.delete_node(node_id);
926 if ok {
927 Ok(())
928 } else {
929 Err(ExecutorError::DeleteNodeWithRelationships { node_id })
930 }
931 }
932 }
933 DeleteTarget::Relationship(rel_id) => {
934 let ok = self.ctx.storage.delete_relationship(rel_id);
935 if ok {
936 Ok(())
937 } else {
938 Err(ExecutorError::DeleteRelationshipFailed { rel_id })
939 }
940 }
941 }
942 }
943
944 fn exec_merge(&mut self, plan: &PhysicalPlan, op: &MergeExec) -> ExecResult<Vec<Row>> {
945 if crate::pull::subtree_is_fully_streaming(plan, op.input) {
950 return self.streaming_apply(plan, op.input, |this, row| {
951 let matched = this.match_merge_pattern(row, &op.pattern_part)?;
952 if !matched {
953 this.apply_create_pattern_part(row, &op.pattern_part)?;
954 }
955 for action in &op.actions {
956 if action.on_match == matched {
957 for item in &action.set.items {
958 this.apply_set_item(row, item)?;
959 }
960 }
961 }
962 Ok(())
963 });
964 }
965
966 let input_rows = self.execute_node(plan, op.input)?;
967 let mut out = Vec::with_capacity(input_rows.len());
968
969 for mut row in input_rows {
970 let matched = self.match_merge_pattern(&mut row, &op.pattern_part)?;
971
972 if !matched {
973 self.apply_create_pattern_part(&mut row, &op.pattern_part)?;
974 }
975
976 for action in &op.actions {
977 if action.on_match == matched {
978 for item in &action.set.items {
979 self.apply_set_item(&row, item)?;
980 }
981 }
982 }
983
984 out.push(row);
985 }
986
987 Ok(out)
988 }
989
990 fn match_merge_pattern(&self, row: &mut Row, part: &ResolvedPatternPart) -> ExecResult<bool> {
995 if self.pattern_part_is_bound(row, part)? {
996 if let (Some(var), Some(path)) = (part.binding, bound_pattern_path(row, part)) {
997 row.insert(var, LoraValue::Path(path));
998 }
999 return Ok(true);
1000 }
1001 self.try_match_merge_pattern(row, part)
1002 }
1003
1004 fn try_match_merge_pattern(
1009 &self,
1010 row: &mut Row,
1011 part: &ResolvedPatternPart,
1012 ) -> ExecResult<bool> {
1013 match &part.element {
1014 ResolvedPatternElement::Node {
1015 var,
1016 labels,
1017 properties,
1018 } => {
1019 let expected_props = self.merge_expected_props(properties.as_ref(), row);
1020 let Some(id) = self
1021 .merge_node_candidates(labels, &expected_props)
1022 .into_iter()
1023 .find(|&id| self.merge_node_matches(id, labels, &expected_props))
1024 else {
1025 return Ok(false);
1026 };
1027 if let Some(var_id) = var {
1028 row.insert(*var_id, LoraValue::Node(id));
1029 }
1030 if let Some(path_var) = part.binding {
1031 let path = LoraPath {
1032 nodes: vec![id],
1033 rels: Vec::new(),
1034 };
1035 row.insert(path_var, LoraValue::Path(path));
1036 }
1037 Ok(true)
1038 }
1039
1040 ResolvedPatternElement::ShortestPath { .. } => {
1041 Ok(false)
1043 }
1044
1045 ResolvedPatternElement::NodeChain { head, chain } => {
1046 let head_candidates = match head.var.and_then(|v| row.get(v)) {
1049 Some(LoraValue::Node(id)) => vec![*id],
1050 _ => {
1051 let expected = self.merge_expected_props(head.properties.as_ref(), row);
1052 self.merge_node_candidates(&head.labels, &expected)
1053 .into_iter()
1054 .filter(|&id| self.merge_node_matches(id, &head.labels, &expected))
1055 .collect()
1056 }
1057 };
1058
1059 for head_id in head_candidates {
1060 let mut trial = row.clone();
1061 if let Some(var_id) = head.var {
1062 trial.insert(var_id, LoraValue::Node(head_id));
1063 }
1064 let mut walked = Vec::with_capacity(chain.len());
1065 if self.match_merge_chain(&mut trial, head_id, chain, &mut walked) {
1066 if let Some(path_var) = part.binding {
1067 let path = LoraPath {
1068 nodes: std::iter::once(head_id)
1069 .chain(walked.iter().map(|&(_, node)| node))
1070 .collect(),
1071 rels: walked.iter().map(|&(rel, _)| rel).collect(),
1072 };
1073 trial.insert(path_var, LoraValue::Path(path));
1074 }
1075 *row = trial;
1076 return Ok(true);
1077 }
1078 }
1079 Ok(false)
1080 }
1081 }
1082 }
1083
1084 fn match_merge_chain(
1090 &self,
1091 row: &mut Row,
1092 current: NodeId,
1093 chain: &[lora_analyzer::ResolvedChain],
1094 walked: &mut Vec<(u64, NodeId)>,
1095 ) -> bool {
1096 let Some((step, rest)) = chain.split_first() else {
1097 return true;
1098 };
1099
1100 let bound_dst = match step.node.var.and_then(|v| row.get(v)) {
1101 Some(LoraValue::Node(id)) => Some(*id),
1102 _ => None,
1103 };
1104 let bound_rel = match step.rel.var.and_then(|v| row.get(v)) {
1105 Some(LoraValue::Relationship(id)) => Some(*id),
1106 _ => None,
1107 };
1108 let expected_node = self.merge_expected_props(step.node.properties.as_ref(), row);
1109 let expected_rel = self.merge_expected_props(step.rel.properties.as_ref(), row);
1110
1111 let edges = self
1112 .ctx
1113 .storage
1114 .expand_ids(current, step.rel.direction, &step.rel.types);
1115 for (rel_id, node_id) in edges {
1116 if bound_dst.is_some_and(|id| id != node_id)
1117 || bound_rel.is_some_and(|id| id != rel_id)
1118 || walked.iter().any(|&(used, _)| used == rel_id)
1119 {
1120 continue;
1121 }
1122 if !self.merge_node_matches(node_id, &step.node.labels, &expected_node) {
1123 continue;
1124 }
1125 if let Some(LoraValue::Map(expected_map)) = &expected_rel {
1126 let rel_ok = self
1127 .ctx
1128 .storage
1129 .with_relationship(rel_id, |rel_rec| {
1130 expected_map.iter().all(|(key, expected_val)| {
1131 rel_rec
1132 .properties
1133 .get(key.as_str())
1134 .map(|actual| value_matches_property_value(expected_val, actual))
1135 .unwrap_or(false)
1136 })
1137 })
1138 .unwrap_or(false);
1139 if !rel_ok {
1140 continue;
1141 }
1142 }
1143
1144 let mut next = row.clone();
1145 if let Some(rel_var) = step.rel.var {
1146 next.insert(rel_var, LoraValue::Relationship(rel_id));
1147 }
1148 if let Some(node_var) = step.node.var {
1149 next.insert(node_var, LoraValue::Node(node_id));
1150 }
1151 walked.push((rel_id, node_id));
1152 if self.match_merge_chain(&mut next, node_id, rest, walked) {
1153 *row = next;
1154 return true;
1155 }
1156 walked.pop();
1157 }
1158 false
1159 }
1160
1161 fn merge_expected_props(
1162 &self,
1163 properties: Option<&ResolvedExpr>,
1164 row: &Row,
1165 ) -> Option<LoraValue> {
1166 let eval_ctx = EvalContext {
1167 storage: &*self.ctx.storage,
1168 params: &self.ctx.params,
1169 };
1170 properties.map(|e| eval_expr(e, row, &eval_ctx))
1171 }
1172
1173 fn merge_node_candidates(
1178 &self,
1179 labels: &[Vec<String>],
1180 expected_props: &Option<LoraValue>,
1181 ) -> Vec<NodeId> {
1182 let indexed = match expected_props {
1183 Some(LoraValue::Map(expected)) => {
1184 merge_candidates_from_index(&*self.ctx.storage, labels, expected)
1185 }
1186 _ => None,
1187 };
1188 match indexed {
1189 Some(ids) => ids,
1190 None if labels.is_empty() => self.ctx.storage.all_node_ids(),
1191 None => scan_node_ids_for_label_groups(&*self.ctx.storage, labels),
1192 }
1193 }
1194
1195 fn merge_node_matches(
1196 &self,
1197 id: NodeId,
1198 labels: &[Vec<String>],
1199 expected_props: &Option<LoraValue>,
1200 ) -> bool {
1201 self.ctx
1202 .storage
1203 .with_node(id, |node| {
1204 if !node_matches_label_groups(&node.labels, labels) {
1205 return false;
1206 }
1207 if let Some(LoraValue::Map(expected)) = expected_props {
1208 return expected.iter().all(|(key, expected_value)| {
1209 node.properties
1210 .get(key.as_str())
1211 .map(|actual| value_matches_property_value(expected_value, actual))
1212 .unwrap_or(false)
1213 });
1214 }
1215 true
1216 })
1217 .unwrap_or(false)
1218 }
1219
1220 fn exec_delete(&mut self, plan: &PhysicalPlan, op: &DeleteExec) -> ExecResult<Vec<Row>> {
1221 let input_rows = self.execute_node(plan, op.input)?;
1222 let mut targets = BTreeSet::new();
1223
1224 for row in &input_rows {
1225 for expr in &op.expressions {
1226 let value = {
1227 let eval_ctx = EvalContext {
1228 storage: &*self.ctx.storage,
1229 params: &self.ctx.params,
1230 };
1231 eval_expr(expr, row, &eval_ctx)
1232 };
1233 self.collect_delete_targets(&value, &mut targets)?;
1234 }
1235 }
1236
1237 self.validate_delete_targets(&targets, op.detach)?;
1238
1239 for target in &targets {
1240 if let DeleteTarget::Relationship(_) = target {
1241 self.delete_target(*target, op.detach)?;
1242 }
1243 }
1244 for target in targets {
1245 if let DeleteTarget::Node(_) = target {
1246 self.delete_target(target, op.detach)?;
1247 }
1248 }
1249
1250 Ok(input_rows)
1251 }
1252
1253 fn exec_set(&mut self, plan: &PhysicalPlan, op: &SetExec) -> ExecResult<Vec<Row>> {
1254 if crate::pull::subtree_is_fully_streaming(plan, op.input) {
1255 return self.streaming_apply(plan, op.input, |this, row| {
1256 for item in &op.items {
1257 this.apply_set_item(row, item)?;
1258 }
1259 Ok(())
1260 });
1261 }
1262
1263 let input_rows = self.execute_node(plan, op.input)?;
1264
1265 for row in &input_rows {
1266 for item in &op.items {
1267 self.apply_set_item(row, item)?;
1268 }
1269 }
1270
1271 Ok(input_rows)
1272 }
1273
1274 fn exec_foreach(&mut self, plan: &PhysicalPlan, op: &ForeachExec) -> ExecResult<Vec<Row>> {
1282 let input_rows = self.execute_node(plan, op.input)?;
1283 let mut out = Vec::with_capacity(input_rows.len());
1284
1285 for row in input_rows {
1286 let list_value = {
1287 let eval_ctx = EvalContext {
1288 storage: &*self.ctx.storage,
1289 params: &self.ctx.params,
1290 };
1291 eval_expr(&op.list, &row, &eval_ctx)
1292 };
1293
1294 let elements: Vec<LoraValue> = match list_value {
1295 LoraValue::List(items) => items,
1296 LoraValue::Null => Vec::new(),
1297 other => {
1298 return Err(ExecutorError::RuntimeError(format!(
1299 "FOREACH expects a list, got {}",
1300 value_kind(&other)
1301 )));
1302 }
1303 };
1304
1305 for element in elements {
1306 let mut iter_row = row.clone();
1309 iter_row.insert(op.variable, element);
1310 for clause in &op.body {
1311 self.apply_foreach_body_clause(&mut iter_row, clause)?;
1312 }
1313 }
1314
1315 out.push(row);
1316 }
1317
1318 Ok(out)
1319 }
1320
1321 fn apply_foreach_body_clause(
1326 &mut self,
1327 row: &mut Row,
1328 clause: &lora_analyzer::ResolvedClause,
1329 ) -> ExecResult<()> {
1330 use lora_analyzer::ResolvedClause;
1331 match clause {
1332 ResolvedClause::Create(c) => self.apply_create_pattern(row, &c.pattern),
1333 ResolvedClause::Set(s) => {
1334 for item in &s.items {
1335 self.apply_set_item(row, item)?;
1336 }
1337 Ok(())
1338 }
1339 ResolvedClause::Remove(r) => {
1340 for item in &r.items {
1341 self.apply_remove_item(row, item)?;
1342 }
1343 Ok(())
1344 }
1345 ResolvedClause::Delete(d) => {
1346 let detach = d.detach;
1347 for expr in &d.expressions {
1348 let value = {
1349 let eval_ctx = EvalContext {
1350 storage: &*self.ctx.storage,
1351 params: &self.ctx.params,
1352 };
1353 eval_expr(expr, row, &eval_ctx)
1354 };
1355 self.delete_value(value, detach)?;
1356 }
1357 Ok(())
1358 }
1359 ResolvedClause::Merge(m) => {
1360 let matched = self.match_merge_pattern(row, &m.pattern_part)?;
1361 if !matched {
1362 self.apply_create_pattern_part(row, &m.pattern_part)?;
1363 }
1364 for action in &m.actions {
1365 if action.on_match == matched {
1366 for item in &action.set.items {
1367 self.apply_set_item(row, item)?;
1368 }
1369 }
1370 }
1371 Ok(())
1372 }
1373 ResolvedClause::Foreach(nested) => {
1374 let list_value = {
1375 let eval_ctx = EvalContext {
1376 storage: &*self.ctx.storage,
1377 params: &self.ctx.params,
1378 };
1379 eval_expr(&nested.list, row, &eval_ctx)
1380 };
1381
1382 let elements: Vec<LoraValue> = match list_value {
1383 LoraValue::List(items) => items,
1384 LoraValue::Null => Vec::new(),
1385 other => {
1386 return Err(ExecutorError::RuntimeError(format!(
1387 "FOREACH expects a list, got {}",
1388 value_kind(&other)
1389 )));
1390 }
1391 };
1392
1393 for element in elements {
1394 let mut iter_row = row.clone();
1395 iter_row.insert(nested.variable, element);
1396 for inner in &nested.body {
1397 self.apply_foreach_body_clause(&mut iter_row, inner)?;
1398 }
1399 }
1400
1401 Ok(())
1402 }
1403 other => Err(ExecutorError::RuntimeError(format!(
1404 "FOREACH body may only contain updating clauses, got {:?}",
1405 std::mem::discriminant(other)
1406 ))),
1407 }
1408 }
1409
1410 fn exec_remove(&mut self, plan: &PhysicalPlan, op: &RemoveExec) -> ExecResult<Vec<Row>> {
1411 if crate::pull::subtree_is_fully_streaming(plan, op.input) {
1412 return self.streaming_apply(plan, op.input, |this, row| {
1413 for item in &op.items {
1414 this.apply_remove_item(row, item)?;
1415 }
1416 Ok(())
1417 });
1418 }
1419
1420 let input_rows = self.execute_node(plan, op.input)?;
1421
1422 for row in &input_rows {
1423 for item in &op.items {
1424 self.apply_remove_item(row, item)?;
1425 }
1426 }
1427
1428 Ok(input_rows)
1429 }
1430
1431 fn apply_set_item(&mut self, row: &Row, item: &ResolvedSetItem) -> ExecResult<()> {
1432 match item {
1433 ResolvedSetItem::SetProperty { target, value } => {
1434 let new_value = {
1435 let eval_ctx = EvalContext {
1436 storage: &*self.ctx.storage,
1437 params: &self.ctx.params,
1438 };
1439 eval_expr(value, row, &eval_ctx)
1440 };
1441
1442 self.set_property_from_expr(row, target, new_value)
1443 }
1444
1445 ResolvedSetItem::SetVariable { variable, value } => {
1446 let entity_ref =
1448 row.get(*variable)
1449 .ok_or(ExecutorError::UnboundVariableForSet {
1450 var: format!("{variable:?}"),
1451 })?;
1452 let entity_target = entity_target_from_value(entity_ref)?;
1453
1454 let new_value = {
1455 let eval_ctx = EvalContext {
1456 storage: &*self.ctx.storage,
1457 params: &self.ctx.params,
1458 };
1459 eval_expr(value, row, &eval_ctx)
1460 };
1461
1462 self.overwrite_entity_target(entity_target, new_value)
1463 }
1464
1465 ResolvedSetItem::MutateVariable { variable, value } => {
1466 let entity_ref =
1467 row.get(*variable)
1468 .ok_or(ExecutorError::UnboundVariableForSet {
1469 var: format!("{variable:?}"),
1470 })?;
1471 let entity_target = entity_target_from_value(entity_ref)?;
1472
1473 let patch = {
1474 let eval_ctx = EvalContext {
1475 storage: &*self.ctx.storage,
1476 params: &self.ctx.params,
1477 };
1478 eval_expr(value, row, &eval_ctx)
1479 };
1480
1481 self.mutate_entity_target(entity_target, patch)
1482 }
1483
1484 ResolvedSetItem::SetLabels { variable, labels } => match row.get(*variable) {
1485 Some(LoraValue::Node(node_id)) => {
1486 let node_id = *node_id;
1487 for label in labels {
1488 let checked = if self.defer_existence {
1492 self.ctx
1493 .storage
1494 .check_node_add_label_deferring_existence(node_id, label)
1495 } else {
1496 self.ctx
1497 .storage
1498 .check_node_add_label_against_constraints(node_id, label)
1499 };
1500 checked.map_err(ExecutorError::ConstraintViolation)?;
1501 self.ctx.storage.add_node_label(node_id, label);
1502 }
1503 if self.defer_existence {
1504 self.pending_existence.push(EntityTarget::Node(node_id));
1505 }
1506 Ok(())
1507 }
1508 Some(other) => Err(ExecutorError::ExpectedNodeForSetLabels {
1509 found: value_kind(other),
1510 }),
1511 None => Err(ExecutorError::UnboundVariableForSet {
1512 var: format!("{variable:?}"),
1513 }),
1514 },
1515 }
1516 }
1517
1518 fn set_property_from_expr(
1519 &mut self,
1520 row: &Row,
1521 target_expr: &ResolvedExpr,
1522 new_value: LoraValue,
1523 ) -> ExecResult<()> {
1524 let ResolvedExpr::Property { expr, property } = target_expr else {
1525 return Err(ExecutorError::UnsupportedSetTarget);
1526 };
1527
1528 let owner = {
1529 let eval_ctx = EvalContext {
1530 storage: &*self.ctx.storage,
1531 params: &self.ctx.params,
1532 };
1533 eval_expr(expr, row, &eval_ctx)
1534 };
1535
1536 if matches!(new_value, LoraValue::Null) {
1538 return match owner {
1539 LoraValue::Node(node_id) => {
1540 self.remove_entity_property(EntityTarget::Node(node_id), property)
1541 }
1542 LoraValue::Relationship(rel_id) => {
1543 self.remove_entity_property(EntityTarget::Relationship(rel_id), property)
1544 }
1545 other => Err(ExecutorError::InvalidSetTarget {
1546 found: value_kind(&other),
1547 }),
1548 };
1549 }
1550
1551 match owner {
1552 LoraValue::Node(node_id) => {
1553 let prop = lora_value_to_property(new_value)
1554 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1555 if let Err(msg) = self
1556 .ctx
1557 .storage
1558 .check_node_set_property_against_constraints(node_id, property, &prop)
1559 {
1560 return Err(ExecutorError::ConstraintViolation(msg));
1561 }
1562 self.ctx
1563 .storage
1564 .set_node_property(node_id, property.clone(), prop);
1565 Ok(())
1566 }
1567 LoraValue::Relationship(rel_id) => {
1568 let prop = lora_value_to_property(new_value)
1569 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1570 if let Err(msg) = self
1571 .ctx
1572 .storage
1573 .check_relationship_set_property_against_constraints(rel_id, property, &prop)
1574 {
1575 return Err(ExecutorError::ConstraintViolation(msg));
1576 }
1577 self.ctx
1578 .storage
1579 .set_relationship_property(rel_id, property.clone(), prop);
1580 Ok(())
1581 }
1582 other => Err(ExecutorError::InvalidSetTarget {
1583 found: value_kind(&other),
1584 }),
1585 }
1586 }
1587
1588 fn remove_entity_property(&mut self, target: EntityTarget, property: &str) -> ExecResult<()> {
1593 if self.defer_existence {
1594 match target {
1595 EntityTarget::Node(node_id) => {
1596 self.ctx.storage.remove_node_property(node_id, property);
1597 }
1598 EntityTarget::Relationship(rel_id) => {
1599 self.ctx
1600 .storage
1601 .remove_relationship_property(rel_id, property);
1602 }
1603 }
1604 self.pending_existence.push(target);
1605 return Ok(());
1606 }
1607 match target {
1608 EntityTarget::Node(node_id) => {
1609 if let Err(msg) = self
1610 .ctx
1611 .storage
1612 .check_node_remove_property_against_constraints(node_id, property)
1613 {
1614 return Err(ExecutorError::ConstraintViolation(msg));
1615 }
1616 self.ctx.storage.remove_node_property(node_id, property);
1617 }
1618 EntityTarget::Relationship(rel_id) => {
1619 if let Err(msg) = self
1620 .ctx
1621 .storage
1622 .check_relationship_remove_property_against_constraints(rel_id, property)
1623 {
1624 return Err(ExecutorError::ConstraintViolation(msg));
1625 }
1626 self.ctx
1627 .storage
1628 .remove_relationship_property(rel_id, property);
1629 }
1630 }
1631 Ok(())
1632 }
1633
1634 fn remove_property_from_expr(&mut self, row: &Row, expr: &ResolvedExpr) -> ExecResult<()> {
1635 let ResolvedExpr::Property {
1636 expr: owner_expr,
1637 property,
1638 } = expr
1639 else {
1640 return Err(ExecutorError::UnsupportedRemoveTarget);
1641 };
1642
1643 let owner = {
1644 let eval_ctx = EvalContext {
1645 storage: &*self.ctx.storage,
1646 params: &self.ctx.params,
1647 };
1648 eval_expr(owner_expr, row, &eval_ctx)
1649 };
1650
1651 match owner {
1652 LoraValue::Node(node_id) => {
1653 self.remove_entity_property(EntityTarget::Node(node_id), property)
1654 }
1655 LoraValue::Relationship(rel_id) => {
1656 self.remove_entity_property(EntityTarget::Relationship(rel_id), property)
1657 }
1658 other => Err(ExecutorError::InvalidRemoveTarget {
1659 found: value_kind(&other),
1660 }),
1661 }
1662 }
1663
1664 fn overwrite_entity_target(
1665 &mut self,
1666 target: EntityTarget,
1667 new_value: LoraValue,
1668 ) -> ExecResult<()> {
1669 let LoraValue::Map(map) = new_value else {
1670 return Err(ExecutorError::ExpectedPropertyMap {
1671 found: value_kind(&new_value),
1672 });
1673 };
1674
1675 let mut props: Properties = Properties::new();
1676 for (k, v) in map {
1677 if matches!(v, LoraValue::Null) {
1679 continue;
1680 }
1681 let prop = lora_value_to_property(v)
1682 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1683 props.insert(lora_store::intern_owned(k), prop);
1684 }
1685
1686 let storage = &*self.ctx.storage;
1690 let checked = match (target, self.defer_existence) {
1691 (EntityTarget::Node(id), true) => {
1692 storage.check_node_replace_properties_deferring_existence(id, &props)
1693 }
1694 (EntityTarget::Node(id), false) => {
1695 storage.check_node_replace_properties_against_constraints(id, &props)
1696 }
1697 (EntityTarget::Relationship(id), true) => {
1698 storage.check_relationship_replace_properties_deferring_existence(id, &props)
1699 }
1700 (EntityTarget::Relationship(id), false) => {
1701 storage.check_relationship_replace_properties_against_constraints(id, &props)
1702 }
1703 };
1704 checked.map_err(ExecutorError::ConstraintViolation)?;
1705 match target {
1706 EntityTarget::Node(node_id) => {
1707 self.ctx.storage.replace_node_properties(node_id, props);
1708 }
1709 EntityTarget::Relationship(rel_id) => {
1710 self.ctx
1711 .storage
1712 .replace_relationship_properties(rel_id, props);
1713 }
1714 }
1715 if self.defer_existence {
1716 self.pending_existence.push(target);
1717 }
1718 Ok(())
1719 }
1720
1721 fn mutate_entity_target(
1722 &mut self,
1723 target: EntityTarget,
1724 patch_value: LoraValue,
1725 ) -> ExecResult<()> {
1726 let LoraValue::Map(map) = patch_value else {
1727 return Err(ExecutorError::ExpectedPropertyMap {
1728 found: value_kind(&patch_value),
1729 });
1730 };
1731
1732 match target {
1733 EntityTarget::Node(node_id) => {
1734 for (k, v) in map {
1735 if matches!(v, LoraValue::Null) {
1737 self.remove_entity_property(target, &k)?;
1738 continue;
1739 }
1740 let prop = lora_value_to_property(v)
1741 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1742 if let Err(msg) = self
1743 .ctx
1744 .storage
1745 .check_node_set_property_against_constraints(node_id, &k, &prop)
1746 {
1747 return Err(ExecutorError::ConstraintViolation(msg));
1748 }
1749 self.ctx.storage.set_node_property(node_id, k, prop);
1750 }
1751 }
1752 EntityTarget::Relationship(rel_id) => {
1753 for (k, v) in map {
1754 if matches!(v, LoraValue::Null) {
1755 self.remove_entity_property(target, &k)?;
1756 continue;
1757 }
1758 let prop = lora_value_to_property(v)
1759 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
1760 if let Err(msg) = self
1761 .ctx
1762 .storage
1763 .check_relationship_set_property_against_constraints(rel_id, &k, &prop)
1764 {
1765 return Err(ExecutorError::ConstraintViolation(msg));
1766 }
1767 self.ctx.storage.set_relationship_property(rel_id, k, prop);
1768 }
1769 }
1770 }
1771 Ok(())
1772 }
1773
1774 pub(crate) fn apply_create_pattern(
1775 &mut self,
1776 row: &mut Row,
1777 pattern: &ResolvedPattern,
1778 ) -> ExecResult<()> {
1779 for part in &pattern.parts {
1780 self.apply_create_pattern_part(row, part)?;
1781 }
1782 Ok(())
1783 }
1784
1785 pub(crate) fn apply_write_op(&mut self, op: &PhysicalOp, row: &mut Row) -> ExecResult<()> {
1791 match op {
1792 PhysicalOp::Create(c) => self.apply_create_pattern(row, &c.pattern),
1793 PhysicalOp::Set(s) => {
1794 for item in &s.items {
1795 self.apply_set_item(row, item)?;
1796 }
1797 Ok(())
1798 }
1799 PhysicalOp::Delete(d) => {
1800 let detach = d.detach;
1801 for expr in &d.expressions {
1802 let value = {
1803 let eval_ctx = EvalContext {
1804 storage: &*self.ctx.storage,
1805 params: &self.ctx.params,
1806 };
1807 eval_expr(expr, row, &eval_ctx)
1808 };
1809 self.delete_value(value, detach)?;
1810 }
1811 Ok(())
1812 }
1813 PhysicalOp::Remove(r) => {
1814 for item in &r.items {
1815 self.apply_remove_item(row, item)?;
1816 }
1817 Ok(())
1818 }
1819 PhysicalOp::Merge(m) => {
1820 let matched = self.match_merge_pattern(row, &m.pattern_part)?;
1821 if !matched {
1822 self.apply_create_pattern_part(row, &m.pattern_part)?;
1823 }
1824 for action in &m.actions {
1825 if action.on_match == matched {
1826 for item in &action.set.items {
1827 self.apply_set_item(row, item)?;
1828 }
1829 }
1830 }
1831 Ok(())
1832 }
1833 other => Err(ExecutorError::RuntimeError(format!(
1834 "apply_write_op called on non-write op: {other:?}"
1835 ))),
1836 }
1837 }
1838
1839 fn apply_create_pattern_part(
1842 &mut self,
1843 row: &mut Row,
1844 part: &ResolvedPatternPart,
1845 ) -> ExecResult<()> {
1846 let path = self.apply_create_pattern_element(row, &part.element)?;
1847 if let (Some(var), Some(path)) = (part.binding, path) {
1848 row.insert(var, LoraValue::Path(path));
1849 }
1850 Ok(())
1851 }
1852
1853 fn apply_create_pattern_element(
1854 &mut self,
1855 row: &mut Row,
1856 element: &ResolvedPatternElement,
1857 ) -> ExecResult<Option<LoraPath>> {
1858 match element {
1859 ResolvedPatternElement::Node {
1860 var,
1861 labels,
1862 properties,
1863 } => {
1864 let node_id =
1865 self.materialize_node_pattern(row, *var, labels, properties.as_ref())?;
1866 Ok(Some(LoraPath {
1867 nodes: vec![node_id],
1868 rels: Vec::new(),
1869 }))
1870 }
1871
1872 ResolvedPatternElement::NodeChain { head, chain } => {
1873 let mut current_node_id = self.materialize_node_pattern(
1874 row,
1875 head.var,
1876 &head.labels,
1877 head.properties.as_ref(),
1878 )?;
1879 let mut path = LoraPath {
1880 nodes: Vec::with_capacity(chain.len() + 1),
1881 rels: Vec::with_capacity(chain.len()),
1882 };
1883 path.nodes.push(current_node_id);
1884
1885 for link in chain {
1886 let next_node_id = self.materialize_node_pattern(
1887 row,
1888 link.node.var,
1889 &link.node.labels,
1890 link.node.properties.as_ref(),
1891 )?;
1892
1893 let rel_id = self.materialize_relationship_pattern(
1894 row,
1895 current_node_id,
1896 next_node_id,
1897 &link.rel,
1898 )?;
1899 path.rels.push(rel_id);
1900 path.nodes.push(next_node_id);
1901
1902 current_node_id = next_node_id;
1903 }
1904
1905 Ok(Some(path))
1906 }
1907
1908 ResolvedPatternElement::ShortestPath { .. } => {
1909 Ok(None)
1911 }
1912 }
1913 }
1914
1915 fn pattern_part_is_bound(&self, row: &Row, part: &ResolvedPatternPart) -> ExecResult<bool> {
1921 fn node_bound(row: &Row, var: Option<VarId>) -> ExecResult<bool> {
1922 let Some(var) = var else { return Ok(false) };
1923 match row.get(var) {
1924 None => Ok(false),
1925 Some(LoraValue::Node(_)) => Ok(true),
1926 Some(other) => Err(ExecutorError::ExpectedNodeForCreate {
1927 var: bound_var_name(row, var),
1928 found: value_kind(other),
1929 }),
1930 }
1931 }
1932 fn rel_bound(row: &Row, var: Option<VarId>) -> ExecResult<bool> {
1933 let Some(var) = var else { return Ok(false) };
1937 match row.get(var) {
1938 None => Ok(false),
1939 Some(LoraValue::Relationship(_)) => Ok(true),
1940 Some(other) => Err(ExecutorError::ExpectedRelationshipForCreate {
1941 var: bound_var_name(row, var),
1942 found: value_kind(other),
1943 }),
1944 }
1945 }
1946
1947 match &part.element {
1948 ResolvedPatternElement::Node { var, .. } => node_bound(row, *var),
1949
1950 ResolvedPatternElement::ShortestPath { .. } => Ok(false),
1951
1952 ResolvedPatternElement::NodeChain { head, chain } => {
1953 let mut all_bound = node_bound(row, head.var)?;
1956 for link in chain {
1957 let node_ok = node_bound(row, link.node.var)?;
1958 let rel_ok = rel_bound(row, link.rel.var)?;
1959 all_bound &= node_ok && rel_ok;
1960 }
1961 Ok(all_bound)
1962 }
1963 }
1964 }
1965
1966 fn materialize_node_pattern(
1967 &mut self,
1968 row: &mut Row,
1969 var: Option<VarId>,
1970 labels: &[Vec<String>],
1971 properties: Option<&ResolvedExpr>,
1972 ) -> ExecResult<u64> {
1973 if let Some(var_id) = var {
1974 match row.get(var_id) {
1975 Some(LoraValue::Node(id)) => return Ok(*id),
1976 Some(other) => {
1981 return Err(ExecutorError::ExpectedNodeForCreate {
1982 var: bound_var_name(row, var_id),
1983 found: value_kind(other),
1984 });
1985 }
1986 None => {}
1987 }
1988 }
1989
1990 let properties = match properties {
1991 Some(expr) => eval_properties_expr(expr, row, &*self.ctx.storage, &self.ctx.params)?,
1992 None => Properties::new(),
1993 };
1994
1995 let flat_labels = flatten_label_groups(labels);
1996 debug!("creating node with labels={flat_labels:?}");
1997 let checked = if self.defer_existence {
1998 self.ctx
1999 .storage
2000 .check_node_create_deferring_existence(&flat_labels, &properties)
2001 } else {
2002 self.ctx
2003 .storage
2004 .check_node_create_against_constraints(&flat_labels, &properties)
2005 };
2006 checked.map_err(ExecutorError::ConstraintViolation)?;
2007 let created = self
2008 .ctx
2009 .storage
2010 .try_create_node(flat_labels, properties)
2011 .ok_or(ExecutorError::NodeCreateFailed)?;
2012 if self.defer_existence {
2013 self.pending_existence.push(EntityTarget::Node(created.id));
2014 }
2015
2016 if let Some(var_id) = var {
2017 row.insert(var_id, LoraValue::Node(created.id));
2018 }
2019
2020 Ok(created.id)
2021 }
2022
2023 fn materialize_relationship_pattern(
2024 &mut self,
2025 row: &mut Row,
2026 left_node_id: u64,
2027 right_node_id: u64,
2028 rel: &lora_analyzer::ResolvedRel,
2029 ) -> ExecResult<u64> {
2030 if let Some(var_id) = rel.var {
2031 if let Some(other) = row
2032 .get(var_id)
2033 .filter(|v| !matches!(v, LoraValue::Relationship(_)))
2034 {
2035 return Err(ExecutorError::ExpectedRelationshipForCreate {
2036 var: bound_var_name(row, var_id),
2037 found: value_kind(other),
2038 });
2039 }
2040 if let Some(LoraValue::Relationship(id)) = row.get(var_id) {
2041 let id = *id;
2042 if let Some((src, dst)) = self.ctx.storage.relationship_endpoints(id) {
2043 let endpoints_match = match rel.direction {
2044 Direction::Right | Direction::Undirected => {
2045 src == left_node_id && dst == right_node_id
2046 }
2047 Direction::Left => src == right_node_id && dst == left_node_id,
2048 };
2049
2050 if endpoints_match {
2051 return Ok(id);
2052 }
2053 }
2054 }
2055 }
2056
2057 if rel.range.is_some() {
2058 return Err(ExecutorError::UnsupportedCreateRelationshipRange);
2059 }
2060
2061 let (src, dst) = match rel.direction {
2062 Direction::Right | Direction::Undirected => (left_node_id, right_node_id),
2063 Direction::Left => (right_node_id, left_node_id),
2064 };
2065
2066 let rel_type = rel
2067 .types
2068 .first()
2069 .ok_or(ExecutorError::MissingRelationshipType)?;
2070
2071 if rel_type.is_empty() {
2072 return Err(ExecutorError::MissingRelationshipType);
2073 }
2074
2075 let properties = match rel.properties.as_ref() {
2076 Some(expr) => eval_properties_expr(expr, row, &*self.ctx.storage, &self.ctx.params)?,
2077 None => Properties::new(),
2078 };
2079
2080 debug!("creating relationship: src={src}, dst={dst}, type={rel_type}");
2081
2082 let checked = if self.defer_existence {
2083 self.ctx
2084 .storage
2085 .check_relationship_create_deferring_existence(rel_type, &properties)
2086 } else {
2087 self.ctx
2088 .storage
2089 .check_relationship_create_against_constraints(rel_type, &properties)
2090 };
2091 checked.map_err(ExecutorError::ConstraintViolation)?;
2092
2093 let created = self
2094 .ctx
2095 .storage
2096 .create_relationship(src, dst, rel_type, properties)
2097 .ok_or_else(|| ExecutorError::RelationshipCreateFailed {
2098 src,
2099 dst,
2100 rel_type: rel_type.clone(),
2101 })?;
2102 if self.defer_existence {
2103 self.pending_existence
2104 .push(EntityTarget::Relationship(created.id));
2105 }
2106
2107 if let Some(var_id) = rel.var {
2108 row.insert(var_id, LoraValue::Relationship(created.id));
2109 }
2110
2111 Ok(created.id)
2112 }
2113}
2114
2115fn bound_pattern_path(row: &Row, part: &ResolvedPatternPart) -> Option<LoraPath> {
2118 let node = |var: Option<VarId>| match var.and_then(|v| row.get(v)) {
2119 Some(LoraValue::Node(id)) => Some(*id),
2120 _ => None,
2121 };
2122 match &part.element {
2123 ResolvedPatternElement::Node { var, .. } => Some(LoraPath {
2124 nodes: vec![node(*var)?],
2125 rels: Vec::new(),
2126 }),
2127 ResolvedPatternElement::NodeChain { head, chain } => {
2128 let mut path = LoraPath {
2129 nodes: vec![node(head.var)?],
2130 rels: Vec::with_capacity(chain.len()),
2131 };
2132 for link in chain {
2133 match link.rel.var.and_then(|v| row.get(v)) {
2134 Some(LoraValue::Relationship(id)) => path.rels.push(*id),
2135 _ => return None,
2136 }
2137 path.nodes.push(node(link.node.var)?);
2138 }
2139 Some(path)
2140 }
2141 ResolvedPatternElement::ShortestPath { .. } => None,
2142 }
2143}
2144
2145pub(crate) fn plan_defers_existence(plan: &PhysicalPlan) -> bool {
2150 let mut creates = 0;
2151 let mut deletes = false;
2152 for op in &plan.nodes {
2153 match op {
2154 PhysicalOp::Create(_) => creates += 1,
2155 PhysicalOp::Delete(_) => deletes = true,
2157 PhysicalOp::Merge(_)
2158 | PhysicalOp::Set(_)
2159 | PhysicalOp::Remove(_)
2160 | PhysicalOp::Foreach(_) => return true,
2161 _ => {}
2162 }
2163 }
2164 creates > 1 || creates == 1 && deletes
2165}
2166
2167pub(crate) fn plan_ends_in_write(plan: &PhysicalPlan) -> bool {
2172 match &plan.nodes[plan.root] {
2173 PhysicalOp::Create(_)
2174 | PhysicalOp::Merge(_)
2175 | PhysicalOp::Set(_)
2176 | PhysicalOp::Delete(_)
2177 | PhysicalOp::Remove(_)
2178 | PhysicalOp::Foreach(_) => true,
2179 PhysicalOp::CallSubquery(op) => op.new_vars.is_empty(),
2181 _ => false,
2182 }
2183}
2184
2185fn merge_candidates_from_index<S: lora_store::GraphStorage>(
2192 storage: &S,
2193 labels: &[Vec<String>],
2194 expected: &std::collections::BTreeMap<String, LoraValue>,
2195) -> Option<Vec<lora_store::NodeId>> {
2196 use lora_store::PropertyValue;
2197
2198 let label = match labels {
2201 [group] if group.len() == 1 => Some(group[0].as_str()),
2202 _ => None,
2203 };
2204 let (key, value) = expected.iter().find(|(_, v)| {
2205 matches!(
2206 v,
2207 LoraValue::String(_) | LoraValue::Bool(_) | LoraValue::Int(_)
2208 ) || matches!(v, LoraValue::Float(f) if f.is_finite() && f.abs() < 9_007_199_254_740_992.0)
2209 })?;
2210 let images: Vec<PropertyValue> = match value {
2211 LoraValue::String(s) => vec![PropertyValue::String(s.clone())],
2212 LoraValue::Bool(b) => vec![PropertyValue::Bool(*b)],
2213 LoraValue::Int(i) => vec![PropertyValue::Int(*i), PropertyValue::Float(*i as f64)],
2214 LoraValue::Float(f) => {
2215 let mut v = vec![PropertyValue::Float(*f)];
2216 if f.fract() == 0.0 {
2217 v.push(PropertyValue::Int(*f as i64));
2218 }
2219 v
2220 }
2221 _ => return None,
2222 };
2223 let mut ids: Vec<lora_store::NodeId> = images
2224 .iter()
2225 .flat_map(|image| storage.find_node_ids_by_property(label, key, image))
2226 .collect();
2227 ids.sort_unstable();
2228 ids.dedup();
2229 Some(ids)
2230}
2231
2232fn bound_var_name(row: &Row, var: VarId) -> String {
2234 row.iter_named()
2235 .find(|(key, _, _)| **key == var)
2236 .map(|(_, name, _)| name.into_owned())
2237 .unwrap_or_else(|| format!("{var:?}"))
2238}