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