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                    // openCypher: collect() ignores null values. Only a
971                    // top-level null is dropped; a list or map that
972                    // merely contains nulls is kept as-is.
973                    values.retain(|v| !matches!(v, LoraValue::Null));
974
975                    if *distinct {
976                        values = dedup_values(values);
977                    }
978
979                    Ok(LoraValue::List(values))
980                }
981
982                Some(kind @ (AggregateFunction::Sum | AggregateFunction::Avg)) => {
983                    if args.is_empty() {
984                        return Ok(LoraValue::Null);
985                    }
986
987                    let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
988
989                    if *distinct {
990                        values = dedup_values(values);
991                    }
992
993                    // One implementation with the streaming fold.
994                    let stream_kind = if kind == AggregateFunction::Sum {
995                        crate::pull::StreamableAggKind::Sum
996                    } else {
997                        crate::pull::StreamableAggKind::Avg
998                    };
999                    let mut state = crate::pull::AggState::seed(stream_kind);
1000                    for value in values {
1001                        state.fold(stream_kind, value);
1002                    }
1003                    state.finalize(stream_kind)
1004                }
1005
1006                Some(AggregateFunction::Min) => {
1007                    if args.is_empty() {
1008                        return Ok(LoraValue::Null);
1009                    }
1010
1011                    let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
1012                    values.retain(|v| !matches!(v, LoraValue::Null));
1013
1014                    if *distinct {
1015                        values = dedup_values(values);
1016                    }
1017
1018                    Ok(values
1019                        .into_iter()
1020                        .min_by(compare_values_total)
1021                        .unwrap_or(LoraValue::Null))
1022                }
1023
1024                Some(AggregateFunction::Max) => {
1025                    if args.is_empty() {
1026                        return Ok(LoraValue::Null);
1027                    }
1028
1029                    let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
1030                    values.retain(|v| !matches!(v, LoraValue::Null));
1031
1032                    if *distinct {
1033                        values = dedup_values(values);
1034                    }
1035
1036                    Ok(values
1037                        .into_iter()
1038                        .max_by(compare_values_total)
1039                        .unwrap_or(LoraValue::Null))
1040                }
1041
1042                Some(AggregateFunction::Stdev | AggregateFunction::Stdevp) => {
1043                    if args.is_empty() {
1044                        return Ok(LoraValue::Null);
1045                    }
1046
1047                    let nums: Vec<f64> = aggregate_numbers(&args[0], *distinct, rows, eval_ctx)?
1048                        .into_iter()
1049                        .map(|(n, _)| n)
1050                        .collect();
1051
1052                    let is_population = matches!(func, Some(AggregateFunction::Stdevp));
1053
1054                    if nums.is_empty() || (!is_population && nums.len() < 2) {
1055                        return Ok(LoraValue::Float(0.0));
1056                    }
1057
1058                    let mean = nums.iter().sum::<f64>() / nums.len() as f64;
1059                    let variance_sum: f64 = nums.iter().map(|x| (x - mean).powi(2)).sum();
1060                    let denom = if is_population {
1061                        nums.len() as f64
1062                    } else {
1063                        (nums.len() - 1) as f64
1064                    };
1065                    Ok(LoraValue::Float((variance_sum / denom).sqrt()))
1066                }
1067
1068                Some(AggregateFunction::PercentileCont) => {
1069                    if args.len() < 2 {
1070                        return Ok(LoraValue::Null);
1071                    }
1072
1073                    let Some(first) = rows.first() else {
1074                        return Ok(LoraValue::Null);
1075                    };
1076
1077                    let percentile = eval_expr_result(&args[1], first, eval_ctx)
1078                        .map_err(ExecutorError::from_eval)?
1079                        .as_f64()
1080                        .map(normalize_percentile)
1081                        .unwrap_or(0.5);
1082                    let mut nums: Vec<f64> =
1083                        aggregate_numbers(&args[0], *distinct, rows, eval_ctx)?
1084                            .into_iter()
1085                            .map(|(n, _)| n)
1086                            .collect();
1087
1088                    if nums.is_empty() {
1089                        return Ok(LoraValue::Null);
1090                    }
1091
1092                    nums.sort_by(|a, b| a.partial_cmp(b).unwrap_or(Ordering::Equal));
1093
1094                    let index = percentile * (nums.len() - 1) as f64;
1095                    let lower = index.floor() as usize;
1096                    let upper = index.ceil() as usize;
1097                    let fraction = index - lower as f64;
1098
1099                    if lower == upper || upper >= nums.len() {
1100                        Ok(LoraValue::Float(nums[lower]))
1101                    } else {
1102                        Ok(LoraValue::Float(
1103                            nums[lower] * (1.0 - fraction) + nums[upper] * fraction,
1104                        ))
1105                    }
1106                }
1107
1108                Some(AggregateFunction::PercentileDisc) => {
1109                    if args.len() < 2 {
1110                        return Ok(LoraValue::Null);
1111                    }
1112
1113                    let Some(first) = rows.first() else {
1114                        return Ok(LoraValue::Null);
1115                    };
1116
1117                    let percentile = eval_expr_result(&args[1], first, eval_ctx)
1118                        .map_err(ExecutorError::from_eval)?
1119                        .as_f64()
1120                        .map(normalize_percentile)
1121                        .unwrap_or(0.5);
1122                    let mut nums = aggregate_numbers(&args[0], *distinct, rows, eval_ctx)?;
1123
1124                    if nums.is_empty() {
1125                        return Ok(LoraValue::Null);
1126                    }
1127
1128                    nums.sort_by(|a, b| a.0.partial_cmp(&b.0).unwrap_or(Ordering::Equal));
1129
1130                    // The value itself, as it was: an integer stays an integer.
1131                    let index = (percentile * (nums.len() - 1) as f64).round() as usize;
1132                    let index = index.min(nums.len() - 1);
1133                    Ok(nums.swap_remove(index).1)
1134                }
1135
1136                _ => eval_first_or_null(expr, rows, eval_ctx),
1137            }
1138        }
1139
1140        _ => eval_first_or_null(expr, rows, eval_ctx),
1141    }
1142}
1143
1144/// The numeric values of an aggregate's argument over `rows`, each with
1145/// its numeric reading, deduplicated for DISTINCT; others are ignored.
1146fn aggregate_numbers<S: GraphStorage>(
1147    expr: &ResolvedExpr,
1148    distinct: bool,
1149    rows: &[Row],
1150    eval_ctx: &EvalContext<'_, S>,
1151) -> ExecResult<Vec<(f64, LoraValue)>> {
1152    let mut values = eval_aggregate_arg_values(expr, rows, eval_ctx)?;
1153    if distinct {
1154        values = dedup_values(values);
1155    }
1156    Ok(values
1157        .into_iter()
1158        .filter_map(|v| match v {
1159            LoraValue::Int(i) => Some((i as f64, v)),
1160            LoraValue::Float(f) => Some((f, v)),
1161            _ => None,
1162        })
1163        .collect())
1164}
1165
1166fn eval_aggregate_arg_values<S: GraphStorage>(
1167    expr: &ResolvedExpr,
1168    rows: &[Row],
1169    eval_ctx: &EvalContext<'_, S>,
1170) -> ExecResult<Vec<LoraValue>> {
1171    rows.iter()
1172        .map(|row| eval_expr_result(expr, row, eval_ctx).map_err(ExecutorError::from_eval))
1173        .collect()
1174}
1175
1176fn normalize_percentile(value: f64) -> f64 {
1177    if value.is_finite() {
1178        value.clamp(0.0, 1.0)
1179    } else {
1180        0.5
1181    }
1182}
1183
1184fn eval_first_or_null<S: GraphStorage>(
1185    expr: &ResolvedExpr,
1186    rows: &[Row],
1187    eval_ctx: &EvalContext<'_, S>,
1188) -> ExecResult<LoraValue> {
1189    match rows.first() {
1190        Some(row) => eval_expr_result(expr, row, eval_ctx).map_err(ExecutorError::from_eval),
1191        None => Ok(LoraValue::Null),
1192    }
1193}
1194
1195fn dedup_values(values: Vec<LoraValue>) -> Vec<LoraValue> {
1196    let mut seen: BTreeSet<GroupValueKey> = BTreeSet::new();
1197    let mut out = Vec::new();
1198
1199    for value in values {
1200        let key = GroupValueKey::from_value(&value);
1201        if seen.insert(key) {
1202            out.push(value);
1203        }
1204    }
1205
1206    out
1207}
1208
1209pub(crate) fn compare_values_total(a: &LoraValue, b: &LoraValue) -> Ordering {
1210    use LoraValue::*;
1211
1212    match (a, b) {
1213        (Bool(x), Bool(y)) => x.cmp(y),
1214        (Int(x), Int(y)) => x.cmp(y),
1215        (Float(x), Float(y)) => x.partial_cmp(y).unwrap_or(Ordering::Equal),
1216        (Int(x), Float(y)) => (*x as f64).partial_cmp(y).unwrap_or(Ordering::Equal),
1217        (Float(x), Int(y)) => x.partial_cmp(&(*y as f64)).unwrap_or(Ordering::Equal),
1218        (String(x), String(y)) => x.cmp(y),
1219        (Binary(x), Binary(y)) => x.segments().cmp(y.segments()),
1220        (Node(x), Node(y)) => x.cmp(y),
1221        (Relationship(x), Relationship(y)) => x.cmp(y),
1222        (Date(x), Date(y)) => x.cmp(y),
1223        (DateTime(x), DateTime(y)) => x.cmp(y),
1224        (LocalDateTime(x), LocalDateTime(y)) => x.cmp(y),
1225        (Time(x), Time(y)) => x.cmp(y),
1226        (LocalTime(x), LocalTime(y)) => x.cmp(y),
1227        (Duration(x), Duration(y)) => x.cmp(y),
1228        (Vector(x), Vector(y)) => x.to_key_string().cmp(&y.to_key_string()),
1229        _ => type_rank(a)
1230            .cmp(&type_rank(b))
1231            .then_with(|| format!("{a:?}").cmp(&format!("{b:?}"))),
1232    }
1233}
1234
1235pub fn value_matches_property_value(expected: &LoraValue, actual: &PropertyValue) -> bool {
1236    match (expected, actual) {
1237        (LoraValue::Null, PropertyValue::Null) => true,
1238        (LoraValue::Bool(a), PropertyValue::Bool(b)) => a == b,
1239        (LoraValue::Int(a), PropertyValue::Int(b)) => a == b,
1240        (LoraValue::Float(a), PropertyValue::Float(b)) => a == b,
1241        (LoraValue::Int(a), PropertyValue::Float(b)) => (*a as f64) == *b,
1242        (LoraValue::Float(a), PropertyValue::Int(b)) => *a == (*b as f64),
1243        (LoraValue::String(a), PropertyValue::String(b)) => a == b,
1244        (LoraValue::Binary(a), PropertyValue::Binary(b)) => a == b,
1245
1246        (LoraValue::List(xs), PropertyValue::List(ys)) => {
1247            xs.len() == ys.len()
1248                && xs
1249                    .iter()
1250                    .zip(ys.iter())
1251                    .all(|(x, y)| value_matches_property_value(x, y))
1252        }
1253
1254        (LoraValue::Map(xm), PropertyValue::Map(ym)) => xm.iter().all(|(k, xv)| {
1255            ym.get(k)
1256                .map(|yv| value_matches_property_value(xv, yv))
1257                .unwrap_or(false)
1258        }),
1259
1260        (LoraValue::Date(a), PropertyValue::Date(b)) => a == b,
1261        (LoraValue::DateTime(a), PropertyValue::DateTime(b)) => a == b,
1262        (LoraValue::LocalDateTime(a), PropertyValue::LocalDateTime(b)) => a == b,
1263        (LoraValue::Time(a), PropertyValue::Time(b)) => a == b,
1264        (LoraValue::LocalTime(a), PropertyValue::LocalTime(b)) => a == b,
1265        (LoraValue::Duration(a), PropertyValue::Duration(b)) => a == b,
1266        (LoraValue::Point(a), PropertyValue::Point(b)) => a == b,
1267        (LoraValue::Vector(a), PropertyValue::Vector(b)) => a == b,
1268
1269        _ => false,
1270    }
1271}
1272
1273pub(crate) fn node_matches_property_filter<S: GraphStorage>(
1274    storage: &S,
1275    node_id: NodeId,
1276    labels: &[Vec<String>],
1277    key: &str,
1278    expected: &LoraValue,
1279) -> bool {
1280    storage
1281        .with_node(node_id, |node| {
1282            node_matches_label_groups(&node.labels, labels)
1283                && node
1284                    .properties
1285                    .get(key)
1286                    .map(|actual| value_matches_property_value(expected, actual))
1287                    .unwrap_or(false)
1288        })
1289        .unwrap_or(false)
1290}
1291
1292fn single_label_hint(labels: &[Vec<String>]) -> Option<&str> {
1293    if labels.len() == 1 && labels[0].len() == 1 {
1294        Some(labels[0][0].as_str())
1295    } else {
1296        None
1297    }
1298}
1299
1300fn property_lookup_values(expected: &LoraValue) -> Option<Vec<PropertyValue>> {
1301    let property = lora_value_to_property(expected.clone()).ok()?;
1302    let mut values = vec![property.clone()];
1303
1304    match property {
1305        PropertyValue::Int(i) => {
1306            values.push(PropertyValue::Float(i as f64));
1307        }
1308        PropertyValue::Float(f)
1309            if f.is_finite()
1310                && f.fract() == 0.0
1311                && f >= i64::MIN as f64
1312                && f <= i64::MAX as f64 =>
1313        {
1314            values.push(PropertyValue::Int(f as i64));
1315        }
1316        _ => {}
1317    }
1318
1319    Some(values)
1320}
1321
1322pub(crate) struct NodePropertyCandidates {
1323    pub(crate) ids: Vec<NodeId>,
1324    pub(crate) prefiltered: bool,
1325}
1326
1327pub(crate) fn node_by_property_range_scan_rows<S: GraphStorage>(
1328    storage: &S,
1329    params: &BTreeMap<String, LoraValue>,
1330    base_rows: Vec<Row>,
1331    op: &lora_compiler::NodeByPropertyRangeScanExec,
1332    deadline: Option<Instant>,
1333) -> ExecResult<Vec<Row>> {
1334    let eval_ctx = EvalContext { storage, params };
1335    let mut out = Vec::new();
1336    let mut other_kinds = OtherKindScan::default();
1337
1338    for row in base_rows {
1339        check_optional_deadline(deadline)?;
1340        let lo_value = op.lo.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
1341        // A null upper bound reads as none: the Filter kept above every
1342        // index scan still judges each row (a STARTS WITH prefix with no
1343        // successor, `string.prefix_end` null, scans to the end).
1344        let hi_value = op
1345            .hi
1346            .as_ref()
1347            .map(|expr| eval_expr(expr, &row, &eval_ctx))
1348            .filter(|v| !matches!(v, LoraValue::Null));
1349        let lo_prop = lo_value
1350            .clone()
1351            .and_then(|v| lora_value_to_property(v).ok());
1352        let hi_prop = hi_value
1353            .clone()
1354            .and_then(|v| lora_value_to_property(v).ok());
1355        let filter = NodeRangeFilter {
1356            labels: &op.labels,
1357            key: &op.key,
1358            lo: lo_value.as_ref(),
1359            lo_inclusive: op.lo_inclusive,
1360            hi: hi_value.as_ref(),
1361            hi_inclusive: op.hi_inclusive,
1362        };
1363        let bounds = [lo_value.as_ref(), hi_value.as_ref()];
1364        let bound_id = bound_node_id_for_expand(&row, op.var)?;
1365
1366        if let Some(existing_id) = bound_id {
1367            if node_matches_range_filter(storage, existing_id, &filter)
1368                || bound_node_has_other_kind(storage, existing_id, op, bounds)
1369            {
1370                out.push(row);
1371            }
1372            continue;
1373        }
1374
1375        // Temporals of another kind first, for the Filter above to judge.
1376        for id in other_kind_node_ids(storage, op, bounds, &mut other_kinds) {
1377            let mut new_row = row.clone();
1378            new_row.insert(op.var, LoraValue::Node(id));
1379            out.push(new_row);
1380        }
1381
1382        if op.order.is_some() {
1383            let mut cursor =
1384                OrderedRangeCursor::new(storage, op, lo_value.clone(), hi_value.clone());
1385            while let Some(id) = cursor.next_id(storage) {
1386                check_optional_deadline(deadline)?;
1387                let mut new_row = row.clone();
1388                new_row.insert(op.var, LoraValue::Node(id));
1389                out.push(new_row);
1390            }
1391            continue;
1392        }
1393
1394        let candidate_ids = match single_label_hint(&op.labels) {
1395            Some(label) => storage
1396                .node_range_candidates(label, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
1397                .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1398            None => scan_node_ids_for_label_groups(storage, &op.labels),
1399        };
1400
1401        for id in candidate_ids {
1402            check_optional_deadline(deadline)?;
1403            if node_matches_range_filter(storage, id, &filter) {
1404                let mut new_row = row.clone();
1405                new_row.insert(op.var, LoraValue::Node(id));
1406                out.push(new_row);
1407            }
1408        }
1409    }
1410
1411    Ok(out)
1412}
1413
1414pub(crate) fn node_by_text_scan_rows<S: GraphStorage>(
1415    storage: &S,
1416    params: &BTreeMap<String, LoraValue>,
1417    base_rows: Vec<Row>,
1418    op: &lora_compiler::NodeByTextScanExec,
1419    deadline: Option<Instant>,
1420) -> ExecResult<Vec<Row>> {
1421    let eval_ctx = EvalContext { storage, params };
1422    let mut out = Vec::new();
1423
1424    for row in base_rows {
1425        check_optional_deadline(deadline)?;
1426        let query = eval_expr(&op.query, &row, &eval_ctx);
1427        let LoraValue::String(query_str) = &query else {
1428            // Non-string query → predicate cannot match anything;
1429            // skip the row entirely (matches scan + filter behaviour).
1430            continue;
1431        };
1432
1433        if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
1434            if node_matches_text_filter(
1435                storage,
1436                existing_id,
1437                &op.labels,
1438                &op.key,
1439                op.predicate,
1440                query_str,
1441            ) {
1442                out.push(row);
1443            }
1444            continue;
1445        }
1446
1447        let candidate_ids = match single_label_hint(&op.labels) {
1448            Some(label) => storage
1449                .node_text_candidates(label, &op.key, query_str)
1450                .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1451            None => scan_node_ids_for_label_groups(storage, &op.labels),
1452        };
1453
1454        for id in candidate_ids {
1455            check_optional_deadline(deadline)?;
1456            if node_matches_text_filter(storage, id, &op.labels, &op.key, op.predicate, query_str) {
1457                let mut new_row = row.clone();
1458                new_row.insert(op.var, LoraValue::Node(id));
1459                out.push(new_row);
1460            }
1461        }
1462    }
1463
1464    Ok(out)
1465}
1466
1467/// The temporal kinds of a range scan's bounds, at most two.
1468fn temporal_bound_kinds(bounds: [Option<&LoraValue>; 2]) -> Vec<(&'static str, &LoraValue)> {
1469    let mut kinds: Vec<(&'static str, &LoraValue)> = Vec::with_capacity(2);
1470    for bound in bounds.into_iter().flatten() {
1471        if let Some(kind) = temporal_kind_name(bound) {
1472            if !kinds.iter().any(|(k, _)| *k == kind) {
1473                kinds.push((kind, bound));
1474            }
1475        }
1476    }
1477    kinds
1478}
1479
1480/// Whether `value` is a temporal of a kind none of `kinds` is.
1481fn is_other_temporal_kind(value: &PropertyValue, kinds: &[(&'static str, &LoraValue)]) -> bool {
1482    temporal_kind_name(&LoraValue::from(value))
1483        .is_some_and(|kind| kinds.iter().any(|(k, _)| *k != kind))
1484}
1485
1486/// Per-call memo of a label or type scan, for when no index can list the
1487/// other-kind ids. Held for one scan call only: in the pull path writes
1488/// can land between rows, so it is never kept across them.
1489#[derive(Default)]
1490pub(crate) struct OtherKindScan<Id> {
1491    temporals: Option<Vec<(Id, &'static str)>>,
1492}
1493
1494impl<Id: Copy> OtherKindScan<Id> {
1495    fn ids(
1496        &mut self,
1497        kinds: &[(&'static str, &LoraValue)],
1498        scan: impl FnOnce() -> Vec<(Id, &'static str)>,
1499    ) -> Vec<Id> {
1500        self.temporals
1501            .get_or_insert_with(scan)
1502            .iter()
1503            .filter(|(_, kind)| kinds.iter().any(|(k, _)| k != kind))
1504            .map(|(id, _)| *id)
1505            .collect()
1506    }
1507}
1508
1509/// Nodes a range scan emits besides its matches: those whose `op.key`
1510/// holds a temporal of another kind than a temporal bound. The index keeps
1511/// each kind apart and would skip them without a word; emitted, they reach
1512/// the `Filter` above the scan, which compares them as an unindexed scan
1513/// does. So an unguarded `x >= date(…)` fails on a DATETIME `x`, while a
1514/// guard such as `type.of(x) = 'DATE' AND x >= date(…)` drops the row
1515/// before the comparison (`AND` short-circuits). Reads only the other
1516/// kinds' runs of the index (two probes plus the ids found); without an
1517/// index, the label once per call.
1518pub(crate) fn other_kind_node_ids<S: GraphStorage>(
1519    storage: &S,
1520    op: &lora_compiler::NodeByPropertyRangeScanExec,
1521    bounds: [Option<&LoraValue>; 2],
1522    fallback: &mut OtherKindScan<NodeId>,
1523) -> Vec<NodeId> {
1524    let kinds = temporal_bound_kinds(bounds);
1525    if kinds.is_empty() {
1526        return Vec::new();
1527    }
1528    // A label every matching node has: its index lists the candidates.
1529    let index_label = op
1530        .labels
1531        .iter()
1532        .find(|group| group.len() == 1)
1533        .map(|group| group[0].as_str());
1534    let mut ids = BTreeSet::new();
1535    let mut need_scan = index_label.is_none();
1536    if let Some(label) = index_label {
1537        for (_, bound) in &kinds {
1538            let listed = lora_value_to_property((*bound).clone())
1539                .ok()
1540                .and_then(|like| storage.node_range_other_temporal_kind_ids(label, &op.key, &like));
1541            match listed {
1542                Some(listed) => ids.extend(listed),
1543                None => need_scan = true,
1544            }
1545        }
1546        // Several bound kinds: a value of one bound's kind is "other" for
1547        // the other bound, which each probe already lists.
1548    }
1549    if need_scan {
1550        ids.extend(fallback.ids(&kinds, || {
1551            scan_node_ids_for_label_groups(storage, &op.labels)
1552                .into_iter()
1553                .filter_map(|id| {
1554                    storage
1555                        .with_node(id, |n| {
1556                            n.properties
1557                                .get(op.key.as_str())
1558                                .and_then(|v| temporal_kind_name(&LoraValue::from(v)))
1559                        })
1560                        .flatten()
1561                        .map(|kind| (id, kind))
1562                })
1563                .collect()
1564        }));
1565        return ids.into_iter().collect();
1566    }
1567    let single = op.labels.len() == 1 && op.labels[0].len() == 1;
1568    ids.into_iter()
1569        .filter(|id| {
1570            single
1571                || storage
1572                    .with_node(*id, |n| node_matches_label_groups(&n.labels, &op.labels))
1573                    .unwrap_or(false)
1574        })
1575        .collect()
1576}
1577
1578/// Whether the bound node `id` matches the scan's labels and holds a
1579/// temporal of another kind than a temporal bound; see
1580/// [`other_kind_node_ids`].
1581pub(crate) fn bound_node_has_other_kind<S: GraphStorage>(
1582    storage: &S,
1583    id: NodeId,
1584    op: &lora_compiler::NodeByPropertyRangeScanExec,
1585    bounds: [Option<&LoraValue>; 2],
1586) -> bool {
1587    let kinds = temporal_bound_kinds(bounds);
1588    !kinds.is_empty()
1589        && storage
1590            .with_node(id, |n| {
1591                node_matches_label_groups(&n.labels, &op.labels)
1592                    && n.properties
1593                        .get(op.key.as_str())
1594                        .is_some_and(|v| is_other_temporal_kind(v, &kinds))
1595            })
1596            .unwrap_or(false)
1597}
1598
1599/// Relationship counterpart of [`other_kind_node_ids`].
1600fn other_kind_rel_ids<S: GraphStorage>(
1601    storage: &S,
1602    op: &lora_compiler::RelByPropertyRangeScanExec,
1603    bounds: [Option<&LoraValue>; 2],
1604    fallback: &mut OtherKindScan<RelationshipId>,
1605) -> Vec<RelationshipId> {
1606    let kinds = temporal_bound_kinds(bounds);
1607    if kinds.is_empty() {
1608        return Vec::new();
1609    }
1610    let mut ids = BTreeSet::new();
1611    let mut need_scan = op.types.is_empty();
1612    for ty in &op.types {
1613        for (_, bound) in &kinds {
1614            let listed = lora_value_to_property((*bound).clone())
1615                .ok()
1616                .and_then(|like| {
1617                    storage.relationship_range_other_temporal_kind_ids(ty, &op.key, &like)
1618                });
1619            match listed {
1620                Some(listed) => ids.extend(listed),
1621                None => need_scan = true,
1622            }
1623        }
1624    }
1625    if need_scan {
1626        ids.extend(fallback.ids(&kinds, || {
1627            rel_candidate_ids(storage, &op.types, |_| None)
1628                .into_iter()
1629                .filter_map(|id| {
1630                    storage
1631                        .with_relationship(id, |r| {
1632                            r.properties
1633                                .get(op.key.as_str())
1634                                .and_then(|v| temporal_kind_name(&LoraValue::from(v)))
1635                        })
1636                        .flatten()
1637                        .map(|kind| (id, kind))
1638                })
1639                .collect()
1640        }));
1641    }
1642    ids.into_iter().collect()
1643}
1644
1645fn temporal_kind_name(value: &LoraValue) -> Option<&'static str> {
1646    Some(match value {
1647        LoraValue::Date(_) => "DATE",
1648        LoraValue::DateTime(_) => "DATETIME",
1649        LoraValue::LocalDateTime(_) => "LOCAL_DATETIME",
1650        LoraValue::Time(_) => "TIME",
1651        LoraValue::LocalTime(_) => "LOCAL_TIME",
1652        _ => return None,
1653    })
1654}
1655
1656pub(crate) struct NodeRangeFilter<'a> {
1657    labels: &'a [Vec<String>],
1658    key: &'a str,
1659    lo: Option<&'a LoraValue>,
1660    lo_inclusive: bool,
1661    hi: Option<&'a LoraValue>,
1662    hi_inclusive: bool,
1663}
1664
1665fn node_matches_range_filter<S: GraphStorage>(
1666    storage: &S,
1667    id: NodeId,
1668    filter: &NodeRangeFilter<'_>,
1669) -> bool {
1670    storage
1671        .with_node(id, |n| {
1672            if !node_matches_label_groups(&n.labels, filter.labels) {
1673                return false;
1674            }
1675            let Some(actual) = n.properties.get(filter.key) else {
1676                return false;
1677            };
1678            let actual_lv = lora_store_property_to_value(actual);
1679            range_predicate_holds(
1680                &actual_lv,
1681                filter.lo,
1682                filter.lo_inclusive,
1683                filter.hi,
1684                filter.hi_inclusive,
1685            )
1686        })
1687        .unwrap_or(false)
1688}
1689
1690fn node_matches_text_filter<S: GraphStorage>(
1691    storage: &S,
1692    id: NodeId,
1693    labels: &[Vec<String>],
1694    key: &str,
1695    predicate: lora_compiler::TextPredicate,
1696    query: &str,
1697) -> bool {
1698    storage
1699        .with_node(id, |n| {
1700            if !node_matches_label_groups(&n.labels, labels) {
1701                return false;
1702            }
1703            let Some(PropertyValue::String(actual)) = n.properties.get(key) else {
1704                return false;
1705            };
1706            text_predicate_holds(actual, predicate, query)
1707        })
1708        .unwrap_or(false)
1709}
1710
1711fn text_predicate_holds(
1712    actual: &str,
1713    predicate: lora_compiler::TextPredicate,
1714    query: &str,
1715) -> bool {
1716    match predicate {
1717        lora_compiler::TextPredicate::StartsWith => actual.starts_with(query),
1718        lora_compiler::TextPredicate::EndsWith => actual.ends_with(query),
1719        lora_compiler::TextPredicate::Contains => actual.contains(query),
1720    }
1721}
1722
1723fn range_predicate_holds(
1724    actual: &LoraValue,
1725    lo: Option<&LoraValue>,
1726    lo_inclusive: bool,
1727    hi: Option<&LoraValue>,
1728    hi_inclusive: bool,
1729) -> bool {
1730    if let Some(lo) = lo {
1731        match range_comparison(actual, lo) {
1732            None => return false,
1733            Some(Ordering::Less) => return false,
1734            Some(Ordering::Equal) if !lo_inclusive => return false,
1735            _ => {}
1736        }
1737    }
1738    if let Some(hi) = hi {
1739        match range_comparison(actual, hi) {
1740            None => return false,
1741            Some(Ordering::Greater) => return false,
1742            Some(Ordering::Equal) if !hi_inclusive => return false,
1743            _ => {}
1744        }
1745    }
1746    true
1747}
1748
1749fn range_comparison(actual: &LoraValue, bound: &LoraValue) -> Option<Ordering> {
1750    match (actual, bound) {
1751        (LoraValue::Null, _) | (_, LoraValue::Null) => None,
1752        (LoraValue::String(a), LoraValue::String(b)) => Some(a.cmp(b)),
1753        (
1754            LoraValue::Date(_)
1755            | LoraValue::DateTime(_)
1756            | LoraValue::LocalDateTime(_)
1757            | LoraValue::Time(_)
1758            | LoraValue::LocalTime(_),
1759            _,
1760        ) => actual.temporal_cmp(bound),
1761        (LoraValue::Duration(a), LoraValue::Duration(b)) => a
1762            .total_seconds_approx()
1763            .partial_cmp(&b.total_seconds_approx()),
1764        _ => actual.as_f64()?.partial_cmp(&bound.as_f64()?),
1765    }
1766}
1767
1768fn lora_store_property_to_value(value: &PropertyValue) -> LoraValue {
1769    LoraValue::from(value)
1770}
1771
1772pub(crate) fn node_by_point_scan_rows<S: GraphStorage>(
1773    storage: &S,
1774    params: &BTreeMap<String, LoraValue>,
1775    base_rows: Vec<Row>,
1776    op: &lora_compiler::NodeByPointScanExec,
1777    deadline: Option<Instant>,
1778) -> ExecResult<Vec<Row>> {
1779    let eval_ctx = EvalContext { storage, params };
1780    let mut out = Vec::new();
1781
1782    for row in base_rows {
1783        check_optional_deadline(deadline)?;
1784
1785        // Resolve the predicate's literal/scalar inputs against the
1786        // current row. These end up as captured `LoraValue`s used both
1787        // to probe the spatial index and to refilter every candidate.
1788        let probe = match &op.predicate {
1789            lora_compiler::PointPredicate::WithinBBox {
1790                lower_left,
1791                upper_right,
1792            } => {
1793                let ll = eval_expr(lower_left, &row, &eval_ctx);
1794                let ur = eval_expr(upper_right, &row, &eval_ctx);
1795                match (ll, ur) {
1796                    (LoraValue::Point(a), LoraValue::Point(b)) => {
1797                        Probe::WithinBBox { ll: a, ur: b }
1798                    }
1799                    _ => continue,
1800                }
1801            }
1802            lora_compiler::PointPredicate::WithinDistance {
1803                center,
1804                max_distance,
1805                inclusive,
1806            } => {
1807                let c = eval_expr(center, &row, &eval_ctx);
1808                let d = eval_expr(max_distance, &row, &eval_ctx);
1809                match (c, d) {
1810                    (LoraValue::Point(c), LoraValue::Float(d)) => Probe::WithinDistance {
1811                        center: c,
1812                        max: d,
1813                        inclusive: *inclusive,
1814                    },
1815                    (LoraValue::Point(c), LoraValue::Int(d)) => Probe::WithinDistance {
1816                        center: c,
1817                        max: d as f64,
1818                        inclusive: *inclusive,
1819                    },
1820                    _ => continue,
1821                }
1822            }
1823        };
1824
1825        if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
1826            if node_matches_point_filter(storage, existing_id, &op.labels, &op.key, &probe) {
1827                out.push(row);
1828            }
1829            continue;
1830        }
1831
1832        let candidate_ids = match single_label_hint(&op.labels) {
1833            Some(label) => match &probe {
1834                Probe::WithinBBox { ll, ur } => {
1835                    // Two seeks when the box crosses the antimeridian.
1836                    let (ranges, n) = lora_store::bbox_x_ranges(ll, ur);
1837                    let (lo_y, hi_y) = (ll.y.min(ur.y), ll.y.max(ur.y));
1838                    let mut ids = Vec::new();
1839                    let mut indexed = true;
1840                    for (lo_x, hi_x) in &ranges[..n] {
1841                        match storage.node_point_within_bbox(
1842                            label,
1843                            &op.key,
1844                            (*lo_x, lo_y),
1845                            (*hi_x, hi_y),
1846                        ) {
1847                            Some(found) => ids.extend(found),
1848                            None => indexed = false,
1849                        }
1850                    }
1851                    if !indexed {
1852                        scan_node_ids_for_label_groups(storage, &op.labels)
1853                    } else {
1854                        if n > 1 {
1855                            ids.sort_unstable();
1856                            ids.dedup();
1857                        }
1858                        ids
1859                    }
1860                }
1861                Probe::WithinDistance { center, max, .. } => storage
1862                    .node_point_within_distance(label, &op.key, (center.x, center.y), *max)
1863                    .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1864            },
1865            None => scan_node_ids_for_label_groups(storage, &op.labels),
1866        };
1867
1868        for id in candidate_ids {
1869            check_optional_deadline(deadline)?;
1870            if node_matches_point_filter(storage, id, &op.labels, &op.key, &probe) {
1871                let mut new_row = row.clone();
1872                new_row.insert(op.var, LoraValue::Node(id));
1873                out.push(new_row);
1874            }
1875        }
1876    }
1877
1878    Ok(out)
1879}
1880
1881fn node_matches_point_filter<S: GraphStorage>(
1882    storage: &S,
1883    id: NodeId,
1884    labels: &[Vec<String>],
1885    key: &str,
1886    probe: &Probe,
1887) -> bool {
1888    storage
1889        .with_node(id, |n| {
1890            if !node_matches_label_groups(&n.labels, labels) {
1891                return false;
1892            }
1893            let Some(PropertyValue::Point(point)) = n.properties.get(key) else {
1894                return false;
1895            };
1896            point_predicate_holds(point, probe)
1897        })
1898        .unwrap_or(false)
1899}
1900
1901enum Probe {
1902    WithinBBox {
1903        ll: lora_store::LoraPoint,
1904        ur: lora_store::LoraPoint,
1905    },
1906    WithinDistance {
1907        center: lora_store::LoraPoint,
1908        max: f64,
1909        inclusive: bool,
1910    },
1911}
1912
1913fn point_predicate_holds(actual: &lora_store::LoraPoint, probe: &Probe) -> bool {
1914    match probe {
1915        Probe::WithinBBox { ll, ur } => lora_store::bbox_contains(actual, ll, ur).unwrap_or(false),
1916        Probe::WithinDistance {
1917            center,
1918            max,
1919            inclusive,
1920        } => {
1921            let Some(d) = lora_store::point_distance(actual, center) else {
1922                return false;
1923            };
1924            if *inclusive {
1925                d <= *max
1926            } else {
1927                d < *max
1928            }
1929        }
1930    }
1931}
1932
1933pub(crate) fn indexed_node_property_candidates<S: GraphStorage>(
1934    storage: &S,
1935    labels: &[Vec<String>],
1936    key: &str,
1937    expected: &LoraValue,
1938) -> NodePropertyCandidates {
1939    let Some(values) = property_lookup_values(expected) else {
1940        return NodePropertyCandidates {
1941            ids: scan_node_ids_for_label_groups(storage, labels),
1942            prefiltered: false,
1943        };
1944    };
1945
1946    let label_hint = single_label_hint(labels);
1947    let mut seen = BTreeSet::new();
1948    let mut out = Vec::new();
1949    for value in values {
1950        for id in storage.find_node_ids_by_property(label_hint, key, &value) {
1951            if seen.insert(id) {
1952                out.push(id);
1953            }
1954        }
1955    }
1956    NodePropertyCandidates {
1957        ids: out,
1958        prefiltered: labels.is_empty() || label_hint.is_some(),
1959    }
1960}
1961
1962/// Candidate ids for a `NodeByPropertyScan`. With `in_list` the scan
1963/// seeks `key IN expected`: one lookup per distinct list element, the
1964/// union deduplicated. A `null` list matches nothing. Any other non-list
1965/// value yields every labelled node unfiltered, so the `Filter` kept
1966/// above the scan evaluates (and reports) the predicate itself.
1967pub(crate) fn property_scan_candidates<S: GraphStorage>(
1968    storage: &S,
1969    labels: &[Vec<String>],
1970    key: &str,
1971    expected: &LoraValue,
1972    in_list: bool,
1973) -> NodePropertyCandidates {
1974    if !in_list {
1975        return indexed_node_property_candidates(storage, labels, key, expected);
1976    }
1977    match expected {
1978        LoraValue::Null => NodePropertyCandidates {
1979            ids: Vec::new(),
1980            prefiltered: true,
1981        },
1982        LoraValue::List(items) => {
1983            let mut seen = BTreeSet::new();
1984            let mut ids = Vec::new();
1985            for item in items {
1986                if matches!(item, LoraValue::Null) {
1987                    continue;
1988                }
1989                let candidates = indexed_node_property_candidates(storage, labels, key, item);
1990                for id in candidates.ids {
1991                    if seen.contains(&id) {
1992                        continue;
1993                    }
1994                    if !candidates.prefiltered
1995                        && !node_matches_property_filter(storage, id, labels, key, item)
1996                    {
1997                        continue;
1998                    }
1999                    seen.insert(id);
2000                    ids.push(id);
2001                }
2002            }
2003            NodePropertyCandidates {
2004                ids,
2005                prefiltered: true,
2006            }
2007        }
2008        _ => NodePropertyCandidates {
2009            ids: scan_node_ids_for_label_groups(storage, labels)
2010                .into_iter()
2011                .filter(|&id| {
2012                    storage
2013                        .with_node(id, |n| node_matches_label_groups(&n.labels, labels))
2014                        .unwrap_or(false)
2015                })
2016                .collect(),
2017            prefiltered: true,
2018        },
2019    }
2020}
2021
2022/// Whether `node_id` passes a `NodeByPropertyScan` (labels plus
2023/// `key = expected`, or `key IN expected` when `in_list`). Mirrors
2024/// [`property_scan_candidates`], including its non-list fallback.
2025pub(crate) fn property_scan_matches<S: GraphStorage>(
2026    storage: &S,
2027    node_id: NodeId,
2028    labels: &[Vec<String>],
2029    key: &str,
2030    expected: &LoraValue,
2031    in_list: bool,
2032) -> bool {
2033    if !in_list {
2034        return node_matches_property_filter(storage, node_id, labels, key, expected);
2035    }
2036    match expected {
2037        LoraValue::Null => false,
2038        LoraValue::List(items) => items.iter().any(|item| {
2039            !matches!(item, LoraValue::Null)
2040                && node_matches_property_filter(storage, node_id, labels, key, item)
2041        }),
2042        _ => storage
2043            .with_node(node_id, |n| node_matches_label_groups(&n.labels, labels))
2044            .unwrap_or(false),
2045    }
2046}
2047
2048/// Build a LoraPath from the node and relationship variables currently in a row.
2049///
2050/// For variable-length relationships (stored as a List of Relationship values),
2051/// intermediate nodes are reconstructed from the storage by walking the
2052/// relationship chain.
2053pub(crate) fn build_path_value<S: GraphStorage>(
2054    row: &Row,
2055    node_vars: &[VarId],
2056    rel_vars: &[VarId],
2057    storage: &S,
2058) -> LoraValue {
2059    let (raw_nodes, rels, has_var_len) = path_bindings(row, node_vars, rel_vars);
2060
2061    let nodes = if has_var_len && !rels.is_empty() && raw_nodes.len() == 2 {
2062        reconstruct_var_len_nodes(raw_nodes[0], &rels, storage)
2063    } else {
2064        raw_nodes
2065    };
2066
2067    LoraValue::Path(LoraPath { nodes, rels })
2068}
2069
2070#[inline]
2071fn path_bindings(
2072    row: &Row,
2073    node_vars: &[VarId],
2074    rel_vars: &[VarId],
2075) -> (Vec<NodeId>, Vec<RelationshipId>, bool) {
2076    let mut raw_nodes = Vec::new();
2077    let mut rels = Vec::new();
2078    let mut has_var_len = false;
2079
2080    for &nv in node_vars {
2081        match row.get(nv) {
2082            Some(LoraValue::Node(id)) => raw_nodes.push(*id),
2083            Some(LoraValue::List(items)) => {
2084                for item in items {
2085                    if let LoraValue::Node(id) = item {
2086                        raw_nodes.push(*id);
2087                    }
2088                }
2089            }
2090            _ => {}
2091        }
2092    }
2093
2094    for &rv in rel_vars {
2095        match row.get(rv) {
2096            Some(LoraValue::Relationship(id)) => rels.push(*id),
2097            Some(LoraValue::List(items)) => {
2098                has_var_len = true;
2099                for item in items {
2100                    if let LoraValue::Relationship(id) = item {
2101                        rels.push(*id);
2102                    }
2103                }
2104            }
2105            _ => {}
2106        }
2107    }
2108
2109    (raw_nodes, rels, has_var_len)
2110}
2111
2112#[inline]
2113fn reconstruct_var_len_nodes<S: GraphStorage>(
2114    start: NodeId,
2115    rels: &[RelationshipId],
2116    storage: &S,
2117) -> Vec<NodeId> {
2118    let mut ordered = Vec::with_capacity(rels.len() + 1);
2119    ordered.push(start);
2120    let mut current = start;
2121    for &rel_id in rels {
2122        if let Some((src, dst)) = storage.relationship_endpoints(rel_id) {
2123            let next = if src == current { dst } else { src };
2124            ordered.push(next);
2125            current = next;
2126        }
2127    }
2128    ordered
2129}
2130
2131fn type_rank(v: &LoraValue) -> u8 {
2132    match v {
2133        LoraValue::Null => 0,
2134        LoraValue::Bool(_) => 1,
2135        LoraValue::Int(_) | LoraValue::Float(_) => 2,
2136        LoraValue::String(_) => 3,
2137        LoraValue::Binary(_) => 4,
2138        LoraValue::Date(_) => 5,
2139        LoraValue::DateTime(_) => 6,
2140        LoraValue::LocalDateTime(_) => 7,
2141        LoraValue::Time(_) => 8,
2142        LoraValue::LocalTime(_) => 9,
2143        LoraValue::Duration(_) => 10,
2144        LoraValue::Point(_) => 11,
2145        LoraValue::Vector(_) => 12,
2146        LoraValue::List(_) => 13,
2147        LoraValue::Map(_) => 14,
2148        LoraValue::Node(_) => 15,
2149        LoraValue::Relationship(_) => 16,
2150        LoraValue::Path(_) => 17,
2151    }
2152}
2153
2154/// Check whether a node's labels satisfy all label groups.
2155/// Each group is a disjunction (OR): the node must have at least one label
2156/// from the group.  Groups are conjunctive (AND): all groups must be satisfied.
2157pub(crate) fn node_matches_label_groups(node_labels: &[String], groups: &[Vec<String>]) -> bool {
2158    groups
2159        .iter()
2160        .all(|group| group.iter().any(|l| node_labels.iter().any(|nl| nl == l)))
2161}
2162
2163/// Scan the graph for candidate node IDs matching the label groups. Uses the
2164/// label index for the pick-first-label phase and avoids cloning NodeRecords.
2165pub(crate) fn scan_node_ids_for_label_groups<S: GraphStorage>(
2166    storage: &S,
2167    groups: &[Vec<String>],
2168) -> Vec<NodeId> {
2169    if groups.is_empty() {
2170        return storage.all_node_ids();
2171    }
2172    if groups.len() == 1 {
2173        return label_group_candidate_ids(storage, &groups[0]);
2174    }
2175
2176    let mut best: Option<Vec<NodeId>> = None;
2177    for group in groups {
2178        let ids = label_group_candidate_ids(storage, group);
2179        if ids.is_empty() {
2180            return Vec::new();
2181        }
2182        if best
2183            .as_ref()
2184            .map(|current| ids.len() < current.len())
2185            .unwrap_or(true)
2186        {
2187            best = Some(ids);
2188        }
2189    }
2190
2191    best.unwrap_or_default()
2192}
2193
2194pub(crate) fn label_group_candidates_prefiltered(groups: &[Vec<String>]) -> bool {
2195    groups.len() <= 1
2196}
2197
2198fn label_group_candidate_ids<S: GraphStorage>(storage: &S, group: &[String]) -> Vec<NodeId> {
2199    match group {
2200        [] => Vec::new(),
2201        [label] => storage.node_ids_by_label(label),
2202        labels => {
2203            let mut seen = BTreeSet::new();
2204            let mut out = Vec::new();
2205            for label in labels {
2206                for id in storage.node_ids_by_label(label) {
2207                    if seen.insert(id) {
2208                        out.push(id);
2209                    }
2210                }
2211            }
2212            out
2213        }
2214    }
2215}
2216
2217pub(crate) fn hydrate_node_record(node: &lora_store::NodeRecord) -> LoraValue {
2218    let mut map = BTreeMap::new();
2219    map.insert("kind".to_string(), LoraValue::String("node".to_string()));
2220    map.insert("id".to_string(), LoraValue::Int(node.id as i64));
2221    map.insert(
2222        "labels".to_string(),
2223        LoraValue::List(
2224            node.labels
2225                .iter()
2226                .map(|s| LoraValue::String(s.clone()))
2227                .collect(),
2228        ),
2229    );
2230    map.insert(
2231        "properties".to_string(),
2232        properties_to_value_map(&node.properties),
2233    );
2234    LoraValue::Map(map)
2235}
2236
2237pub(crate) fn hydrate_relationship_record(rel: &lora_store::RelationshipRecord) -> LoraValue {
2238    let mut map = BTreeMap::new();
2239    map.insert(
2240        "kind".to_string(),
2241        LoraValue::String("relationship".to_string()),
2242    );
2243    map.insert("id".to_string(), LoraValue::Int(rel.id as i64));
2244    map.insert("startId".to_string(), LoraValue::Int(rel.src as i64));
2245    map.insert("endId".to_string(), LoraValue::Int(rel.dst as i64));
2246    map.insert("type".to_string(), LoraValue::String(rel.rel_type.clone()));
2247    map.insert(
2248        "properties".to_string(),
2249        properties_to_value_map(&rel.properties),
2250    );
2251    LoraValue::Map(map)
2252}
2253
2254/// Flatten label groups into a simple Vec<String> (for CREATE/MERGE where
2255/// disjunction doesn't apply — all labels are created).
2256pub(super) fn flatten_label_groups(groups: &[Vec<String>]) -> Vec<String> {
2257    groups.iter().flat_map(|g| g.iter().cloned()).collect()
2258}
2259
2260#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
2261pub(crate) enum GroupValueKey {
2262    Null,
2263    Bool(bool),
2264    Int(i64),
2265    Float(String),
2266    String(String),
2267    Binary(Vec<Vec<u8>>),
2268    List(Vec<GroupValueKey>),
2269    Map(Vec<(String, GroupValueKey)>),
2270    Node(u64),
2271    Relationship(u64),
2272}
2273
2274impl GroupValueKey {
2275    pub(crate) fn from_value(v: &LoraValue) -> Self {
2276        match v {
2277            LoraValue::Null => Self::Null,
2278            LoraValue::Bool(x) => Self::Bool(*x),
2279            LoraValue::Int(x) => Self::Int(*x),
2280            LoraValue::Float(x) => Self::Float(x.to_string()),
2281            LoraValue::String(x) => Self::String(x.clone()),
2282            LoraValue::Binary(x) => Self::Binary(x.segments().to_vec()),
2283            LoraValue::List(xs) => Self::List(xs.iter().map(Self::from_value).collect()),
2284            LoraValue::Map(m) => Self::Map(
2285                m.iter()
2286                    .map(|(k, v)| (k.clone(), Self::from_value(v)))
2287                    .collect(),
2288            ),
2289            LoraValue::Node(id) => Self::Node(*id),
2290            LoraValue::Relationship(id) => Self::Relationship(*id),
2291            LoraValue::Path(_) => Self::Null,
2292            // Temporal types: use their string representation as group key
2293            LoraValue::Date(d) => Self::String(d.to_string()),
2294            LoraValue::DateTime(dt) => Self::String(dt.to_string()),
2295            LoraValue::LocalDateTime(dt) => Self::String(dt.to_string()),
2296            LoraValue::Time(t) => Self::String(t.to_string()),
2297            LoraValue::LocalTime(t) => Self::String(t.to_string()),
2298            LoraValue::Duration(dur) => Self::String(dur.to_string()),
2299            LoraValue::Point(p) => Self::String(p.to_string()),
2300            LoraValue::Vector(v) => Self::String(format!("vector:{}", v.to_key_string())),
2301        }
2302    }
2303}
2304
2305/// Compute effective (min_hops, max_hops) from a `RangeLiteral`.
2306///
2307/// Lora semantics:
2308/// - `*`       → 1..∞   (start=None, end=None)
2309/// - `*2..5`   → 2..5   (start=Some(2), end=Some(5))
2310/// - `*..3`    → 1..3   (start=None, end=Some(3))
2311/// - `*2..`    → 2..∞   (start=Some(2), end=None)
2312/// - `*3`      → 3..3   (start=Some(3), end=None, no dots → exactly 3)
2313/// - `*0..1`   → 0..1
2314///
2315/// For unbounded upper, we cap at `MAX_VAR_LEN_HOPS` to prevent runaway.
2316const MAX_VAR_LEN_HOPS: u64 = 100;
2317
2318pub(crate) fn resolve_range(range: &RangeLiteral) -> (u64, u64) {
2319    let min_hops = range.start.unwrap_or(1);
2320    let max_hops = range.end.unwrap_or(MAX_VAR_LEN_HOPS);
2321    (min_hops, max_hops)
2322}
2323
2324/// An entry produced during BFS variable-length expansion.
2325pub(crate) struct VarLenResult {
2326    /// The destination node at the end of this path.
2327    pub(crate) dst_node_id: NodeId,
2328    /// The relationship IDs traversed (in order).
2329    pub(crate) rel_ids: Vec<u64>,
2330}
2331
2332/// Perform variable-length expansion from `start_node_id` following
2333/// relationships of the given `types` and `direction`, collecting all
2334/// reachable nodes at hop distances in `[min_hops, max_hops]`.
2335///
2336/// Uses BFS with relationship-uniqueness per path (each path does not
2337/// reuse the same relationship, but may revisit nodes).
2338pub(crate) fn variable_length_expand<S: GraphStorage>(
2339    storage: &S,
2340    start_node_id: NodeId,
2341    direction: Direction,
2342    types: &[String],
2343    min_hops: u64,
2344    max_hops: u64,
2345    bind_relationships: bool,
2346) -> Vec<VarLenResult> {
2347    let mut results = Vec::new();
2348
2349    // Each frontier entry: (current_node_id, relationships_used_so_far)
2350    let mut frontier: Vec<(NodeId, Vec<u64>)> = vec![(start_node_id, Vec::new())];
2351
2352    for depth in 1..=max_hops {
2353        // On the final hop we don't need to build next_frontier at all; every
2354        // path gets recorded and then the loop terminates. Avoids one full
2355        // pass of Vec clones on deep traversals.
2356        let is_last_hop = depth == max_hops;
2357        let mut next_frontier: Vec<(NodeId, Vec<u64>)> = Vec::new();
2358
2359        for (current_node, rels_used) in &frontier {
2360            // ID-only expand avoids cloning full records/properties for every
2361            // neighbour on every hop.
2362            for (rel_id, neighbor_id) in storage.expand_ids(*current_node, direction, types) {
2363                // Relationship-uniqueness: skip if this relationship was already
2364                // traversed on this particular path.
2365                if rels_used.contains(&rel_id) {
2366                    continue;
2367                }
2368
2369                if is_last_hop {
2370                    // Terminal hop: just record the result. Allocate rel_ids
2371                    // once (no duplicate clone) by extending a fresh copy.
2372                    if depth >= min_hops {
2373                        let mut rel_ids = Vec::with_capacity(rels_used.len() + 1);
2374                        rel_ids.extend_from_slice(rels_used);
2375                        rel_ids.push(rel_id);
2376                        results.push(VarLenResult {
2377                            dst_node_id: neighbor_id,
2378                            rel_ids: if bind_relationships {
2379                                rel_ids
2380                            } else {
2381                                Vec::new()
2382                            },
2383                        });
2384                    }
2385                    continue;
2386                }
2387
2388                let mut new_rels = Vec::with_capacity(rels_used.len() + 1);
2389                new_rels.extend_from_slice(rels_used);
2390                new_rels.push(rel_id);
2391
2392                if depth >= min_hops {
2393                    results.push(VarLenResult {
2394                        dst_node_id: neighbor_id,
2395                        rel_ids: if bind_relationships {
2396                            new_rels.clone()
2397                        } else {
2398                            Vec::new()
2399                        },
2400                    });
2401                }
2402
2403                next_frontier.push((neighbor_id, new_rels));
2404            }
2405        }
2406
2407        if is_last_hop || next_frontier.is_empty() {
2408            break;
2409        }
2410
2411        frontier = next_frontier;
2412    }
2413
2414    // Handle min_hops == 0: include the start node itself at depth 0.
2415    if min_hops == 0 {
2416        results.insert(
2417            0,
2418            VarLenResult {
2419                dst_node_id: start_node_id,
2420                rel_ids: Vec::new(),
2421            },
2422        );
2423    }
2424
2425    results
2426}
2427
2428/// Filter rows to keep only shortest paths.
2429/// `all` = false → keep one shortest path; `all` = true → keep all shortest.
2430pub(crate) fn filter_shortest_paths(rows: Vec<Row>, path_var: VarId, all: bool) -> Vec<Row> {
2431    if rows.is_empty() {
2432        return rows;
2433    }
2434
2435    // Compute path length for each row
2436    let lengths: Vec<usize> = rows
2437        .iter()
2438        .map(|row| match row.get(path_var) {
2439            Some(LoraValue::Path(p)) => p.rels.len(),
2440            _ => usize::MAX,
2441        })
2442        .collect();
2443
2444    let min_len = lengths.iter().copied().min().unwrap_or(usize::MAX);
2445
2446    let mut result: Vec<Row> = rows
2447        .into_iter()
2448        .zip(lengths.iter())
2449        .filter(|(_, len)| **len == min_len)
2450        .map(|(row, _)| row)
2451        .collect();
2452
2453    if !all && result.len() > 1 {
2454        result.truncate(1);
2455    }
2456
2457    result
2458}
2459
2460// ---------- Rel scans ----------
2461//
2462// Mirror of the `node_by_*_scan_rows` helpers above for the
2463// relationship-targeted index operators. Each helper resolves the
2464// candidate set via the corresponding `relationship_*_candidates`
2465// trait method (falling back to `rel_ids_by_type` when no scope is
2466// active), refilters with the precise predicate, and emits one row
2467// per matching relationship per input row.
2468//
2469// Direction handling: the optimizer only emits a Rel*Scan when both
2470// endpoints of the pattern have no upstream constraint. With a
2471// directed pattern (Direction::Right), each rel is bound as
2472// (src=stored_src, rel, dst=stored_dst). With Direction::Left those
2473// are swapped. With Direction::Undirected we emit both orientations
2474// to preserve parity with the Expand-based plan it replaces.
2475
2476fn emit_rel_rows(
2477    direction: Direction,
2478    src_var: VarId,
2479    rel_var: VarId,
2480    dst_var: VarId,
2481    rel: &lora_store::RelationshipRecord,
2482    base: &Row,
2483    out: &mut Vec<Row>,
2484) -> ExecResult<()> {
2485    match direction {
2486        Direction::Right => {
2487            emit_one_rel_row(
2488                src_var, rel_var, dst_var, rel.src, rel.id, rel.dst, base, out,
2489            )?;
2490        }
2491        Direction::Left => {
2492            emit_one_rel_row(
2493                src_var, rel_var, dst_var, rel.dst, rel.id, rel.src, base, out,
2494            )?;
2495        }
2496        Direction::Undirected => {
2497            emit_one_rel_row(
2498                src_var, rel_var, dst_var, rel.src, rel.id, rel.dst, base, out,
2499            )?;
2500            // Self-loops produce a single row even under undirected
2501            // semantics — emitting two copies of (a=a, b=a) for
2502            // a=stored_src=stored_dst would double-count.
2503            if rel.src != rel.dst {
2504                emit_one_rel_row(
2505                    src_var, rel_var, dst_var, rel.dst, rel.id, rel.src, base, out,
2506                )?;
2507            }
2508        }
2509    }
2510    Ok(())
2511}
2512
2513#[allow(clippy::too_many_arguments)]
2514fn emit_one_rel_row(
2515    src_var: VarId,
2516    rel_var: VarId,
2517    dst_var: VarId,
2518    src_id: NodeId,
2519    rel_id: RelationshipId,
2520    dst_id: NodeId,
2521    base: &Row,
2522    out: &mut Vec<Row>,
2523) -> ExecResult<()> {
2524    let mut row = base.clone();
2525    if bind_node_value(&mut row, src_var, src_id)?
2526        && bind_relationship_value(&mut row, rel_var, rel_id)?
2527        && bind_node_value(&mut row, dst_var, dst_id)?
2528    {
2529        out.push(row);
2530    }
2531    Ok(())
2532}
2533
2534fn bind_node_value(row: &mut Row, var: VarId, id: NodeId) -> ExecResult<bool> {
2535    match row.get(var) {
2536        Some(LoraValue::Node(existing)) => Ok(*existing == id),
2537        Some(other) => Err(ExecutorError::ExpectedNodeForExpand {
2538            var: format!("{var:?}"),
2539            found: value_kind(other),
2540        }),
2541        None => {
2542            row.insert(var, LoraValue::Node(id));
2543            Ok(true)
2544        }
2545    }
2546}
2547
2548fn bind_relationship_value(row: &mut Row, var: VarId, id: RelationshipId) -> ExecResult<bool> {
2549    match row.get(var) {
2550        Some(LoraValue::Relationship(existing)) => Ok(*existing == id),
2551        Some(other) => Err(ExecutorError::ExpectedRelationshipForExpand {
2552            var: format!("{var:?}"),
2553            found: value_kind(other),
2554        }),
2555        None => {
2556            row.insert(var, LoraValue::Relationship(id));
2557            Ok(true)
2558        }
2559    }
2560}
2561
2562fn rel_candidate_ids<S, F>(storage: &S, types: &[String], indexed: F) -> Vec<RelationshipId>
2563where
2564    S: GraphStorage,
2565    F: Fn(&str) -> Option<Vec<RelationshipId>>,
2566{
2567    if types.is_empty() {
2568        // No type constraint → no rel-typed scope to probe; fall back
2569        // to scanning every relationship.
2570        return storage.all_rel_ids();
2571    }
2572    let mut all = Vec::new();
2573    let mut seen = BTreeSet::new();
2574    for ty in types {
2575        let ids = indexed(ty).unwrap_or_else(|| storage.rel_ids_by_type(ty));
2576        for id in ids {
2577            if seen.insert(id) {
2578                all.push(id);
2579            }
2580        }
2581    }
2582    all
2583}
2584
2585pub(crate) fn rel_by_property_range_scan_rows<S: GraphStorage>(
2586    storage: &S,
2587    params: &BTreeMap<String, LoraValue>,
2588    base_rows: Vec<Row>,
2589    op: &lora_compiler::RelByPropertyRangeScanExec,
2590    deadline: Option<Instant>,
2591) -> ExecResult<Vec<Row>> {
2592    let eval_ctx = EvalContext { storage, params };
2593    let mut out = Vec::new();
2594    let mut other_kinds = OtherKindScan::default();
2595
2596    for row in base_rows {
2597        check_optional_deadline(deadline)?;
2598        let lo_value = op.lo.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
2599        // A null upper bound reads as none: the Filter kept above every
2600        // index scan still judges each row (a STARTS WITH prefix with no
2601        // successor, `string.prefix_end` null, scans to the end).
2602        let hi_value = op
2603            .hi
2604            .as_ref()
2605            .map(|expr| eval_expr(expr, &row, &eval_ctx))
2606            .filter(|v| !matches!(v, LoraValue::Null));
2607        let lo_prop = lo_value
2608            .clone()
2609            .and_then(|v| lora_value_to_property(v).ok());
2610        let hi_prop = hi_value
2611            .clone()
2612            .and_then(|v| lora_value_to_property(v).ok());
2613
2614        // Temporals of another kind first, for the Filter above to judge
2615        // (see `other_kind_node_ids`).
2616        let bounds = [lo_value.as_ref(), hi_value.as_ref()];
2617        for rel_id in other_kind_rel_ids(storage, op, bounds, &mut other_kinds) {
2618            if let Some(result) = storage.with_relationship(rel_id, |rel| {
2619                if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
2620                    return Ok(());
2621                }
2622                emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2623            }) {
2624                result?;
2625            }
2626        }
2627        let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| {
2628            storage.relationship_range_candidates(ty, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
2629        });
2630
2631        for rel_id in candidate_ids {
2632            check_optional_deadline(deadline)?;
2633            if let Some(result) = storage.with_relationship(rel_id, |rel| {
2634                if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
2635                    return Ok(());
2636                }
2637                let Some(actual) = rel.properties.get(op.key.as_str()) else {
2638                    return Ok(());
2639                };
2640                let actual_lv = LoraValue::from(actual);
2641                if !range_predicate_holds(
2642                    &actual_lv,
2643                    lo_value.as_ref(),
2644                    op.lo_inclusive,
2645                    hi_value.as_ref(),
2646                    op.hi_inclusive,
2647                ) {
2648                    return Ok(());
2649                }
2650                emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2651            }) {
2652                result?;
2653            }
2654        }
2655    }
2656
2657    Ok(out)
2658}
2659
2660pub(crate) fn rel_by_text_scan_rows<S: GraphStorage>(
2661    storage: &S,
2662    params: &BTreeMap<String, LoraValue>,
2663    base_rows: Vec<Row>,
2664    op: &lora_compiler::RelByTextScanExec,
2665    deadline: Option<Instant>,
2666) -> ExecResult<Vec<Row>> {
2667    let eval_ctx = EvalContext { storage, params };
2668    let mut out = Vec::new();
2669
2670    for row in base_rows {
2671        check_optional_deadline(deadline)?;
2672        let query = eval_expr(&op.query, &row, &eval_ctx);
2673        let LoraValue::String(query_str) = &query else {
2674            continue;
2675        };
2676
2677        let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| {
2678            storage.relationship_text_candidates(ty, &op.key, query_str)
2679        });
2680
2681        for rel_id in candidate_ids {
2682            check_optional_deadline(deadline)?;
2683            if let Some(result) = storage.with_relationship(rel_id, |rel| {
2684                if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
2685                    return Ok(());
2686                }
2687                let Some(PropertyValue::String(actual)) = rel.properties.get(op.key.as_str())
2688                else {
2689                    return Ok(());
2690                };
2691                if !text_predicate_holds(actual, op.predicate, query_str) {
2692                    return Ok(());
2693                }
2694                emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2695            }) {
2696                result?;
2697            }
2698        }
2699    }
2700
2701    Ok(out)
2702}
2703
2704pub(crate) fn rel_by_point_scan_rows<S: GraphStorage>(
2705    storage: &S,
2706    params: &BTreeMap<String, LoraValue>,
2707    base_rows: Vec<Row>,
2708    op: &lora_compiler::RelByPointScanExec,
2709    deadline: Option<Instant>,
2710) -> ExecResult<Vec<Row>> {
2711    let eval_ctx = EvalContext { storage, params };
2712    let mut out = Vec::new();
2713
2714    for row in base_rows {
2715        check_optional_deadline(deadline)?;
2716
2717        let probe = match &op.predicate {
2718            lora_compiler::PointPredicate::WithinBBox {
2719                lower_left,
2720                upper_right,
2721            } => {
2722                let ll = eval_expr(lower_left, &row, &eval_ctx);
2723                let ur = eval_expr(upper_right, &row, &eval_ctx);
2724                match (ll, ur) {
2725                    (LoraValue::Point(a), LoraValue::Point(b)) => {
2726                        Probe::WithinBBox { ll: a, ur: b }
2727                    }
2728                    _ => continue,
2729                }
2730            }
2731            lora_compiler::PointPredicate::WithinDistance {
2732                center,
2733                max_distance,
2734                inclusive,
2735            } => {
2736                let c = eval_expr(center, &row, &eval_ctx);
2737                let d = eval_expr(max_distance, &row, &eval_ctx);
2738                match (c, d) {
2739                    (LoraValue::Point(c), LoraValue::Float(d)) => Probe::WithinDistance {
2740                        center: c,
2741                        max: d,
2742                        inclusive: *inclusive,
2743                    },
2744                    (LoraValue::Point(c), LoraValue::Int(d)) => Probe::WithinDistance {
2745                        center: c,
2746                        max: d as f64,
2747                        inclusive: *inclusive,
2748                    },
2749                    _ => continue,
2750                }
2751            }
2752        };
2753
2754        let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| match &probe {
2755            Probe::WithinBBox { ll, ur } => {
2756                // Two seeks when the box crosses the antimeridian.
2757                let (ranges, n) = lora_store::bbox_x_ranges(ll, ur);
2758                let (lo_y, hi_y) = (ll.y.min(ur.y), ll.y.max(ur.y));
2759                let mut ids = Vec::new();
2760                for (lo_x, hi_x) in &ranges[..n] {
2761                    ids.extend(storage.relationship_point_within_bbox(
2762                        ty,
2763                        &op.key,
2764                        (*lo_x, lo_y),
2765                        (*hi_x, hi_y),
2766                    )?);
2767                }
2768                if n > 1 {
2769                    ids.sort_unstable();
2770                    ids.dedup();
2771                }
2772                Some(ids)
2773            }
2774            Probe::WithinDistance { center, max, .. } => {
2775                storage.relationship_point_within_distance(ty, &op.key, (center.x, center.y), *max)
2776            }
2777        });
2778
2779        for rel_id in candidate_ids {
2780            check_optional_deadline(deadline)?;
2781            if let Some(result) = storage.with_relationship(rel_id, |rel| {
2782                if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
2783                    return Ok(());
2784                }
2785                let Some(PropertyValue::Point(actual)) = rel.properties.get(op.key.as_str()) else {
2786                    return Ok(());
2787                };
2788                if !point_predicate_holds(actual, &probe) {
2789                    return Ok(());
2790                }
2791                emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2792            }) {
2793                result?;
2794            }
2795        }
2796    }
2797
2798    Ok(out)
2799}
2800
2801/// Node ids of an ordered range scan (`op.order` set), in `ORDER BY op.key`
2802/// order, each already checked against the exact bounds and labels.
2803///
2804/// The index serves the order when every bound is a string: a string
2805/// comparison only lets strings through, and among strings the index's
2806/// order is Cypher's order. Ids are then pulled from the index in chunks
2807/// that grow as rows are consumed, so a `LIMIT` above stops the scan
2808/// after a few small index reads. For any other bound the index order is
2809/// not Cypher's (it orders every integer before every float), so the
2810/// cursor collects the candidates and sorts them by value itself.
2811pub(crate) struct OrderedRangeCursor {
2812    lo: Option<LoraValue>,
2813    hi: Option<LoraValue>,
2814    mode: OrderedMode,
2815    pending: std::vec::IntoIter<NodeId>,
2816    labels: Vec<Vec<String>>,
2817    key: String,
2818    lo_inclusive: bool,
2819    hi_inclusive: bool,
2820}
2821
2822enum OrderedMode {
2823    Index {
2824        label: String,
2825        descending: bool,
2826        after: Option<(PropertyValue, NodeId)>,
2827        chunk: usize,
2828        exhausted: bool,
2829    },
2830    Buffered,
2831}
2832
2833impl OrderedRangeCursor {
2834    pub(crate) fn new<S: GraphStorage>(
2835        storage: &S,
2836        op: &lora_compiler::NodeByPropertyRangeScanExec,
2837        lo: Option<LoraValue>,
2838        hi: Option<LoraValue>,
2839    ) -> Self {
2840        let descending = matches!(op.order, Some(SortDirection::Desc));
2841        // Strings and temporals of one kind have a single contiguous,
2842        // correctly ordered run in the index, so the walk can stream from
2843        // it. Numbers cannot: integers and floats are keyed separately.
2844        let mut bounds = lo.iter().chain(hi.iter());
2845        let string_bounds = match bounds.next() {
2846            Some(LoraValue::String(_)) => bounds.all(|v| matches!(v, LoraValue::String(_))),
2847            Some(
2848                first @ (LoraValue::Date(_)
2849                | LoraValue::DateTime(_)
2850                | LoraValue::LocalDateTime(_)
2851                | LoraValue::Time(_)
2852                | LoraValue::LocalTime(_)),
2853            ) => bounds.all(|v| std::mem::discriminant(v) == std::mem::discriminant(first)),
2854            _ => false,
2855        };
2856        let index_label = single_label_hint(&op.labels).filter(|label| {
2857            string_bounds
2858                && storage
2859                    .node_range_ordered_chunk(label, &op.key, None, None, descending, None, 0)
2860                    .is_some()
2861        });
2862
2863        let mut cursor = Self {
2864            lo,
2865            hi,
2866            mode: OrderedMode::Buffered,
2867            pending: Vec::new().into_iter(),
2868            labels: op.labels.clone(),
2869            key: op.key.clone(),
2870            lo_inclusive: op.lo_inclusive,
2871            hi_inclusive: op.hi_inclusive,
2872        };
2873        match index_label {
2874            Some(label) => {
2875                cursor.mode = OrderedMode::Index {
2876                    label: label.to_string(),
2877                    descending,
2878                    after: None,
2879                    chunk: 64,
2880                    exhausted: false,
2881                }
2882            }
2883            None => {
2884                cursor.pending = cursor
2885                    .sorted_candidates(storage, op, descending)
2886                    .into_iter()
2887            }
2888        }
2889        cursor
2890    }
2891
2892    pub(crate) fn filter(&self) -> NodeRangeFilter<'_> {
2893        NodeRangeFilter {
2894            labels: &self.labels,
2895            key: &self.key,
2896            lo: self.lo.as_ref(),
2897            lo_inclusive: self.lo_inclusive,
2898            hi: self.hi.as_ref(),
2899            hi_inclusive: self.hi_inclusive,
2900        }
2901    }
2902
2903    /// Fallback: every matching id, sorted by the property value in Cypher
2904    /// order (ties by id).
2905    fn sorted_candidates<S: GraphStorage>(
2906        &self,
2907        storage: &S,
2908        op: &lora_compiler::NodeByPropertyRangeScanExec,
2909        descending: bool,
2910    ) -> Vec<NodeId> {
2911        let lo_prop = self.lo.clone().and_then(|v| lora_value_to_property(v).ok());
2912        let hi_prop = self.hi.clone().and_then(|v| lora_value_to_property(v).ok());
2913        let candidates = match single_label_hint(&op.labels) {
2914            Some(label) => storage
2915                .node_range_candidates(label, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
2916                .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
2917            None => scan_node_ids_for_label_groups(storage, &op.labels),
2918        };
2919        let filter = self.filter();
2920        let mut keyed: Vec<(LoraValue, NodeId)> = candidates
2921            .into_iter()
2922            .filter(|&id| node_matches_range_filter(storage, id, &filter))
2923            .map(|id| {
2924                let value = storage
2925                    .with_node(id, |n| n.properties.get(op.key.as_str()).cloned())
2926                    .flatten()
2927                    .map(LoraValue::from)
2928                    .unwrap_or(LoraValue::Null);
2929                (value, id)
2930            })
2931            .collect();
2932        keyed.sort_by(|(a, ai), (b, bi)| {
2933            let ord = compare_values_total(a, b).then(ai.cmp(bi));
2934            if descending {
2935                ord.reverse()
2936            } else {
2937                ord
2938            }
2939        });
2940        keyed.into_iter().map(|(_, id)| id).collect()
2941    }
2942
2943    pub(crate) fn next_id<S: GraphStorage>(&mut self, storage: &S) -> Option<NodeId> {
2944        loop {
2945            if let Some(id) = self.pending.next() {
2946                if matches!(self.mode, OrderedMode::Buffered)
2947                    || node_matches_range_filter(storage, id, &self.filter())
2948                {
2949                    return Some(id);
2950                }
2951                continue;
2952            }
2953            let lo_prop = self.lo.clone().and_then(|v| lora_value_to_property(v).ok());
2954            let hi_prop = self.hi.clone().and_then(|v| lora_value_to_property(v).ok());
2955            let OrderedMode::Index {
2956                label,
2957                descending,
2958                after,
2959                chunk,
2960                exhausted,
2961            } = &mut self.mode
2962            else {
2963                return None;
2964            };
2965            if *exhausted {
2966                return None;
2967            }
2968            let ids = storage.node_range_ordered_chunk(
2969                label,
2970                &self.key,
2971                lo_prop.as_ref(),
2972                hi_prop.as_ref(),
2973                *descending,
2974                after.as_ref().map(|(v, id)| (v, *id)),
2975                *chunk,
2976            )?;
2977            if ids.len() < *chunk {
2978                *exhausted = true;
2979            }
2980            let last = *ids.last()?;
2981            let last_value = storage
2982                .with_node(last, |n| n.properties.get(self.key.as_str()).cloned())
2983                .flatten()?;
2984            *after = Some((last_value, last));
2985            *chunk = (*chunk * 2).min(4096);
2986            self.pending = ids.into_iter();
2987        }
2988    }
2989}
2990
2991// ---------------------------------------------------------------------------
2992// Seek by internal id
2993// ---------------------------------------------------------------------------
2994
2995/// Ids an id seek visits for `value`, sorted and distinct. A superset of
2996/// the ids `id(x) = value` (or, with `in_list`, `id(x) IN value`) holds
2997/// for: the `Filter` kept above the seek re-judges each row, so only the
2998/// ids that could make the predicate true need to be here. `id()` is an
2999/// integer, which equals an integer of the same value or a float that
3000/// converts to it exactly (`1 = 1.0`); null, a negative or fractional
3001/// number, or any other kind matches nothing. `all_ids` lists every id;
3002/// it is only called for a float of 2^53 or more, which `i64 as f64`
3003/// rounding makes equal to several ids.
3004pub(crate) fn id_seek_candidates(
3005    value: &LoraValue,
3006    in_list: bool,
3007    all_ids: impl Fn() -> Vec<u64>,
3008) -> Vec<u64> {
3009    fn push_candidates(value: &LoraValue, all_ids: &dyn Fn() -> Vec<u64>, out: &mut Vec<u64>) {
3010        /// 2^53: below it every integer converts to `f64` exactly.
3011        const EXACT_F64_LIMIT: f64 = 9_007_199_254_740_992.0;
3012        match value {
3013            LoraValue::Int(i) if *i >= 0 => out.push(*i as u64),
3014            LoraValue::Float(f) if f.is_finite() && *f >= 0.0 && f.fract() == 0.0 => {
3015                if *f < EXACT_F64_LIMIT {
3016                    out.push(*f as u64);
3017                } else {
3018                    out.extend(all_ids().into_iter().filter(|id| (*id as i64) as f64 == *f));
3019                }
3020            }
3021            _ => {}
3022        }
3023    }
3024
3025    let mut out = Vec::new();
3026    if in_list {
3027        if let LoraValue::List(items) = value {
3028            for item in items {
3029                push_candidates(item, &all_ids, &mut out);
3030            }
3031        }
3032    } else {
3033        push_candidates(value, &all_ids, &mut out);
3034    }
3035    out.sort_unstable();
3036    out.dedup();
3037    out
3038}
3039
3040/// Rows of a `NodeByIdSeek`: for each base row, one row per existing node
3041/// among the sought ids whose labels match. A base row that already binds
3042/// `op.var` passes through iff that node is among them.
3043pub(crate) fn node_by_id_seek_rows<S: GraphStorage>(
3044    storage: &S,
3045    params: &BTreeMap<String, LoraValue>,
3046    base_rows: Vec<Row>,
3047    op: &lora_compiler::NodeByIdSeekExec,
3048    deadline: Option<Instant>,
3049) -> ExecResult<Vec<Row>> {
3050    let eval_ctx = EvalContext { storage, params };
3051    let node_ok = |id: NodeId| {
3052        if op.labels.is_empty() {
3053            storage.contains_node(id)
3054        } else {
3055            storage
3056                .with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
3057                .unwrap_or(false)
3058        }
3059    };
3060    let mut out = Vec::new();
3061    for row in base_rows {
3062        check_optional_deadline(deadline)?;
3063        let value = eval_expr(&op.ids, &row, &eval_ctx);
3064        let ids = id_seek_candidates(&value, op.in_list, || storage.all_node_ids());
3065        if let Some(bound) = bound_node_id_for_expand(&row, op.var)? {
3066            if ids.binary_search(&bound).is_ok() && node_ok(bound) {
3067                out.push(row);
3068            }
3069            continue;
3070        }
3071        for id in ids {
3072            if node_ok(id) {
3073                let mut new_row = row.clone();
3074                new_row.insert(op.var, LoraValue::Node(id));
3075                out.push(new_row);
3076            }
3077        }
3078    }
3079    Ok(out)
3080}
3081
3082/// Rows of a `RelByIdSeek`: for each base row and each existing sought
3083/// relationship of one of `op.types`, the rows the replaced
3084/// `Expand(NodeScan(src))` would produce for it: endpoints oriented by
3085/// `op.direction` (both orientations when undirected, once for a
3086/// self-loop), bound through any variable the row already binds, and the
3087/// source checked against `op.src_labels`.
3088pub(crate) fn rel_by_id_seek_rows<S: GraphStorage>(
3089    storage: &S,
3090    params: &BTreeMap<String, LoraValue>,
3091    base_rows: Vec<Row>,
3092    op: &lora_compiler::RelByIdSeekExec,
3093    deadline: Option<Instant>,
3094) -> ExecResult<Vec<Row>> {
3095    let eval_ctx = EvalContext { storage, params };
3096    let mut out = Vec::new();
3097    let mut emitted = Vec::new();
3098    for row in base_rows {
3099        check_optional_deadline(deadline)?;
3100        let value = eval_expr(&op.ids, &row, &eval_ctx);
3101        let ids = id_seek_candidates(&value, op.in_list, || storage.all_rel_ids());
3102        for rel_id in ids {
3103            emitted.clear();
3104            if let Some(result) = storage.with_relationship(rel_id, |rel| {
3105                if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
3106                    return Ok(());
3107                }
3108                emit_rel_rows(
3109                    op.direction,
3110                    op.src,
3111                    op.rel,
3112                    op.dst,
3113                    rel,
3114                    &row,
3115                    &mut emitted,
3116                )
3117            }) {
3118                result?;
3119            }
3120            for new_row in emitted.drain(..) {
3121                if !op.src_labels.is_empty() {
3122                    let Some(src) = bound_node_id_for_expand(&new_row, op.src)? else {
3123                        continue;
3124                    };
3125                    let labels_ok = storage
3126                        .with_node(src, |n| {
3127                            node_matches_label_groups(&n.labels, &op.src_labels)
3128                        })
3129                        .unwrap_or(false);
3130                    if !labels_ok {
3131                        continue;
3132                    }
3133                }
3134                out.push(new_row);
3135            }
3136        }
3137    }
3138    Ok(out)
3139}