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