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