Skip to main content

lora_executor/executor/
helpers.rs

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