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