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