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 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 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 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 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
1480fn 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
1493fn 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#[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
1522pub(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 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 }
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
1591pub(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
1612fn 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 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 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
1975pub(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
2035pub(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
2061pub(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
2167pub(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
2176pub(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
2267pub(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 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
2318const 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
2337pub(crate) struct VarLenResult {
2339 pub(crate) dst_node_id: NodeId,
2341 pub(crate) rel_ids: Vec<u64>,
2343}
2344
2345pub(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 let mut frontier: Vec<(NodeId, Vec<u64>)> = vec![(start_node_id, Vec::new())];
2364
2365 for depth in 1..=max_hops {
2366 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 for (rel_id, neighbor_id) in storage.expand_ids(*current_node, direction, types) {
2376 if rels_used.contains(&rel_id) {
2379 continue;
2380 }
2381
2382 if is_last_hop {
2383 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 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
2441pub(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 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
2473fn 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 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 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 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 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 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
2814pub(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 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 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}