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