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