1use crate::errors::{value_kind, ExecResult, ExecutorError};
12use crate::eval::{clear_eval_error, eval_expr, eval_expr_result, EvalContext};
13use crate::value::{LoraValue, Row};
14use crate::{project_rows, ExecuteOptions, QueryResult};
15
16use lora_analyzer::{ResolvedExpr, ResolvedProjection};
17use lora_ast::RangeLiteral;
18use lora_compiler::physical::*;
19use lora_compiler::CompiledQuery;
20use lora_store::{GraphStorage, Properties};
21
22use std::cmp::Ordering;
23use std::collections::BTreeMap;
24use std::time::Instant;
25use tracing::{error, trace};
26
27use super::helpers::{
28 build_path_value, check_deadline_at, compare_sort_item, compute_aggregate_expr, dedup_rows,
29 filter_rows_checked, filter_shortest_paths, hydrate_node_record, hydrate_relationship_record,
30 indexed_node_property_candidates, label_group_candidates_prefiltered,
31 node_matches_label_groups, node_matches_property_filter, project_rows_checked, resolve_range,
32 scan_node_ids_for_label_groups, value_matches_property_value, variable_length_expand,
33 GroupValueKey,
34};
35
36pub struct ExecutionContext<'a, S: GraphStorage> {
37 pub storage: &'a S,
38 pub params: BTreeMap<String, LoraValue>,
39}
40
41pub struct Executor<'a, S: GraphStorage> {
42 ctx: ExecutionContext<'a, S>,
43 deadline: Option<Instant>,
44}
45
46impl<'a, S: GraphStorage> Executor<'a, S> {
47 pub fn new(ctx: ExecutionContext<'a, S>) -> Self {
48 Self {
49 ctx,
50 deadline: None,
51 }
52 }
53
54 pub fn with_deadline(ctx: ExecutionContext<'a, S>, deadline: Option<Instant>) -> Self {
55 Self { ctx, deadline }
56 }
57
58 #[inline]
59 fn check_deadline(&self) -> ExecResult<()> {
60 if let Some(deadline) = self.deadline {
61 check_deadline_at(deadline)
62 } else {
63 Ok(())
64 }
65 }
66
67 #[inline]
68 fn check_loop_deadline(deadline: Option<Instant>) -> ExecResult<()> {
69 if let Some(deadline) = deadline {
70 check_deadline_at(deadline)
71 } else {
72 Ok(())
73 }
74 }
75}
76
77impl<'a, S: GraphStorage> Executor<'a, S> {
78 pub fn execute(
79 &self,
80 plan: &PhysicalPlan,
81 options: Option<ExecuteOptions>,
82 ) -> ExecResult<QueryResult> {
83 let rows = self.execute_rows(plan)?;
84 Ok(project_rows(rows, options.unwrap_or_default()))
85 }
86
87 pub fn execute_compiled(
88 &self,
89 compiled: &CompiledQuery,
90 options: Option<ExecuteOptions>,
91 ) -> ExecResult<QueryResult> {
92 let rows = self.execute_compiled_rows(compiled)?;
93 Ok(project_rows(rows, options.unwrap_or_default()))
94 }
95
96 pub fn execute_compiled_rows(&self, compiled: &CompiledQuery) -> ExecResult<Vec<Row>> {
97 self.check_deadline()?;
98 if compiled.unions.is_empty() {
99 return self.execute_rows(&compiled.physical);
100 }
101
102 clear_eval_error();
103
104 let mut all_rows = self.execute_rows(&compiled.physical)?;
105 let mut needs_dedup = false;
106
107 for branch in &compiled.unions {
108 self.check_deadline()?;
109 let branch_rows = self.execute_rows(&branch.physical)?;
110 all_rows.extend(branch_rows);
111
112 if !branch.all {
113 needs_dedup = true;
114 }
115 }
116
117 if needs_dedup {
118 all_rows = dedup_rows(all_rows);
119 }
120
121 Ok(all_rows)
122 }
123
124 pub fn execute_rows(&self, plan: &PhysicalPlan) -> ExecResult<Vec<Row>> {
125 self.check_deadline()?;
126 clear_eval_error();
129
130 let rows = self.execute_node(plan, plan.root)?;
131 Ok(rows
132 .into_iter()
133 .map(|row| self.hydrate_row(row))
134 .collect::<Vec<_>>())
135 }
136
137 fn hydrate_row(&self, row: Row) -> Row {
138 let mut out = Row::new();
139
140 for (var, name, value) in row.into_iter_named() {
141 out.insert_named(var, name, self.hydrate_value(value));
142 }
143
144 out
145 }
146
147 pub(crate) fn execute_subtree(
151 &self,
152 plan: &PhysicalPlan,
153 node_id: PhysicalNodeId,
154 ) -> ExecResult<Vec<Row>> {
155 self.execute_node(plan, node_id)
156 }
157
158 fn execute_node(&self, plan: &PhysicalPlan, node_id: PhysicalNodeId) -> ExecResult<Vec<Row>> {
159 self.check_deadline()?;
160 trace!("read-only execute_node start: node_id={node_id:?}");
161
162 let result = match &plan.nodes[node_id] {
163 PhysicalOp::Argument(op) => self.exec_argument(op),
164 PhysicalOp::NodeScan(op) => self.exec_node_scan(plan, op),
165 PhysicalOp::NodeByLabelScan(op) => self.exec_node_by_label_scan(plan, op),
166 PhysicalOp::NodeByPropertyScan(op) => self.exec_node_by_property_scan(plan, op),
167 PhysicalOp::Expand(op) => self.exec_expand(plan, op),
168 PhysicalOp::Filter(op) => self.exec_filter(plan, op),
169 PhysicalOp::Projection(op) => self.exec_projection(plan, op),
170 PhysicalOp::Unwind(op) => self.exec_unwind(plan, op),
171 PhysicalOp::HashAggregation(op) => self.exec_hash_aggregation(plan, op),
172 PhysicalOp::Sort(op) => self.exec_sort(plan, op),
173 PhysicalOp::Limit(op) => self.exec_limit(plan, op),
174 PhysicalOp::OptionalMatch(op) => self.exec_optional_match(plan, op),
175 PhysicalOp::PathBuild(op) => self.exec_path_build(plan, op),
176 PhysicalOp::Create(_) => Err(ExecutorError::ReadOnlyCreate { node_id }),
177 PhysicalOp::Merge(_) => Err(ExecutorError::ReadOnlyMerge { node_id }),
178 PhysicalOp::Delete(_) => Err(ExecutorError::ReadOnlyDelete { node_id }),
179 PhysicalOp::Set(_) => Err(ExecutorError::ReadOnlySet { node_id }),
180 PhysicalOp::Remove(_) => Err(ExecutorError::ReadOnlyRemove { node_id }),
181 };
182
183 match &result {
184 Ok(rows) => trace!(
185 "read-only execute_node ok: node_id={node_id:?}, rows={}",
186 rows.len()
187 ),
188 Err(err) => error!("read-only execute_node failed: node_id={node_id:?}, error={err}"),
189 }
190
191 result
192 }
193
194 fn exec_argument(&self, _op: &ArgumentExec) -> ExecResult<Vec<Row>> {
195 Ok(vec![Row::new()])
196 }
197
198 fn exec_node_scan(&self, plan: &PhysicalPlan, op: &NodeScanExec) -> ExecResult<Vec<Row>> {
199 let base_rows = match op.input {
200 Some(input) => self.execute_node(plan, input)?,
201 None => vec![Row::new()],
202 };
203
204 let node_ids = self.ctx.storage.all_node_ids();
205 let mut out = Vec::new();
206
207 let deadline = self.deadline;
208 for row in base_rows {
209 Self::check_loop_deadline(deadline)?;
210 if let Some(existing) = row.get(op.var) {
211 match existing {
212 LoraValue::Node(existing_id) => {
213 if self.ctx.storage.has_node(*existing_id) {
214 out.push(row);
215 }
216 }
217 other => {
218 return Err(ExecutorError::ExpectedNodeForExpand {
219 var: format!("{:?}", op.var),
220 found: value_kind(other),
221 });
222 }
223 }
224 continue;
225 }
226
227 for &id in &node_ids {
228 Self::check_loop_deadline(deadline)?;
229 let mut new_row = row.clone();
230 new_row.insert(op.var, LoraValue::Node(id));
231 out.push(new_row);
232 }
233 }
234
235 Ok(out)
236 }
237
238 fn exec_node_by_label_scan(
239 &self,
240 plan: &PhysicalPlan,
241 op: &NodeByLabelScanExec,
242 ) -> ExecResult<Vec<Row>> {
243 let base_rows = match op.input {
244 Some(input) => self.execute_node(plan, input)?,
245 None => vec![Row::new()],
246 };
247
248 let candidate_ids = scan_node_ids_for_label_groups(self.ctx.storage, &op.labels);
249 let candidates_prefiltered = label_group_candidates_prefiltered(&op.labels);
250 let mut out = Vec::new();
251
252 match self.deadline {
253 Some(deadline) => {
254 for row in base_rows {
255 check_deadline_at(deadline)?;
256 if let Some(existing) = row.get(op.var) {
257 match existing {
258 LoraValue::Node(existing_id) => {
259 let labels_ok = self
260 .ctx
261 .storage
262 .with_node(*existing_id, |n| {
263 node_matches_label_groups(&n.labels, &op.labels)
264 })
265 .unwrap_or(false);
266 if labels_ok {
267 out.push(row);
268 }
269 }
270 other => {
271 return Err(ExecutorError::ExpectedNodeForExpand {
272 var: format!("{:?}", op.var),
273 found: value_kind(other),
274 });
275 }
276 }
277 continue;
278 }
279
280 for &id in &candidate_ids {
281 check_deadline_at(deadline)?;
282 if !candidates_prefiltered {
283 let labels_ok = self
284 .ctx
285 .storage
286 .with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
287 .unwrap_or(false);
288 if !labels_ok {
289 continue;
290 }
291 }
292 let mut new_row = row.clone();
293 new_row.insert(op.var, LoraValue::Node(id));
294 out.push(new_row);
295 }
296 }
297 }
298 None => {
299 for row in base_rows {
300 if let Some(existing) = row.get(op.var) {
301 match existing {
302 LoraValue::Node(existing_id) => {
303 let labels_ok = self
304 .ctx
305 .storage
306 .with_node(*existing_id, |n| {
307 node_matches_label_groups(&n.labels, &op.labels)
308 })
309 .unwrap_or(false);
310 if labels_ok {
311 out.push(row);
312 }
313 }
314 other => {
315 return Err(ExecutorError::ExpectedNodeForExpand {
316 var: format!("{:?}", op.var),
317 found: value_kind(other),
318 });
319 }
320 }
321 continue;
322 }
323
324 for &id in &candidate_ids {
325 if !candidates_prefiltered {
326 let labels_ok = self
327 .ctx
328 .storage
329 .with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
330 .unwrap_or(false);
331 if !labels_ok {
332 continue;
333 }
334 }
335 let mut new_row = row.clone();
336 new_row.insert(op.var, LoraValue::Node(id));
337 out.push(new_row);
338 }
339 }
340 }
341 }
342
343 Ok(out)
344 }
345
346 fn exec_node_by_property_scan(
347 &self,
348 plan: &PhysicalPlan,
349 op: &NodeByPropertyScanExec,
350 ) -> ExecResult<Vec<Row>> {
351 let base_rows = match op.input {
352 Some(input) => self.execute_node(plan, input)?,
353 None => vec![Row::new()],
354 };
355
356 let eval_ctx = EvalContext {
357 storage: self.ctx.storage,
358 params: &self.ctx.params,
359 };
360 let mut out = Vec::new();
361
362 let deadline = self.deadline;
363 for row in base_rows {
364 Self::check_loop_deadline(deadline)?;
365 let expected = eval_expr(&op.value, &row, &eval_ctx);
366
367 if let Some(existing) = row.get(op.var) {
368 match existing {
369 LoraValue::Node(existing_id) => {
370 if node_matches_property_filter(
371 self.ctx.storage,
372 *existing_id,
373 &op.labels,
374 &op.key,
375 &expected,
376 ) {
377 out.push(row);
378 }
379 }
380 other => {
381 return Err(ExecutorError::ExpectedNodeForExpand {
382 var: format!("{:?}", op.var),
383 found: value_kind(other),
384 });
385 }
386 }
387 continue;
388 }
389
390 let candidates =
391 indexed_node_property_candidates(self.ctx.storage, &op.labels, &op.key, &expected);
392 for id in candidates.ids {
393 Self::check_loop_deadline(deadline)?;
394 if !candidates.prefiltered
395 && !node_matches_property_filter(
396 self.ctx.storage,
397 id,
398 &op.labels,
399 &op.key,
400 &expected,
401 )
402 {
403 continue;
404 }
405 let mut new_row = row.clone();
406 new_row.insert(op.var, LoraValue::Node(id));
407 out.push(new_row);
408 }
409 }
410
411 Ok(out)
412 }
413
414 fn exec_expand(&self, plan: &PhysicalPlan, op: &ExpandExec) -> ExecResult<Vec<Row>> {
415 if let Some(range) = &op.range {
417 return self.exec_expand_var_len(plan, op, range);
418 }
419
420 let input_rows = self.execute_node(plan, op.input)?;
421 let mut out = Vec::new();
422
423 for row in input_rows {
424 let src_node_id = match row.get(op.src) {
425 Some(LoraValue::Node(id)) => *id,
426 Some(other) => {
427 return Err(ExecutorError::ExpectedNodeForExpand {
428 var: format!("{:?}", op.src),
429 found: value_kind(other),
430 });
431 }
432 None => continue,
433 };
434
435 for (rel_id, dst_id) in
436 self.ctx
437 .storage
438 .expand_ids(src_node_id, op.direction, &op.types)
439 {
440 if let Some(expr) = op.rel_properties.as_ref() {
441 let actual_props = self
442 .ctx
443 .storage
444 .with_relationship(rel_id, |rel| rel.properties.clone());
445 let matches = match actual_props {
446 Some(props) => {
447 self.relationship_matches_properties(&props, Some(expr), &row)?
448 }
449 None => false,
450 };
451 if !matches {
452 continue;
453 }
454 }
455
456 if let Some(existing_dst) = row.get(op.dst) {
457 match existing_dst {
458 LoraValue::Node(existing_id) if *existing_id == dst_id => {}
459 LoraValue::Node(_) => continue,
460 other => {
461 return Err(ExecutorError::ExpectedNodeForExpand {
462 var: format!("{:?}", op.dst),
463 found: value_kind(other),
464 });
465 }
466 }
467 }
468
469 if let Some(rel_var) = op.rel {
470 if let Some(existing_rel) = row.get(rel_var) {
471 match existing_rel {
472 LoraValue::Relationship(existing_id) if *existing_id == rel_id => {}
473 LoraValue::Relationship(_) => continue,
474 other => {
475 return Err(ExecutorError::ExpectedRelationshipForExpand {
476 var: format!("{:?}", rel_var),
477 found: value_kind(other),
478 });
479 }
480 }
481 }
482 }
483
484 let mut new_row = row.clone();
485
486 if !new_row.contains_key(op.dst) {
487 new_row.insert(op.dst, LoraValue::Node(dst_id));
488 }
489
490 if let Some(rel_var) = op.rel {
491 if !new_row.contains_key(rel_var) {
492 new_row.insert(rel_var, LoraValue::Relationship(rel_id));
493 }
494 }
495
496 out.push(new_row);
497 }
498 }
499
500 Ok(out)
501 }
502
503 fn exec_expand_var_len(
504 &self,
505 plan: &PhysicalPlan,
506 op: &ExpandExec,
507 range: &RangeLiteral,
508 ) -> ExecResult<Vec<Row>> {
509 let input_rows = self.execute_node(plan, op.input)?;
510 let (min_hops, max_hops) = resolve_range(range);
511 let mut out = Vec::new();
512
513 for row in input_rows {
514 let src_node_id = match row.get(op.src) {
515 Some(LoraValue::Node(id)) => *id,
516 Some(other) => {
517 return Err(ExecutorError::ExpectedNodeForExpand {
518 var: format!("{:?}", op.src),
519 found: value_kind(other),
520 });
521 }
522 None => continue,
523 };
524
525 let expansions = variable_length_expand(
526 self.ctx.storage,
527 src_node_id,
528 op.direction,
529 &op.types,
530 min_hops,
531 max_hops,
532 );
533
534 for result in expansions {
535 let mut new_row = row.clone();
536 new_row.insert(op.dst, LoraValue::Node(result.dst_node_id));
537
538 if let Some(rel_var) = op.rel {
541 let rel_list = LoraValue::List(
543 result
544 .rel_ids
545 .into_iter()
546 .map(LoraValue::Relationship)
547 .collect(),
548 );
549 new_row.insert(rel_var, rel_list);
550 }
551
552 out.push(new_row);
553 }
554 }
555
556 Ok(out)
557 }
558
559 fn relationship_matches_properties(
560 &self,
561 actual: &Properties,
562 expected_expr: Option<&ResolvedExpr>,
563 row: &Row,
564 ) -> ExecResult<bool> {
565 let Some(expr) = expected_expr else {
566 return Ok(true);
567 };
568
569 let eval_ctx = EvalContext {
570 storage: self.ctx.storage,
571 params: &self.ctx.params,
572 };
573
574 let expected = eval_expr(expr, row, &eval_ctx);
575
576 let LoraValue::Map(expected_map) = expected else {
577 return Err(ExecutorError::ExpectedPropertyMap {
578 found: value_kind(&expected),
579 });
580 };
581
582 Ok(expected_map.iter().all(|(key, expected_value)| {
583 actual
584 .get(key)
585 .map(|actual_value| value_matches_property_value(expected_value, actual_value))
586 .unwrap_or(false)
587 }))
588 }
589
590 fn exec_filter(&self, plan: &PhysicalPlan, op: &FilterExec) -> ExecResult<Vec<Row>> {
591 let input_rows = self.execute_node(plan, op.input)?;
592 let eval_ctx = EvalContext {
593 storage: self.ctx.storage,
594 params: &self.ctx.params,
595 };
596
597 filter_rows_checked(input_rows, &op.predicate, &eval_ctx)
598 }
599
600 fn exec_projection(&self, plan: &PhysicalPlan, op: &ProjectionExec) -> ExecResult<Vec<Row>> {
601 let input_rows = self.execute_node(plan, op.input)?;
602 let eval_ctx = EvalContext {
603 storage: self.ctx.storage,
604 params: &self.ctx.params,
605 };
606
607 project_rows_checked(input_rows, op, &eval_ctx)
608 }
609
610 fn hydrate_value(&self, value: LoraValue) -> LoraValue {
611 match value {
612 LoraValue::Node(id) => self.hydrate_node(id),
613 LoraValue::Relationship(id) => self.hydrate_relationship(id),
614 LoraValue::List(values) => {
615 LoraValue::List(values.into_iter().map(|v| self.hydrate_value(v)).collect())
616 }
617 LoraValue::Map(map) => LoraValue::Map(
618 map.into_iter()
619 .map(|(k, v)| (k, self.hydrate_value(v)))
620 .collect(),
621 ),
622 other => other,
623 }
624 }
625
626 fn hydrate_node(&self, id: u64) -> LoraValue {
627 self.ctx
628 .storage
629 .with_node(id, hydrate_node_record)
630 .unwrap_or(LoraValue::Null)
631 }
632
633 fn hydrate_relationship(&self, id: u64) -> LoraValue {
634 self.ctx
635 .storage
636 .with_relationship(id, hydrate_relationship_record)
637 .unwrap_or(LoraValue::Null)
638 }
639
640 fn exec_unwind(&self, plan: &PhysicalPlan, op: &UnwindExec) -> ExecResult<Vec<Row>> {
641 let input_rows = self.execute_node(plan, op.input)?;
642 let eval_ctx = EvalContext {
643 storage: self.ctx.storage,
644 params: &self.ctx.params,
645 };
646
647 let mut out = Vec::new();
648
649 for row in input_rows {
650 match eval_expr(&op.expr, &row, &eval_ctx) {
651 LoraValue::List(values) => {
652 for value in values {
653 let mut new_row = row.clone();
654 new_row.insert(op.alias, value);
655 out.push(new_row);
656 }
657 }
658 LoraValue::Null => {}
659 other => {
660 let mut new_row = row;
661 new_row.insert(op.alias, other);
662 out.push(new_row);
663 }
664 }
665 }
666
667 Ok(out)
668 }
669
670 fn exec_hash_aggregation(
671 &self,
672 plan: &PhysicalPlan,
673 op: &HashAggregationExec,
674 ) -> ExecResult<Vec<Row>> {
675 let input_rows = self.execute_node(plan, op.input)?;
676 let eval_ctx = EvalContext {
677 storage: self.ctx.storage,
678 params: &self.ctx.params,
679 };
680
681 if let Some(specs) = crate::pull::classify_streamable_aggregates(&op.aggregates) {
687 return self.exec_hash_aggregation_streaming(
688 input_rows,
689 &op.group_by,
690 &op.aggregates,
691 &specs,
692 &eval_ctx,
693 );
694 }
695
696 let mut groups: BTreeMap<Vec<GroupValueKey>, Vec<Row>> = BTreeMap::new();
697
698 if op.group_by.is_empty() {
699 groups.insert(Vec::new(), input_rows);
700 } else {
701 for row in input_rows {
702 let mut key = Vec::with_capacity(op.group_by.len());
703 for proj in &op.group_by {
704 let value = eval_expr_result(&proj.expr, &row, &eval_ctx)
705 .map_err(ExecutorError::RuntimeError)?;
706 key.push(GroupValueKey::from_value(&value));
707 }
708
709 groups.entry(key).or_default().push(row);
710 }
711 }
712
713 let mut out = Vec::new();
714
715 for rows in groups.into_values() {
716 let mut result = Row::new();
717
718 if let Some(first) = rows.first() {
719 for proj in &op.group_by {
720 let value = eval_expr_result(&proj.expr, first, &eval_ctx)
721 .map_err(ExecutorError::RuntimeError)?;
722 let value = self.hydrate_value(value);
723 result.insert_named(proj.output, proj.name.clone(), value);
724 }
725 }
726
727 for proj in &op.aggregates {
728 let value = compute_aggregate_expr(&proj.expr, &rows, &eval_ctx)?;
729 result.insert_named(proj.output, proj.name.clone(), value);
730 }
731
732 out.push(result);
733 }
734
735 Ok(out)
736 }
737
738 fn exec_hash_aggregation_streaming(
739 &self,
740 input_rows: Vec<Row>,
741 group_by: &[ResolvedProjection],
742 aggregates: &[ResolvedProjection],
743 specs: &[crate::pull::StreamableAggSpec],
744 eval_ctx: &EvalContext<'_, S>,
745 ) -> ExecResult<Vec<Row>> {
746 if group_by.is_empty() {
748 if specs
753 .iter()
754 .all(|s| matches!(s.kind, crate::pull::StreamableAggKind::CountAll))
755 {
756 let count = LoraValue::Int(input_rows.len() as i64);
757 let mut result = Row::new();
758 for proj in aggregates {
759 result.insert_named(proj.output, proj.name.clone(), count.clone());
760 }
761 return Ok(vec![result]);
762 }
763
764 let mut aggs: Vec<crate::pull::AggState> = specs
765 .iter()
766 .map(|s| crate::pull::AggState::seed(s.kind))
767 .collect();
768 for row in &input_rows {
769 for (i, spec) in specs.iter().enumerate() {
770 let value = match &spec.arg {
771 Some(arg) => eval_expr_result(arg, row, eval_ctx)
772 .map_err(ExecutorError::RuntimeError)?,
773 None => LoraValue::Null,
774 };
775 aggs[i].fold(spec.kind, value);
776 }
777 }
778 let mut result = Row::new();
779 for (i, proj) in aggregates.iter().enumerate() {
780 let value =
781 std::mem::replace(&mut aggs[i], crate::pull::AggState::seed(specs[i].kind))
782 .finalize(specs[i].kind);
783 result.insert_named(proj.output, proj.name.clone(), value);
784 }
785 return Ok(vec![result]);
786 }
787
788 let mut groups: BTreeMap<Vec<GroupValueKey>, (Row, Vec<crate::pull::AggState>)> =
792 BTreeMap::new();
793
794 for row in input_rows {
795 let mut key = Vec::with_capacity(group_by.len());
796 for proj in group_by {
797 let value = eval_expr_result(&proj.expr, &row, eval_ctx)
798 .map_err(ExecutorError::RuntimeError)?;
799 key.push(GroupValueKey::from_value(&value));
800 }
801
802 let entry = groups.entry(key).or_insert_with(|| {
803 (
804 row.clone(),
805 specs
806 .iter()
807 .map(|s| crate::pull::AggState::seed(s.kind))
808 .collect(),
809 )
810 });
811
812 for (i, spec) in specs.iter().enumerate() {
813 let value = match &spec.arg {
814 Some(arg) => eval_expr_result(arg, &row, eval_ctx)
815 .map_err(ExecutorError::RuntimeError)?,
816 None => LoraValue::Null,
817 };
818 entry.1[i].fold(spec.kind, value);
819 }
820 }
821
822 let mut out = Vec::with_capacity(groups.len());
823 for (_, (first_row, mut aggs)) in groups {
824 let mut result = Row::new();
825 for proj in group_by {
826 let value = eval_expr_result(&proj.expr, &first_row, eval_ctx)
827 .map_err(ExecutorError::RuntimeError)?;
828 let value = self.hydrate_value(value);
829 result.insert_named(proj.output, proj.name.clone(), value);
830 }
831 for (i, proj) in aggregates.iter().enumerate() {
832 let value =
833 std::mem::replace(&mut aggs[i], crate::pull::AggState::seed(specs[i].kind))
834 .finalize(specs[i].kind);
835 result.insert_named(proj.output, proj.name.clone(), value);
836 }
837 out.push(result);
838 }
839 Ok(out)
840 }
841
842 fn exec_sort(&self, plan: &PhysicalPlan, op: &SortExec) -> ExecResult<Vec<Row>> {
843 let mut rows = self.execute_node(plan, op.input)?;
844 let eval_ctx = EvalContext {
845 storage: self.ctx.storage,
846 params: &self.ctx.params,
847 };
848
849 rows.sort_by(|a, b| {
850 for item in &op.items {
851 let ord = compare_sort_item(item, a, b, &eval_ctx);
852 if ord != Ordering::Equal {
853 return ord;
854 }
855 }
856 Ordering::Equal
857 });
858
859 Ok(rows)
860 }
861
862 fn exec_limit(&self, plan: &PhysicalPlan, op: &LimitExec) -> ExecResult<Vec<Row>> {
863 let mut rows = self.execute_node(plan, op.input)?;
864 let eval_ctx = EvalContext {
865 storage: self.ctx.storage,
866 params: &self.ctx.params,
867 };
868
869 let limit = op
870 .limit
871 .as_ref()
872 .and_then(|e| eval_expr(e, &Row::new(), &eval_ctx).as_i64())
873 .unwrap_or(rows.len() as i64)
874 .max(0) as usize;
875
876 let skip = op
877 .skip
878 .as_ref()
879 .and_then(|e| eval_expr(e, &Row::new(), &eval_ctx).as_i64())
880 .unwrap_or(0)
881 .max(0) as usize;
882
883 if skip >= rows.len() {
884 return Ok(Vec::new());
885 }
886
887 rows.drain(0..skip);
888 rows.truncate(limit);
889 Ok(rows)
890 }
891
892 fn exec_optional_match(
893 &self,
894 plan: &PhysicalPlan,
895 op: &OptionalMatchExec,
896 ) -> ExecResult<Vec<Row>> {
897 let input_rows = self.execute_node(plan, op.input)?;
898
899 let inner_rows = self.execute_node(plan, op.inner)?;
904
905 let mut out = Vec::new();
906
907 for input_row in input_rows {
908 let mut matched = false;
909
910 for inner_row in &inner_rows {
911 let compatible = input_row
913 .iter()
914 .all(|(var, val)| match inner_row.get(*var) {
915 Some(inner_val) => inner_val == val,
916 None => true,
917 });
918 if !compatible {
919 continue;
920 }
921
922 let mut merged = input_row.clone();
923 for (var, name, val) in inner_row.iter_named() {
924 if !merged.contains_key(*var) {
925 merged.insert_named(*var, name.into_owned(), val.clone());
926 }
927 }
928 out.push(merged);
929 matched = true;
930 }
931
932 if !matched {
933 let mut null_row = input_row;
934 for &var_id in &op.new_vars {
935 if !null_row.contains_key(var_id) {
936 null_row.insert(var_id, LoraValue::Null);
937 }
938 }
939 out.push(null_row);
940 }
941 }
942
943 Ok(out)
944 }
945
946 fn exec_path_build(&self, plan: &PhysicalPlan, op: &PathBuildExec) -> ExecResult<Vec<Row>> {
947 let input_rows = self.execute_node(plan, op.input)?;
948 let mut rows: Vec<Row> = input_rows
949 .into_iter()
950 .map(|mut row| {
951 let path = build_path_value(&row, &op.node_vars, &op.rel_vars, self.ctx.storage);
952 row.insert(op.output, path);
953 row
954 })
955 .collect();
956
957 if let Some(all) = op.shortest_path_all {
958 rows = filter_shortest_paths(rows, op.output, all);
959 }
960 Ok(rows)
961 }
962}