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