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