Skip to main content

lora_executor/executor/
immutable.rs

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