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