1use crate::errors::{ExecResult, ExecutorError};
12use crate::eval::{clear_eval_error, EvalContext};
13#[cfg(all(feature = "parallel", not(target_arch = "wasm32")))]
14use crate::eval::{eval_expr, eval_expr_result, eval_truthy_result};
15use crate::value::{LoraValue, Row};
16use crate::{project_rows, ExecuteOptions, QueryResult};
17
18use lora_compiler::physical::*;
19use lora_compiler::CompiledQuery;
20use lora_store::GraphStorage;
21
22use std::collections::BTreeMap;
23use std::time::Instant;
24use tracing::{error, trace};
25
26use super::aggregate_rows;
27use super::helpers::{
28 build_path_value, check_deadline_at, dedup_rows, expand_rows, expand_var_len_rows,
29 filter_rows_checked, filter_shortest_paths, hydrate_node_record, hydrate_relationship_record,
30 limit_rows, node_by_label_scan_rows, node_by_property_scan_rows, node_scan_rows,
31 plan_may_need_hydration, project_rows_checked, unwind_rows,
32};
33#[cfg(all(feature = "parallel", not(target_arch = "wasm32")))]
34use super::helpers::{
35 dedup_rows_by_vars, indexed_node_property_candidates, label_group_candidates_prefiltered,
36 node_matches_label_groups, node_matches_property_filter, scan_node_ids_for_label_groups,
37};
38use super::sort_rows_with_top_k;
39use super::{merge_optional_rows, optional_match_rows};
40
41#[cfg(all(feature = "parallel", not(target_arch = "wasm32")))]
42const PARALLEL_ROW_THRESHOLD: usize = 20_000;
43
44pub struct ExecutionContext<'a, S: GraphStorage> {
45 pub storage: &'a S,
46 pub params: BTreeMap<String, LoraValue>,
47}
48
49pub struct Executor<'a, S: GraphStorage> {
50 ctx: ExecutionContext<'a, S>,
51 deadline: Option<Instant>,
52}
53
54impl<'a, S: GraphStorage> Executor<'a, S> {
55 pub fn new(ctx: ExecutionContext<'a, S>) -> Self {
56 Self {
57 ctx,
58 deadline: None,
59 }
60 }
61
62 pub fn with_deadline(ctx: ExecutionContext<'a, S>, deadline: Option<Instant>) -> Self {
63 Self { ctx, deadline }
64 }
65
66 #[inline]
67 fn check_deadline(&self) -> ExecResult<()> {
68 if let Some(deadline) = self.deadline {
69 check_deadline_at(deadline)
70 } else {
71 Ok(())
72 }
73 }
74}
75
76impl<'a, S: GraphStorage> Executor<'a, S> {
77 pub fn execute(
78 &self,
79 plan: &PhysicalPlan,
80 options: Option<ExecuteOptions>,
81 ) -> ExecResult<QueryResult> {
82 let rows = self.execute_rows(plan)?;
83 Ok(project_rows(rows, options.unwrap_or_default()))
84 }
85
86 pub fn execute_compiled(
87 &self,
88 compiled: &CompiledQuery,
89 options: Option<ExecuteOptions>,
90 ) -> ExecResult<QueryResult> {
91 let rows = self.execute_compiled_rows(compiled)?;
92 Ok(project_rows(rows, options.unwrap_or_default()))
93 }
94
95 pub fn execute_compiled_rows(&self, compiled: &CompiledQuery) -> ExecResult<Vec<Row>> {
96 self.check_deadline()?;
97 if compiled.unions.is_empty() {
98 return self.execute_rows(&compiled.physical);
99 }
100
101 clear_eval_error();
102
103 let mut all_rows = self.execute_rows(&compiled.physical)?;
104 let mut needs_dedup = false;
105
106 for branch in &compiled.unions {
107 self.check_deadline()?;
108 let branch_rows = self.execute_rows(&branch.physical)?;
109 all_rows.extend(branch_rows);
110
111 if !branch.all {
112 needs_dedup = true;
113 }
114 }
115
116 if needs_dedup {
117 all_rows = dedup_rows(all_rows);
118 }
119
120 Ok(all_rows)
121 }
122
123 pub fn execute_compiled_rows_parallel_safe(
128 &self,
129 compiled: &CompiledQuery,
130 ) -> ExecResult<Vec<Row>>
131 where
132 S: Sync,
133 {
134 #[cfg(all(feature = "parallel", not(target_arch = "wasm32")))]
135 {
136 self.check_deadline()?;
137 if compiled.unions.is_empty() && plan_is_parallel_safe(&compiled.physical) {
138 return self.execute_rows_parallel_safe(&compiled.physical);
139 }
140 }
141
142 self.execute_compiled_rows(compiled)
143 }
144
145 pub fn execute_rows(&self, plan: &PhysicalPlan) -> ExecResult<Vec<Row>> {
146 self.check_deadline()?;
147 clear_eval_error();
150
151 let rows = self.execute_node(plan, plan.root)?;
152 if !plan_may_need_hydration(plan) {
153 return Ok(rows);
154 }
155 Ok(rows
156 .into_iter()
157 .map(|row| self.hydrate_row(row))
158 .collect::<Vec<_>>())
159 }
160
161 fn hydrate_row(&self, row: Row) -> Row {
162 let mut out = Row::new();
163
164 for (var, name, value) in row.into_iter_named() {
165 out.insert_named(var, name, self.hydrate_value(value));
166 }
167
168 out
169 }
170
171 pub(crate) fn execute_subtree(
175 &self,
176 plan: &PhysicalPlan,
177 node_id: PhysicalNodeId,
178 ) -> ExecResult<Vec<Row>> {
179 self.execute_node(plan, node_id)
180 }
181
182 fn execute_node(&self, plan: &PhysicalPlan, node_id: PhysicalNodeId) -> ExecResult<Vec<Row>> {
183 self.check_deadline()?;
184 trace!("read-only execute_node start: node_id={node_id:?}");
185
186 let result = match &plan.nodes[node_id] {
187 PhysicalOp::Argument(op) => self.exec_argument(op),
188 PhysicalOp::NodeScan(op) => self.exec_node_scan(plan, op),
189 PhysicalOp::NodeByLabelScan(op) => self.exec_node_by_label_scan(plan, op),
190 PhysicalOp::NodeByPropertyScan(op) => self.exec_node_by_property_scan(plan, op),
191 PhysicalOp::NodeByPropertyRangeScan(op) => {
192 self.exec_node_by_property_range_scan(plan, op)
193 }
194 PhysicalOp::NodeByTextScan(op) => self.exec_node_by_text_scan(plan, op),
195 PhysicalOp::NodeByPointScan(op) => self.exec_node_by_point_scan(plan, op),
196 PhysicalOp::RelByPropertyRangeScan(op) => {
197 self.exec_rel_by_property_range_scan(plan, op)
198 }
199 PhysicalOp::RelByTextScan(op) => self.exec_rel_by_text_scan(plan, op),
200 PhysicalOp::RelByPointScan(op) => self.exec_rel_by_point_scan(plan, op),
201 PhysicalOp::Expand(op) => self.exec_expand(plan, op),
202 PhysicalOp::Filter(op) => self.exec_filter(plan, op),
203 PhysicalOp::Projection(op) => self.exec_projection(plan, op),
204 PhysicalOp::Unwind(op) => self.exec_unwind(plan, op),
205 PhysicalOp::HashAggregation(op) => self.exec_hash_aggregation(plan, op),
206 PhysicalOp::Sort(op) => self.exec_sort(plan, op),
207 PhysicalOp::Limit(op) => self.exec_limit(plan, op),
208 PhysicalOp::OptionalMatch(op) => self.exec_optional_match(plan, op),
209 PhysicalOp::CallSubquery(op) => self.exec_call_subquery(plan, op),
210 PhysicalOp::PathBuild(op) => self.exec_path_build(plan, op),
211 PhysicalOp::Create(_) => Err(ExecutorError::ReadOnlyCreate { node_id }),
212 PhysicalOp::Merge(_) => Err(ExecutorError::ReadOnlyMerge { node_id }),
213 PhysicalOp::Delete(_) => Err(ExecutorError::ReadOnlyDelete { node_id }),
214 PhysicalOp::Set(_) => Err(ExecutorError::ReadOnlySet { node_id }),
215 PhysicalOp::Remove(_) => Err(ExecutorError::ReadOnlyRemove { node_id }),
216 };
217
218 match &result {
219 Ok(rows) => trace!(
220 "read-only execute_node ok: node_id={node_id:?}, rows={}",
221 rows.len()
222 ),
223 Err(err) => error!("read-only execute_node failed: node_id={node_id:?}, error={err}"),
224 }
225
226 result
227 }
228
229 #[cfg(all(feature = "parallel", not(target_arch = "wasm32")))]
230 fn execute_rows_parallel_safe(&self, plan: &PhysicalPlan) -> ExecResult<Vec<Row>>
231 where
232 S: Sync,
233 {
234 self.check_deadline()?;
235 clear_eval_error();
236
237 let rows = self.execute_node_parallel_safe(plan, plan.root)?;
238 if !plan_may_need_hydration(plan) {
239 return Ok(rows);
240 }
241 if rows.len() < PARALLEL_ROW_THRESHOLD {
242 return Ok(rows
243 .into_iter()
244 .map(|row| self.hydrate_row(row))
245 .collect::<Vec<_>>());
246 }
247
248 use rayon::prelude::*;
249 Ok(rows
250 .into_par_iter()
251 .map(|row| self.hydrate_row(row))
252 .collect::<Vec<_>>())
253 }
254
255 #[cfg(all(feature = "parallel", not(target_arch = "wasm32")))]
256 fn execute_node_parallel_safe(
257 &self,
258 plan: &PhysicalPlan,
259 node_id: PhysicalNodeId,
260 ) -> ExecResult<Vec<Row>>
261 where
262 S: Sync,
263 {
264 self.check_deadline()?;
265 match &plan.nodes[node_id] {
266 PhysicalOp::Argument(op) => self.exec_argument(op),
267 PhysicalOp::NodeScan(op) => self.exec_node_scan_parallel_safe(plan, op),
268 PhysicalOp::NodeByLabelScan(op) => self.exec_node_by_label_scan_parallel_safe(plan, op),
269 PhysicalOp::NodeByPropertyScan(op) => {
270 self.exec_node_by_property_scan_parallel_safe(plan, op)
271 }
272 PhysicalOp::Filter(op) => self.exec_filter_parallel_safe(plan, op),
273 PhysicalOp::Projection(op) => self.exec_projection_parallel_safe(plan, op),
274 _ => Err(ExecutorError::RuntimeError(
275 "parallel-safe executor called with unsupported operator".into(),
276 )),
277 }
278 }
279
280 #[cfg(all(feature = "parallel", not(target_arch = "wasm32")))]
281 fn exec_node_scan_parallel_safe(
282 &self,
283 plan: &PhysicalPlan,
284 op: &NodeScanExec,
285 ) -> ExecResult<Vec<Row>>
286 where
287 S: Sync,
288 {
289 let base_rows = match op.input {
290 Some(input) => self.execute_node_parallel_safe(plan, input)?,
291 None => vec![Row::new()],
292 };
293 let node_ids = self.ctx.storage.all_node_ids();
294 if base_rows.len().saturating_mul(node_ids.len()) < PARALLEL_ROW_THRESHOLD {
295 return node_scan_rows(self.ctx.storage, base_rows, op, self.deadline);
296 }
297
298 use rayon::prelude::*;
299 if base_rows.len() == 1 {
300 let Some(row) = base_rows.into_iter().next() else {
301 return Err(ExecutorError::RuntimeError(
302 "parallel node scan expected one base row".into(),
303 ));
304 };
305 if let Some(existing_id) = super::helpers::bound_node_id_for_expand(&row, op.var)? {
306 return Ok(if self.ctx.storage.has_node(existing_id) {
307 vec![row]
308 } else {
309 Vec::new()
310 });
311 }
312
313 return node_ids
314 .into_par_iter()
315 .map(|id| {
316 if let Some(deadline) = self.deadline {
317 check_deadline_at(deadline)?;
318 }
319 let mut new_row = row.clone();
320 new_row.insert(op.var, LoraValue::Node(id));
321 Ok(new_row)
322 })
323 .collect();
324 }
325
326 let chunks: ExecResult<Vec<Vec<Row>>> = base_rows
327 .into_par_iter()
328 .map(|row| {
329 if let Some(deadline) = self.deadline {
330 check_deadline_at(deadline)?;
331 }
332 if let Some(existing_id) = super::helpers::bound_node_id_for_expand(&row, op.var)? {
333 return Ok(if self.ctx.storage.has_node(existing_id) {
334 vec![row]
335 } else {
336 Vec::new()
337 });
338 }
339
340 let mut out = Vec::with_capacity(node_ids.len());
341 for &id in &node_ids {
342 if let Some(deadline) = self.deadline {
343 check_deadline_at(deadline)?;
344 }
345 let mut new_row = row.clone();
346 new_row.insert(op.var, LoraValue::Node(id));
347 out.push(new_row);
348 }
349 Ok(out)
350 })
351 .collect();
352 Ok(chunks?.into_iter().flatten().collect())
353 }
354
355 #[cfg(all(feature = "parallel", not(target_arch = "wasm32")))]
356 fn exec_node_by_label_scan_parallel_safe(
357 &self,
358 plan: &PhysicalPlan,
359 op: &NodeByLabelScanExec,
360 ) -> ExecResult<Vec<Row>>
361 where
362 S: Sync,
363 {
364 let base_rows = match op.input {
365 Some(input) => self.execute_node_parallel_safe(plan, input)?,
366 None => vec![Row::new()],
367 };
368 let candidate_ids = scan_node_ids_for_label_groups(self.ctx.storage, &op.labels);
369 if base_rows.len().saturating_mul(candidate_ids.len()) < PARALLEL_ROW_THRESHOLD {
370 return node_by_label_scan_rows(self.ctx.storage, base_rows, op, self.deadline);
371 }
372
373 let candidates_prefiltered = label_group_candidates_prefiltered(&op.labels);
374 use rayon::prelude::*;
375 if base_rows.len() == 1 {
376 let Some(row) = base_rows.into_iter().next() else {
377 return Err(ExecutorError::RuntimeError(
378 "parallel label scan expected one base row".into(),
379 ));
380 };
381 if let Some(existing_id) = super::helpers::bound_node_id_for_expand(&row, op.var)? {
382 let labels_ok = self
383 .ctx
384 .storage
385 .with_node(existing_id, |n| {
386 node_matches_label_groups(&n.labels, &op.labels)
387 })
388 .unwrap_or(false);
389 return Ok(if labels_ok { vec![row] } else { Vec::new() });
390 }
391
392 return candidate_ids
393 .into_par_iter()
394 .filter_map(|id| {
395 if let Some(deadline) = self.deadline {
396 if let Err(err) = check_deadline_at(deadline) {
397 return Some(Err(err));
398 }
399 }
400 if !candidates_prefiltered {
401 let labels_ok = self
402 .ctx
403 .storage
404 .with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
405 .unwrap_or(false);
406 if !labels_ok {
407 return None;
408 }
409 }
410 let mut new_row = row.clone();
411 new_row.insert(op.var, LoraValue::Node(id));
412 Some(Ok(new_row))
413 })
414 .collect();
415 }
416
417 let chunks: ExecResult<Vec<Vec<Row>>> = base_rows
418 .into_par_iter()
419 .map(|row| {
420 if let Some(deadline) = self.deadline {
421 check_deadline_at(deadline)?;
422 }
423 if let Some(existing_id) = super::helpers::bound_node_id_for_expand(&row, op.var)? {
424 let labels_ok = self
425 .ctx
426 .storage
427 .with_node(existing_id, |n| {
428 node_matches_label_groups(&n.labels, &op.labels)
429 })
430 .unwrap_or(false);
431 return Ok(if labels_ok { vec![row] } else { Vec::new() });
432 }
433
434 let mut out = Vec::with_capacity(candidate_ids.len());
435 for &id in &candidate_ids {
436 if let Some(deadline) = self.deadline {
437 check_deadline_at(deadline)?;
438 }
439 if !candidates_prefiltered {
440 let labels_ok = self
441 .ctx
442 .storage
443 .with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
444 .unwrap_or(false);
445 if !labels_ok {
446 continue;
447 }
448 }
449 let mut new_row = row.clone();
450 new_row.insert(op.var, LoraValue::Node(id));
451 out.push(new_row);
452 }
453 Ok(out)
454 })
455 .collect();
456 Ok(chunks?.into_iter().flatten().collect())
457 }
458
459 #[cfg(all(feature = "parallel", not(target_arch = "wasm32")))]
460 fn exec_node_by_property_scan_parallel_safe(
461 &self,
462 plan: &PhysicalPlan,
463 op: &NodeByPropertyScanExec,
464 ) -> ExecResult<Vec<Row>>
465 where
466 S: Sync,
467 {
468 let base_rows = match op.input {
469 Some(input) => self.execute_node_parallel_safe(plan, input)?,
470 None => vec![Row::new()],
471 };
472 let eval_ctx = EvalContext {
473 storage: self.ctx.storage,
474 params: &self.ctx.params,
475 };
476 use rayon::prelude::*;
477
478 if base_rows.len() == 1 {
479 let Some(row) = base_rows.into_iter().next() else {
480 return Err(ExecutorError::RuntimeError(
481 "parallel property scan expected one base row".into(),
482 ));
483 };
484 if let Some(deadline) = self.deadline {
485 check_deadline_at(deadline)?;
486 }
487 let expected = eval_expr(&op.value, &row, &eval_ctx);
488 if let Some(existing_id) = super::helpers::bound_node_id_for_expand(&row, op.var)? {
489 return Ok(
490 if node_matches_property_filter(
491 self.ctx.storage,
492 existing_id,
493 &op.labels,
494 &op.key,
495 &expected,
496 ) {
497 vec![row]
498 } else {
499 Vec::new()
500 },
501 );
502 }
503
504 let candidates =
505 indexed_node_property_candidates(self.ctx.storage, &op.labels, &op.key, &expected);
506 if candidates.ids.len() < PARALLEL_ROW_THRESHOLD {
507 let mut out = Vec::with_capacity(candidates.ids.len());
508 for id in candidates.ids {
509 if !candidates.prefiltered
510 && !node_matches_property_filter(
511 self.ctx.storage,
512 id,
513 &op.labels,
514 &op.key,
515 &expected,
516 )
517 {
518 continue;
519 }
520 let mut new_row = row.clone();
521 new_row.insert(op.var, LoraValue::Node(id));
522 out.push(new_row);
523 }
524 return Ok(out);
525 }
526
527 return candidates
528 .ids
529 .into_par_iter()
530 .filter_map(|id| {
531 if let Some(deadline) = self.deadline {
532 if let Err(err) = check_deadline_at(deadline) {
533 return Some(Err(err));
534 }
535 }
536 if !candidates.prefiltered
537 && !node_matches_property_filter(
538 self.ctx.storage,
539 id,
540 &op.labels,
541 &op.key,
542 &expected,
543 )
544 {
545 return None;
546 }
547 let mut new_row = row.clone();
548 new_row.insert(op.var, LoraValue::Node(id));
549 Some(Ok(new_row))
550 })
551 .collect();
552 }
553
554 if base_rows.len() < PARALLEL_ROW_THRESHOLD {
555 return node_by_property_scan_rows(
556 self.ctx.storage,
557 &self.ctx.params,
558 base_rows,
559 op,
560 self.deadline,
561 );
562 }
563
564 let chunks: ExecResult<Vec<Vec<Row>>> = base_rows
565 .into_par_iter()
566 .map(|row| {
567 if let Some(deadline) = self.deadline {
568 check_deadline_at(deadline)?;
569 }
570 let expected = eval_expr(&op.value, &row, &eval_ctx);
571
572 if let Some(existing_id) = super::helpers::bound_node_id_for_expand(&row, op.var)? {
573 return Ok(
574 if node_matches_property_filter(
575 self.ctx.storage,
576 existing_id,
577 &op.labels,
578 &op.key,
579 &expected,
580 ) {
581 vec![row]
582 } else {
583 Vec::new()
584 },
585 );
586 }
587
588 let candidates = indexed_node_property_candidates(
589 self.ctx.storage,
590 &op.labels,
591 &op.key,
592 &expected,
593 );
594 let mut out = Vec::with_capacity(candidates.ids.len());
595 for id in candidates.ids {
596 if let Some(deadline) = self.deadline {
597 check_deadline_at(deadline)?;
598 }
599 if !candidates.prefiltered
600 && !node_matches_property_filter(
601 self.ctx.storage,
602 id,
603 &op.labels,
604 &op.key,
605 &expected,
606 )
607 {
608 continue;
609 }
610 let mut new_row = row.clone();
611 new_row.insert(op.var, LoraValue::Node(id));
612 out.push(new_row);
613 }
614 Ok(out)
615 })
616 .collect();
617 Ok(chunks?.into_iter().flatten().collect())
618 }
619
620 #[cfg(all(feature = "parallel", not(target_arch = "wasm32")))]
621 fn exec_filter_parallel_safe(
622 &self,
623 plan: &PhysicalPlan,
624 op: &FilterExec,
625 ) -> ExecResult<Vec<Row>>
626 where
627 S: Sync,
628 {
629 let input_rows = self.execute_node_parallel_safe(plan, op.input)?;
630 if input_rows.len() < PARALLEL_ROW_THRESHOLD {
631 let eval_ctx = EvalContext {
632 storage: self.ctx.storage,
633 params: &self.ctx.params,
634 };
635 return filter_rows_checked(input_rows, &op.predicate, &eval_ctx);
636 }
637
638 let eval_ctx = EvalContext {
639 storage: self.ctx.storage,
640 params: &self.ctx.params,
641 };
642 use rayon::prelude::*;
643 let filtered: ExecResult<Vec<Option<Row>>> = input_rows
644 .into_par_iter()
645 .map(|row| {
646 if let Some(deadline) = self.deadline {
647 check_deadline_at(deadline)?;
648 }
649 let keep = eval_truthy_result(&op.predicate, &row, &eval_ctx)
650 .map_err(ExecutorError::RuntimeError)?;
651 Ok(if keep { Some(row) } else { None })
652 })
653 .collect();
654 Ok(filtered?.into_iter().flatten().collect())
655 }
656
657 #[cfg(all(feature = "parallel", not(target_arch = "wasm32")))]
658 fn exec_projection_parallel_safe(
659 &self,
660 plan: &PhysicalPlan,
661 op: &ProjectionExec,
662 ) -> ExecResult<Vec<Row>>
663 where
664 S: Sync,
665 {
666 let input_rows = self.execute_node_parallel_safe(plan, op.input)?;
667 if input_rows.len() < PARALLEL_ROW_THRESHOLD {
668 let eval_ctx = EvalContext {
669 storage: self.ctx.storage,
670 params: &self.ctx.params,
671 };
672 return project_rows_checked(input_rows, op, &eval_ctx);
673 }
674
675 let eval_ctx = EvalContext {
676 storage: self.ctx.storage,
677 params: &self.ctx.params,
678 };
679 use rayon::prelude::*;
680 let projected: ExecResult<Vec<Row>> = input_rows
681 .into_par_iter()
682 .map(|row| {
683 if let Some(deadline) = self.deadline {
684 check_deadline_at(deadline)?;
685 }
686 if op.include_existing {
687 let mut projected = row;
688 for item in &op.items {
689 let value = eval_expr_result(&item.expr, &projected, &eval_ctx)
690 .map_err(ExecutorError::RuntimeError)?;
691 projected.insert_named(item.output, item.name.clone(), value);
692 }
693 Ok(projected)
694 } else {
695 let mut projected = Row::new();
696 for item in &op.items {
697 let value = eval_expr_result(&item.expr, &row, &eval_ctx)
698 .map_err(ExecutorError::RuntimeError)?;
699 projected.insert_named(item.output, item.name.clone(), value);
700 }
701 Ok(projected)
702 }
703 })
704 .collect();
705 let rows = projected?;
706 Ok(if op.distinct {
707 dedup_rows_by_vars(rows)
708 } else {
709 rows
710 })
711 }
712
713 fn exec_argument(&self, _op: &ArgumentExec) -> ExecResult<Vec<Row>> {
714 Ok(vec![Row::new()])
715 }
716
717 fn exec_node_scan(&self, plan: &PhysicalPlan, op: &NodeScanExec) -> ExecResult<Vec<Row>> {
718 let base_rows = match op.input {
719 Some(input) => self.execute_node(plan, input)?,
720 None => vec![Row::new()],
721 };
722
723 node_scan_rows(self.ctx.storage, base_rows, op, self.deadline)
724 }
725
726 fn exec_node_by_label_scan(
727 &self,
728 plan: &PhysicalPlan,
729 op: &NodeByLabelScanExec,
730 ) -> ExecResult<Vec<Row>> {
731 let base_rows = match op.input {
732 Some(input) => self.execute_node(plan, input)?,
733 None => vec![Row::new()],
734 };
735
736 node_by_label_scan_rows(self.ctx.storage, base_rows, op, self.deadline)
737 }
738
739 fn exec_node_by_property_scan(
740 &self,
741 plan: &PhysicalPlan,
742 op: &NodeByPropertyScanExec,
743 ) -> ExecResult<Vec<Row>> {
744 let base_rows = match op.input {
745 Some(input) => self.execute_node(plan, input)?,
746 None => vec![Row::new()],
747 };
748
749 node_by_property_scan_rows(
750 self.ctx.storage,
751 &self.ctx.params,
752 base_rows,
753 op,
754 self.deadline,
755 )
756 }
757
758 fn exec_node_by_property_range_scan(
759 &self,
760 plan: &PhysicalPlan,
761 op: &lora_compiler::NodeByPropertyRangeScanExec,
762 ) -> ExecResult<Vec<Row>> {
763 let base_rows = match op.input {
764 Some(input) => self.execute_node(plan, input)?,
765 None => vec![Row::new()],
766 };
767 super::helpers::node_by_property_range_scan_rows(
768 self.ctx.storage,
769 &self.ctx.params,
770 base_rows,
771 op,
772 self.deadline,
773 )
774 }
775
776 fn exec_node_by_text_scan(
777 &self,
778 plan: &PhysicalPlan,
779 op: &lora_compiler::NodeByTextScanExec,
780 ) -> ExecResult<Vec<Row>> {
781 let base_rows = match op.input {
782 Some(input) => self.execute_node(plan, input)?,
783 None => vec![Row::new()],
784 };
785 super::helpers::node_by_text_scan_rows(
786 self.ctx.storage,
787 &self.ctx.params,
788 base_rows,
789 op,
790 self.deadline,
791 )
792 }
793
794 fn exec_node_by_point_scan(
795 &self,
796 plan: &PhysicalPlan,
797 op: &lora_compiler::NodeByPointScanExec,
798 ) -> ExecResult<Vec<Row>> {
799 let base_rows = match op.input {
800 Some(input) => self.execute_node(plan, input)?,
801 None => vec![Row::new()],
802 };
803 super::helpers::node_by_point_scan_rows(
804 self.ctx.storage,
805 &self.ctx.params,
806 base_rows,
807 op,
808 self.deadline,
809 )
810 }
811
812 fn exec_rel_by_property_range_scan(
813 &self,
814 plan: &PhysicalPlan,
815 op: &lora_compiler::RelByPropertyRangeScanExec,
816 ) -> ExecResult<Vec<Row>> {
817 let base_rows = match op.input {
818 Some(input) => self.execute_node(plan, input)?,
819 None => vec![Row::new()],
820 };
821 super::helpers::rel_by_property_range_scan_rows(
822 self.ctx.storage,
823 &self.ctx.params,
824 base_rows,
825 op,
826 self.deadline,
827 )
828 }
829
830 fn exec_rel_by_text_scan(
831 &self,
832 plan: &PhysicalPlan,
833 op: &lora_compiler::RelByTextScanExec,
834 ) -> ExecResult<Vec<Row>> {
835 let base_rows = match op.input {
836 Some(input) => self.execute_node(plan, input)?,
837 None => vec![Row::new()],
838 };
839 super::helpers::rel_by_text_scan_rows(
840 self.ctx.storage,
841 &self.ctx.params,
842 base_rows,
843 op,
844 self.deadline,
845 )
846 }
847
848 fn exec_rel_by_point_scan(
849 &self,
850 plan: &PhysicalPlan,
851 op: &lora_compiler::RelByPointScanExec,
852 ) -> ExecResult<Vec<Row>> {
853 let base_rows = match op.input {
854 Some(input) => self.execute_node(plan, input)?,
855 None => vec![Row::new()],
856 };
857 super::helpers::rel_by_point_scan_rows(
858 self.ctx.storage,
859 &self.ctx.params,
860 base_rows,
861 op,
862 self.deadline,
863 )
864 }
865
866 fn exec_expand(&self, plan: &PhysicalPlan, op: &ExpandExec) -> ExecResult<Vec<Row>> {
867 let input_rows = self.execute_node(plan, op.input)?;
868 if let Some(range) = &op.range {
869 expand_var_len_rows(self.ctx.storage, input_rows, op, range)
870 } else {
871 expand_rows(self.ctx.storage, &self.ctx.params, input_rows, op)
872 }
873 }
874
875 fn exec_filter(&self, plan: &PhysicalPlan, op: &FilterExec) -> ExecResult<Vec<Row>> {
876 let input_rows = self.execute_node(plan, op.input)?;
877 let eval_ctx = EvalContext {
878 storage: self.ctx.storage,
879 params: &self.ctx.params,
880 };
881
882 filter_rows_checked(input_rows, &op.predicate, &eval_ctx)
883 }
884
885 fn exec_projection(&self, plan: &PhysicalPlan, op: &ProjectionExec) -> ExecResult<Vec<Row>> {
886 let input_rows = self.execute_node(plan, op.input)?;
887 let eval_ctx = EvalContext {
888 storage: self.ctx.storage,
889 params: &self.ctx.params,
890 };
891
892 project_rows_checked(input_rows, op, &eval_ctx)
893 }
894
895 fn hydrate_value(&self, value: LoraValue) -> LoraValue {
896 match value {
897 LoraValue::Node(id) => self.hydrate_node(id),
898 LoraValue::Relationship(id) => self.hydrate_relationship(id),
899 LoraValue::List(values) => {
900 LoraValue::List(values.into_iter().map(|v| self.hydrate_value(v)).collect())
901 }
902 LoraValue::Map(map) => LoraValue::Map(
903 map.into_iter()
904 .map(|(k, v)| (k, self.hydrate_value(v)))
905 .collect(),
906 ),
907 other => other,
908 }
909 }
910
911 fn hydrate_node(&self, id: u64) -> LoraValue {
912 self.ctx
913 .storage
914 .with_node(id, hydrate_node_record)
915 .unwrap_or(LoraValue::Null)
916 }
917
918 fn hydrate_relationship(&self, id: u64) -> LoraValue {
919 self.ctx
920 .storage
921 .with_relationship(id, hydrate_relationship_record)
922 .unwrap_or(LoraValue::Null)
923 }
924
925 fn exec_unwind(&self, plan: &PhysicalPlan, op: &UnwindExec) -> ExecResult<Vec<Row>> {
926 let input_rows = self.execute_node(plan, op.input)?;
927 let eval_ctx = EvalContext {
928 storage: self.ctx.storage,
929 params: &self.ctx.params,
930 };
931
932 Ok(unwind_rows(input_rows, op, &eval_ctx))
933 }
934
935 fn exec_hash_aggregation(
936 &self,
937 plan: &PhysicalPlan,
938 op: &HashAggregationExec,
939 ) -> ExecResult<Vec<Row>> {
940 if let Some(rows) =
941 super::helpers::count_all_scan_aggregation_rows(self.ctx.storage, plan, op)
942 {
943 return Ok(rows);
944 }
945
946 let input_rows = self.execute_node(plan, op.input)?;
947 let eval_ctx = EvalContext {
948 storage: self.ctx.storage,
949 params: &self.ctx.params,
950 };
951
952 aggregate_rows(
953 input_rows,
954 &op.group_by,
955 &op.aggregates,
956 &eval_ctx,
957 |value| self.hydrate_value(value),
958 )
959 }
960
961 fn exec_sort(&self, plan: &PhysicalPlan, op: &SortExec) -> ExecResult<Vec<Row>> {
962 let mut rows = self.execute_node(plan, op.input)?;
963 let eval_ctx = EvalContext {
964 storage: self.ctx.storage,
965 params: &self.ctx.params,
966 };
967
968 sort_rows_with_top_k(&mut rows, &op.items, &eval_ctx, op.top_k);
969
970 Ok(rows)
971 }
972
973 fn exec_limit(&self, plan: &PhysicalPlan, op: &LimitExec) -> ExecResult<Vec<Row>> {
974 let rows = self.execute_node(plan, op.input)?;
975 let eval_ctx = EvalContext {
976 storage: self.ctx.storage,
977 params: &self.ctx.params,
978 };
979
980 Ok(limit_rows(rows, op, &eval_ctx))
981 }
982
983 fn exec_call_subquery(
984 &self,
985 plan: &PhysicalPlan,
986 op: &CallSubqueryExec,
987 ) -> ExecResult<Vec<Row>> {
988 let input_rows = self.execute_node(plan, op.input)?;
989 let mut out = Vec::with_capacity(input_rows.len());
990 let params = std::sync::Arc::new(self.ctx.params.clone());
991 for outer_row in input_rows {
992 let mut inner_source = crate::pull::build_streaming_seeded(
993 plan,
994 op.inner,
995 self.ctx.storage,
996 params.clone(),
997 outer_row.clone(),
998 )?;
999 let inner_rows = crate::pull::drain(inner_source.as_mut())?;
1000 for inner_row in inner_rows {
1001 out.push(merge_optional_rows(&outer_row, &inner_row));
1002 }
1003 }
1004 Ok(out)
1005 }
1006
1007 fn exec_optional_match(
1008 &self,
1009 plan: &PhysicalPlan,
1010 op: &OptionalMatchExec,
1011 ) -> ExecResult<Vec<Row>> {
1012 let input_rows = self.execute_node(plan, op.input)?;
1013
1014 let inner_rows = self.execute_node(plan, op.inner)?;
1019
1020 Ok(optional_match_rows(input_rows, &inner_rows, &op.new_vars))
1021 }
1022
1023 fn exec_path_build(&self, plan: &PhysicalPlan, op: &PathBuildExec) -> ExecResult<Vec<Row>> {
1024 let input_rows = self.execute_node(plan, op.input)?;
1025 let mut rows: Vec<Row> = input_rows
1026 .into_iter()
1027 .map(|mut row| {
1028 let path = build_path_value(&row, &op.node_vars, &op.rel_vars, self.ctx.storage);
1029 row.insert(op.output, path);
1030 row
1031 })
1032 .collect();
1033
1034 if let Some(all) = op.shortest_path_all {
1035 rows = filter_shortest_paths(rows, op.output, all);
1036 }
1037 Ok(rows)
1038 }
1039}
1040
1041#[cfg(all(feature = "parallel", not(target_arch = "wasm32")))]
1042fn plan_is_parallel_safe(plan: &PhysicalPlan) -> bool {
1043 subtree_is_parallel_safe(plan, plan.root)
1044}
1045
1046#[cfg(all(feature = "parallel", not(target_arch = "wasm32")))]
1047fn subtree_is_parallel_safe(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
1048 match &plan.nodes[node_id] {
1049 PhysicalOp::Argument(_) => true,
1050 PhysicalOp::NodeScan(op) => op
1051 .input
1052 .map(|input| subtree_is_parallel_safe(plan, input))
1053 .unwrap_or(true),
1054 PhysicalOp::NodeByLabelScan(op) => op
1055 .input
1056 .map(|input| subtree_is_parallel_safe(plan, input))
1057 .unwrap_or(true),
1058 PhysicalOp::NodeByPropertyScan(op) => op
1059 .input
1060 .map(|input| subtree_is_parallel_safe(plan, input))
1061 .unwrap_or(true),
1062 PhysicalOp::Filter(op) => subtree_is_parallel_safe(plan, op.input),
1063 PhysicalOp::Projection(op) => subtree_is_parallel_safe(plan, op.input),
1064 PhysicalOp::NodeByPropertyRangeScan(_)
1065 | PhysicalOp::NodeByTextScan(_)
1066 | PhysicalOp::NodeByPointScan(_)
1067 | PhysicalOp::RelByPropertyRangeScan(_)
1068 | PhysicalOp::RelByTextScan(_)
1069 | PhysicalOp::RelByPointScan(_) => false,
1070 PhysicalOp::Expand(_)
1071 | PhysicalOp::Unwind(_)
1072 | PhysicalOp::HashAggregation(_)
1073 | PhysicalOp::Sort(_)
1074 | PhysicalOp::Limit(_)
1075 | PhysicalOp::Create(_)
1076 | PhysicalOp::Merge(_)
1077 | PhysicalOp::Delete(_)
1078 | PhysicalOp::Set(_)
1079 | PhysicalOp::Remove(_)
1080 | PhysicalOp::OptionalMatch(_)
1081 | PhysicalOp::PathBuild(_)
1082 | PhysicalOp::CallSubquery(_) => false,
1083 }
1084}