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::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            // An IN seek touches one index bucket per element; the
485            // sequential path is already cheap.
486            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                // Rayon workers do not inherit the caller's thread-local
669                // deadline; expression-level checks need it.
670                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, |n| hydrate_node_record(n))
973            .unwrap_or(LoraValue::Null)
974    }
975
976    fn hydrate_relationship(&self, id: u64) -> LoraValue {
977        self.ctx
978            .storage
979            .with_relationship(id, |r| hydrate_relationship_record(r))
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        // Fallback for inner plans the pull pipeline cannot seed: execute
1079        // the pattern once, uncorrelated, and join it against every input
1080        // row.
1081        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}