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