1use std::cmp::Ordering;
36use std::collections::{BTreeMap, BTreeSet};
37use std::time::Instant;
38
39use lora_analyzer::symbols::VarId;
40use lora_analyzer::{ResolvedExpr, ResolvedMapSelector};
41use lora_ast::{Direction, RangeLiteral};
42use lora_compiler::physical::{
43 ExpandExec, HashAggregationExec, LimitExec, NodeByLabelScanExec, NodeByPropertyScanExec,
44 NodeScanExec, PhysicalNodeId, PhysicalOp, PhysicalPlan, ProjectionExec, UnwindExec,
45};
46use lora_store::{GraphStorage, NodeId, Properties, PropertyValue, RelationshipId};
47
48use crate::errors::{value_kind, ExecResult, ExecutorError};
49use crate::eval::{eval_expr, eval_expr_result, eval_truthy_result, EvalContext};
50use crate::value::{lora_value_to_property, LoraPath, LoraValue, Row};
51
52#[inline]
56pub(super) fn check_deadline_at(deadline: Instant) -> ExecResult<()> {
57 if Instant::now() >= deadline {
58 Err(ExecutorError::QueryTimeout)
59 } else {
60 Ok(())
61 }
62}
63
64pub(super) fn filter_rows_checked<S: GraphStorage>(
65 input_rows: Vec<Row>,
66 predicate: &ResolvedExpr,
67 eval_ctx: &EvalContext<'_, S>,
68) -> ExecResult<Vec<Row>> {
69 let mut out = Vec::with_capacity(input_rows.len());
70 for row in input_rows {
71 if eval_truthy_result(predicate, &row, eval_ctx).map_err(ExecutorError::RuntimeError)? {
72 out.push(row);
73 }
74 }
75 Ok(out)
76}
77
78pub(super) fn project_rows_checked<S: GraphStorage>(
79 input_rows: Vec<Row>,
80 op: &ProjectionExec,
81 eval_ctx: &EvalContext<'_, S>,
82) -> ExecResult<Vec<Row>> {
83 let mut out = Vec::with_capacity(input_rows.len());
84
85 for row in input_rows {
86 if op.include_existing {
87 let mut projected = row;
88 for item in &op.items {
89 let value = eval_expr_result(&item.expr, &projected, eval_ctx)
90 .map_err(ExecutorError::RuntimeError)?;
91 projected.insert_named(item.output, item.name.clone(), value);
92 }
93 out.push(projected);
94 } else {
95 let mut projected = Row::new();
96 for item in &op.items {
97 let value = eval_expr_result(&item.expr, &row, eval_ctx)
98 .map_err(ExecutorError::RuntimeError)?;
99 projected.insert_named(item.output, item.name.clone(), value);
100 }
101 out.push(projected);
102 }
103 }
104
105 Ok(if op.distinct {
106 dedup_rows_by_vars(out)
107 } else {
108 out
109 })
110}
111
112pub(super) fn unwind_rows<S: GraphStorage>(
113 input_rows: Vec<Row>,
114 op: &UnwindExec,
115 eval_ctx: &EvalContext<'_, S>,
116) -> Vec<Row> {
117 let mut out = Vec::new();
118
119 for row in input_rows {
120 match eval_expr(&op.expr, &row, eval_ctx) {
121 LoraValue::List(values) => {
122 for value in values {
123 let mut new_row = row.clone();
124 new_row.insert(op.alias, value);
125 out.push(new_row);
126 }
127 }
128 LoraValue::Null => {}
129 other => {
130 let mut new_row = row;
131 new_row.insert(op.alias, other);
132 out.push(new_row);
133 }
134 }
135 }
136
137 out
138}
139
140pub(super) fn limit_rows<S: GraphStorage>(
141 mut rows: Vec<Row>,
142 op: &LimitExec,
143 eval_ctx: &EvalContext<'_, S>,
144) -> Vec<Row> {
145 let limit = op
146 .limit
147 .as_ref()
148 .and_then(|e| eval_expr(e, &Row::new(), eval_ctx).as_i64())
149 .unwrap_or(rows.len() as i64)
150 .max(0) as usize;
151
152 let skip = op
153 .skip
154 .as_ref()
155 .and_then(|e| eval_expr(e, &Row::new(), eval_ctx).as_i64())
156 .unwrap_or(0)
157 .max(0) as usize;
158
159 if skip >= rows.len() {
160 return Vec::new();
161 }
162
163 rows.drain(0..skip);
164 rows.truncate(limit);
165 rows
166}
167
168#[inline]
169pub(crate) fn bound_node_id_for_expand(row: &Row, var: VarId) -> ExecResult<Option<NodeId>> {
170 match row.get(var) {
171 Some(LoraValue::Node(id)) => Ok(Some(*id)),
172 Some(other) => Err(ExecutorError::ExpectedNodeForExpand {
173 var: format!("{var:?}"),
174 found: value_kind(other),
175 }),
176 None => Ok(None),
177 }
178}
179
180#[inline]
181pub(crate) fn bound_relationship_id_for_expand(
182 row: &Row,
183 var: VarId,
184) -> ExecResult<Option<RelationshipId>> {
185 match row.get(var) {
186 Some(LoraValue::Relationship(id)) => Ok(Some(*id)),
187 Some(other) => Err(ExecutorError::ExpectedRelationshipForExpand {
188 var: format!("{var:?}"),
189 found: value_kind(other),
190 }),
191 None => Ok(None),
192 }
193}
194
195pub(super) fn node_scan_rows<S: GraphStorage>(
196 storage: &S,
197 base_rows: Vec<Row>,
198 op: &NodeScanExec,
199 deadline: Option<Instant>,
200) -> ExecResult<Vec<Row>> {
201 let node_ids = storage.all_node_ids();
202 let mut out = Vec::with_capacity(base_rows.len().saturating_mul(node_ids.len()));
203
204 if deadline.is_none() {
205 for row in base_rows {
206 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
207 if storage.has_node(existing_id) {
208 out.push(row);
209 }
210 continue;
211 }
212
213 for &id in &node_ids {
214 let mut new_row = row.clone();
215 new_row.insert(op.var, LoraValue::Node(id));
216 out.push(new_row);
217 }
218 }
219 return Ok(out);
220 }
221
222 for row in base_rows {
223 check_optional_deadline(deadline)?;
224 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
225 if storage.has_node(existing_id) {
226 out.push(row);
227 }
228 continue;
229 }
230
231 for &id in &node_ids {
232 check_optional_deadline(deadline)?;
233 let mut new_row = row.clone();
234 new_row.insert(op.var, LoraValue::Node(id));
235 out.push(new_row);
236 }
237 }
238
239 Ok(out)
240}
241
242pub(super) fn node_by_label_scan_rows<S: GraphStorage>(
243 storage: &S,
244 base_rows: Vec<Row>,
245 op: &NodeByLabelScanExec,
246 deadline: Option<Instant>,
247) -> ExecResult<Vec<Row>> {
248 let candidate_ids = scan_node_ids_for_label_groups(storage, &op.labels);
249 let candidates_prefiltered = label_group_candidates_prefiltered(&op.labels);
250 let mut out = Vec::with_capacity(base_rows.len().saturating_mul(candidate_ids.len()));
251
252 if deadline.is_none() {
253 for row in base_rows {
254 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
255 let labels_ok = storage
256 .with_node(existing_id, |n| {
257 node_matches_label_groups(&n.labels, &op.labels)
258 })
259 .unwrap_or(false);
260 if labels_ok {
261 out.push(row);
262 }
263 continue;
264 }
265
266 for &id in &candidate_ids {
267 if !candidates_prefiltered {
268 let labels_ok = storage
269 .with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
270 .unwrap_or(false);
271 if !labels_ok {
272 continue;
273 }
274 }
275 let mut new_row = row.clone();
276 new_row.insert(op.var, LoraValue::Node(id));
277 out.push(new_row);
278 }
279 }
280 return Ok(out);
281 }
282
283 for row in base_rows {
284 check_optional_deadline(deadline)?;
285 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
286 let labels_ok = storage
287 .with_node(existing_id, |n| {
288 node_matches_label_groups(&n.labels, &op.labels)
289 })
290 .unwrap_or(false);
291 if labels_ok {
292 out.push(row);
293 }
294 continue;
295 }
296
297 for &id in &candidate_ids {
298 check_optional_deadline(deadline)?;
299 if !candidates_prefiltered {
300 let labels_ok = storage
301 .with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
302 .unwrap_or(false);
303 if !labels_ok {
304 continue;
305 }
306 }
307 let mut new_row = row.clone();
308 new_row.insert(op.var, LoraValue::Node(id));
309 out.push(new_row);
310 }
311 }
312
313 Ok(out)
314}
315
316pub(super) fn node_by_property_scan_rows<S: GraphStorage>(
317 storage: &S,
318 params: &BTreeMap<String, LoraValue>,
319 base_rows: Vec<Row>,
320 op: &NodeByPropertyScanExec,
321 deadline: Option<Instant>,
322) -> ExecResult<Vec<Row>> {
323 let eval_ctx = EvalContext { storage, params };
324 let mut out = Vec::new();
325
326 if deadline.is_none() {
327 for row in base_rows {
328 let expected = eval_expr(&op.value, &row, &eval_ctx);
329
330 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
331 if node_matches_property_filter(
332 storage,
333 existing_id,
334 &op.labels,
335 &op.key,
336 &expected,
337 ) {
338 out.push(row);
339 }
340 continue;
341 }
342
343 let candidates =
344 indexed_node_property_candidates(storage, &op.labels, &op.key, &expected);
345 for id in candidates.ids {
346 if !candidates.prefiltered
347 && !node_matches_property_filter(storage, id, &op.labels, &op.key, &expected)
348 {
349 continue;
350 }
351 let mut new_row = row.clone();
352 new_row.insert(op.var, LoraValue::Node(id));
353 out.push(new_row);
354 }
355 }
356 return Ok(out);
357 }
358
359 for row in base_rows {
360 check_optional_deadline(deadline)?;
361 let expected = eval_expr(&op.value, &row, &eval_ctx);
362
363 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
364 if node_matches_property_filter(storage, existing_id, &op.labels, &op.key, &expected) {
365 out.push(row);
366 }
367 continue;
368 }
369
370 let candidates = indexed_node_property_candidates(storage, &op.labels, &op.key, &expected);
371 for id in candidates.ids {
372 check_optional_deadline(deadline)?;
373 if !candidates.prefiltered
374 && !node_matches_property_filter(storage, id, &op.labels, &op.key, &expected)
375 {
376 continue;
377 }
378 let mut new_row = row.clone();
379 new_row.insert(op.var, LoraValue::Node(id));
380 out.push(new_row);
381 }
382 }
383
384 Ok(out)
385}
386
387#[inline]
388fn check_optional_deadline(deadline: Option<Instant>) -> ExecResult<()> {
389 match deadline {
390 Some(deadline) => check_deadline_at(deadline),
391 None => Ok(()),
392 }
393}
394
395pub(crate) fn plan_may_need_hydration(plan: &PhysicalPlan) -> bool {
396 op_may_need_hydration(plan, plan.root)
397}
398
399pub(crate) fn count_all_scan_aggregation_rows<S: GraphStorage>(
400 storage: &S,
401 plan: &PhysicalPlan,
402 op: &HashAggregationExec,
403) -> Option<Vec<Row>> {
404 if !op.group_by.is_empty() {
405 return None;
406 }
407 let specs = crate::pull::classify_streamable_aggregates(&op.aggregates)?;
408 if !specs
409 .iter()
410 .all(|spec| matches!(spec.kind, crate::pull::StreamableAggKind::CountAll))
411 {
412 return None;
413 }
414
415 let count = count_rows_for_scan_subtree(storage, plan, op.input)? as i64;
416 let value = LoraValue::Int(count);
417 let mut row = Row::new();
418 for proj in &op.aggregates {
419 row.insert_named(proj.output, proj.name.clone(), value.clone());
420 }
421 Some(vec![row])
422}
423
424fn count_rows_for_scan_subtree<S: GraphStorage>(
425 storage: &S,
426 plan: &PhysicalPlan,
427 node_id: PhysicalNodeId,
428) -> Option<usize> {
429 match &plan.nodes[node_id] {
430 PhysicalOp::NodeScan(op) if scan_input_is_argument(plan, op.input) => {
431 Some(storage.node_count())
432 }
433 PhysicalOp::NodeByLabelScan(op) if scan_input_is_argument(plan, op.input) => {
434 let ids = scan_node_ids_for_label_groups(storage, &op.labels);
435 if label_group_candidates_prefiltered(&op.labels) {
436 return Some(ids.len());
437 }
438 Some(
439 ids.into_iter()
440 .filter(|&id| {
441 storage
442 .with_node(id, |node| {
443 node_matches_label_groups(&node.labels, &op.labels)
444 })
445 .unwrap_or(false)
446 })
447 .count(),
448 )
449 }
450 _ => None,
451 }
452}
453
454fn scan_input_is_argument(plan: &PhysicalPlan, input: Option<PhysicalNodeId>) -> bool {
455 match input {
456 None => true,
457 Some(id) => matches!(plan.nodes.get(id), Some(PhysicalOp::Argument(_))),
458 }
459}
460
461fn op_may_need_hydration(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
462 match &plan.nodes[node_id] {
463 PhysicalOp::Projection(op) if !op.include_existing => op
464 .items
465 .iter()
466 .any(|item| expr_may_produce_hydratable_value(&item.expr)),
467 PhysicalOp::HashAggregation(op) => {
468 op.group_by
469 .iter()
470 .any(|item| expr_may_produce_hydratable_value(&item.expr))
471 || op
472 .aggregates
473 .iter()
474 .any(|item| expr_may_produce_hydratable_value(&item.expr))
475 }
476 PhysicalOp::Sort(op) => op_may_need_hydration(plan, op.input),
477 PhysicalOp::Limit(op) => op_may_need_hydration(plan, op.input),
478 PhysicalOp::Filter(op) => op_may_need_hydration(plan, op.input),
479 PhysicalOp::Unwind(op) => op_may_need_hydration(plan, op.input),
480 PhysicalOp::PathBuild(op) => op_may_need_hydration(plan, op.input),
481 _ => true,
482 }
483}
484
485fn expr_may_produce_hydratable_value(expr: &ResolvedExpr) -> bool {
486 match expr {
487 ResolvedExpr::Variable(_) | ResolvedExpr::Parameter(_) => true,
488 ResolvedExpr::Literal(_)
489 | ResolvedExpr::Property { .. }
490 | ResolvedExpr::ExistsSubquery { .. }
491 | ResolvedExpr::Binary { .. }
492 | ResolvedExpr::Unary { .. }
493 | ResolvedExpr::ListPredicate { .. } => false,
494 ResolvedExpr::Function { name, args, .. } => {
495 function_may_produce_hydratable_value(name, args)
496 }
497 ResolvedExpr::List(items) => items.iter().any(expr_may_produce_hydratable_value),
498 ResolvedExpr::Map(items) => items
499 .iter()
500 .any(|(_, value)| expr_may_produce_hydratable_value(value)),
501 ResolvedExpr::Case {
502 alternatives,
503 else_expr,
504 ..
505 } => {
506 alternatives
507 .iter()
508 .any(|(_, value)| expr_may_produce_hydratable_value(value))
509 || else_expr
510 .as_deref()
511 .is_some_and(expr_may_produce_hydratable_value)
512 }
513 ResolvedExpr::ListComprehension { map_expr, .. } => map_expr
514 .as_deref()
515 .map(expr_may_produce_hydratable_value)
516 .unwrap_or(true),
517 ResolvedExpr::Reduce { expr, .. } => expr_may_produce_hydratable_value(expr),
518 ResolvedExpr::MapProjection { selectors, .. } => selectors.iter().any(|selector| {
519 matches!(selector, ResolvedMapSelector::Literal(_, expr) if expr_may_produce_hydratable_value(expr))
520 }),
521 ResolvedExpr::Index { expr, .. } | ResolvedExpr::Slice { expr, .. } => {
522 expr_may_produce_hydratable_value(expr)
523 }
524 ResolvedExpr::PatternComprehension { map_expr, .. } => {
525 expr_may_produce_hydratable_value(map_expr)
526 }
527 }
528}
529
530fn function_may_produce_hydratable_value(name: &str, args: &[ResolvedExpr]) -> bool {
531 match name.to_ascii_lowercase().as_str() {
532 "nodes" | "relationships" | "head" | "last" => true,
533 "coalesce" | "collect" => args.iter().any(expr_may_produce_hydratable_value),
534 "tail" | "reverse" => args.first().is_some_and(expr_may_produce_hydratable_value),
535 _ => false,
536 }
537}
538
539pub(super) fn expand_rows<S: GraphStorage>(
540 storage: &S,
541 params: &BTreeMap<String, LoraValue>,
542 input_rows: Vec<Row>,
543 op: &ExpandExec,
544) -> ExecResult<Vec<Row>> {
545 let eval_ctx = EvalContext { storage, params };
546 let mut out = Vec::new();
547
548 for row in input_rows {
549 let Some(src_node_id) = bound_node_id_for_expand(&row, op.src)? else {
550 continue;
551 };
552
553 let mut rel_property_filter = None;
554
555 storage.try_for_each_expand_id(
556 src_node_id,
557 op.direction,
558 &op.types,
559 |rel_id, dst_id| {
560 if let Some(expr) = op.rel_properties.as_ref() {
561 if rel_property_filter.is_none() {
562 let expected = eval_expr(expr, &row, &eval_ctx);
563 let LoraValue::Map(map) = expected else {
564 return Err(ExecutorError::ExpectedPropertyMap {
565 found: value_kind(&expected),
566 });
567 };
568 rel_property_filter = Some(map);
569 }
570
571 let Some(map) = rel_property_filter.as_ref() else {
572 return Ok(());
573 };
574 let matches = storage
575 .with_relationship(rel_id, |rel| {
576 map.iter().all(|(key, expected)| {
577 rel.properties
578 .get(key)
579 .map(|actual| value_matches_property_value(expected, actual))
580 .unwrap_or(false)
581 })
582 })
583 .unwrap_or(false);
584 if !matches {
585 return Ok(());
586 }
587 }
588
589 if let Some(existing_id) = bound_node_id_for_expand(&row, op.dst)? {
590 if existing_id != dst_id {
591 return Ok(());
592 }
593 }
594
595 if let Some(rel_var) = op.rel {
596 if let Some(existing_id) = bound_relationship_id_for_expand(&row, rel_var)? {
597 if existing_id != rel_id {
598 return Ok(());
599 }
600 }
601 }
602
603 let mut new_row = row.clone();
604 if !new_row.contains_key(op.dst) {
605 new_row.insert(op.dst, LoraValue::Node(dst_id));
606 }
607 if let Some(rel_var) = op.rel {
608 if !new_row.contains_key(rel_var) {
609 new_row.insert(rel_var, LoraValue::Relationship(rel_id));
610 }
611 }
612 out.push(new_row);
613 Ok(())
614 },
615 )?;
616 }
617
618 Ok(out)
619}
620
621pub(super) fn expand_var_len_rows<S: GraphStorage>(
622 storage: &S,
623 input_rows: Vec<Row>,
624 op: &ExpandExec,
625 range: &RangeLiteral,
626) -> ExecResult<Vec<Row>> {
627 let (min_hops, max_hops) = resolve_range(range);
628 let bind_relationships = op.rel.is_some();
629 let mut out = Vec::new();
630
631 for row in input_rows {
632 let Some(src_node_id) = bound_node_id_for_expand(&row, op.src)? else {
633 continue;
634 };
635
636 let expansions = variable_length_expand(
637 storage,
638 src_node_id,
639 op.direction,
640 &op.types,
641 min_hops,
642 max_hops,
643 bind_relationships,
644 );
645
646 for result in expansions {
647 let mut new_row = row.clone();
648 new_row.insert(op.dst, LoraValue::Node(result.dst_node_id));
649
650 if let Some(rel_var) = op.rel {
651 let rel_list = LoraValue::List(
652 result
653 .rel_ids
654 .into_iter()
655 .map(LoraValue::Relationship)
656 .collect(),
657 );
658 new_row.insert(rel_var, rel_list);
659 }
660
661 out.push(new_row);
662 }
663 }
664
665 Ok(out)
666}
667
668pub(super) fn properties_to_value_map(props: &Properties) -> LoraValue {
669 let mut map = BTreeMap::new();
670 for (k, v) in props.iter() {
671 map.insert(k.clone(), LoraValue::from(v));
672 }
673 LoraValue::Map(map)
674}
675
676pub(crate) fn dedup_rows_by_vars(rows: Vec<Row>) -> Vec<Row> {
680 let mut seen: BTreeSet<Vec<GroupValueKey>> = BTreeSet::new();
681 let mut out = Vec::new();
682
683 for row in rows {
684 let key: Vec<GroupValueKey> = row
685 .iter()
686 .map(|(_, val)| GroupValueKey::from_value(val))
687 .collect();
688 if seen.insert(key) {
689 out.push(row);
690 }
691 }
692
693 out
694}
695
696pub(crate) fn dedup_rows(rows: Vec<Row>) -> Vec<Row> {
700 let mut seen: BTreeSet<Vec<(String, GroupValueKey)>> = BTreeSet::new();
701 let mut out = Vec::new();
702
703 for row in rows {
704 let key: Vec<(String, GroupValueKey)> = row
705 .iter_named()
706 .map(|(_, name, val)| (name.into_owned(), GroupValueKey::from_value(val)))
707 .collect();
708 if seen.insert(key) {
709 out.push(row);
710 }
711 }
712
713 out
714}
715
716pub(super) fn eval_properties_expr<S: GraphStorage>(
717 expr: &ResolvedExpr,
718 row: &Row,
719 storage: &S,
720 params: &BTreeMap<String, LoraValue>,
721) -> ExecResult<Properties> {
722 let eval_ctx = EvalContext { storage, params };
723
724 match eval_expr(expr, row, &eval_ctx) {
725 LoraValue::Map(map) => {
726 let mut out = Properties::new();
727 for (k, v) in map {
728 let prop = lora_value_to_property(v)
729 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
730 out.insert(k, prop);
731 }
732 Ok(out)
733 }
734 other => Err(ExecutorError::ExpectedPropertyMap {
735 found: value_kind(&other),
736 }),
737 }
738}
739
740pub(crate) fn compute_aggregate_expr<S: GraphStorage>(
741 expr: &ResolvedExpr,
742 rows: &[Row],
743 eval_ctx: &EvalContext<'_, S>,
744) -> ExecResult<LoraValue> {
745 match expr {
746 ResolvedExpr::Function {
747 name,
748 distinct,
749 args,
750 } => {
751 let func = name.to_ascii_lowercase();
752
753 match func.as_str() {
754 "count" => {
755 if args.is_empty() {
756 return Ok(LoraValue::Int(rows.len() as i64));
757 }
758
759 let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
760 values.retain(|v| !matches!(v, LoraValue::Null));
761
762 if *distinct {
763 values = dedup_values(values);
764 }
765
766 Ok(LoraValue::Int(values.len() as i64))
767 }
768
769 "collect" => {
770 if args.is_empty() {
771 return Ok(LoraValue::List(Vec::new()));
772 }
773
774 let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
775
776 if *distinct {
777 values = dedup_values(values);
778 }
779
780 Ok(LoraValue::List(values))
781 }
782
783 "sum" => {
784 if args.is_empty() {
785 return Ok(LoraValue::Null);
786 }
787
788 let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
789
790 if *distinct {
791 values = dedup_values(values);
792 }
793
794 let nums = values
795 .into_iter()
796 .filter_map(as_f64_lossy)
797 .collect::<Vec<_>>();
798
799 if nums.is_empty() {
800 Ok(LoraValue::Null)
801 } else if nums.iter().all(|n| n.fract() == 0.0) {
802 Ok(LoraValue::Int(nums.iter().sum::<f64>() as i64))
803 } else {
804 Ok(LoraValue::Float(nums.iter().sum::<f64>()))
805 }
806 }
807
808 "avg" => {
809 if args.is_empty() {
810 return Ok(LoraValue::Null);
811 }
812
813 let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
814
815 if *distinct {
816 values = dedup_values(values);
817 }
818
819 let nums = values
820 .into_iter()
821 .filter_map(as_f64_lossy)
822 .collect::<Vec<_>>();
823
824 if nums.is_empty() {
825 Ok(LoraValue::Null)
826 } else {
827 Ok(LoraValue::Float(
828 nums.iter().sum::<f64>() / nums.len() as f64,
829 ))
830 }
831 }
832
833 "min" => {
834 if args.is_empty() {
835 return Ok(LoraValue::Null);
836 }
837
838 let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
839 values.retain(|v| !matches!(v, LoraValue::Null));
840
841 if *distinct {
842 values = dedup_values(values);
843 }
844
845 Ok(values
846 .into_iter()
847 .min_by(compare_values_total)
848 .unwrap_or(LoraValue::Null))
849 }
850
851 "max" => {
852 if args.is_empty() {
853 return Ok(LoraValue::Null);
854 }
855
856 let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
857 values.retain(|v| !matches!(v, LoraValue::Null));
858
859 if *distinct {
860 values = dedup_values(values);
861 }
862
863 Ok(values
864 .into_iter()
865 .max_by(compare_values_total)
866 .unwrap_or(LoraValue::Null))
867 }
868
869 "stdev" | "stdevp" => {
870 if args.is_empty() {
871 return Ok(LoraValue::Null);
872 }
873
874 let nums: Vec<f64> = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?
875 .into_iter()
876 .filter_map(as_f64_lossy)
877 .collect();
878
879 let is_population = func == "stdevp";
880
881 if nums.is_empty() || (!is_population && nums.len() < 2) {
882 return Ok(LoraValue::Float(0.0));
883 }
884
885 let mean = nums.iter().sum::<f64>() / nums.len() as f64;
886 let variance_sum: f64 = nums.iter().map(|x| (x - mean).powi(2)).sum();
887 let denom = if is_population {
888 nums.len() as f64
889 } else {
890 (nums.len() - 1) as f64
891 };
892 Ok(LoraValue::Float((variance_sum / denom).sqrt()))
893 }
894
895 "percentilecont" => {
896 if args.len() < 2 {
897 return Ok(LoraValue::Null);
898 }
899
900 let Some(first) = rows.first() else {
901 return Ok(LoraValue::Null);
902 };
903
904 let percentile = eval_expr_result(&args[1], first, eval_ctx)
905 .map_err(ExecutorError::RuntimeError)?
906 .as_f64()
907 .map(normalize_percentile)
908 .unwrap_or(0.5);
909 let mut nums: Vec<f64> = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?
910 .into_iter()
911 .filter_map(as_f64_lossy)
912 .collect();
913
914 if nums.is_empty() {
915 return Ok(LoraValue::Null);
916 }
917
918 nums.sort_by(|a, b| a.partial_cmp(b).unwrap_or(Ordering::Equal));
919
920 let index = percentile * (nums.len() - 1) as f64;
921 let lower = index.floor() as usize;
922 let upper = index.ceil() as usize;
923 let fraction = index - lower as f64;
924
925 if lower == upper || upper >= nums.len() {
926 Ok(LoraValue::Float(nums[lower]))
927 } else {
928 Ok(LoraValue::Float(
929 nums[lower] * (1.0 - fraction) + nums[upper] * fraction,
930 ))
931 }
932 }
933
934 "percentiledisc" => {
935 if args.len() < 2 {
936 return Ok(LoraValue::Null);
937 }
938
939 let Some(first) = rows.first() else {
940 return Ok(LoraValue::Null);
941 };
942
943 let percentile = eval_expr_result(&args[1], first, eval_ctx)
944 .map_err(ExecutorError::RuntimeError)?
945 .as_f64()
946 .map(normalize_percentile)
947 .unwrap_or(0.5);
948 let mut nums: Vec<f64> = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?
949 .into_iter()
950 .filter_map(as_f64_lossy)
951 .collect();
952
953 if nums.is_empty() {
954 return Ok(LoraValue::Null);
955 }
956
957 nums.sort_by(|a, b| a.partial_cmp(b).unwrap_or(Ordering::Equal));
958
959 let index = (percentile * (nums.len() - 1) as f64).round() as usize;
960 let index = index.min(nums.len() - 1);
961 Ok(LoraValue::Float(nums[index]))
962 }
963
964 _ => eval_first_or_null(expr, rows, eval_ctx),
965 }
966 }
967
968 _ => eval_first_or_null(expr, rows, eval_ctx),
969 }
970}
971
972fn eval_aggregate_arg_values<S: GraphStorage>(
973 expr: &ResolvedExpr,
974 rows: &[Row],
975 eval_ctx: &EvalContext<'_, S>,
976) -> ExecResult<Vec<LoraValue>> {
977 rows.iter()
978 .map(|row| eval_expr_result(expr, row, eval_ctx).map_err(ExecutorError::RuntimeError))
979 .collect()
980}
981
982fn normalize_percentile(value: f64) -> f64 {
983 if value.is_finite() {
984 value.clamp(0.0, 1.0)
985 } else {
986 0.5
987 }
988}
989
990fn eval_first_or_null<S: GraphStorage>(
991 expr: &ResolvedExpr,
992 rows: &[Row],
993 eval_ctx: &EvalContext<'_, S>,
994) -> ExecResult<LoraValue> {
995 match rows.first() {
996 Some(row) => eval_expr_result(expr, row, eval_ctx).map_err(ExecutorError::RuntimeError),
997 None => Ok(LoraValue::Null),
998 }
999}
1000
1001fn dedup_values(values: Vec<LoraValue>) -> Vec<LoraValue> {
1002 let mut seen: BTreeSet<GroupValueKey> = BTreeSet::new();
1003 let mut out = Vec::new();
1004
1005 for value in values {
1006 let key = GroupValueKey::from_value(&value);
1007 if seen.insert(key) {
1008 out.push(value);
1009 }
1010 }
1011
1012 out
1013}
1014
1015fn as_f64_lossy(v: LoraValue) -> Option<f64> {
1016 match v {
1017 LoraValue::Int(i) => Some(i as f64),
1018 LoraValue::Float(f) => Some(f),
1019 _ => None,
1020 }
1021}
1022
1023pub(super) fn compare_values_total(a: &LoraValue, b: &LoraValue) -> Ordering {
1024 use LoraValue::*;
1025
1026 match (a, b) {
1027 (Bool(x), Bool(y)) => x.cmp(y),
1028 (Int(x), Int(y)) => x.cmp(y),
1029 (Float(x), Float(y)) => x.partial_cmp(y).unwrap_or(Ordering::Equal),
1030 (Int(x), Float(y)) => (*x as f64).partial_cmp(y).unwrap_or(Ordering::Equal),
1031 (Float(x), Int(y)) => x.partial_cmp(&(*y as f64)).unwrap_or(Ordering::Equal),
1032 (String(x), String(y)) => x.cmp(y),
1033 (Binary(x), Binary(y)) => x.segments().cmp(y.segments()),
1034 (Node(x), Node(y)) => x.cmp(y),
1035 (Relationship(x), Relationship(y)) => x.cmp(y),
1036 (Date(x), Date(y)) => x.cmp(y),
1037 (DateTime(x), DateTime(y)) => x.cmp(y),
1038 (Duration(x), Duration(y)) => x.cmp(y),
1039 (Vector(x), Vector(y)) => x.to_key_string().cmp(&y.to_key_string()),
1040 _ => type_rank(a)
1041 .cmp(&type_rank(b))
1042 .then_with(|| format!("{a:?}").cmp(&format!("{b:?}"))),
1043 }
1044}
1045
1046pub fn value_matches_property_value(expected: &LoraValue, actual: &PropertyValue) -> bool {
1047 match (expected, actual) {
1048 (LoraValue::Null, PropertyValue::Null) => true,
1049 (LoraValue::Bool(a), PropertyValue::Bool(b)) => a == b,
1050 (LoraValue::Int(a), PropertyValue::Int(b)) => a == b,
1051 (LoraValue::Float(a), PropertyValue::Float(b)) => a == b,
1052 (LoraValue::Int(a), PropertyValue::Float(b)) => (*a as f64) == *b,
1053 (LoraValue::Float(a), PropertyValue::Int(b)) => *a == (*b as f64),
1054 (LoraValue::String(a), PropertyValue::String(b)) => a == b,
1055 (LoraValue::Binary(a), PropertyValue::Binary(b)) => a == b,
1056
1057 (LoraValue::List(xs), PropertyValue::List(ys)) => {
1058 xs.len() == ys.len()
1059 && xs
1060 .iter()
1061 .zip(ys.iter())
1062 .all(|(x, y)| value_matches_property_value(x, y))
1063 }
1064
1065 (LoraValue::Map(xm), PropertyValue::Map(ym)) => xm.iter().all(|(k, xv)| {
1066 ym.get(k)
1067 .map(|yv| value_matches_property_value(xv, yv))
1068 .unwrap_or(false)
1069 }),
1070
1071 (LoraValue::Date(a), PropertyValue::Date(b)) => a == b,
1072 (LoraValue::DateTime(a), PropertyValue::DateTime(b)) => a == b,
1073 (LoraValue::LocalDateTime(a), PropertyValue::LocalDateTime(b)) => a == b,
1074 (LoraValue::Time(a), PropertyValue::Time(b)) => a == b,
1075 (LoraValue::LocalTime(a), PropertyValue::LocalTime(b)) => a == b,
1076 (LoraValue::Duration(a), PropertyValue::Duration(b)) => a == b,
1077 (LoraValue::Point(a), PropertyValue::Point(b)) => a == b,
1078 (LoraValue::Vector(a), PropertyValue::Vector(b)) => a == b,
1079
1080 _ => false,
1081 }
1082}
1083
1084pub(crate) fn node_matches_property_filter<S: GraphStorage>(
1085 storage: &S,
1086 node_id: NodeId,
1087 labels: &[Vec<String>],
1088 key: &str,
1089 expected: &LoraValue,
1090) -> bool {
1091 storage
1092 .with_node(node_id, |node| {
1093 node_matches_label_groups(&node.labels, labels)
1094 && node
1095 .properties
1096 .get(key)
1097 .map(|actual| value_matches_property_value(expected, actual))
1098 .unwrap_or(false)
1099 })
1100 .unwrap_or(false)
1101}
1102
1103fn single_label_hint(labels: &[Vec<String>]) -> Option<&str> {
1104 if labels.len() == 1 && labels[0].len() == 1 {
1105 Some(labels[0][0].as_str())
1106 } else {
1107 None
1108 }
1109}
1110
1111fn property_lookup_values(expected: &LoraValue) -> Option<Vec<PropertyValue>> {
1112 let property = lora_value_to_property(expected.clone()).ok()?;
1113 let mut values = vec![property.clone()];
1114
1115 match property {
1116 PropertyValue::Int(i) => {
1117 values.push(PropertyValue::Float(i as f64));
1118 }
1119 PropertyValue::Float(f)
1120 if f.is_finite()
1121 && f.fract() == 0.0
1122 && f >= i64::MIN as f64
1123 && f <= i64::MAX as f64 =>
1124 {
1125 values.push(PropertyValue::Int(f as i64));
1126 }
1127 _ => {}
1128 }
1129
1130 Some(values)
1131}
1132
1133pub(crate) struct NodePropertyCandidates {
1134 pub(crate) ids: Vec<NodeId>,
1135 pub(crate) prefiltered: bool,
1136}
1137
1138pub(crate) fn node_by_property_range_scan_rows<S: GraphStorage>(
1139 storage: &S,
1140 params: &BTreeMap<String, LoraValue>,
1141 base_rows: Vec<Row>,
1142 op: &lora_compiler::NodeByPropertyRangeScanExec,
1143 deadline: Option<Instant>,
1144) -> ExecResult<Vec<Row>> {
1145 let eval_ctx = EvalContext { storage, params };
1146 let mut out = Vec::new();
1147
1148 for row in base_rows {
1149 check_optional_deadline(deadline)?;
1150 let lo_value = op.lo.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
1151 let hi_value = op.hi.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
1152 let lo_prop = lo_value
1153 .clone()
1154 .and_then(|v| lora_value_to_property(v).ok());
1155 let hi_prop = hi_value
1156 .clone()
1157 .and_then(|v| lora_value_to_property(v).ok());
1158 let filter = NodeRangeFilter {
1159 labels: &op.labels,
1160 key: &op.key,
1161 lo: lo_value.as_ref(),
1162 lo_inclusive: op.lo_inclusive,
1163 hi: hi_value.as_ref(),
1164 hi_inclusive: op.hi_inclusive,
1165 };
1166
1167 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
1168 if node_matches_range_filter(storage, existing_id, &filter) {
1169 out.push(row);
1170 }
1171 continue;
1172 }
1173
1174 let candidate_ids = match single_label_hint(&op.labels) {
1175 Some(label) => storage
1176 .node_range_candidates(label, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
1177 .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1178 None => scan_node_ids_for_label_groups(storage, &op.labels),
1179 };
1180
1181 for id in candidate_ids {
1182 check_optional_deadline(deadline)?;
1183 if node_matches_range_filter(storage, id, &filter) {
1184 let mut new_row = row.clone();
1185 new_row.insert(op.var, LoraValue::Node(id));
1186 out.push(new_row);
1187 }
1188 }
1189 }
1190
1191 Ok(out)
1192}
1193
1194pub(crate) fn node_by_text_scan_rows<S: GraphStorage>(
1195 storage: &S,
1196 params: &BTreeMap<String, LoraValue>,
1197 base_rows: Vec<Row>,
1198 op: &lora_compiler::NodeByTextScanExec,
1199 deadline: Option<Instant>,
1200) -> ExecResult<Vec<Row>> {
1201 let eval_ctx = EvalContext { storage, params };
1202 let mut out = Vec::new();
1203
1204 for row in base_rows {
1205 check_optional_deadline(deadline)?;
1206 let query = eval_expr(&op.query, &row, &eval_ctx);
1207 let LoraValue::String(query_str) = &query else {
1208 continue;
1211 };
1212
1213 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
1214 if node_matches_text_filter(
1215 storage,
1216 existing_id,
1217 &op.labels,
1218 &op.key,
1219 op.predicate,
1220 query_str,
1221 ) {
1222 out.push(row);
1223 }
1224 continue;
1225 }
1226
1227 let candidate_ids = match single_label_hint(&op.labels) {
1228 Some(label) => storage
1229 .node_text_candidates(label, &op.key, query_str)
1230 .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1231 None => scan_node_ids_for_label_groups(storage, &op.labels),
1232 };
1233
1234 for id in candidate_ids {
1235 check_optional_deadline(deadline)?;
1236 if node_matches_text_filter(storage, id, &op.labels, &op.key, op.predicate, query_str) {
1237 let mut new_row = row.clone();
1238 new_row.insert(op.var, LoraValue::Node(id));
1239 out.push(new_row);
1240 }
1241 }
1242 }
1243
1244 Ok(out)
1245}
1246
1247struct NodeRangeFilter<'a> {
1248 labels: &'a [Vec<String>],
1249 key: &'a str,
1250 lo: Option<&'a LoraValue>,
1251 lo_inclusive: bool,
1252 hi: Option<&'a LoraValue>,
1253 hi_inclusive: bool,
1254}
1255
1256fn node_matches_range_filter<S: GraphStorage>(
1257 storage: &S,
1258 id: NodeId,
1259 filter: &NodeRangeFilter<'_>,
1260) -> bool {
1261 storage
1262 .with_node(id, |n| {
1263 if !node_matches_label_groups(&n.labels, filter.labels) {
1264 return false;
1265 }
1266 let Some(actual) = n.properties.get(filter.key) else {
1267 return false;
1268 };
1269 let actual_lv = lora_store_property_to_value(actual);
1270 range_predicate_holds(
1271 &actual_lv,
1272 filter.lo,
1273 filter.lo_inclusive,
1274 filter.hi,
1275 filter.hi_inclusive,
1276 )
1277 })
1278 .unwrap_or(false)
1279}
1280
1281fn node_matches_text_filter<S: GraphStorage>(
1282 storage: &S,
1283 id: NodeId,
1284 labels: &[Vec<String>],
1285 key: &str,
1286 predicate: lora_compiler::TextPredicate,
1287 query: &str,
1288) -> bool {
1289 storage
1290 .with_node(id, |n| {
1291 if !node_matches_label_groups(&n.labels, labels) {
1292 return false;
1293 }
1294 let Some(PropertyValue::String(actual)) = n.properties.get(key) else {
1295 return false;
1296 };
1297 text_predicate_holds(actual, predicate, query)
1298 })
1299 .unwrap_or(false)
1300}
1301
1302fn text_predicate_holds(
1303 actual: &str,
1304 predicate: lora_compiler::TextPredicate,
1305 query: &str,
1306) -> bool {
1307 match predicate {
1308 lora_compiler::TextPredicate::StartsWith => actual.starts_with(query),
1309 lora_compiler::TextPredicate::EndsWith => actual.ends_with(query),
1310 lora_compiler::TextPredicate::Contains => actual.contains(query),
1311 }
1312}
1313
1314fn range_predicate_holds(
1315 actual: &LoraValue,
1316 lo: Option<&LoraValue>,
1317 lo_inclusive: bool,
1318 hi: Option<&LoraValue>,
1319 hi_inclusive: bool,
1320) -> bool {
1321 if let Some(lo) = lo {
1322 match range_comparison(actual, lo) {
1323 None => return false,
1324 Some(Ordering::Less) => return false,
1325 Some(Ordering::Equal) if !lo_inclusive => return false,
1326 _ => {}
1327 }
1328 }
1329 if let Some(hi) = hi {
1330 match range_comparison(actual, hi) {
1331 None => return false,
1332 Some(Ordering::Greater) => return false,
1333 Some(Ordering::Equal) if !hi_inclusive => return false,
1334 _ => {}
1335 }
1336 }
1337 true
1338}
1339
1340fn range_comparison(actual: &LoraValue, bound: &LoraValue) -> Option<Ordering> {
1341 match (actual, bound) {
1342 (LoraValue::Null, _) | (_, LoraValue::Null) => None,
1343 (LoraValue::String(a), LoraValue::String(b)) => Some(a.cmp(b)),
1344 (LoraValue::Date(a), LoraValue::Date(b)) => Some(a.to_epoch_days().cmp(&b.to_epoch_days())),
1345 (LoraValue::DateTime(a), LoraValue::DateTime(b)) => {
1346 Some(a.to_epoch_millis().cmp(&b.to_epoch_millis()))
1347 }
1348 (LoraValue::Duration(a), LoraValue::Duration(b)) => a
1349 .total_seconds_approx()
1350 .partial_cmp(&b.total_seconds_approx()),
1351 _ => actual.as_f64()?.partial_cmp(&bound.as_f64()?),
1352 }
1353}
1354
1355fn lora_store_property_to_value(value: &PropertyValue) -> LoraValue {
1356 LoraValue::from(value)
1357}
1358
1359pub(crate) fn node_by_point_scan_rows<S: GraphStorage>(
1360 storage: &S,
1361 params: &BTreeMap<String, LoraValue>,
1362 base_rows: Vec<Row>,
1363 op: &lora_compiler::NodeByPointScanExec,
1364 deadline: Option<Instant>,
1365) -> ExecResult<Vec<Row>> {
1366 let eval_ctx = EvalContext { storage, params };
1367 let mut out = Vec::new();
1368
1369 for row in base_rows {
1370 check_optional_deadline(deadline)?;
1371
1372 let probe = match &op.predicate {
1376 lora_compiler::PointPredicate::WithinBBox {
1377 lower_left,
1378 upper_right,
1379 } => {
1380 let ll = eval_expr(lower_left, &row, &eval_ctx);
1381 let ur = eval_expr(upper_right, &row, &eval_ctx);
1382 match (ll, ur) {
1383 (LoraValue::Point(a), LoraValue::Point(b)) => {
1384 Probe::WithinBBox { ll: a, ur: b }
1385 }
1386 _ => continue,
1387 }
1388 }
1389 lora_compiler::PointPredicate::WithinDistance {
1390 center,
1391 max_distance,
1392 inclusive,
1393 } => {
1394 let c = eval_expr(center, &row, &eval_ctx);
1395 let d = eval_expr(max_distance, &row, &eval_ctx);
1396 match (c, d) {
1397 (LoraValue::Point(c), LoraValue::Float(d)) => Probe::WithinDistance {
1398 center: c,
1399 max: d,
1400 inclusive: *inclusive,
1401 },
1402 (LoraValue::Point(c), LoraValue::Int(d)) => Probe::WithinDistance {
1403 center: c,
1404 max: d as f64,
1405 inclusive: *inclusive,
1406 },
1407 _ => continue,
1408 }
1409 }
1410 };
1411
1412 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
1413 if node_matches_point_filter(storage, existing_id, &op.labels, &op.key, &probe) {
1414 out.push(row);
1415 }
1416 continue;
1417 }
1418
1419 let candidate_ids = match single_label_hint(&op.labels) {
1420 Some(label) => match &probe {
1421 Probe::WithinBBox { ll, ur } => storage
1422 .node_point_within_bbox(label, &op.key, (ll.x, ll.y), (ur.x, ur.y))
1423 .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1424 Probe::WithinDistance { center, max, .. } => storage
1425 .node_point_within_distance(label, &op.key, (center.x, center.y), *max)
1426 .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1427 },
1428 None => scan_node_ids_for_label_groups(storage, &op.labels),
1429 };
1430
1431 for id in candidate_ids {
1432 check_optional_deadline(deadline)?;
1433 if node_matches_point_filter(storage, id, &op.labels, &op.key, &probe) {
1434 let mut new_row = row.clone();
1435 new_row.insert(op.var, LoraValue::Node(id));
1436 out.push(new_row);
1437 }
1438 }
1439 }
1440
1441 Ok(out)
1442}
1443
1444fn node_matches_point_filter<S: GraphStorage>(
1445 storage: &S,
1446 id: NodeId,
1447 labels: &[Vec<String>],
1448 key: &str,
1449 probe: &Probe,
1450) -> bool {
1451 storage
1452 .with_node(id, |n| {
1453 if !node_matches_label_groups(&n.labels, labels) {
1454 return false;
1455 }
1456 let Some(PropertyValue::Point(point)) = n.properties.get(key) else {
1457 return false;
1458 };
1459 point_predicate_holds(point, probe)
1460 })
1461 .unwrap_or(false)
1462}
1463
1464enum Probe {
1465 WithinBBox {
1466 ll: lora_store::LoraPoint,
1467 ur: lora_store::LoraPoint,
1468 },
1469 WithinDistance {
1470 center: lora_store::LoraPoint,
1471 max: f64,
1472 inclusive: bool,
1473 },
1474}
1475
1476fn point_predicate_holds(actual: &lora_store::LoraPoint, probe: &Probe) -> bool {
1477 match probe {
1478 Probe::WithinBBox { ll, ur } => {
1479 if actual.srid != ll.srid || actual.srid != ur.srid {
1480 return false;
1481 }
1482 let in_x = actual.x >= ll.x.min(ur.x) && actual.x <= ll.x.max(ur.x);
1483 let in_y = actual.y >= ll.y.min(ur.y) && actual.y <= ll.y.max(ur.y);
1484 let in_z = match (actual.z, ll.z, ur.z) {
1485 (Some(pz), Some(lz), Some(uz)) => pz >= lz.min(uz) && pz <= lz.max(uz),
1486 (None, None, None) => true,
1487 _ => return false,
1488 };
1489 in_x && in_y && in_z
1490 }
1491 Probe::WithinDistance {
1492 center,
1493 max,
1494 inclusive,
1495 } => {
1496 let Some(d) = lora_store::point_distance(actual, center) else {
1497 return false;
1498 };
1499 if *inclusive {
1500 d <= *max
1501 } else {
1502 d < *max
1503 }
1504 }
1505 }
1506}
1507
1508pub(crate) fn indexed_node_property_candidates<S: GraphStorage>(
1509 storage: &S,
1510 labels: &[Vec<String>],
1511 key: &str,
1512 expected: &LoraValue,
1513) -> NodePropertyCandidates {
1514 let Some(values) = property_lookup_values(expected) else {
1515 return NodePropertyCandidates {
1516 ids: scan_node_ids_for_label_groups(storage, labels),
1517 prefiltered: false,
1518 };
1519 };
1520
1521 let label_hint = single_label_hint(labels);
1522 let mut seen = BTreeSet::new();
1523 let mut out = Vec::new();
1524 for value in values {
1525 for id in storage.find_node_ids_by_property(label_hint, key, &value) {
1526 if seen.insert(id) {
1527 out.push(id);
1528 }
1529 }
1530 }
1531 NodePropertyCandidates {
1532 ids: out,
1533 prefiltered: labels.is_empty() || label_hint.is_some(),
1534 }
1535}
1536
1537pub(crate) fn build_path_value<S: GraphStorage>(
1543 row: &Row,
1544 node_vars: &[VarId],
1545 rel_vars: &[VarId],
1546 storage: &S,
1547) -> LoraValue {
1548 let (raw_nodes, rels, has_var_len) = path_bindings(row, node_vars, rel_vars);
1549
1550 let nodes = if has_var_len && !rels.is_empty() && raw_nodes.len() == 2 {
1551 reconstruct_var_len_nodes(raw_nodes[0], &rels, storage)
1552 } else {
1553 raw_nodes
1554 };
1555
1556 LoraValue::Path(LoraPath { nodes, rels })
1557}
1558
1559#[inline]
1560fn path_bindings(
1561 row: &Row,
1562 node_vars: &[VarId],
1563 rel_vars: &[VarId],
1564) -> (Vec<NodeId>, Vec<RelationshipId>, bool) {
1565 let mut raw_nodes = Vec::new();
1566 let mut rels = Vec::new();
1567 let mut has_var_len = false;
1568
1569 for &nv in node_vars {
1570 match row.get(nv) {
1571 Some(LoraValue::Node(id)) => raw_nodes.push(*id),
1572 Some(LoraValue::List(items)) => {
1573 for item in items {
1574 if let LoraValue::Node(id) = item {
1575 raw_nodes.push(*id);
1576 }
1577 }
1578 }
1579 _ => {}
1580 }
1581 }
1582
1583 for &rv in rel_vars {
1584 match row.get(rv) {
1585 Some(LoraValue::Relationship(id)) => rels.push(*id),
1586 Some(LoraValue::List(items)) => {
1587 has_var_len = true;
1588 for item in items {
1589 if let LoraValue::Relationship(id) = item {
1590 rels.push(*id);
1591 }
1592 }
1593 }
1594 _ => {}
1595 }
1596 }
1597
1598 (raw_nodes, rels, has_var_len)
1599}
1600
1601#[inline]
1602fn reconstruct_var_len_nodes<S: GraphStorage>(
1603 start: NodeId,
1604 rels: &[RelationshipId],
1605 storage: &S,
1606) -> Vec<NodeId> {
1607 let mut ordered = Vec::with_capacity(rels.len() + 1);
1608 ordered.push(start);
1609 let mut current = start;
1610 for &rel_id in rels {
1611 if let Some((src, dst)) = storage.relationship_endpoints(rel_id) {
1612 let next = if src == current { dst } else { src };
1613 ordered.push(next);
1614 current = next;
1615 }
1616 }
1617 ordered
1618}
1619
1620fn type_rank(v: &LoraValue) -> u8 {
1621 match v {
1622 LoraValue::Null => 0,
1623 LoraValue::Bool(_) => 1,
1624 LoraValue::Int(_) | LoraValue::Float(_) => 2,
1625 LoraValue::String(_) => 3,
1626 LoraValue::Binary(_) => 4,
1627 LoraValue::Date(_) => 5,
1628 LoraValue::DateTime(_) => 6,
1629 LoraValue::LocalDateTime(_) => 7,
1630 LoraValue::Time(_) => 8,
1631 LoraValue::LocalTime(_) => 9,
1632 LoraValue::Duration(_) => 10,
1633 LoraValue::Point(_) => 11,
1634 LoraValue::Vector(_) => 12,
1635 LoraValue::List(_) => 13,
1636 LoraValue::Map(_) => 14,
1637 LoraValue::Node(_) => 15,
1638 LoraValue::Relationship(_) => 16,
1639 LoraValue::Path(_) => 17,
1640 }
1641}
1642
1643pub(crate) fn node_matches_label_groups(node_labels: &[String], groups: &[Vec<String>]) -> bool {
1647 groups
1648 .iter()
1649 .all(|group| group.iter().any(|l| node_labels.iter().any(|nl| nl == l)))
1650}
1651
1652pub(crate) fn scan_node_ids_for_label_groups<S: GraphStorage>(
1655 storage: &S,
1656 groups: &[Vec<String>],
1657) -> Vec<NodeId> {
1658 if groups.is_empty() {
1659 return storage.all_node_ids();
1660 }
1661 if groups.len() == 1 {
1662 return label_group_candidate_ids(storage, &groups[0]);
1663 }
1664
1665 let mut best: Option<Vec<NodeId>> = None;
1666 for group in groups {
1667 let ids = label_group_candidate_ids(storage, group);
1668 if ids.is_empty() {
1669 return Vec::new();
1670 }
1671 if best
1672 .as_ref()
1673 .map(|current| ids.len() < current.len())
1674 .unwrap_or(true)
1675 {
1676 best = Some(ids);
1677 }
1678 }
1679
1680 best.unwrap_or_default()
1681}
1682
1683pub(crate) fn label_group_candidates_prefiltered(groups: &[Vec<String>]) -> bool {
1684 groups.len() <= 1
1685}
1686
1687fn label_group_candidate_ids<S: GraphStorage>(storage: &S, group: &[String]) -> Vec<NodeId> {
1688 match group {
1689 [] => Vec::new(),
1690 [label] => storage.node_ids_by_label(label),
1691 labels => {
1692 let mut seen = BTreeSet::new();
1693 let mut out = Vec::new();
1694 for label in labels {
1695 for id in storage.node_ids_by_label(label) {
1696 if seen.insert(id) {
1697 out.push(id);
1698 }
1699 }
1700 }
1701 out
1702 }
1703 }
1704}
1705
1706pub(crate) fn hydrate_node_record(node: &lora_store::NodeRecord) -> LoraValue {
1707 let mut map = BTreeMap::new();
1708 map.insert("kind".to_string(), LoraValue::String("node".to_string()));
1709 map.insert("id".to_string(), LoraValue::Int(node.id as i64));
1710 map.insert(
1711 "labels".to_string(),
1712 LoraValue::List(
1713 node.labels
1714 .iter()
1715 .map(|s| LoraValue::String(s.clone()))
1716 .collect(),
1717 ),
1718 );
1719 map.insert(
1720 "properties".to_string(),
1721 properties_to_value_map(&node.properties),
1722 );
1723 LoraValue::Map(map)
1724}
1725
1726pub(crate) fn hydrate_relationship_record(rel: &lora_store::RelationshipRecord) -> LoraValue {
1727 let mut map = BTreeMap::new();
1728 map.insert(
1729 "kind".to_string(),
1730 LoraValue::String("relationship".to_string()),
1731 );
1732 map.insert("id".to_string(), LoraValue::Int(rel.id as i64));
1733 map.insert("startId".to_string(), LoraValue::Int(rel.src as i64));
1734 map.insert("endId".to_string(), LoraValue::Int(rel.dst as i64));
1735 map.insert("type".to_string(), LoraValue::String(rel.rel_type.clone()));
1736 map.insert(
1737 "properties".to_string(),
1738 properties_to_value_map(&rel.properties),
1739 );
1740 LoraValue::Map(map)
1741}
1742
1743pub(super) fn flatten_label_groups(groups: &[Vec<String>]) -> Vec<String> {
1746 groups.iter().flat_map(|g| g.iter().cloned()).collect()
1747}
1748
1749#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
1750pub(crate) enum GroupValueKey {
1751 Null,
1752 Bool(bool),
1753 Int(i64),
1754 Float(String),
1755 String(String),
1756 Binary(Vec<Vec<u8>>),
1757 List(Vec<GroupValueKey>),
1758 Map(Vec<(String, GroupValueKey)>),
1759 Node(u64),
1760 Relationship(u64),
1761}
1762
1763impl GroupValueKey {
1764 pub(crate) fn from_value(v: &LoraValue) -> Self {
1765 match v {
1766 LoraValue::Null => Self::Null,
1767 LoraValue::Bool(x) => Self::Bool(*x),
1768 LoraValue::Int(x) => Self::Int(*x),
1769 LoraValue::Float(x) => Self::Float(x.to_string()),
1770 LoraValue::String(x) => Self::String(x.clone()),
1771 LoraValue::Binary(x) => Self::Binary(x.segments().to_vec()),
1772 LoraValue::List(xs) => Self::List(xs.iter().map(Self::from_value).collect()),
1773 LoraValue::Map(m) => Self::Map(
1774 m.iter()
1775 .map(|(k, v)| (k.clone(), Self::from_value(v)))
1776 .collect(),
1777 ),
1778 LoraValue::Node(id) => Self::Node(*id),
1779 LoraValue::Relationship(id) => Self::Relationship(*id),
1780 LoraValue::Path(_) => Self::Null,
1781 LoraValue::Date(d) => Self::String(d.to_string()),
1783 LoraValue::DateTime(dt) => Self::String(dt.to_string()),
1784 LoraValue::LocalDateTime(dt) => Self::String(dt.to_string()),
1785 LoraValue::Time(t) => Self::String(t.to_string()),
1786 LoraValue::LocalTime(t) => Self::String(t.to_string()),
1787 LoraValue::Duration(dur) => Self::String(dur.to_string()),
1788 LoraValue::Point(p) => Self::String(p.to_string()),
1789 LoraValue::Vector(v) => Self::String(format!("vector:{}", v.to_key_string())),
1790 }
1791 }
1792}
1793
1794const MAX_VAR_LEN_HOPS: u64 = 100;
1806
1807pub(crate) fn resolve_range(range: &RangeLiteral) -> (u64, u64) {
1808 let min_hops = range.start.unwrap_or(1);
1809 let max_hops = range.end.unwrap_or(MAX_VAR_LEN_HOPS);
1810 (min_hops, max_hops)
1811}
1812
1813pub(crate) struct VarLenResult {
1815 pub(crate) dst_node_id: NodeId,
1817 pub(crate) rel_ids: Vec<u64>,
1819}
1820
1821pub(crate) fn variable_length_expand<S: GraphStorage>(
1828 storage: &S,
1829 start_node_id: NodeId,
1830 direction: Direction,
1831 types: &[String],
1832 min_hops: u64,
1833 max_hops: u64,
1834 bind_relationships: bool,
1835) -> Vec<VarLenResult> {
1836 let mut results = Vec::new();
1837
1838 let mut frontier: Vec<(NodeId, Vec<u64>)> = vec![(start_node_id, Vec::new())];
1840
1841 for depth in 1..=max_hops {
1842 let is_last_hop = depth == max_hops;
1846 let mut next_frontier: Vec<(NodeId, Vec<u64>)> = Vec::new();
1847
1848 for (current_node, rels_used) in &frontier {
1849 for (rel_id, neighbor_id) in storage.expand_ids(*current_node, direction, types) {
1852 if rels_used.contains(&rel_id) {
1855 continue;
1856 }
1857
1858 if is_last_hop {
1859 if depth >= min_hops {
1862 let mut rel_ids = Vec::with_capacity(rels_used.len() + 1);
1863 rel_ids.extend_from_slice(rels_used);
1864 rel_ids.push(rel_id);
1865 results.push(VarLenResult {
1866 dst_node_id: neighbor_id,
1867 rel_ids: if bind_relationships {
1868 rel_ids
1869 } else {
1870 Vec::new()
1871 },
1872 });
1873 }
1874 continue;
1875 }
1876
1877 let mut new_rels = Vec::with_capacity(rels_used.len() + 1);
1878 new_rels.extend_from_slice(rels_used);
1879 new_rels.push(rel_id);
1880
1881 if depth >= min_hops {
1882 results.push(VarLenResult {
1883 dst_node_id: neighbor_id,
1884 rel_ids: if bind_relationships {
1885 new_rels.clone()
1886 } else {
1887 Vec::new()
1888 },
1889 });
1890 }
1891
1892 next_frontier.push((neighbor_id, new_rels));
1893 }
1894 }
1895
1896 if is_last_hop || next_frontier.is_empty() {
1897 break;
1898 }
1899
1900 frontier = next_frontier;
1901 }
1902
1903 if min_hops == 0 {
1905 results.insert(
1906 0,
1907 VarLenResult {
1908 dst_node_id: start_node_id,
1909 rel_ids: Vec::new(),
1910 },
1911 );
1912 }
1913
1914 results
1915}
1916
1917pub(crate) fn filter_shortest_paths(rows: Vec<Row>, path_var: VarId, all: bool) -> Vec<Row> {
1920 if rows.is_empty() {
1921 return rows;
1922 }
1923
1924 let lengths: Vec<usize> = rows
1926 .iter()
1927 .map(|row| match row.get(path_var) {
1928 Some(LoraValue::Path(p)) => p.rels.len(),
1929 _ => usize::MAX,
1930 })
1931 .collect();
1932
1933 let min_len = lengths.iter().copied().min().unwrap_or(usize::MAX);
1934
1935 let mut result: Vec<Row> = rows
1936 .into_iter()
1937 .zip(lengths.iter())
1938 .filter(|(_, len)| **len == min_len)
1939 .map(|(row, _)| row)
1940 .collect();
1941
1942 if !all && result.len() > 1 {
1943 result.truncate(1);
1944 }
1945
1946 result
1947}
1948
1949fn emit_rel_rows(
1966 direction: Direction,
1967 src_var: VarId,
1968 rel_var: VarId,
1969 dst_var: VarId,
1970 rel: &lora_store::RelationshipRecord,
1971 base: &Row,
1972 out: &mut Vec<Row>,
1973) -> ExecResult<()> {
1974 match direction {
1975 Direction::Right => {
1976 emit_one_rel_row(
1977 src_var, rel_var, dst_var, rel.src, rel.id, rel.dst, base, out,
1978 )?;
1979 }
1980 Direction::Left => {
1981 emit_one_rel_row(
1982 src_var, rel_var, dst_var, rel.dst, rel.id, rel.src, base, out,
1983 )?;
1984 }
1985 Direction::Undirected => {
1986 emit_one_rel_row(
1987 src_var, rel_var, dst_var, rel.src, rel.id, rel.dst, base, out,
1988 )?;
1989 if rel.src != rel.dst {
1993 emit_one_rel_row(
1994 src_var, rel_var, dst_var, rel.dst, rel.id, rel.src, base, out,
1995 )?;
1996 }
1997 }
1998 }
1999 Ok(())
2000}
2001
2002#[allow(clippy::too_many_arguments)]
2003fn emit_one_rel_row(
2004 src_var: VarId,
2005 rel_var: VarId,
2006 dst_var: VarId,
2007 src_id: NodeId,
2008 rel_id: RelationshipId,
2009 dst_id: NodeId,
2010 base: &Row,
2011 out: &mut Vec<Row>,
2012) -> ExecResult<()> {
2013 let mut row = base.clone();
2014 if bind_node_value(&mut row, src_var, src_id)?
2015 && bind_relationship_value(&mut row, rel_var, rel_id)?
2016 && bind_node_value(&mut row, dst_var, dst_id)?
2017 {
2018 out.push(row);
2019 }
2020 Ok(())
2021}
2022
2023fn bind_node_value(row: &mut Row, var: VarId, id: NodeId) -> ExecResult<bool> {
2024 match row.get(var) {
2025 Some(LoraValue::Node(existing)) => Ok(*existing == id),
2026 Some(other) => Err(ExecutorError::ExpectedNodeForExpand {
2027 var: format!("{var:?}"),
2028 found: value_kind(other),
2029 }),
2030 None => {
2031 row.insert(var, LoraValue::Node(id));
2032 Ok(true)
2033 }
2034 }
2035}
2036
2037fn bind_relationship_value(row: &mut Row, var: VarId, id: RelationshipId) -> ExecResult<bool> {
2038 match row.get(var) {
2039 Some(LoraValue::Relationship(existing)) => Ok(*existing == id),
2040 Some(other) => Err(ExecutorError::ExpectedRelationshipForExpand {
2041 var: format!("{var:?}"),
2042 found: value_kind(other),
2043 }),
2044 None => {
2045 row.insert(var, LoraValue::Relationship(id));
2046 Ok(true)
2047 }
2048 }
2049}
2050
2051fn rel_candidate_ids<S, F>(storage: &S, types: &[String], indexed: F) -> Vec<RelationshipId>
2052where
2053 S: GraphStorage,
2054 F: Fn(&str) -> Option<Vec<RelationshipId>>,
2055{
2056 if types.is_empty() {
2057 return storage.all_rel_ids();
2060 }
2061 let mut all = Vec::new();
2062 let mut seen = BTreeSet::new();
2063 for ty in types {
2064 let ids = indexed(ty).unwrap_or_else(|| storage.rel_ids_by_type(ty));
2065 for id in ids {
2066 if seen.insert(id) {
2067 all.push(id);
2068 }
2069 }
2070 }
2071 all
2072}
2073
2074pub(crate) fn rel_by_property_range_scan_rows<S: GraphStorage>(
2075 storage: &S,
2076 params: &BTreeMap<String, LoraValue>,
2077 base_rows: Vec<Row>,
2078 op: &lora_compiler::RelByPropertyRangeScanExec,
2079 deadline: Option<Instant>,
2080) -> ExecResult<Vec<Row>> {
2081 let eval_ctx = EvalContext { storage, params };
2082 let mut out = Vec::new();
2083
2084 for row in base_rows {
2085 check_optional_deadline(deadline)?;
2086 let lo_value = op.lo.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
2087 let hi_value = op.hi.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
2088 let lo_prop = lo_value
2089 .clone()
2090 .and_then(|v| lora_value_to_property(v).ok());
2091 let hi_prop = hi_value
2092 .clone()
2093 .and_then(|v| lora_value_to_property(v).ok());
2094
2095 let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| {
2096 storage.relationship_range_candidates(ty, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
2097 });
2098
2099 for rel_id in candidate_ids {
2100 check_optional_deadline(deadline)?;
2101 if let Some(result) = storage.with_relationship(rel_id, |rel| {
2102 if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
2103 return Ok(());
2104 }
2105 let Some(actual) = rel.properties.get(&op.key) else {
2106 return Ok(());
2107 };
2108 let actual_lv = LoraValue::from(actual);
2109 if !range_predicate_holds(
2110 &actual_lv,
2111 lo_value.as_ref(),
2112 op.lo_inclusive,
2113 hi_value.as_ref(),
2114 op.hi_inclusive,
2115 ) {
2116 return Ok(());
2117 }
2118 emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2119 }) {
2120 result?;
2121 }
2122 }
2123 }
2124
2125 Ok(out)
2126}
2127
2128pub(crate) fn rel_by_text_scan_rows<S: GraphStorage>(
2129 storage: &S,
2130 params: &BTreeMap<String, LoraValue>,
2131 base_rows: Vec<Row>,
2132 op: &lora_compiler::RelByTextScanExec,
2133 deadline: Option<Instant>,
2134) -> ExecResult<Vec<Row>> {
2135 let eval_ctx = EvalContext { storage, params };
2136 let mut out = Vec::new();
2137
2138 for row in base_rows {
2139 check_optional_deadline(deadline)?;
2140 let query = eval_expr(&op.query, &row, &eval_ctx);
2141 let LoraValue::String(query_str) = &query else {
2142 continue;
2143 };
2144
2145 let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| {
2146 storage.relationship_text_candidates(ty, &op.key, query_str)
2147 });
2148
2149 for rel_id in candidate_ids {
2150 check_optional_deadline(deadline)?;
2151 if let Some(result) = storage.with_relationship(rel_id, |rel| {
2152 if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
2153 return Ok(());
2154 }
2155 let Some(PropertyValue::String(actual)) = rel.properties.get(&op.key) else {
2156 return Ok(());
2157 };
2158 if !text_predicate_holds(actual, op.predicate, query_str) {
2159 return Ok(());
2160 }
2161 emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2162 }) {
2163 result?;
2164 }
2165 }
2166 }
2167
2168 Ok(out)
2169}
2170
2171pub(crate) fn rel_by_point_scan_rows<S: GraphStorage>(
2172 storage: &S,
2173 params: &BTreeMap<String, LoraValue>,
2174 base_rows: Vec<Row>,
2175 op: &lora_compiler::RelByPointScanExec,
2176 deadline: Option<Instant>,
2177) -> ExecResult<Vec<Row>> {
2178 let eval_ctx = EvalContext { storage, params };
2179 let mut out = Vec::new();
2180
2181 for row in base_rows {
2182 check_optional_deadline(deadline)?;
2183
2184 let probe = match &op.predicate {
2185 lora_compiler::PointPredicate::WithinBBox {
2186 lower_left,
2187 upper_right,
2188 } => {
2189 let ll = eval_expr(lower_left, &row, &eval_ctx);
2190 let ur = eval_expr(upper_right, &row, &eval_ctx);
2191 match (ll, ur) {
2192 (LoraValue::Point(a), LoraValue::Point(b)) => {
2193 Probe::WithinBBox { ll: a, ur: b }
2194 }
2195 _ => continue,
2196 }
2197 }
2198 lora_compiler::PointPredicate::WithinDistance {
2199 center,
2200 max_distance,
2201 inclusive,
2202 } => {
2203 let c = eval_expr(center, &row, &eval_ctx);
2204 let d = eval_expr(max_distance, &row, &eval_ctx);
2205 match (c, d) {
2206 (LoraValue::Point(c), LoraValue::Float(d)) => Probe::WithinDistance {
2207 center: c,
2208 max: d,
2209 inclusive: *inclusive,
2210 },
2211 (LoraValue::Point(c), LoraValue::Int(d)) => Probe::WithinDistance {
2212 center: c,
2213 max: d as f64,
2214 inclusive: *inclusive,
2215 },
2216 _ => continue,
2217 }
2218 }
2219 };
2220
2221 let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| match &probe {
2222 Probe::WithinBBox { ll, ur } => {
2223 storage.relationship_point_within_bbox(ty, &op.key, (ll.x, ll.y), (ur.x, ur.y))
2224 }
2225 Probe::WithinDistance { center, max, .. } => {
2226 storage.relationship_point_within_distance(ty, &op.key, (center.x, center.y), *max)
2227 }
2228 });
2229
2230 for rel_id in candidate_ids {
2231 check_optional_deadline(deadline)?;
2232 if let Some(result) = storage.with_relationship(rel_id, |rel| {
2233 if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
2234 return Ok(());
2235 }
2236 let Some(PropertyValue::Point(actual)) = rel.properties.get(&op.key) else {
2237 return Ok(());
2238 };
2239 if !point_predicate_holds(actual, &probe) {
2240 return Ok(());
2241 }
2242 emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2243 }) {
2244 result?;
2245 }
2246 }
2247 }
2248
2249 Ok(out)
2250}