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