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