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