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