Skip to main content

lora_executor/executor/
immutable.rs

1//! Read-only buffered executor: lower a [`PhysicalPlan`] into a fully
2//! materialized `Vec<Row>` without touching the store.
3//!
4//! [`Executor`] mirrors the operator set of the streaming pipeline in
5//! `crate::pull` so write operators can fall back to it for subtrees
6//! that are not fully streamable. Aggregation and DISTINCT projection
7//! reuse the streaming `StreamableAggSpec` / `AggState` machinery
8//! (re-exported as `crate::pull::*` for that purpose) on the
9//! fold-only fast path; everything else materializes.
10
11use 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    /// Execute a compiled read-only query, using rayon only for the narrow
127    /// materialized-safe subset (scan/filter/projection) and only above a
128    /// measured threshold. Smaller plans and unsupported operators fall back
129    /// to the normal buffered executor.
130    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 any error residue that a previous query on this thread may have
153        // left in the thread-local eval-error slot.
154        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    /// Buffered execution of an arbitrary subplan. Public to the
177    /// crate so the pull pipeline can fall back to materialized
178    /// execution for operators that have no streaming source yet.
179    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            // An IN seek touches one index bucket per element; the
481            // sequential path is already cheap.
482            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        // Fallback for inner plans the pull pipeline cannot seed: execute
1044        // the pattern once, uncorrelated, and join it against every input
1045        // row.
1046        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}