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