Skip to main content

lora_executor/executor/
helpers.rs

1//! Cross-cutting executor helpers used by both the read-only and the
2//! mutable executor (and, via `pub(crate)` re-exports through
3//! `super::mod`, by the streaming pull pipeline in `crate::pull`).
4//!
5//! Roughly four groups:
6//!
7//! 1. Row-set primitives: `dedup_rows` / `dedup_rows_by_vars` for
8//!    UNION / DISTINCT, [`compute_aggregate_expr`] for the buffered
9//!    aggregation path, [`compare_sort_item`] for the buffered Sort
10//!    operator. The streaming pipeline in `crate::pull` calls
11//!    [`compute_aggregate_expr`] when the streamable-fold fast-path
12//!    classifier rejects a projection.
13//! 2. Label / property scans: [`scan_node_ids_for_label_groups`],
14//!    [`indexed_node_property_candidates`],
15//!    [`node_matches_label_groups`], [`node_matches_property_filter`],
16//!    [`label_group_candidates_prefiltered`]. Both NodeByLabelScan and
17//!    NodeByPropertyScan share these helpers across the buffered and
18//!    streaming pipelines.
19//! 3. Path construction: [`build_path_value`] for `PathBuild`,
20//!    [`variable_length_expand`] (the buffered BFS path used as the
21//!    fallback under non-streaming variable-length expansions) and
22//!    [`filter_shortest_paths`] for SHORTEST PATH.
23//! 4. Value classification: [`value_matches_property_value`] (used
24//!    by every property prefilter), [`hydrate_node_record`] /
25//!    [`hydrate_relationship_record`] (single-record hydration), and
26//!    the [`GroupValueKey`] dedup / group key.
27//!
28//! Also hosts the small `eval_properties_expr`, `eval_aggregate_arg_values`,
29//! `eval_first_or_null`, `dedup_values`, `as_f64_lossy`,
30//! `compare_values_for_sort`, `compare_values_total`,
31//! `single_label_hint`, `property_lookup_values`, `type_rank`,
32//! `flatten_label_groups` private helpers, and the
33//! `MAX_VAR_LEN_HOPS` cap on unbounded variable-length expansion.
34
35use std::cmp::Ordering;
36use std::collections::{BTreeMap, BTreeSet};
37
38use web_time::Instant;
39
40use lora_analyzer::symbols::VarId;
41use lora_analyzer::{AggregateFunction, FunctionId, ResolvedExpr, ResolvedMapSelector};
42use lora_ast::{Direction, RangeLiteral, SortDirection};
43use lora_compiler::physical::{
44    ExpandExec, HashAggregationExec, LimitExec, NodeByLabelScanExec, NodeByPropertyScanExec,
45    NodeScanExec, PhysicalNodeId, PhysicalOp, PhysicalPlan, ProjectionExec, UnwindExec,
46};
47use lora_store::{GraphStorage, NodeId, Properties, PropertyValue, RelationshipId};
48
49use crate::errors::{value_kind, ExecResult, ExecutorError};
50use crate::eval::{eval_expr, eval_expr_result, eval_truthy_result, EvalContext};
51use crate::value::{lora_value_to_property, LoraPath, LoraValue, Row};
52
53/// Deadline guard. Returns `QueryTimeout` once the deadline has
54/// elapsed; both executors call this every operator-level recursion
55/// step and from inside per-row inner loops.
56#[inline]
57/// Pre-size for a scan that emits `rows x candidates` rows. Capped: a
58/// cartesian product (`MATCH (a:A), (b:B)`) would otherwise reserve its
59/// full size up front, hundreds of gigabytes at 20k x 20k, before the
60/// first deadline check or the first row. Past the cap the Vec grows on
61/// demand.
62fn scan_output_capacity(rows: usize, candidates: usize) -> usize {
63    const MAX_PRESIZE_ROWS: usize = 1 << 16;
64    rows.saturating_mul(candidates).min(MAX_PRESIZE_ROWS)
65}
66
67pub(super) fn check_deadline_at(deadline: Instant) -> ExecResult<()> {
68    // Also true when the query was cancelled; see `crate::cancel`.
69    if crate::cancel::deadline_reached(deadline) {
70        Err(ExecutorError::QueryTimeout)
71    } else {
72        Ok(())
73    }
74}
75
76pub(super) fn filter_rows_checked<S: GraphStorage>(
77    input_rows: Vec<Row>,
78    predicate: &ResolvedExpr,
79    eval_ctx: &EvalContext<'_, S>,
80) -> ExecResult<Vec<Row>> {
81    let mut out = Vec::with_capacity(input_rows.len());
82    for row in input_rows {
83        if eval_truthy_result(predicate, &row, eval_ctx).map_err(ExecutorError::RuntimeError)? {
84            out.push(row);
85        }
86    }
87    Ok(out)
88}
89
90pub(super) fn project_rows_checked<S: GraphStorage>(
91    input_rows: Vec<Row>,
92    op: &ProjectionExec,
93    eval_ctx: &EvalContext<'_, S>,
94) -> ExecResult<Vec<Row>> {
95    let mut out = Vec::with_capacity(input_rows.len());
96
97    for row in input_rows {
98        if op.include_existing {
99            let mut projected = row;
100            for item in &op.items {
101                let value = eval_expr_result(&item.expr, &projected, eval_ctx)
102                    .map_err(ExecutorError::RuntimeError)?;
103                projected.insert_named(item.output, item.name.clone(), value);
104            }
105            out.push(projected);
106        } else {
107            let mut projected = Row::new();
108            for item in &op.items {
109                let value = eval_expr_result(&item.expr, &row, eval_ctx)
110                    .map_err(ExecutorError::RuntimeError)?;
111                projected.insert_named(item.output, item.name.clone(), value);
112            }
113            out.push(projected);
114        }
115    }
116
117    Ok(if op.distinct {
118        dedup_rows_by_vars(out)
119    } else {
120        out
121    })
122}
123
124pub(super) fn unwind_rows<S: GraphStorage>(
125    input_rows: Vec<Row>,
126    op: &UnwindExec,
127    eval_ctx: &EvalContext<'_, S>,
128) -> ExecResult<Vec<Row>> {
129    let mut out = Vec::with_capacity(input_rows.len());
130
131    for row in input_rows {
132        // `eval_expr_result` so a failing list expression (a bad function
133        // argument, a missing index) is an error, not an empty UNWIND.
134        match eval_expr_result(&op.expr, &row, eval_ctx).map_err(ExecutorError::RuntimeError)? {
135            LoraValue::List(values) => {
136                let mut values = values.into_iter();
137                let last = values.next_back();
138                for value in values {
139                    let mut new_row = row.clone();
140                    new_row.insert(op.alias, value);
141                    out.push(new_row);
142                }
143                // The input row moves into the last element's output.
144                if let Some(value) = last {
145                    let mut new_row = row;
146                    new_row.insert(op.alias, value);
147                    out.push(new_row);
148                }
149            }
150            LoraValue::Null => {}
151            other => {
152                let mut new_row = row;
153                new_row.insert(op.alias, other);
154                out.push(new_row);
155            }
156        }
157    }
158
159    Ok(out)
160}
161
162pub(super) fn limit_rows<S: GraphStorage>(
163    mut rows: Vec<Row>,
164    op: &LimitExec,
165    eval_ctx: &EvalContext<'_, S>,
166) -> Vec<Row> {
167    let limit = op
168        .limit
169        .as_ref()
170        .and_then(|e| eval_expr(e, &Row::new(), eval_ctx).as_i64())
171        .unwrap_or(rows.len() as i64)
172        .max(0) as usize;
173
174    let skip = op
175        .skip
176        .as_ref()
177        .and_then(|e| eval_expr(e, &Row::new(), eval_ctx).as_i64())
178        .unwrap_or(0)
179        .max(0) as usize;
180
181    if skip >= rows.len() {
182        return Vec::new();
183    }
184
185    rows.drain(0..skip);
186    rows.truncate(limit);
187    rows
188}
189
190#[inline]
191pub(crate) fn bound_node_id_for_expand(row: &Row, var: VarId) -> ExecResult<Option<NodeId>> {
192    match row.get(var) {
193        Some(LoraValue::Node(id)) => Ok(Some(*id)),
194        Some(other) => Err(ExecutorError::ExpectedNodeForExpand {
195            var: format!("{var:?}"),
196            found: value_kind(other),
197        }),
198        None => Ok(None),
199    }
200}
201
202#[inline]
203pub(crate) fn bound_relationship_id_for_expand(
204    row: &Row,
205    var: VarId,
206) -> ExecResult<Option<RelationshipId>> {
207    match row.get(var) {
208        Some(LoraValue::Relationship(id)) => Ok(Some(*id)),
209        Some(other) => Err(ExecutorError::ExpectedRelationshipForExpand {
210            var: format!("{var:?}"),
211            found: value_kind(other),
212        }),
213        None => Ok(None),
214    }
215}
216
217pub(super) fn node_scan_rows<S: GraphStorage>(
218    storage: &S,
219    base_rows: Vec<Row>,
220    op: &NodeScanExec,
221    deadline: Option<Instant>,
222) -> ExecResult<Vec<Row>> {
223    let node_ids = storage.all_node_ids();
224    let mut out = Vec::with_capacity(scan_output_capacity(base_rows.len(), node_ids.len()));
225
226    if deadline.is_none() {
227        for row in base_rows {
228            if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
229                if storage.has_node(existing_id) {
230                    out.push(row);
231                }
232                continue;
233            }
234
235            for &id in &node_ids {
236                let mut new_row = row.clone();
237                new_row.insert(op.var, LoraValue::Node(id));
238                out.push(new_row);
239            }
240        }
241        return Ok(out);
242    }
243
244    for row in base_rows {
245        check_optional_deadline(deadline)?;
246        if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
247            if storage.has_node(existing_id) {
248                out.push(row);
249            }
250            continue;
251        }
252
253        for &id in &node_ids {
254            check_optional_deadline(deadline)?;
255            let mut new_row = row.clone();
256            new_row.insert(op.var, LoraValue::Node(id));
257            out.push(new_row);
258        }
259    }
260
261    Ok(out)
262}
263
264pub(super) fn node_by_label_scan_rows<S: GraphStorage>(
265    storage: &S,
266    base_rows: Vec<Row>,
267    op: &NodeByLabelScanExec,
268    deadline: Option<Instant>,
269) -> ExecResult<Vec<Row>> {
270    let candidate_ids = scan_node_ids_for_label_groups(storage, &op.labels);
271    let candidates_prefiltered = label_group_candidates_prefiltered(&op.labels);
272    let mut out = Vec::with_capacity(scan_output_capacity(base_rows.len(), candidate_ids.len()));
273
274    if deadline.is_none() {
275        for row in base_rows {
276            if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
277                let labels_ok = storage
278                    .with_node(existing_id, |n| {
279                        node_matches_label_groups(&n.labels, &op.labels)
280                    })
281                    .unwrap_or(false);
282                if labels_ok {
283                    out.push(row);
284                }
285                continue;
286            }
287
288            for &id in &candidate_ids {
289                if !candidates_prefiltered {
290                    let labels_ok = storage
291                        .with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
292                        .unwrap_or(false);
293                    if !labels_ok {
294                        continue;
295                    }
296                }
297                let mut new_row = row.clone();
298                new_row.insert(op.var, LoraValue::Node(id));
299                out.push(new_row);
300            }
301        }
302        return Ok(out);
303    }
304
305    for row in base_rows {
306        check_optional_deadline(deadline)?;
307        if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
308            let labels_ok = storage
309                .with_node(existing_id, |n| {
310                    node_matches_label_groups(&n.labels, &op.labels)
311                })
312                .unwrap_or(false);
313            if labels_ok {
314                out.push(row);
315            }
316            continue;
317        }
318
319        for &id in &candidate_ids {
320            check_optional_deadline(deadline)?;
321            if !candidates_prefiltered {
322                let labels_ok = storage
323                    .with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
324                    .unwrap_or(false);
325                if !labels_ok {
326                    continue;
327                }
328            }
329            let mut new_row = row.clone();
330            new_row.insert(op.var, LoraValue::Node(id));
331            out.push(new_row);
332        }
333    }
334
335    Ok(out)
336}
337
338pub(super) fn node_by_property_scan_rows<S: GraphStorage>(
339    storage: &S,
340    params: &BTreeMap<String, LoraValue>,
341    base_rows: Vec<Row>,
342    op: &NodeByPropertyScanExec,
343    deadline: Option<Instant>,
344) -> ExecResult<Vec<Row>> {
345    let eval_ctx = EvalContext { storage, params };
346    let mut out = Vec::new();
347
348    for row in base_rows {
349        check_optional_deadline(deadline)?;
350        let expected = eval_expr(&op.value, &row, &eval_ctx);
351
352        if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
353            if property_scan_matches(
354                storage,
355                existing_id,
356                &op.labels,
357                &op.key,
358                &expected,
359                op.in_list,
360            ) {
361                out.push(row);
362            }
363            continue;
364        }
365
366        let candidates =
367            property_scan_candidates(storage, &op.labels, &op.key, &expected, op.in_list);
368        for id in candidates.ids {
369            check_optional_deadline(deadline)?;
370            if !candidates.prefiltered
371                && !property_scan_matches(storage, id, &op.labels, &op.key, &expected, op.in_list)
372            {
373                continue;
374            }
375            let mut new_row = row.clone();
376            new_row.insert(op.var, LoraValue::Node(id));
377            out.push(new_row);
378        }
379    }
380
381    Ok(out)
382}
383
384#[inline]
385fn check_optional_deadline(deadline: Option<Instant>) -> ExecResult<()> {
386    match deadline {
387        Some(deadline) => check_deadline_at(deadline),
388        None => Ok(()),
389    }
390}
391
392pub(crate) fn plan_may_need_hydration(plan: &PhysicalPlan) -> bool {
393    op_may_need_hydration(plan, plan.root)
394}
395
396pub(crate) fn count_all_scan_aggregation_rows<S: GraphStorage>(
397    storage: &S,
398    plan: &PhysicalPlan,
399    op: &HashAggregationExec,
400) -> Option<Vec<Row>> {
401    if !op.group_by.is_empty() {
402        return None;
403    }
404    let specs = crate::pull::classify_streamable_aggregates(&op.aggregates)?;
405    let scan_var = scan_subtree_var(plan, op.input)?;
406    // `count(*)` counts rows; `count(n)` of the scanned node counts the
407    // same rows, since a scan never binds `n` to null.
408    let counts_rows = specs.iter().all(|spec| match spec.kind {
409        crate::pull::StreamableAggKind::CountAll => true,
410        crate::pull::StreamableAggKind::CountField => {
411            matches!(&spec.arg, Some(ResolvedExpr::Variable(v)) if *v == scan_var)
412        }
413        _ => false,
414    });
415    if !counts_rows {
416        return None;
417    }
418
419    let count = count_rows_for_scan_subtree(storage, plan, op.input)? as i64;
420    let value = LoraValue::Int(count);
421    let mut row = Row::new();
422    for proj in &op.aggregates {
423        row.insert_named(proj.output, proj.name.clone(), value.clone());
424    }
425    Some(vec![row])
426}
427
428fn scan_subtree_var(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> Option<VarId> {
429    match &plan.nodes[node_id] {
430        PhysicalOp::NodeScan(op) => Some(op.var),
431        PhysicalOp::NodeByLabelScan(op) => Some(op.var),
432        _ => None,
433    }
434}
435
436fn count_rows_for_scan_subtree<S: GraphStorage>(
437    storage: &S,
438    plan: &PhysicalPlan,
439    node_id: PhysicalNodeId,
440) -> Option<usize> {
441    match &plan.nodes[node_id] {
442        PhysicalOp::NodeScan(op) if scan_input_is_argument(plan, op.input) => {
443            Some(storage.node_count())
444        }
445        PhysicalOp::NodeByLabelScan(op) if scan_input_is_argument(plan, op.input) => {
446            // One label: the label index already holds the answer.
447            if let [group] = op.labels.as_slice() {
448                if let [label] = group.as_slice() {
449                    return Some(storage.node_count_by_label(label));
450                }
451            }
452            let ids = scan_node_ids_for_label_groups(storage, &op.labels);
453            if label_group_candidates_prefiltered(&op.labels) {
454                return Some(ids.len());
455            }
456            Some(
457                ids.into_iter()
458                    .filter(|&id| {
459                        storage
460                            .with_node(id, |node| {
461                                node_matches_label_groups(&node.labels, &op.labels)
462                            })
463                            .unwrap_or(false)
464                    })
465                    .count(),
466            )
467        }
468        _ => None,
469    }
470}
471
472fn scan_input_is_argument(plan: &PhysicalPlan, input: Option<PhysicalNodeId>) -> bool {
473    match input {
474        None => true,
475        Some(id) => matches!(plan.nodes.get(id), Some(PhysicalOp::Argument(_))),
476    }
477}
478
479fn op_may_need_hydration(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
480    match &plan.nodes[node_id] {
481        PhysicalOp::Projection(op) if !op.include_existing => op
482            .items
483            .iter()
484            .any(|item| expr_may_produce_hydratable_value(&item.expr)),
485        PhysicalOp::HashAggregation(op) => {
486            op.group_by
487                .iter()
488                .any(|item| expr_may_produce_hydratable_value(&item.expr))
489                || op
490                    .aggregates
491                    .iter()
492                    .any(|item| expr_may_produce_hydratable_value(&item.expr))
493        }
494        PhysicalOp::Sort(op) => op_may_need_hydration(plan, op.input),
495        PhysicalOp::Limit(op) => op_may_need_hydration(plan, op.input),
496        PhysicalOp::Filter(op) => op_may_need_hydration(plan, op.input),
497        PhysicalOp::Unwind(op) => op_may_need_hydration(plan, op.input),
498        PhysicalOp::PathBuild(op) => op_may_need_hydration(plan, op.input),
499        _ => true,
500    }
501}
502
503fn expr_may_produce_hydratable_value(expr: &ResolvedExpr) -> bool {
504    match expr {
505        ResolvedExpr::Variable(_) | ResolvedExpr::Parameter(_) => true,
506        ResolvedExpr::Literal(_)
507        | ResolvedExpr::Property { .. }
508        | ResolvedExpr::ExistsSubquery { .. }
509        | ResolvedExpr::Binary { .. }
510        | ResolvedExpr::Unary { .. }
511        | ResolvedExpr::ListPredicate { .. } => false,
512        ResolvedExpr::Function { function, args, .. } => {
513            function_may_produce_hydratable_value(*function, args)
514        }
515        ResolvedExpr::List(items) => items.iter().any(expr_may_produce_hydratable_value),
516        ResolvedExpr::Map(items) => items
517            .iter()
518            .any(|(_, value)| expr_may_produce_hydratable_value(value)),
519        ResolvedExpr::Case {
520            alternatives,
521            else_expr,
522            ..
523        } => {
524            alternatives
525                .iter()
526                .any(|(_, value)| expr_may_produce_hydratable_value(value))
527                || else_expr
528                    .as_deref()
529                    .is_some_and(expr_may_produce_hydratable_value)
530        }
531        ResolvedExpr::ListComprehension { map_expr, .. } => map_expr
532            .as_deref()
533            .map(expr_may_produce_hydratable_value)
534            .unwrap_or(true),
535        ResolvedExpr::Reduce { expr, .. } => expr_may_produce_hydratable_value(expr),
536        ResolvedExpr::MapProjection { selectors, .. } => selectors.iter().any(|selector| {
537            matches!(selector, ResolvedMapSelector::Literal(_, expr) if expr_may_produce_hydratable_value(expr))
538        }),
539        ResolvedExpr::Index { expr, .. } | ResolvedExpr::Slice { expr, .. } => {
540            expr_may_produce_hydratable_value(expr)
541        }
542        ResolvedExpr::PatternComprehension { map_expr, .. } => {
543            expr_may_produce_hydratable_value(map_expr)
544        }
545    }
546}
547
548fn function_may_produce_hydratable_value(function: FunctionId, args: &[ResolvedExpr]) -> bool {
549    match function.name() {
550        "path.nodes" | "path.edges" | "path.first" | "path.last" | "list.first" | "list.last" => {
551            true
552        }
553        "value.coalesce" | "value.first_non_null" | "collect" => {
554            args.iter().any(expr_may_produce_hydratable_value)
555        }
556        "list.rest" | "value.reverse" | "list.reverse" => {
557            args.first().is_some_and(expr_may_produce_hydratable_value)
558        }
559        _ => false,
560    }
561}
562
563pub(super) fn expand_rows<S: GraphStorage>(
564    storage: &S,
565    params: &BTreeMap<String, LoraValue>,
566    input_rows: Vec<Row>,
567    op: &ExpandExec,
568) -> ExecResult<Vec<Row>> {
569    let eval_ctx = EvalContext { storage, params };
570    // Most expansions emit at least one row per input; sizing up front
571    // avoids re-copying every row as a large output buffer grows.
572    let mut out = Vec::with_capacity(input_rows.len());
573    // Matching (rel, dst) pairs for the current input row. Collected
574    // first so the input row can be *moved* into its last output
575    // instead of cloned: a degree-k node costs k-1 row clones, and a
576    // degree-1 hop (chains, trees walked upward) costs none.
577    let mut matches: Vec<(u64, u64)> = Vec::new();
578
579    for row in input_rows {
580        let Some(src_node_id) = bound_node_id_for_expand(&row, op.src)? else {
581            continue;
582        };
583
584        let mut rel_property_filter = None;
585        matches.clear();
586
587        storage.try_for_each_expand_id(
588            src_node_id,
589            op.direction,
590            &op.types,
591            |rel_id, dst_id| {
592                if let Some(expr) = op.rel_properties.as_ref() {
593                    if rel_property_filter.is_none() {
594                        let expected = eval_expr(expr, &row, &eval_ctx);
595                        let LoraValue::Map(map) = expected else {
596                            return Err(ExecutorError::ExpectedPropertyMap {
597                                found: value_kind(&expected),
598                            });
599                        };
600                        rel_property_filter = Some(map);
601                    }
602
603                    let Some(map) = rel_property_filter.as_ref() else {
604                        return Ok(());
605                    };
606                    let matches = storage
607                        .with_relationship(rel_id, |rel| {
608                            map.iter().all(|(key, expected)| {
609                                rel.properties
610                                    .get(key.as_str())
611                                    .map(|actual| value_matches_property_value(expected, actual))
612                                    .unwrap_or(false)
613                            })
614                        })
615                        .unwrap_or(false);
616                    if !matches {
617                        return Ok(());
618                    }
619                }
620
621                if let Some(existing_id) = bound_node_id_for_expand(&row, op.dst)? {
622                    if existing_id != dst_id {
623                        return Ok(());
624                    }
625                }
626
627                if let Some(rel_var) = op.rel {
628                    if let Some(existing_id) = bound_relationship_id_for_expand(&row, rel_var)? {
629                        if existing_id != rel_id {
630                            return Ok(());
631                        }
632                    }
633                }
634
635                matches.push((rel_id, dst_id));
636                Ok(())
637            },
638        )?;
639
640        let Some((&last, rest)) = matches.split_last() else {
641            continue;
642        };
643        for &(rel_id, dst_id) in rest {
644            let mut new_row = row.clone();
645            bind_expand_target(&mut new_row, op, rel_id, dst_id);
646            out.push(new_row);
647        }
648        let mut new_row = row;
649        bind_expand_target(&mut new_row, op, last.0, last.1);
650        out.push(new_row);
651    }
652
653    Ok(out)
654}
655
656#[inline]
657fn bind_expand_target(row: &mut Row, op: &ExpandExec, rel_id: u64, dst_id: u64) {
658    if !row.contains_key(op.dst) {
659        row.insert(op.dst, LoraValue::Node(dst_id));
660    }
661    if let Some(rel_var) = op.rel {
662        if !row.contains_key(rel_var) {
663            row.insert(rel_var, LoraValue::Relationship(rel_id));
664        }
665    }
666}
667
668pub(super) fn expand_var_len_rows<S: GraphStorage>(
669    storage: &S,
670    input_rows: Vec<Row>,
671    op: &ExpandExec,
672    range: &RangeLiteral,
673) -> ExecResult<Vec<Row>> {
674    let (min_hops, max_hops) = resolve_range(range);
675    let bind_relationships = op.rel.is_some();
676    let mut out = Vec::new();
677
678    for row in input_rows {
679        let Some(src_node_id) = bound_node_id_for_expand(&row, op.src)? else {
680            continue;
681        };
682
683        let expansions = variable_length_expand(
684            storage,
685            src_node_id,
686            op.direction,
687            &op.types,
688            min_hops,
689            max_hops,
690            bind_relationships,
691        );
692
693        for result in expansions {
694            let mut new_row = row.clone();
695            new_row.insert(op.dst, LoraValue::Node(result.dst_node_id));
696
697            if let Some(rel_var) = op.rel {
698                let rel_list = LoraValue::List(
699                    result
700                        .rel_ids
701                        .into_iter()
702                        .map(LoraValue::Relationship)
703                        .collect(),
704                );
705                new_row.insert(rel_var, rel_list);
706            }
707
708            out.push(new_row);
709        }
710    }
711
712    Ok(out)
713}
714
715pub(super) fn properties_to_value_map(props: &Properties) -> LoraValue {
716    let mut map = BTreeMap::new();
717    for (k, v) in props.iter() {
718        map.insert(k.to_string(), LoraValue::from(v));
719    }
720    LoraValue::Map(map)
721}
722
723/// Dedup rows that share the same schema (same VarId set). Compares rows by
724/// a Vec<GroupValueKey> keyed on VarId iteration order — avoids the per-row
725/// column-name String clones of `dedup_rows`. Used by DISTINCT projection.
726pub(crate) fn dedup_rows_by_vars(rows: Vec<Row>) -> Vec<Row> {
727    let mut seen: BTreeSet<Vec<GroupValueKey>> = BTreeSet::new();
728    let mut out = Vec::new();
729
730    for row in rows {
731        let key: Vec<GroupValueKey> = row
732            .iter()
733            .map(|(_, val)| GroupValueKey::from_value(val))
734            .collect();
735        if seen.insert(key) {
736            out.push(row);
737        }
738    }
739
740    out
741}
742
743/// Dedup rows using named entries so rows with different VarIds but the same
744/// column name + value are collapsed. Needed for UNION where each branch has
745/// its own VarIds.
746pub(crate) fn dedup_rows(rows: Vec<Row>) -> Vec<Row> {
747    let mut seen: BTreeSet<Vec<(String, GroupValueKey)>> = BTreeSet::new();
748    let mut out = Vec::new();
749
750    for row in rows {
751        let key: Vec<(String, GroupValueKey)> = row
752            .iter_named()
753            .map(|(_, name, val)| (name.into_owned(), GroupValueKey::from_value(val)))
754            .collect();
755        if seen.insert(key) {
756            out.push(row);
757        }
758    }
759
760    out
761}
762
763pub(super) fn eval_properties_expr<S: GraphStorage>(
764    expr: &ResolvedExpr,
765    row: &Row,
766    storage: &S,
767    params: &BTreeMap<String, LoraValue>,
768) -> ExecResult<Properties> {
769    let eval_ctx = EvalContext { storage, params };
770
771    // A `null` value in a property map means "no property": it is never
772    // stored, so `keys(n)` and existence checks do not see it.
773    if let ResolvedExpr::Map(items) = expr {
774        let mut out = Properties::new();
775        for (k, v) in items {
776            let value = eval_expr(v, row, &eval_ctx);
777            if matches!(value, LoraValue::Null) {
778                continue;
779            }
780            let prop = lora_value_to_property(value)
781                .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
782            out.insert(lora_store::intern(k), prop);
783        }
784        return Ok(out);
785    }
786
787    match eval_expr(expr, row, &eval_ctx) {
788        LoraValue::Map(map) => {
789            let mut out = Properties::new();
790            for (k, v) in map {
791                if matches!(v, LoraValue::Null) {
792                    continue;
793                }
794                let prop = lora_value_to_property(v)
795                    .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
796                // Route every CREATE/SET property key through the
797                // process-wide intern table so all nodes sharing a
798                // column name share one heap allocation. This is the
799                // hot path for `UNWIND $rows AS r CREATE (...)` and
800                // is where the bulk CSV-import memory win actually
801                // lands.
802                out.insert(lora_store::intern_owned(k), prop);
803            }
804            Ok(out)
805        }
806        other => Err(ExecutorError::ExpectedPropertyMap {
807            found: value_kind(&other),
808        }),
809    }
810}
811
812pub(crate) fn compute_aggregate_expr<S: GraphStorage>(
813    expr: &ResolvedExpr,
814    rows: &[Row],
815    eval_ctx: &EvalContext<'_, S>,
816) -> ExecResult<LoraValue> {
817    match expr {
818        ResolvedExpr::Function {
819            function,
820            distinct,
821            args,
822        } => {
823            let func = function.as_aggregate();
824            match func {
825                Some(AggregateFunction::Count) => {
826                    if args.is_empty() {
827                        return Ok(LoraValue::Int(rows.len() as i64));
828                    }
829
830                    let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
831                    values.retain(|v| !matches!(v, LoraValue::Null));
832
833                    if *distinct {
834                        values = dedup_values(values);
835                    }
836
837                    Ok(LoraValue::Int(values.len() as i64))
838                }
839
840                Some(AggregateFunction::Collect) => {
841                    if args.is_empty() {
842                        return Ok(LoraValue::List(Vec::new()));
843                    }
844
845                    let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
846
847                    if *distinct {
848                        values = dedup_values(values);
849                    }
850
851                    Ok(LoraValue::List(values))
852                }
853
854                Some(AggregateFunction::Sum) => {
855                    if args.is_empty() {
856                        return Ok(LoraValue::Null);
857                    }
858
859                    let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
860
861                    if *distinct {
862                        values = dedup_values(values);
863                    }
864
865                    let nums = values
866                        .into_iter()
867                        .filter_map(as_f64_lossy)
868                        .collect::<Vec<_>>();
869
870                    if nums.is_empty() {
871                        Ok(LoraValue::Null)
872                    } else if nums.iter().all(|n| n.fract() == 0.0) {
873                        Ok(LoraValue::Int(nums.iter().sum::<f64>() as i64))
874                    } else {
875                        Ok(LoraValue::Float(nums.iter().sum::<f64>()))
876                    }
877                }
878
879                Some(AggregateFunction::Avg) => {
880                    if args.is_empty() {
881                        return Ok(LoraValue::Null);
882                    }
883
884                    let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
885
886                    if *distinct {
887                        values = dedup_values(values);
888                    }
889
890                    let nums = values
891                        .into_iter()
892                        .filter_map(as_f64_lossy)
893                        .collect::<Vec<_>>();
894
895                    if nums.is_empty() {
896                        Ok(LoraValue::Null)
897                    } else {
898                        Ok(LoraValue::Float(
899                            nums.iter().sum::<f64>() / nums.len() as f64,
900                        ))
901                    }
902                }
903
904                Some(AggregateFunction::Min) => {
905                    if args.is_empty() {
906                        return Ok(LoraValue::Null);
907                    }
908
909                    let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
910                    values.retain(|v| !matches!(v, LoraValue::Null));
911
912                    if *distinct {
913                        values = dedup_values(values);
914                    }
915
916                    Ok(values
917                        .into_iter()
918                        .min_by(compare_values_total)
919                        .unwrap_or(LoraValue::Null))
920                }
921
922                Some(AggregateFunction::Max) => {
923                    if args.is_empty() {
924                        return Ok(LoraValue::Null);
925                    }
926
927                    let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
928                    values.retain(|v| !matches!(v, LoraValue::Null));
929
930                    if *distinct {
931                        values = dedup_values(values);
932                    }
933
934                    Ok(values
935                        .into_iter()
936                        .max_by(compare_values_total)
937                        .unwrap_or(LoraValue::Null))
938                }
939
940                Some(AggregateFunction::Stdev | AggregateFunction::Stdevp) => {
941                    if args.is_empty() {
942                        return Ok(LoraValue::Null);
943                    }
944
945                    let nums: Vec<f64> = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?
946                        .into_iter()
947                        .filter_map(as_f64_lossy)
948                        .collect();
949
950                    let is_population = matches!(func, Some(AggregateFunction::Stdevp));
951
952                    if nums.is_empty() || (!is_population && nums.len() < 2) {
953                        return Ok(LoraValue::Float(0.0));
954                    }
955
956                    let mean = nums.iter().sum::<f64>() / nums.len() as f64;
957                    let variance_sum: f64 = nums.iter().map(|x| (x - mean).powi(2)).sum();
958                    let denom = if is_population {
959                        nums.len() as f64
960                    } else {
961                        (nums.len() - 1) as f64
962                    };
963                    Ok(LoraValue::Float((variance_sum / denom).sqrt()))
964                }
965
966                Some(AggregateFunction::PercentileCont) => {
967                    if args.len() < 2 {
968                        return Ok(LoraValue::Null);
969                    }
970
971                    let Some(first) = rows.first() else {
972                        return Ok(LoraValue::Null);
973                    };
974
975                    let percentile = eval_expr_result(&args[1], first, eval_ctx)
976                        .map_err(ExecutorError::RuntimeError)?
977                        .as_f64()
978                        .map(normalize_percentile)
979                        .unwrap_or(0.5);
980                    let mut nums: Vec<f64> = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?
981                        .into_iter()
982                        .filter_map(as_f64_lossy)
983                        .collect();
984
985                    if nums.is_empty() {
986                        return Ok(LoraValue::Null);
987                    }
988
989                    nums.sort_by(|a, b| a.partial_cmp(b).unwrap_or(Ordering::Equal));
990
991                    let index = percentile * (nums.len() - 1) as f64;
992                    let lower = index.floor() as usize;
993                    let upper = index.ceil() as usize;
994                    let fraction = index - lower as f64;
995
996                    if lower == upper || upper >= nums.len() {
997                        Ok(LoraValue::Float(nums[lower]))
998                    } else {
999                        Ok(LoraValue::Float(
1000                            nums[lower] * (1.0 - fraction) + nums[upper] * fraction,
1001                        ))
1002                    }
1003                }
1004
1005                Some(AggregateFunction::PercentileDisc) => {
1006                    if args.len() < 2 {
1007                        return Ok(LoraValue::Null);
1008                    }
1009
1010                    let Some(first) = rows.first() else {
1011                        return Ok(LoraValue::Null);
1012                    };
1013
1014                    let percentile = eval_expr_result(&args[1], first, eval_ctx)
1015                        .map_err(ExecutorError::RuntimeError)?
1016                        .as_f64()
1017                        .map(normalize_percentile)
1018                        .unwrap_or(0.5);
1019                    let mut nums: Vec<f64> = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?
1020                        .into_iter()
1021                        .filter_map(as_f64_lossy)
1022                        .collect();
1023
1024                    if nums.is_empty() {
1025                        return Ok(LoraValue::Null);
1026                    }
1027
1028                    nums.sort_by(|a, b| a.partial_cmp(b).unwrap_or(Ordering::Equal));
1029
1030                    let index = (percentile * (nums.len() - 1) as f64).round() as usize;
1031                    let index = index.min(nums.len() - 1);
1032                    Ok(LoraValue::Float(nums[index]))
1033                }
1034
1035                _ => eval_first_or_null(expr, rows, eval_ctx),
1036            }
1037        }
1038
1039        _ => eval_first_or_null(expr, rows, eval_ctx),
1040    }
1041}
1042
1043fn eval_aggregate_arg_values<S: GraphStorage>(
1044    expr: &ResolvedExpr,
1045    rows: &[Row],
1046    eval_ctx: &EvalContext<'_, S>,
1047) -> ExecResult<Vec<LoraValue>> {
1048    rows.iter()
1049        .map(|row| eval_expr_result(expr, row, eval_ctx).map_err(ExecutorError::RuntimeError))
1050        .collect()
1051}
1052
1053fn normalize_percentile(value: f64) -> f64 {
1054    if value.is_finite() {
1055        value.clamp(0.0, 1.0)
1056    } else {
1057        0.5
1058    }
1059}
1060
1061fn eval_first_or_null<S: GraphStorage>(
1062    expr: &ResolvedExpr,
1063    rows: &[Row],
1064    eval_ctx: &EvalContext<'_, S>,
1065) -> ExecResult<LoraValue> {
1066    match rows.first() {
1067        Some(row) => eval_expr_result(expr, row, eval_ctx).map_err(ExecutorError::RuntimeError),
1068        None => Ok(LoraValue::Null),
1069    }
1070}
1071
1072fn dedup_values(values: Vec<LoraValue>) -> Vec<LoraValue> {
1073    let mut seen: BTreeSet<GroupValueKey> = BTreeSet::new();
1074    let mut out = Vec::new();
1075
1076    for value in values {
1077        let key = GroupValueKey::from_value(&value);
1078        if seen.insert(key) {
1079            out.push(value);
1080        }
1081    }
1082
1083    out
1084}
1085
1086fn as_f64_lossy(v: LoraValue) -> Option<f64> {
1087    match v {
1088        LoraValue::Int(i) => Some(i as f64),
1089        LoraValue::Float(f) => Some(f),
1090        _ => None,
1091    }
1092}
1093
1094pub(super) fn compare_values_total(a: &LoraValue, b: &LoraValue) -> Ordering {
1095    use LoraValue::*;
1096
1097    match (a, b) {
1098        (Bool(x), Bool(y)) => x.cmp(y),
1099        (Int(x), Int(y)) => x.cmp(y),
1100        (Float(x), Float(y)) => x.partial_cmp(y).unwrap_or(Ordering::Equal),
1101        (Int(x), Float(y)) => (*x as f64).partial_cmp(y).unwrap_or(Ordering::Equal),
1102        (Float(x), Int(y)) => x.partial_cmp(&(*y as f64)).unwrap_or(Ordering::Equal),
1103        (String(x), String(y)) => x.cmp(y),
1104        (Binary(x), Binary(y)) => x.segments().cmp(y.segments()),
1105        (Node(x), Node(y)) => x.cmp(y),
1106        (Relationship(x), Relationship(y)) => x.cmp(y),
1107        (Date(x), Date(y)) => x.cmp(y),
1108        (DateTime(x), DateTime(y)) => x.cmp(y),
1109        (LocalDateTime(x), LocalDateTime(y)) => x.cmp(y),
1110        (Time(x), Time(y)) => x.cmp(y),
1111        (LocalTime(x), LocalTime(y)) => x.cmp(y),
1112        (Duration(x), Duration(y)) => x.cmp(y),
1113        (Vector(x), Vector(y)) => x.to_key_string().cmp(&y.to_key_string()),
1114        _ => type_rank(a)
1115            .cmp(&type_rank(b))
1116            .then_with(|| format!("{a:?}").cmp(&format!("{b:?}"))),
1117    }
1118}
1119
1120pub fn value_matches_property_value(expected: &LoraValue, actual: &PropertyValue) -> bool {
1121    match (expected, actual) {
1122        (LoraValue::Null, PropertyValue::Null) => true,
1123        (LoraValue::Bool(a), PropertyValue::Bool(b)) => a == b,
1124        (LoraValue::Int(a), PropertyValue::Int(b)) => a == b,
1125        (LoraValue::Float(a), PropertyValue::Float(b)) => a == b,
1126        (LoraValue::Int(a), PropertyValue::Float(b)) => (*a as f64) == *b,
1127        (LoraValue::Float(a), PropertyValue::Int(b)) => *a == (*b as f64),
1128        (LoraValue::String(a), PropertyValue::String(b)) => a == b,
1129        (LoraValue::Binary(a), PropertyValue::Binary(b)) => a == b,
1130
1131        (LoraValue::List(xs), PropertyValue::List(ys)) => {
1132            xs.len() == ys.len()
1133                && xs
1134                    .iter()
1135                    .zip(ys.iter())
1136                    .all(|(x, y)| value_matches_property_value(x, y))
1137        }
1138
1139        (LoraValue::Map(xm), PropertyValue::Map(ym)) => xm.iter().all(|(k, xv)| {
1140            ym.get(k)
1141                .map(|yv| value_matches_property_value(xv, yv))
1142                .unwrap_or(false)
1143        }),
1144
1145        (LoraValue::Date(a), PropertyValue::Date(b)) => a == b,
1146        (LoraValue::DateTime(a), PropertyValue::DateTime(b)) => a == b,
1147        (LoraValue::LocalDateTime(a), PropertyValue::LocalDateTime(b)) => a == b,
1148        (LoraValue::Time(a), PropertyValue::Time(b)) => a == b,
1149        (LoraValue::LocalTime(a), PropertyValue::LocalTime(b)) => a == b,
1150        (LoraValue::Duration(a), PropertyValue::Duration(b)) => a == b,
1151        (LoraValue::Point(a), PropertyValue::Point(b)) => a == b,
1152        (LoraValue::Vector(a), PropertyValue::Vector(b)) => a == b,
1153
1154        _ => false,
1155    }
1156}
1157
1158pub(crate) fn node_matches_property_filter<S: GraphStorage>(
1159    storage: &S,
1160    node_id: NodeId,
1161    labels: &[Vec<String>],
1162    key: &str,
1163    expected: &LoraValue,
1164) -> bool {
1165    storage
1166        .with_node(node_id, |node| {
1167            node_matches_label_groups(&node.labels, labels)
1168                && node
1169                    .properties
1170                    .get(key)
1171                    .map(|actual| value_matches_property_value(expected, actual))
1172                    .unwrap_or(false)
1173        })
1174        .unwrap_or(false)
1175}
1176
1177fn single_label_hint(labels: &[Vec<String>]) -> Option<&str> {
1178    if labels.len() == 1 && labels[0].len() == 1 {
1179        Some(labels[0][0].as_str())
1180    } else {
1181        None
1182    }
1183}
1184
1185fn property_lookup_values(expected: &LoraValue) -> Option<Vec<PropertyValue>> {
1186    let property = lora_value_to_property(expected.clone()).ok()?;
1187    let mut values = vec![property.clone()];
1188
1189    match property {
1190        PropertyValue::Int(i) => {
1191            values.push(PropertyValue::Float(i as f64));
1192        }
1193        PropertyValue::Float(f)
1194            if f.is_finite()
1195                && f.fract() == 0.0
1196                && f >= i64::MIN as f64
1197                && f <= i64::MAX as f64 =>
1198        {
1199            values.push(PropertyValue::Int(f as i64));
1200        }
1201        _ => {}
1202    }
1203
1204    Some(values)
1205}
1206
1207pub(crate) struct NodePropertyCandidates {
1208    pub(crate) ids: Vec<NodeId>,
1209    pub(crate) prefiltered: bool,
1210}
1211
1212pub(crate) fn node_by_property_range_scan_rows<S: GraphStorage>(
1213    storage: &S,
1214    params: &BTreeMap<String, LoraValue>,
1215    base_rows: Vec<Row>,
1216    op: &lora_compiler::NodeByPropertyRangeScanExec,
1217    deadline: Option<Instant>,
1218) -> ExecResult<Vec<Row>> {
1219    let eval_ctx = EvalContext { storage, params };
1220    let mut out = Vec::new();
1221
1222    for row in base_rows {
1223        check_optional_deadline(deadline)?;
1224        let lo_value = op.lo.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
1225        let hi_value = op.hi.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
1226        let lo_prop = lo_value
1227            .clone()
1228            .and_then(|v| lora_value_to_property(v).ok());
1229        let hi_prop = hi_value
1230            .clone()
1231            .and_then(|v| lora_value_to_property(v).ok());
1232        let filter = NodeRangeFilter {
1233            labels: &op.labels,
1234            key: &op.key,
1235            lo: lo_value.as_ref(),
1236            lo_inclusive: op.lo_inclusive,
1237            hi: hi_value.as_ref(),
1238            hi_inclusive: op.hi_inclusive,
1239        };
1240
1241        if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
1242            if node_matches_range_filter(storage, existing_id, &filter) {
1243                out.push(row);
1244            }
1245            continue;
1246        }
1247
1248        if op.order.is_some() {
1249            let mut cursor =
1250                OrderedRangeCursor::new(storage, op, lo_value.clone(), hi_value.clone());
1251            while let Some(id) = cursor.next_id(storage) {
1252                check_optional_deadline(deadline)?;
1253                let mut new_row = row.clone();
1254                new_row.insert(op.var, LoraValue::Node(id));
1255                out.push(new_row);
1256            }
1257            continue;
1258        }
1259
1260        let candidate_ids = match single_label_hint(&op.labels) {
1261            Some(label) => storage
1262                .node_range_candidates(label, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
1263                .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1264            None => scan_node_ids_for_label_groups(storage, &op.labels),
1265        };
1266
1267        for id in candidate_ids {
1268            check_optional_deadline(deadline)?;
1269            if node_matches_range_filter(storage, id, &filter) {
1270                let mut new_row = row.clone();
1271                new_row.insert(op.var, LoraValue::Node(id));
1272                out.push(new_row);
1273            }
1274        }
1275    }
1276
1277    Ok(out)
1278}
1279
1280pub(crate) fn node_by_text_scan_rows<S: GraphStorage>(
1281    storage: &S,
1282    params: &BTreeMap<String, LoraValue>,
1283    base_rows: Vec<Row>,
1284    op: &lora_compiler::NodeByTextScanExec,
1285    deadline: Option<Instant>,
1286) -> ExecResult<Vec<Row>> {
1287    let eval_ctx = EvalContext { storage, params };
1288    let mut out = Vec::new();
1289
1290    for row in base_rows {
1291        check_optional_deadline(deadline)?;
1292        let query = eval_expr(&op.query, &row, &eval_ctx);
1293        let LoraValue::String(query_str) = &query else {
1294            // Non-string query → predicate cannot match anything;
1295            // skip the row entirely (matches scan + filter behaviour).
1296            continue;
1297        };
1298
1299        if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
1300            if node_matches_text_filter(
1301                storage,
1302                existing_id,
1303                &op.labels,
1304                &op.key,
1305                op.predicate,
1306                query_str,
1307            ) {
1308                out.push(row);
1309            }
1310            continue;
1311        }
1312
1313        let candidate_ids = match single_label_hint(&op.labels) {
1314            Some(label) => storage
1315                .node_text_candidates(label, &op.key, query_str)
1316                .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1317            None => scan_node_ids_for_label_groups(storage, &op.labels),
1318        };
1319
1320        for id in candidate_ids {
1321            check_optional_deadline(deadline)?;
1322            if node_matches_text_filter(storage, id, &op.labels, &op.key, op.predicate, query_str) {
1323                let mut new_row = row.clone();
1324                new_row.insert(op.var, LoraValue::Node(id));
1325                out.push(new_row);
1326            }
1327        }
1328    }
1329
1330    Ok(out)
1331}
1332
1333struct NodeRangeFilter<'a> {
1334    labels: &'a [Vec<String>],
1335    key: &'a str,
1336    lo: Option<&'a LoraValue>,
1337    lo_inclusive: bool,
1338    hi: Option<&'a LoraValue>,
1339    hi_inclusive: bool,
1340}
1341
1342fn node_matches_range_filter<S: GraphStorage>(
1343    storage: &S,
1344    id: NodeId,
1345    filter: &NodeRangeFilter<'_>,
1346) -> bool {
1347    storage
1348        .with_node(id, |n| {
1349            if !node_matches_label_groups(&n.labels, filter.labels) {
1350                return false;
1351            }
1352            let Some(actual) = n.properties.get(filter.key) else {
1353                return false;
1354            };
1355            let actual_lv = lora_store_property_to_value(actual);
1356            range_predicate_holds(
1357                &actual_lv,
1358                filter.lo,
1359                filter.lo_inclusive,
1360                filter.hi,
1361                filter.hi_inclusive,
1362            )
1363        })
1364        .unwrap_or(false)
1365}
1366
1367fn node_matches_text_filter<S: GraphStorage>(
1368    storage: &S,
1369    id: NodeId,
1370    labels: &[Vec<String>],
1371    key: &str,
1372    predicate: lora_compiler::TextPredicate,
1373    query: &str,
1374) -> bool {
1375    storage
1376        .with_node(id, |n| {
1377            if !node_matches_label_groups(&n.labels, labels) {
1378                return false;
1379            }
1380            let Some(PropertyValue::String(actual)) = n.properties.get(key) else {
1381                return false;
1382            };
1383            text_predicate_holds(actual, predicate, query)
1384        })
1385        .unwrap_or(false)
1386}
1387
1388fn text_predicate_holds(
1389    actual: &str,
1390    predicate: lora_compiler::TextPredicate,
1391    query: &str,
1392) -> bool {
1393    match predicate {
1394        lora_compiler::TextPredicate::StartsWith => actual.starts_with(query),
1395        lora_compiler::TextPredicate::EndsWith => actual.ends_with(query),
1396        lora_compiler::TextPredicate::Contains => actual.contains(query),
1397    }
1398}
1399
1400fn range_predicate_holds(
1401    actual: &LoraValue,
1402    lo: Option<&LoraValue>,
1403    lo_inclusive: bool,
1404    hi: Option<&LoraValue>,
1405    hi_inclusive: bool,
1406) -> bool {
1407    if let Some(lo) = lo {
1408        match range_comparison(actual, lo) {
1409            None => return false,
1410            Some(Ordering::Less) => return false,
1411            Some(Ordering::Equal) if !lo_inclusive => return false,
1412            _ => {}
1413        }
1414    }
1415    if let Some(hi) = hi {
1416        match range_comparison(actual, hi) {
1417            None => return false,
1418            Some(Ordering::Greater) => return false,
1419            Some(Ordering::Equal) if !hi_inclusive => return false,
1420            _ => {}
1421        }
1422    }
1423    true
1424}
1425
1426fn range_comparison(actual: &LoraValue, bound: &LoraValue) -> Option<Ordering> {
1427    match (actual, bound) {
1428        (LoraValue::Null, _) | (_, LoraValue::Null) => None,
1429        (LoraValue::String(a), LoraValue::String(b)) => Some(a.cmp(b)),
1430        (
1431            LoraValue::Date(_)
1432            | LoraValue::DateTime(_)
1433            | LoraValue::LocalDateTime(_)
1434            | LoraValue::Time(_)
1435            | LoraValue::LocalTime(_),
1436            _,
1437        ) => actual.temporal_cmp(bound),
1438        (LoraValue::Duration(a), LoraValue::Duration(b)) => a
1439            .total_seconds_approx()
1440            .partial_cmp(&b.total_seconds_approx()),
1441        _ => actual.as_f64()?.partial_cmp(&bound.as_f64()?),
1442    }
1443}
1444
1445fn lora_store_property_to_value(value: &PropertyValue) -> LoraValue {
1446    LoraValue::from(value)
1447}
1448
1449pub(crate) fn node_by_point_scan_rows<S: GraphStorage>(
1450    storage: &S,
1451    params: &BTreeMap<String, LoraValue>,
1452    base_rows: Vec<Row>,
1453    op: &lora_compiler::NodeByPointScanExec,
1454    deadline: Option<Instant>,
1455) -> ExecResult<Vec<Row>> {
1456    let eval_ctx = EvalContext { storage, params };
1457    let mut out = Vec::new();
1458
1459    for row in base_rows {
1460        check_optional_deadline(deadline)?;
1461
1462        // Resolve the predicate's literal/scalar inputs against the
1463        // current row. These end up as captured `LoraValue`s used both
1464        // to probe the spatial index and to refilter every candidate.
1465        let probe = match &op.predicate {
1466            lora_compiler::PointPredicate::WithinBBox {
1467                lower_left,
1468                upper_right,
1469            } => {
1470                let ll = eval_expr(lower_left, &row, &eval_ctx);
1471                let ur = eval_expr(upper_right, &row, &eval_ctx);
1472                match (ll, ur) {
1473                    (LoraValue::Point(a), LoraValue::Point(b)) => {
1474                        Probe::WithinBBox { ll: a, ur: b }
1475                    }
1476                    _ => continue,
1477                }
1478            }
1479            lora_compiler::PointPredicate::WithinDistance {
1480                center,
1481                max_distance,
1482                inclusive,
1483            } => {
1484                let c = eval_expr(center, &row, &eval_ctx);
1485                let d = eval_expr(max_distance, &row, &eval_ctx);
1486                match (c, d) {
1487                    (LoraValue::Point(c), LoraValue::Float(d)) => Probe::WithinDistance {
1488                        center: c,
1489                        max: d,
1490                        inclusive: *inclusive,
1491                    },
1492                    (LoraValue::Point(c), LoraValue::Int(d)) => Probe::WithinDistance {
1493                        center: c,
1494                        max: d as f64,
1495                        inclusive: *inclusive,
1496                    },
1497                    _ => continue,
1498                }
1499            }
1500        };
1501
1502        if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
1503            if node_matches_point_filter(storage, existing_id, &op.labels, &op.key, &probe) {
1504                out.push(row);
1505            }
1506            continue;
1507        }
1508
1509        let candidate_ids = match single_label_hint(&op.labels) {
1510            Some(label) => match &probe {
1511                Probe::WithinBBox { ll, ur } => storage
1512                    .node_point_within_bbox(label, &op.key, (ll.x, ll.y), (ur.x, ur.y))
1513                    .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1514                Probe::WithinDistance { center, max, .. } => storage
1515                    .node_point_within_distance(label, &op.key, (center.x, center.y), *max)
1516                    .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1517            },
1518            None => scan_node_ids_for_label_groups(storage, &op.labels),
1519        };
1520
1521        for id in candidate_ids {
1522            check_optional_deadline(deadline)?;
1523            if node_matches_point_filter(storage, id, &op.labels, &op.key, &probe) {
1524                let mut new_row = row.clone();
1525                new_row.insert(op.var, LoraValue::Node(id));
1526                out.push(new_row);
1527            }
1528        }
1529    }
1530
1531    Ok(out)
1532}
1533
1534fn node_matches_point_filter<S: GraphStorage>(
1535    storage: &S,
1536    id: NodeId,
1537    labels: &[Vec<String>],
1538    key: &str,
1539    probe: &Probe,
1540) -> bool {
1541    storage
1542        .with_node(id, |n| {
1543            if !node_matches_label_groups(&n.labels, labels) {
1544                return false;
1545            }
1546            let Some(PropertyValue::Point(point)) = n.properties.get(key) else {
1547                return false;
1548            };
1549            point_predicate_holds(point, probe)
1550        })
1551        .unwrap_or(false)
1552}
1553
1554enum Probe {
1555    WithinBBox {
1556        ll: lora_store::LoraPoint,
1557        ur: lora_store::LoraPoint,
1558    },
1559    WithinDistance {
1560        center: lora_store::LoraPoint,
1561        max: f64,
1562        inclusive: bool,
1563    },
1564}
1565
1566fn point_predicate_holds(actual: &lora_store::LoraPoint, probe: &Probe) -> bool {
1567    match probe {
1568        Probe::WithinBBox { ll, ur } => {
1569            if actual.srid != ll.srid || actual.srid != ur.srid {
1570                return false;
1571            }
1572            let in_x = actual.x >= ll.x.min(ur.x) && actual.x <= ll.x.max(ur.x);
1573            let in_y = actual.y >= ll.y.min(ur.y) && actual.y <= ll.y.max(ur.y);
1574            let in_z = match (actual.z, ll.z, ur.z) {
1575                (Some(pz), Some(lz), Some(uz)) => pz >= lz.min(uz) && pz <= lz.max(uz),
1576                (None, None, None) => true,
1577                _ => return false,
1578            };
1579            in_x && in_y && in_z
1580        }
1581        Probe::WithinDistance {
1582            center,
1583            max,
1584            inclusive,
1585        } => {
1586            let Some(d) = lora_store::point_distance(actual, center) else {
1587                return false;
1588            };
1589            if *inclusive {
1590                d <= *max
1591            } else {
1592                d < *max
1593            }
1594        }
1595    }
1596}
1597
1598pub(crate) fn indexed_node_property_candidates<S: GraphStorage>(
1599    storage: &S,
1600    labels: &[Vec<String>],
1601    key: &str,
1602    expected: &LoraValue,
1603) -> NodePropertyCandidates {
1604    let Some(values) = property_lookup_values(expected) else {
1605        return NodePropertyCandidates {
1606            ids: scan_node_ids_for_label_groups(storage, labels),
1607            prefiltered: false,
1608        };
1609    };
1610
1611    let label_hint = single_label_hint(labels);
1612    let mut seen = BTreeSet::new();
1613    let mut out = Vec::new();
1614    for value in values {
1615        for id in storage.find_node_ids_by_property(label_hint, key, &value) {
1616            if seen.insert(id) {
1617                out.push(id);
1618            }
1619        }
1620    }
1621    NodePropertyCandidates {
1622        ids: out,
1623        prefiltered: labels.is_empty() || label_hint.is_some(),
1624    }
1625}
1626
1627/// Candidate ids for a `NodeByPropertyScan`. With `in_list` the scan
1628/// seeks `key IN expected`: one lookup per distinct list element, the
1629/// union deduplicated. A `null` list matches nothing. Any other non-list
1630/// value yields every labelled node unfiltered, so the `Filter` kept
1631/// above the scan evaluates (and reports) the predicate itself.
1632pub(crate) fn property_scan_candidates<S: GraphStorage>(
1633    storage: &S,
1634    labels: &[Vec<String>],
1635    key: &str,
1636    expected: &LoraValue,
1637    in_list: bool,
1638) -> NodePropertyCandidates {
1639    if !in_list {
1640        return indexed_node_property_candidates(storage, labels, key, expected);
1641    }
1642    match expected {
1643        LoraValue::Null => NodePropertyCandidates {
1644            ids: Vec::new(),
1645            prefiltered: true,
1646        },
1647        LoraValue::List(items) => {
1648            let mut seen = BTreeSet::new();
1649            let mut ids = Vec::new();
1650            for item in items {
1651                if matches!(item, LoraValue::Null) {
1652                    continue;
1653                }
1654                let candidates = indexed_node_property_candidates(storage, labels, key, item);
1655                for id in candidates.ids {
1656                    if seen.contains(&id) {
1657                        continue;
1658                    }
1659                    if !candidates.prefiltered
1660                        && !node_matches_property_filter(storage, id, labels, key, item)
1661                    {
1662                        continue;
1663                    }
1664                    seen.insert(id);
1665                    ids.push(id);
1666                }
1667            }
1668            NodePropertyCandidates {
1669                ids,
1670                prefiltered: true,
1671            }
1672        }
1673        _ => NodePropertyCandidates {
1674            ids: scan_node_ids_for_label_groups(storage, labels)
1675                .into_iter()
1676                .filter(|&id| {
1677                    storage
1678                        .with_node(id, |n| node_matches_label_groups(&n.labels, labels))
1679                        .unwrap_or(false)
1680                })
1681                .collect(),
1682            prefiltered: true,
1683        },
1684    }
1685}
1686
1687/// Whether `node_id` passes a `NodeByPropertyScan` (labels plus
1688/// `key = expected`, or `key IN expected` when `in_list`). Mirrors
1689/// [`property_scan_candidates`], including its non-list fallback.
1690pub(crate) fn property_scan_matches<S: GraphStorage>(
1691    storage: &S,
1692    node_id: NodeId,
1693    labels: &[Vec<String>],
1694    key: &str,
1695    expected: &LoraValue,
1696    in_list: bool,
1697) -> bool {
1698    if !in_list {
1699        return node_matches_property_filter(storage, node_id, labels, key, expected);
1700    }
1701    match expected {
1702        LoraValue::Null => false,
1703        LoraValue::List(items) => items.iter().any(|item| {
1704            !matches!(item, LoraValue::Null)
1705                && node_matches_property_filter(storage, node_id, labels, key, item)
1706        }),
1707        _ => storage
1708            .with_node(node_id, |n| node_matches_label_groups(&n.labels, labels))
1709            .unwrap_or(false),
1710    }
1711}
1712
1713/// Build a LoraPath from the node and relationship variables currently in a row.
1714///
1715/// For variable-length relationships (stored as a List of Relationship values),
1716/// intermediate nodes are reconstructed from the storage by walking the
1717/// relationship chain.
1718pub(crate) fn build_path_value<S: GraphStorage>(
1719    row: &Row,
1720    node_vars: &[VarId],
1721    rel_vars: &[VarId],
1722    storage: &S,
1723) -> LoraValue {
1724    let (raw_nodes, rels, has_var_len) = path_bindings(row, node_vars, rel_vars);
1725
1726    let nodes = if has_var_len && !rels.is_empty() && raw_nodes.len() == 2 {
1727        reconstruct_var_len_nodes(raw_nodes[0], &rels, storage)
1728    } else {
1729        raw_nodes
1730    };
1731
1732    LoraValue::Path(LoraPath { nodes, rels })
1733}
1734
1735#[inline]
1736fn path_bindings(
1737    row: &Row,
1738    node_vars: &[VarId],
1739    rel_vars: &[VarId],
1740) -> (Vec<NodeId>, Vec<RelationshipId>, bool) {
1741    let mut raw_nodes = Vec::new();
1742    let mut rels = Vec::new();
1743    let mut has_var_len = false;
1744
1745    for &nv in node_vars {
1746        match row.get(nv) {
1747            Some(LoraValue::Node(id)) => raw_nodes.push(*id),
1748            Some(LoraValue::List(items)) => {
1749                for item in items {
1750                    if let LoraValue::Node(id) = item {
1751                        raw_nodes.push(*id);
1752                    }
1753                }
1754            }
1755            _ => {}
1756        }
1757    }
1758
1759    for &rv in rel_vars {
1760        match row.get(rv) {
1761            Some(LoraValue::Relationship(id)) => rels.push(*id),
1762            Some(LoraValue::List(items)) => {
1763                has_var_len = true;
1764                for item in items {
1765                    if let LoraValue::Relationship(id) = item {
1766                        rels.push(*id);
1767                    }
1768                }
1769            }
1770            _ => {}
1771        }
1772    }
1773
1774    (raw_nodes, rels, has_var_len)
1775}
1776
1777#[inline]
1778fn reconstruct_var_len_nodes<S: GraphStorage>(
1779    start: NodeId,
1780    rels: &[RelationshipId],
1781    storage: &S,
1782) -> Vec<NodeId> {
1783    let mut ordered = Vec::with_capacity(rels.len() + 1);
1784    ordered.push(start);
1785    let mut current = start;
1786    for &rel_id in rels {
1787        if let Some((src, dst)) = storage.relationship_endpoints(rel_id) {
1788            let next = if src == current { dst } else { src };
1789            ordered.push(next);
1790            current = next;
1791        }
1792    }
1793    ordered
1794}
1795
1796fn type_rank(v: &LoraValue) -> u8 {
1797    match v {
1798        LoraValue::Null => 0,
1799        LoraValue::Bool(_) => 1,
1800        LoraValue::Int(_) | LoraValue::Float(_) => 2,
1801        LoraValue::String(_) => 3,
1802        LoraValue::Binary(_) => 4,
1803        LoraValue::Date(_) => 5,
1804        LoraValue::DateTime(_) => 6,
1805        LoraValue::LocalDateTime(_) => 7,
1806        LoraValue::Time(_) => 8,
1807        LoraValue::LocalTime(_) => 9,
1808        LoraValue::Duration(_) => 10,
1809        LoraValue::Point(_) => 11,
1810        LoraValue::Vector(_) => 12,
1811        LoraValue::List(_) => 13,
1812        LoraValue::Map(_) => 14,
1813        LoraValue::Node(_) => 15,
1814        LoraValue::Relationship(_) => 16,
1815        LoraValue::Path(_) => 17,
1816    }
1817}
1818
1819/// Check whether a node's labels satisfy all label groups.
1820/// Each group is a disjunction (OR): the node must have at least one label
1821/// from the group.  Groups are conjunctive (AND): all groups must be satisfied.
1822pub(crate) fn node_matches_label_groups(node_labels: &[String], groups: &[Vec<String>]) -> bool {
1823    groups
1824        .iter()
1825        .all(|group| group.iter().any(|l| node_labels.iter().any(|nl| nl == l)))
1826}
1827
1828/// Scan the graph for candidate node IDs matching the label groups. Uses the
1829/// label index for the pick-first-label phase and avoids cloning NodeRecords.
1830pub(crate) fn scan_node_ids_for_label_groups<S: GraphStorage>(
1831    storage: &S,
1832    groups: &[Vec<String>],
1833) -> Vec<NodeId> {
1834    if groups.is_empty() {
1835        return storage.all_node_ids();
1836    }
1837    if groups.len() == 1 {
1838        return label_group_candidate_ids(storage, &groups[0]);
1839    }
1840
1841    let mut best: Option<Vec<NodeId>> = None;
1842    for group in groups {
1843        let ids = label_group_candidate_ids(storage, group);
1844        if ids.is_empty() {
1845            return Vec::new();
1846        }
1847        if best
1848            .as_ref()
1849            .map(|current| ids.len() < current.len())
1850            .unwrap_or(true)
1851        {
1852            best = Some(ids);
1853        }
1854    }
1855
1856    best.unwrap_or_default()
1857}
1858
1859pub(crate) fn label_group_candidates_prefiltered(groups: &[Vec<String>]) -> bool {
1860    groups.len() <= 1
1861}
1862
1863fn label_group_candidate_ids<S: GraphStorage>(storage: &S, group: &[String]) -> Vec<NodeId> {
1864    match group {
1865        [] => Vec::new(),
1866        [label] => storage.node_ids_by_label(label),
1867        labels => {
1868            let mut seen = BTreeSet::new();
1869            let mut out = Vec::new();
1870            for label in labels {
1871                for id in storage.node_ids_by_label(label) {
1872                    if seen.insert(id) {
1873                        out.push(id);
1874                    }
1875                }
1876            }
1877            out
1878        }
1879    }
1880}
1881
1882pub(crate) fn hydrate_node_record(node: &lora_store::NodeRecord) -> LoraValue {
1883    let mut map = BTreeMap::new();
1884    map.insert("kind".to_string(), LoraValue::String("node".to_string()));
1885    map.insert("id".to_string(), LoraValue::Int(node.id as i64));
1886    map.insert(
1887        "labels".to_string(),
1888        LoraValue::List(
1889            node.labels
1890                .iter()
1891                .map(|s| LoraValue::String(s.clone()))
1892                .collect(),
1893        ),
1894    );
1895    map.insert(
1896        "properties".to_string(),
1897        properties_to_value_map(&node.properties),
1898    );
1899    LoraValue::Map(map)
1900}
1901
1902pub(crate) fn hydrate_relationship_record(rel: &lora_store::RelationshipRecord) -> LoraValue {
1903    let mut map = BTreeMap::new();
1904    map.insert(
1905        "kind".to_string(),
1906        LoraValue::String("relationship".to_string()),
1907    );
1908    map.insert("id".to_string(), LoraValue::Int(rel.id as i64));
1909    map.insert("startId".to_string(), LoraValue::Int(rel.src as i64));
1910    map.insert("endId".to_string(), LoraValue::Int(rel.dst as i64));
1911    map.insert("type".to_string(), LoraValue::String(rel.rel_type.clone()));
1912    map.insert(
1913        "properties".to_string(),
1914        properties_to_value_map(&rel.properties),
1915    );
1916    LoraValue::Map(map)
1917}
1918
1919/// Flatten label groups into a simple Vec<String> (for CREATE/MERGE where
1920/// disjunction doesn't apply — all labels are created).
1921pub(super) fn flatten_label_groups(groups: &[Vec<String>]) -> Vec<String> {
1922    groups.iter().flat_map(|g| g.iter().cloned()).collect()
1923}
1924
1925#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
1926pub(crate) enum GroupValueKey {
1927    Null,
1928    Bool(bool),
1929    Int(i64),
1930    Float(String),
1931    String(String),
1932    Binary(Vec<Vec<u8>>),
1933    List(Vec<GroupValueKey>),
1934    Map(Vec<(String, GroupValueKey)>),
1935    Node(u64),
1936    Relationship(u64),
1937}
1938
1939impl GroupValueKey {
1940    pub(crate) fn from_value(v: &LoraValue) -> Self {
1941        match v {
1942            LoraValue::Null => Self::Null,
1943            LoraValue::Bool(x) => Self::Bool(*x),
1944            LoraValue::Int(x) => Self::Int(*x),
1945            LoraValue::Float(x) => Self::Float(x.to_string()),
1946            LoraValue::String(x) => Self::String(x.clone()),
1947            LoraValue::Binary(x) => Self::Binary(x.segments().to_vec()),
1948            LoraValue::List(xs) => Self::List(xs.iter().map(Self::from_value).collect()),
1949            LoraValue::Map(m) => Self::Map(
1950                m.iter()
1951                    .map(|(k, v)| (k.clone(), Self::from_value(v)))
1952                    .collect(),
1953            ),
1954            LoraValue::Node(id) => Self::Node(*id),
1955            LoraValue::Relationship(id) => Self::Relationship(*id),
1956            LoraValue::Path(_) => Self::Null,
1957            // Temporal types: use their string representation as group key
1958            LoraValue::Date(d) => Self::String(d.to_string()),
1959            LoraValue::DateTime(dt) => Self::String(dt.to_string()),
1960            LoraValue::LocalDateTime(dt) => Self::String(dt.to_string()),
1961            LoraValue::Time(t) => Self::String(t.to_string()),
1962            LoraValue::LocalTime(t) => Self::String(t.to_string()),
1963            LoraValue::Duration(dur) => Self::String(dur.to_string()),
1964            LoraValue::Point(p) => Self::String(p.to_string()),
1965            LoraValue::Vector(v) => Self::String(format!("vector:{}", v.to_key_string())),
1966        }
1967    }
1968}
1969
1970/// Compute effective (min_hops, max_hops) from a `RangeLiteral`.
1971///
1972/// Lora semantics:
1973/// - `*`       → 1..∞   (start=None, end=None)
1974/// - `*2..5`   → 2..5   (start=Some(2), end=Some(5))
1975/// - `*..3`    → 1..3   (start=None, end=Some(3))
1976/// - `*2..`    → 2..∞   (start=Some(2), end=None)
1977/// - `*3`      → 3..3   (start=Some(3), end=None, no dots → exactly 3)
1978/// - `*0..1`   → 0..1
1979///
1980/// For unbounded upper, we cap at `MAX_VAR_LEN_HOPS` to prevent runaway.
1981const MAX_VAR_LEN_HOPS: u64 = 100;
1982
1983pub(crate) fn resolve_range(range: &RangeLiteral) -> (u64, u64) {
1984    let min_hops = range.start.unwrap_or(1);
1985    let max_hops = range.end.unwrap_or(MAX_VAR_LEN_HOPS);
1986    (min_hops, max_hops)
1987}
1988
1989/// An entry produced during BFS variable-length expansion.
1990pub(crate) struct VarLenResult {
1991    /// The destination node at the end of this path.
1992    pub(crate) dst_node_id: NodeId,
1993    /// The relationship IDs traversed (in order).
1994    pub(crate) rel_ids: Vec<u64>,
1995}
1996
1997/// Perform variable-length expansion from `start_node_id` following
1998/// relationships of the given `types` and `direction`, collecting all
1999/// reachable nodes at hop distances in `[min_hops, max_hops]`.
2000///
2001/// Uses BFS with relationship-uniqueness per path (each path does not
2002/// reuse the same relationship, but may revisit nodes).
2003pub(crate) fn variable_length_expand<S: GraphStorage>(
2004    storage: &S,
2005    start_node_id: NodeId,
2006    direction: Direction,
2007    types: &[String],
2008    min_hops: u64,
2009    max_hops: u64,
2010    bind_relationships: bool,
2011) -> Vec<VarLenResult> {
2012    let mut results = Vec::new();
2013
2014    // Each frontier entry: (current_node_id, relationships_used_so_far)
2015    let mut frontier: Vec<(NodeId, Vec<u64>)> = vec![(start_node_id, Vec::new())];
2016
2017    for depth in 1..=max_hops {
2018        // On the final hop we don't need to build next_frontier at all; every
2019        // path gets recorded and then the loop terminates. Avoids one full
2020        // pass of Vec clones on deep traversals.
2021        let is_last_hop = depth == max_hops;
2022        let mut next_frontier: Vec<(NodeId, Vec<u64>)> = Vec::new();
2023
2024        for (current_node, rels_used) in &frontier {
2025            // ID-only expand avoids cloning full records/properties for every
2026            // neighbour on every hop.
2027            for (rel_id, neighbor_id) in storage.expand_ids(*current_node, direction, types) {
2028                // Relationship-uniqueness: skip if this relationship was already
2029                // traversed on this particular path.
2030                if rels_used.contains(&rel_id) {
2031                    continue;
2032                }
2033
2034                if is_last_hop {
2035                    // Terminal hop: just record the result. Allocate rel_ids
2036                    // once (no duplicate clone) by extending a fresh copy.
2037                    if depth >= min_hops {
2038                        let mut rel_ids = Vec::with_capacity(rels_used.len() + 1);
2039                        rel_ids.extend_from_slice(rels_used);
2040                        rel_ids.push(rel_id);
2041                        results.push(VarLenResult {
2042                            dst_node_id: neighbor_id,
2043                            rel_ids: if bind_relationships {
2044                                rel_ids
2045                            } else {
2046                                Vec::new()
2047                            },
2048                        });
2049                    }
2050                    continue;
2051                }
2052
2053                let mut new_rels = Vec::with_capacity(rels_used.len() + 1);
2054                new_rels.extend_from_slice(rels_used);
2055                new_rels.push(rel_id);
2056
2057                if depth >= min_hops {
2058                    results.push(VarLenResult {
2059                        dst_node_id: neighbor_id,
2060                        rel_ids: if bind_relationships {
2061                            new_rels.clone()
2062                        } else {
2063                            Vec::new()
2064                        },
2065                    });
2066                }
2067
2068                next_frontier.push((neighbor_id, new_rels));
2069            }
2070        }
2071
2072        if is_last_hop || next_frontier.is_empty() {
2073            break;
2074        }
2075
2076        frontier = next_frontier;
2077    }
2078
2079    // Handle min_hops == 0: include the start node itself at depth 0.
2080    if min_hops == 0 {
2081        results.insert(
2082            0,
2083            VarLenResult {
2084                dst_node_id: start_node_id,
2085                rel_ids: Vec::new(),
2086            },
2087        );
2088    }
2089
2090    results
2091}
2092
2093/// Filter rows to keep only shortest paths.
2094/// `all` = false → keep one shortest path; `all` = true → keep all shortest.
2095pub(crate) fn filter_shortest_paths(rows: Vec<Row>, path_var: VarId, all: bool) -> Vec<Row> {
2096    if rows.is_empty() {
2097        return rows;
2098    }
2099
2100    // Compute path length for each row
2101    let lengths: Vec<usize> = rows
2102        .iter()
2103        .map(|row| match row.get(path_var) {
2104            Some(LoraValue::Path(p)) => p.rels.len(),
2105            _ => usize::MAX,
2106        })
2107        .collect();
2108
2109    let min_len = lengths.iter().copied().min().unwrap_or(usize::MAX);
2110
2111    let mut result: Vec<Row> = rows
2112        .into_iter()
2113        .zip(lengths.iter())
2114        .filter(|(_, len)| **len == min_len)
2115        .map(|(row, _)| row)
2116        .collect();
2117
2118    if !all && result.len() > 1 {
2119        result.truncate(1);
2120    }
2121
2122    result
2123}
2124
2125// ---------- Rel scans ----------
2126//
2127// Mirror of the `node_by_*_scan_rows` helpers above for the
2128// relationship-targeted index operators. Each helper resolves the
2129// candidate set via the corresponding `relationship_*_candidates`
2130// trait method (falling back to `rel_ids_by_type` when no scope is
2131// active), refilters with the precise predicate, and emits one row
2132// per matching relationship per input row.
2133//
2134// Direction handling: the optimizer only emits a Rel*Scan when both
2135// endpoints of the pattern have no upstream constraint. With a
2136// directed pattern (Direction::Right), each rel is bound as
2137// (src=stored_src, rel, dst=stored_dst). With Direction::Left those
2138// are swapped. With Direction::Undirected we emit both orientations
2139// to preserve parity with the Expand-based plan it replaces.
2140
2141fn emit_rel_rows(
2142    direction: Direction,
2143    src_var: VarId,
2144    rel_var: VarId,
2145    dst_var: VarId,
2146    rel: &lora_store::RelationshipRecord,
2147    base: &Row,
2148    out: &mut Vec<Row>,
2149) -> ExecResult<()> {
2150    match direction {
2151        Direction::Right => {
2152            emit_one_rel_row(
2153                src_var, rel_var, dst_var, rel.src, rel.id, rel.dst, base, out,
2154            )?;
2155        }
2156        Direction::Left => {
2157            emit_one_rel_row(
2158                src_var, rel_var, dst_var, rel.dst, rel.id, rel.src, base, out,
2159            )?;
2160        }
2161        Direction::Undirected => {
2162            emit_one_rel_row(
2163                src_var, rel_var, dst_var, rel.src, rel.id, rel.dst, base, out,
2164            )?;
2165            // Self-loops produce a single row even under undirected
2166            // semantics — emitting two copies of (a=a, b=a) for
2167            // a=stored_src=stored_dst would double-count.
2168            if rel.src != rel.dst {
2169                emit_one_rel_row(
2170                    src_var, rel_var, dst_var, rel.dst, rel.id, rel.src, base, out,
2171                )?;
2172            }
2173        }
2174    }
2175    Ok(())
2176}
2177
2178#[allow(clippy::too_many_arguments)]
2179fn emit_one_rel_row(
2180    src_var: VarId,
2181    rel_var: VarId,
2182    dst_var: VarId,
2183    src_id: NodeId,
2184    rel_id: RelationshipId,
2185    dst_id: NodeId,
2186    base: &Row,
2187    out: &mut Vec<Row>,
2188) -> ExecResult<()> {
2189    let mut row = base.clone();
2190    if bind_node_value(&mut row, src_var, src_id)?
2191        && bind_relationship_value(&mut row, rel_var, rel_id)?
2192        && bind_node_value(&mut row, dst_var, dst_id)?
2193    {
2194        out.push(row);
2195    }
2196    Ok(())
2197}
2198
2199fn bind_node_value(row: &mut Row, var: VarId, id: NodeId) -> ExecResult<bool> {
2200    match row.get(var) {
2201        Some(LoraValue::Node(existing)) => Ok(*existing == id),
2202        Some(other) => Err(ExecutorError::ExpectedNodeForExpand {
2203            var: format!("{var:?}"),
2204            found: value_kind(other),
2205        }),
2206        None => {
2207            row.insert(var, LoraValue::Node(id));
2208            Ok(true)
2209        }
2210    }
2211}
2212
2213fn bind_relationship_value(row: &mut Row, var: VarId, id: RelationshipId) -> ExecResult<bool> {
2214    match row.get(var) {
2215        Some(LoraValue::Relationship(existing)) => Ok(*existing == id),
2216        Some(other) => Err(ExecutorError::ExpectedRelationshipForExpand {
2217            var: format!("{var:?}"),
2218            found: value_kind(other),
2219        }),
2220        None => {
2221            row.insert(var, LoraValue::Relationship(id));
2222            Ok(true)
2223        }
2224    }
2225}
2226
2227fn rel_candidate_ids<S, F>(storage: &S, types: &[String], indexed: F) -> Vec<RelationshipId>
2228where
2229    S: GraphStorage,
2230    F: Fn(&str) -> Option<Vec<RelationshipId>>,
2231{
2232    if types.is_empty() {
2233        // No type constraint → no rel-typed scope to probe; fall back
2234        // to scanning every relationship.
2235        return storage.all_rel_ids();
2236    }
2237    let mut all = Vec::new();
2238    let mut seen = BTreeSet::new();
2239    for ty in types {
2240        let ids = indexed(ty).unwrap_or_else(|| storage.rel_ids_by_type(ty));
2241        for id in ids {
2242            if seen.insert(id) {
2243                all.push(id);
2244            }
2245        }
2246    }
2247    all
2248}
2249
2250pub(crate) fn rel_by_property_range_scan_rows<S: GraphStorage>(
2251    storage: &S,
2252    params: &BTreeMap<String, LoraValue>,
2253    base_rows: Vec<Row>,
2254    op: &lora_compiler::RelByPropertyRangeScanExec,
2255    deadline: Option<Instant>,
2256) -> ExecResult<Vec<Row>> {
2257    let eval_ctx = EvalContext { storage, params };
2258    let mut out = Vec::new();
2259
2260    for row in base_rows {
2261        check_optional_deadline(deadline)?;
2262        let lo_value = op.lo.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
2263        let hi_value = op.hi.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
2264        let lo_prop = lo_value
2265            .clone()
2266            .and_then(|v| lora_value_to_property(v).ok());
2267        let hi_prop = hi_value
2268            .clone()
2269            .and_then(|v| lora_value_to_property(v).ok());
2270
2271        let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| {
2272            storage.relationship_range_candidates(ty, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
2273        });
2274
2275        for rel_id in candidate_ids {
2276            check_optional_deadline(deadline)?;
2277            if let Some(result) = storage.with_relationship(rel_id, |rel| {
2278                if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
2279                    return Ok(());
2280                }
2281                let Some(actual) = rel.properties.get(op.key.as_str()) else {
2282                    return Ok(());
2283                };
2284                let actual_lv = LoraValue::from(actual);
2285                if !range_predicate_holds(
2286                    &actual_lv,
2287                    lo_value.as_ref(),
2288                    op.lo_inclusive,
2289                    hi_value.as_ref(),
2290                    op.hi_inclusive,
2291                ) {
2292                    return Ok(());
2293                }
2294                emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2295            }) {
2296                result?;
2297            }
2298        }
2299    }
2300
2301    Ok(out)
2302}
2303
2304pub(crate) fn rel_by_text_scan_rows<S: GraphStorage>(
2305    storage: &S,
2306    params: &BTreeMap<String, LoraValue>,
2307    base_rows: Vec<Row>,
2308    op: &lora_compiler::RelByTextScanExec,
2309    deadline: Option<Instant>,
2310) -> ExecResult<Vec<Row>> {
2311    let eval_ctx = EvalContext { storage, params };
2312    let mut out = Vec::new();
2313
2314    for row in base_rows {
2315        check_optional_deadline(deadline)?;
2316        let query = eval_expr(&op.query, &row, &eval_ctx);
2317        let LoraValue::String(query_str) = &query else {
2318            continue;
2319        };
2320
2321        let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| {
2322            storage.relationship_text_candidates(ty, &op.key, query_str)
2323        });
2324
2325        for rel_id in candidate_ids {
2326            check_optional_deadline(deadline)?;
2327            if let Some(result) = storage.with_relationship(rel_id, |rel| {
2328                if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
2329                    return Ok(());
2330                }
2331                let Some(PropertyValue::String(actual)) = rel.properties.get(op.key.as_str())
2332                else {
2333                    return Ok(());
2334                };
2335                if !text_predicate_holds(actual, op.predicate, query_str) {
2336                    return Ok(());
2337                }
2338                emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2339            }) {
2340                result?;
2341            }
2342        }
2343    }
2344
2345    Ok(out)
2346}
2347
2348pub(crate) fn rel_by_point_scan_rows<S: GraphStorage>(
2349    storage: &S,
2350    params: &BTreeMap<String, LoraValue>,
2351    base_rows: Vec<Row>,
2352    op: &lora_compiler::RelByPointScanExec,
2353    deadline: Option<Instant>,
2354) -> ExecResult<Vec<Row>> {
2355    let eval_ctx = EvalContext { storage, params };
2356    let mut out = Vec::new();
2357
2358    for row in base_rows {
2359        check_optional_deadline(deadline)?;
2360
2361        let probe = match &op.predicate {
2362            lora_compiler::PointPredicate::WithinBBox {
2363                lower_left,
2364                upper_right,
2365            } => {
2366                let ll = eval_expr(lower_left, &row, &eval_ctx);
2367                let ur = eval_expr(upper_right, &row, &eval_ctx);
2368                match (ll, ur) {
2369                    (LoraValue::Point(a), LoraValue::Point(b)) => {
2370                        Probe::WithinBBox { ll: a, ur: b }
2371                    }
2372                    _ => continue,
2373                }
2374            }
2375            lora_compiler::PointPredicate::WithinDistance {
2376                center,
2377                max_distance,
2378                inclusive,
2379            } => {
2380                let c = eval_expr(center, &row, &eval_ctx);
2381                let d = eval_expr(max_distance, &row, &eval_ctx);
2382                match (c, d) {
2383                    (LoraValue::Point(c), LoraValue::Float(d)) => Probe::WithinDistance {
2384                        center: c,
2385                        max: d,
2386                        inclusive: *inclusive,
2387                    },
2388                    (LoraValue::Point(c), LoraValue::Int(d)) => Probe::WithinDistance {
2389                        center: c,
2390                        max: d as f64,
2391                        inclusive: *inclusive,
2392                    },
2393                    _ => continue,
2394                }
2395            }
2396        };
2397
2398        let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| match &probe {
2399            Probe::WithinBBox { ll, ur } => {
2400                storage.relationship_point_within_bbox(ty, &op.key, (ll.x, ll.y), (ur.x, ur.y))
2401            }
2402            Probe::WithinDistance { center, max, .. } => {
2403                storage.relationship_point_within_distance(ty, &op.key, (center.x, center.y), *max)
2404            }
2405        });
2406
2407        for rel_id in candidate_ids {
2408            check_optional_deadline(deadline)?;
2409            if let Some(result) = storage.with_relationship(rel_id, |rel| {
2410                if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
2411                    return Ok(());
2412                }
2413                let Some(PropertyValue::Point(actual)) = rel.properties.get(op.key.as_str()) else {
2414                    return Ok(());
2415                };
2416                if !point_predicate_holds(actual, &probe) {
2417                    return Ok(());
2418                }
2419                emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2420            }) {
2421                result?;
2422            }
2423        }
2424    }
2425
2426    Ok(out)
2427}
2428
2429/// Node ids of an ordered range scan (`op.order` set), in `ORDER BY op.key`
2430/// order, each already checked against the exact bounds and labels.
2431///
2432/// The index serves the order when every bound is a string: a string
2433/// comparison only lets strings through, and among strings the index's
2434/// order is Cypher's order. Ids are then pulled from the index in chunks
2435/// that grow as rows are consumed, so a `LIMIT` above stops the scan
2436/// after a few small index reads. For any other bound the index order is
2437/// not Cypher's (it orders every integer before every float), so the
2438/// cursor collects the candidates and sorts them by value itself.
2439pub(crate) struct OrderedRangeCursor {
2440    lo: Option<LoraValue>,
2441    hi: Option<LoraValue>,
2442    mode: OrderedMode,
2443    pending: std::vec::IntoIter<NodeId>,
2444    labels: Vec<Vec<String>>,
2445    key: String,
2446    lo_inclusive: bool,
2447    hi_inclusive: bool,
2448}
2449
2450enum OrderedMode {
2451    Index {
2452        label: String,
2453        descending: bool,
2454        after: Option<(PropertyValue, NodeId)>,
2455        chunk: usize,
2456        exhausted: bool,
2457    },
2458    Buffered,
2459}
2460
2461impl OrderedRangeCursor {
2462    pub(crate) fn new<S: GraphStorage>(
2463        storage: &S,
2464        op: &lora_compiler::NodeByPropertyRangeScanExec,
2465        lo: Option<LoraValue>,
2466        hi: Option<LoraValue>,
2467    ) -> Self {
2468        let descending = matches!(op.order, Some(SortDirection::Desc));
2469        // Strings and temporals of one kind have a single contiguous,
2470        // correctly ordered run in the index, so the walk can stream from
2471        // it. Numbers cannot: integers and floats are keyed separately.
2472        let mut bounds = lo.iter().chain(hi.iter());
2473        let string_bounds = match bounds.next() {
2474            Some(LoraValue::String(_)) => bounds.all(|v| matches!(v, LoraValue::String(_))),
2475            Some(
2476                first @ (LoraValue::Date(_)
2477                | LoraValue::DateTime(_)
2478                | LoraValue::LocalDateTime(_)
2479                | LoraValue::Time(_)
2480                | LoraValue::LocalTime(_)),
2481            ) => bounds.all(|v| std::mem::discriminant(v) == std::mem::discriminant(first)),
2482            _ => false,
2483        };
2484        let index_label = single_label_hint(&op.labels).filter(|label| {
2485            string_bounds
2486                && storage
2487                    .node_range_ordered_chunk(label, &op.key, None, None, descending, None, 0)
2488                    .is_some()
2489        });
2490
2491        let mut cursor = Self {
2492            lo,
2493            hi,
2494            mode: OrderedMode::Buffered,
2495            pending: Vec::new().into_iter(),
2496            labels: op.labels.clone(),
2497            key: op.key.clone(),
2498            lo_inclusive: op.lo_inclusive,
2499            hi_inclusive: op.hi_inclusive,
2500        };
2501        match index_label {
2502            Some(label) => {
2503                cursor.mode = OrderedMode::Index {
2504                    label: label.to_string(),
2505                    descending,
2506                    after: None,
2507                    chunk: 64,
2508                    exhausted: false,
2509                }
2510            }
2511            None => {
2512                cursor.pending = cursor
2513                    .sorted_candidates(storage, op, descending)
2514                    .into_iter()
2515            }
2516        }
2517        cursor
2518    }
2519
2520    fn filter(&self) -> NodeRangeFilter<'_> {
2521        NodeRangeFilter {
2522            labels: &self.labels,
2523            key: &self.key,
2524            lo: self.lo.as_ref(),
2525            lo_inclusive: self.lo_inclusive,
2526            hi: self.hi.as_ref(),
2527            hi_inclusive: self.hi_inclusive,
2528        }
2529    }
2530
2531    /// Fallback: every matching id, sorted by the property value in Cypher
2532    /// order (ties by id).
2533    fn sorted_candidates<S: GraphStorage>(
2534        &self,
2535        storage: &S,
2536        op: &lora_compiler::NodeByPropertyRangeScanExec,
2537        descending: bool,
2538    ) -> Vec<NodeId> {
2539        let lo_prop = self.lo.clone().and_then(|v| lora_value_to_property(v).ok());
2540        let hi_prop = self.hi.clone().and_then(|v| lora_value_to_property(v).ok());
2541        let candidates = match single_label_hint(&op.labels) {
2542            Some(label) => storage
2543                .node_range_candidates(label, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
2544                .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
2545            None => scan_node_ids_for_label_groups(storage, &op.labels),
2546        };
2547        let filter = self.filter();
2548        let mut keyed: Vec<(LoraValue, NodeId)> = candidates
2549            .into_iter()
2550            .filter(|&id| node_matches_range_filter(storage, id, &filter))
2551            .map(|id| {
2552                let value = storage
2553                    .with_node(id, |n| n.properties.get(op.key.as_str()).cloned())
2554                    .flatten()
2555                    .map(LoraValue::from)
2556                    .unwrap_or(LoraValue::Null);
2557                (value, id)
2558            })
2559            .collect();
2560        keyed.sort_by(|(a, ai), (b, bi)| {
2561            let ord = compare_values_total(a, b).then(ai.cmp(bi));
2562            if descending {
2563                ord.reverse()
2564            } else {
2565                ord
2566            }
2567        });
2568        keyed.into_iter().map(|(_, id)| id).collect()
2569    }
2570
2571    pub(crate) fn next_id<S: GraphStorage>(&mut self, storage: &S) -> Option<NodeId> {
2572        loop {
2573            if let Some(id) = self.pending.next() {
2574                if matches!(self.mode, OrderedMode::Buffered)
2575                    || node_matches_range_filter(storage, id, &self.filter())
2576                {
2577                    return Some(id);
2578                }
2579                continue;
2580            }
2581            let lo_prop = self.lo.clone().and_then(|v| lora_value_to_property(v).ok());
2582            let hi_prop = self.hi.clone().and_then(|v| lora_value_to_property(v).ok());
2583            let OrderedMode::Index {
2584                label,
2585                descending,
2586                after,
2587                chunk,
2588                exhausted,
2589            } = &mut self.mode
2590            else {
2591                return None;
2592            };
2593            if *exhausted {
2594                return None;
2595            }
2596            let ids = storage.node_range_ordered_chunk(
2597                label,
2598                &self.key,
2599                lo_prop.as_ref(),
2600                hi_prop.as_ref(),
2601                *descending,
2602                after.as_ref().map(|(v, id)| (v, *id)),
2603                *chunk,
2604            )?;
2605            if ids.len() < *chunk {
2606                *exhausted = true;
2607            }
2608            let last = *ids.last()?;
2609            let last_value = storage
2610                .with_node(last, |n| n.properties.get(self.key.as_str()).cloned())
2611                .flatten()?;
2612            *after = Some((last_value, last));
2613            *chunk = (*chunk * 2).min(4096);
2614            self.pending = ids.into_iter();
2615        }
2616    }
2617}