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