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
67pub(super) fn check_deadline_at(deadline: Instant) -> ExecResult<()> {
68 if crate::cancel::deadline_reached(deadline) {
70 Err(ExecutorError::QueryTimeout)
71 } else {
72 Ok(())
73 }
74}
75
76#[inline]
79pub(crate) fn project_item<S: GraphStorage>(
80 projected: &mut Row,
81 source: &Row,
82 item: &lora_analyzer::ResolvedProjection,
83 eval_ctx: &EvalContext<'_, S>,
84) -> ExecResult<()> {
85 if let ResolvedExpr::Variable(var) = &item.expr {
86 if projected.insert_named_from(item.output, item.name.clone(), source, *var) {
87 return Ok(());
88 }
89 }
90 let value =
91 eval_expr_result(&item.expr, source, eval_ctx).map_err(ExecutorError::RuntimeError)?;
92 projected.insert_named(item.output, item.name.clone(), value);
93 Ok(())
94}
95
96#[inline]
98pub(crate) fn project_item_in_place<S: GraphStorage>(
99 row: &mut Row,
100 item: &lora_analyzer::ResolvedProjection,
101 eval_ctx: &EvalContext<'_, S>,
102) -> ExecResult<()> {
103 if let ResolvedExpr::Variable(var) = &item.expr {
104 if row.insert_named_from_self(item.output, item.name.clone(), *var) {
105 return Ok(());
106 }
107 }
108 let value = eval_expr_result(&item.expr, row, eval_ctx).map_err(ExecutorError::RuntimeError)?;
109 row.insert_named(item.output, item.name.clone(), value);
110 Ok(())
111}
112
113pub(super) fn filter_rows_checked<S: GraphStorage>(
114 input_rows: Vec<Row>,
115 predicate: &ResolvedExpr,
116 eval_ctx: &EvalContext<'_, S>,
117) -> ExecResult<Vec<Row>> {
118 let mut out = Vec::with_capacity(input_rows.len());
119 for row in input_rows {
120 if eval_truthy_result(predicate, &row, eval_ctx).map_err(ExecutorError::RuntimeError)? {
121 out.push(row);
122 }
123 }
124 Ok(out)
125}
126
127pub(super) fn project_rows_checked<S: GraphStorage>(
128 input_rows: Vec<Row>,
129 op: &ProjectionExec,
130 eval_ctx: &EvalContext<'_, S>,
131) -> ExecResult<Vec<Row>> {
132 let mut out = Vec::with_capacity(input_rows.len());
133
134 for row in input_rows {
135 if op.include_existing {
136 let mut projected = row;
137 for item in &op.items {
138 project_item_in_place(&mut projected, item, eval_ctx)?;
139 }
140 out.push(projected);
141 } else {
142 let mut projected = Row::new();
143 for item in &op.items {
144 project_item(&mut projected, &row, item, eval_ctx)?;
145 }
146 out.push(projected);
147 }
148 }
149
150 Ok(if op.distinct {
151 dedup_rows_by_vars(out)
152 } else {
153 out
154 })
155}
156
157pub(super) fn unwind_rows<S: GraphStorage>(
158 input_rows: Vec<Row>,
159 op: &UnwindExec,
160 eval_ctx: &EvalContext<'_, S>,
161) -> ExecResult<Vec<Row>> {
162 let mut out = Vec::with_capacity(input_rows.len());
163
164 for row in input_rows {
165 match eval_expr_result(&op.expr, &row, eval_ctx).map_err(ExecutorError::RuntimeError)? {
168 LoraValue::List(values) => {
169 let mut values = values.into_iter();
170 let last = values.next_back();
171 for value in values {
172 let mut new_row = row.clone();
173 new_row.insert(op.alias, value);
174 out.push(new_row);
175 }
176 if let Some(value) = last {
178 let mut new_row = row;
179 new_row.insert(op.alias, value);
180 out.push(new_row);
181 }
182 }
183 LoraValue::Null => {}
184 other => {
185 let mut new_row = row;
186 new_row.insert(op.alias, other);
187 out.push(new_row);
188 }
189 }
190 }
191
192 Ok(out)
193}
194
195pub(crate) fn eval_row_count<S: GraphStorage>(
200 clause: &str,
201 expr: &ResolvedExpr,
202 eval_ctx: &EvalContext<'_, S>,
203) -> ExecResult<usize> {
204 let value =
205 eval_expr_result(expr, &Row::new(), eval_ctx).map_err(ExecutorError::RuntimeError)?;
206 let n = match value {
207 LoraValue::Int(n) => Some(n),
208 LoraValue::Float(f) if f.fract() == 0.0 && f.is_finite() => Some(f as i64),
209 _ => None,
210 };
211 match n {
212 Some(n) if n >= 0 => Ok(n as usize),
213 _ => Err(ExecutorError::RuntimeError(format!(
214 "{clause} expects a non-negative integer, got {}",
215 describe_count(&value)
216 ))),
217 }
218}
219
220fn describe_count(value: &LoraValue) -> String {
221 match value {
222 LoraValue::Null => "null".to_string(),
223 LoraValue::Int(n) => n.to_string(),
224 LoraValue::Float(f) => f.to_string(),
225 other => value_kind(other),
226 }
227}
228
229pub(super) fn limit_rows<S: GraphStorage>(
230 mut rows: Vec<Row>,
231 op: &LimitExec,
232 eval_ctx: &EvalContext<'_, S>,
233) -> ExecResult<Vec<Row>> {
234 let limit = match op.limit.as_ref() {
235 Some(e) => eval_row_count("LIMIT", e, eval_ctx)?,
236 None => rows.len(),
237 };
238 let skip = match op.skip.as_ref() {
239 Some(e) => eval_row_count("SKIP", e, eval_ctx)?,
240 None => 0,
241 };
242
243 if skip >= rows.len() {
244 return Ok(Vec::new());
245 }
246
247 rows.drain(0..skip);
248 rows.truncate(limit);
249 Ok(rows)
250}
251
252#[inline]
253pub(crate) fn bound_node_id_for_expand(row: &Row, var: VarId) -> ExecResult<Option<NodeId>> {
254 match row.get(var) {
255 Some(LoraValue::Node(id)) => Ok(Some(*id)),
256 Some(other) => Err(ExecutorError::ExpectedNodeForExpand {
257 var: format!("{var:?}"),
258 found: value_kind(other),
259 }),
260 None => Ok(None),
261 }
262}
263
264#[inline]
265pub(crate) fn bound_relationship_id_for_expand(
266 row: &Row,
267 var: VarId,
268) -> ExecResult<Option<RelationshipId>> {
269 match row.get(var) {
270 Some(LoraValue::Relationship(id)) => Ok(Some(*id)),
271 Some(other) => Err(ExecutorError::ExpectedRelationshipForExpand {
272 var: format!("{var:?}"),
273 found: value_kind(other),
274 }),
275 None => Ok(None),
276 }
277}
278
279pub(super) fn node_scan_rows<S: GraphStorage>(
280 storage: &S,
281 base_rows: Vec<Row>,
282 op: &NodeScanExec,
283 deadline: Option<Instant>,
284) -> ExecResult<Vec<Row>> {
285 let node_ids = storage.all_node_ids();
286 let mut out = Vec::with_capacity(scan_output_capacity(base_rows.len(), node_ids.len()));
287
288 if deadline.is_none() {
289 for row in base_rows {
290 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
291 if storage.has_node(existing_id) {
292 out.push(row);
293 }
294 continue;
295 }
296
297 for &id in &node_ids {
298 let mut new_row = row.clone();
299 new_row.insert(op.var, LoraValue::Node(id));
300 out.push(new_row);
301 }
302 }
303 return Ok(out);
304 }
305
306 for row in base_rows {
307 check_optional_deadline(deadline)?;
308 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
309 if storage.has_node(existing_id) {
310 out.push(row);
311 }
312 continue;
313 }
314
315 for &id in &node_ids {
316 check_optional_deadline(deadline)?;
317 let mut new_row = row.clone();
318 new_row.insert(op.var, LoraValue::Node(id));
319 out.push(new_row);
320 }
321 }
322
323 Ok(out)
324}
325
326pub(super) fn node_by_label_scan_rows<S: GraphStorage>(
327 storage: &S,
328 base_rows: Vec<Row>,
329 op: &NodeByLabelScanExec,
330 deadline: Option<Instant>,
331) -> ExecResult<Vec<Row>> {
332 let candidate_ids = scan_node_ids_for_label_groups(storage, &op.labels);
333 let candidates_prefiltered = label_group_candidates_prefiltered(&op.labels);
334 let mut out = Vec::with_capacity(scan_output_capacity(base_rows.len(), candidate_ids.len()));
335
336 if deadline.is_none() {
337 for row in base_rows {
338 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
339 let labels_ok = storage
340 .with_node(existing_id, |n| {
341 node_matches_label_groups(&n.labels, &op.labels)
342 })
343 .unwrap_or(false);
344 if labels_ok {
345 out.push(row);
346 }
347 continue;
348 }
349
350 for &id in &candidate_ids {
351 if !candidates_prefiltered {
352 let labels_ok = storage
353 .with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
354 .unwrap_or(false);
355 if !labels_ok {
356 continue;
357 }
358 }
359 let mut new_row = row.clone();
360 new_row.insert(op.var, LoraValue::Node(id));
361 out.push(new_row);
362 }
363 }
364 return Ok(out);
365 }
366
367 for row in base_rows {
368 check_optional_deadline(deadline)?;
369 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
370 let labels_ok = storage
371 .with_node(existing_id, |n| {
372 node_matches_label_groups(&n.labels, &op.labels)
373 })
374 .unwrap_or(false);
375 if labels_ok {
376 out.push(row);
377 }
378 continue;
379 }
380
381 for &id in &candidate_ids {
382 check_optional_deadline(deadline)?;
383 if !candidates_prefiltered {
384 let labels_ok = storage
385 .with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
386 .unwrap_or(false);
387 if !labels_ok {
388 continue;
389 }
390 }
391 let mut new_row = row.clone();
392 new_row.insert(op.var, LoraValue::Node(id));
393 out.push(new_row);
394 }
395 }
396
397 Ok(out)
398}
399
400pub(super) fn node_by_property_scan_rows<S: GraphStorage>(
401 storage: &S,
402 params: &BTreeMap<String, LoraValue>,
403 base_rows: Vec<Row>,
404 op: &NodeByPropertyScanExec,
405 deadline: Option<Instant>,
406) -> ExecResult<Vec<Row>> {
407 let eval_ctx = EvalContext { storage, params };
408 let mut out = Vec::new();
409
410 for row in base_rows {
411 check_optional_deadline(deadline)?;
412 let expected = eval_expr(&op.value, &row, &eval_ctx);
413
414 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
415 if property_scan_matches(
416 storage,
417 existing_id,
418 &op.labels,
419 &op.key,
420 &expected,
421 op.in_list,
422 ) {
423 out.push(row);
424 }
425 continue;
426 }
427
428 let candidates =
429 property_scan_candidates(storage, &op.labels, &op.key, &expected, op.in_list);
430 for id in candidates.ids {
431 check_optional_deadline(deadline)?;
432 if !candidates.prefiltered
433 && !property_scan_matches(storage, id, &op.labels, &op.key, &expected, op.in_list)
434 {
435 continue;
436 }
437 let mut new_row = row.clone();
438 new_row.insert(op.var, LoraValue::Node(id));
439 out.push(new_row);
440 }
441 }
442
443 Ok(out)
444}
445
446#[inline]
447fn check_optional_deadline(deadline: Option<Instant>) -> ExecResult<()> {
448 match deadline {
449 Some(deadline) => check_deadline_at(deadline),
450 None => Ok(()),
451 }
452}
453
454pub(crate) fn plan_may_need_hydration(plan: &PhysicalPlan) -> bool {
455 op_may_need_hydration(plan, plan.root)
456}
457
458pub(crate) fn count_all_scan_aggregation_rows<S: GraphStorage>(
459 storage: &S,
460 plan: &PhysicalPlan,
461 op: &HashAggregationExec,
462) -> Option<Vec<Row>> {
463 if !op.group_by.is_empty() {
464 return None;
465 }
466 let specs = crate::pull::classify_streamable_aggregates(&op.aggregates)?;
467 let scan_var = scan_subtree_var(plan, op.input)?;
468 let counts_rows = specs.iter().all(|spec| match spec.kind {
471 crate::pull::StreamableAggKind::CountAll => true,
472 crate::pull::StreamableAggKind::CountField => {
473 matches!(&spec.arg, Some(ResolvedExpr::Variable(v)) if *v == scan_var)
474 }
475 _ => false,
476 });
477 if !counts_rows {
478 return None;
479 }
480
481 let count = count_rows_for_scan_subtree(storage, plan, op.input)? as i64;
482 let value = LoraValue::Int(count);
483 let mut row = Row::new();
484 for proj in &op.aggregates {
485 row.insert_named(proj.output, proj.name.clone(), value.clone());
486 }
487 Some(vec![row])
488}
489
490fn scan_subtree_var(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> Option<VarId> {
491 match &plan.nodes[node_id] {
492 PhysicalOp::NodeScan(op) => Some(op.var),
493 PhysicalOp::NodeByLabelScan(op) => Some(op.var),
494 _ => None,
495 }
496}
497
498fn count_rows_for_scan_subtree<S: GraphStorage>(
499 storage: &S,
500 plan: &PhysicalPlan,
501 node_id: PhysicalNodeId,
502) -> Option<usize> {
503 match &plan.nodes[node_id] {
504 PhysicalOp::NodeScan(op) if scan_input_is_argument(plan, op.input) => {
505 Some(storage.node_count())
506 }
507 PhysicalOp::NodeByLabelScan(op) if scan_input_is_argument(plan, op.input) => {
508 if let [group] = op.labels.as_slice() {
510 if let [label] = group.as_slice() {
511 return Some(storage.node_count_by_label(label));
512 }
513 }
514 let ids = scan_node_ids_for_label_groups(storage, &op.labels);
515 if label_group_candidates_prefiltered(&op.labels) {
516 return Some(ids.len());
517 }
518 Some(
519 ids.into_iter()
520 .filter(|&id| {
521 storage
522 .with_node(id, |node| {
523 node_matches_label_groups(&node.labels, &op.labels)
524 })
525 .unwrap_or(false)
526 })
527 .count(),
528 )
529 }
530 _ => None,
531 }
532}
533
534fn scan_input_is_argument(plan: &PhysicalPlan, input: Option<PhysicalNodeId>) -> bool {
535 match input {
536 None => true,
537 Some(id) => matches!(plan.nodes.get(id), Some(PhysicalOp::Argument(_))),
538 }
539}
540
541fn op_may_need_hydration(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
542 match &plan.nodes[node_id] {
543 PhysicalOp::Projection(op) if !op.include_existing => op
544 .items
545 .iter()
546 .any(|item| expr_may_produce_hydratable_value(&item.expr)),
547 PhysicalOp::HashAggregation(op) => {
548 op.group_by
549 .iter()
550 .any(|item| expr_may_produce_hydratable_value(&item.expr))
551 || op
552 .aggregates
553 .iter()
554 .any(|item| expr_may_produce_hydratable_value(&item.expr))
555 }
556 PhysicalOp::Sort(op) => op_may_need_hydration(plan, op.input),
557 PhysicalOp::Limit(op) => op_may_need_hydration(plan, op.input),
558 PhysicalOp::Filter(op) => op_may_need_hydration(plan, op.input),
559 PhysicalOp::Unwind(op) => op_may_need_hydration(plan, op.input),
560 PhysicalOp::PathBuild(op) => op_may_need_hydration(plan, op.input),
561 _ => true,
562 }
563}
564
565fn expr_may_produce_hydratable_value(expr: &ResolvedExpr) -> bool {
566 match expr {
567 ResolvedExpr::Variable(_) | ResolvedExpr::Parameter(_) => true,
568 ResolvedExpr::Literal(_)
569 | ResolvedExpr::Property { .. }
570 | ResolvedExpr::ExistsSubquery { .. }
571 | ResolvedExpr::Binary { .. }
572 | ResolvedExpr::Unary { .. }
573 | ResolvedExpr::ListPredicate { .. } => false,
574 ResolvedExpr::Function { function, args, .. } => {
575 function_may_produce_hydratable_value(*function, args)
576 }
577 ResolvedExpr::List(items) => items.iter().any(expr_may_produce_hydratable_value),
578 ResolvedExpr::Map(items) => items
579 .iter()
580 .any(|(_, value)| expr_may_produce_hydratable_value(value)),
581 ResolvedExpr::Case {
582 alternatives,
583 else_expr,
584 ..
585 } => {
586 alternatives
587 .iter()
588 .any(|(_, value)| expr_may_produce_hydratable_value(value))
589 || else_expr
590 .as_deref()
591 .is_some_and(expr_may_produce_hydratable_value)
592 }
593 ResolvedExpr::ListComprehension { map_expr, .. } => map_expr
594 .as_deref()
595 .map(expr_may_produce_hydratable_value)
596 .unwrap_or(true),
597 ResolvedExpr::Reduce { expr, .. } => expr_may_produce_hydratable_value(expr),
598 ResolvedExpr::MapProjection { selectors, .. } => selectors.iter().any(|selector| {
599 matches!(selector, ResolvedMapSelector::Literal(_, expr) if expr_may_produce_hydratable_value(expr))
600 }),
601 ResolvedExpr::Index { expr, .. } | ResolvedExpr::Slice { expr, .. } => {
602 expr_may_produce_hydratable_value(expr)
603 }
604 ResolvedExpr::PatternComprehension { map_expr, .. } => {
605 expr_may_produce_hydratable_value(map_expr)
606 }
607 }
608}
609
610fn function_may_produce_hydratable_value(function: FunctionId, args: &[ResolvedExpr]) -> bool {
611 match function.name() {
612 "path.nodes" | "path.edges" | "path.first" | "path.last" | "list.first" | "list.last" => {
613 true
614 }
615 "value.coalesce" | "value.first_non_null" | "collect" => {
616 args.iter().any(expr_may_produce_hydratable_value)
617 }
618 "list.rest" | "value.reverse" | "list.reverse" => {
619 args.first().is_some_and(expr_may_produce_hydratable_value)
620 }
621 _ => false,
622 }
623}
624
625pub(super) fn expand_rows<S: GraphStorage>(
626 storage: &S,
627 params: &BTreeMap<String, LoraValue>,
628 input_rows: Vec<Row>,
629 op: &ExpandExec,
630) -> ExecResult<Vec<Row>> {
631 let eval_ctx = EvalContext { storage, params };
632 let mut out = Vec::with_capacity(input_rows.len());
635 let mut matches: Vec<(u64, u64)> = Vec::new();
640
641 for row in input_rows {
642 let Some(src_node_id) = bound_node_id_for_expand(&row, op.src)? else {
643 continue;
644 };
645
646 let mut rel_property_filter = None;
647 matches.clear();
648
649 storage.try_for_each_expand_id(
650 src_node_id,
651 op.direction,
652 &op.types,
653 |rel_id, dst_id| {
654 if let Some(expr) = op.rel_properties.as_ref() {
655 if rel_property_filter.is_none() {
656 let expected = eval_expr(expr, &row, &eval_ctx);
657 let LoraValue::Map(map) = expected else {
658 return Err(ExecutorError::ExpectedPropertyMap {
659 found: value_kind(&expected),
660 });
661 };
662 rel_property_filter = Some(map);
663 }
664
665 let Some(map) = rel_property_filter.as_ref() else {
666 return Ok(());
667 };
668 let matches = storage
669 .with_relationship(rel_id, |rel| {
670 map.iter().all(|(key, expected)| {
671 rel.properties
672 .get(key.as_str())
673 .map(|actual| value_matches_property_value(expected, actual))
674 .unwrap_or(false)
675 })
676 })
677 .unwrap_or(false);
678 if !matches {
679 return Ok(());
680 }
681 }
682
683 if let Some(existing_id) = bound_node_id_for_expand(&row, op.dst)? {
684 if existing_id != dst_id {
685 return Ok(());
686 }
687 }
688
689 if let Some(rel_var) = op.rel {
690 if let Some(existing_id) = bound_relationship_id_for_expand(&row, rel_var)? {
691 if existing_id != rel_id {
692 return Ok(());
693 }
694 }
695 }
696
697 matches.push((rel_id, dst_id));
698 Ok(())
699 },
700 )?;
701
702 let Some((&last, rest)) = matches.split_last() else {
703 continue;
704 };
705 for &(rel_id, dst_id) in rest {
706 let mut new_row = row.clone();
707 bind_expand_target(&mut new_row, op, rel_id, dst_id);
708 out.push(new_row);
709 }
710 let mut new_row = row;
711 bind_expand_target(&mut new_row, op, last.0, last.1);
712 out.push(new_row);
713 }
714
715 Ok(out)
716}
717
718#[inline]
719fn bind_expand_target(row: &mut Row, op: &ExpandExec, rel_id: u64, dst_id: u64) {
720 if !row.contains_key(op.dst) {
721 row.insert(op.dst, LoraValue::Node(dst_id));
722 }
723 if let Some(rel_var) = op.rel {
724 if !row.contains_key(rel_var) {
725 row.insert(rel_var, LoraValue::Relationship(rel_id));
726 }
727 }
728}
729
730pub(super) fn expand_var_len_rows<S: GraphStorage>(
731 storage: &S,
732 input_rows: Vec<Row>,
733 op: &ExpandExec,
734 range: &RangeLiteral,
735) -> ExecResult<Vec<Row>> {
736 let (min_hops, max_hops) = resolve_range(range);
737 let bind_relationships = op.rel.is_some();
738 let mut out = Vec::new();
739
740 for row in input_rows {
741 let Some(src_node_id) = bound_node_id_for_expand(&row, op.src)? else {
742 continue;
743 };
744 let bound_dst = bound_node_id_for_expand(&row, op.dst)?;
746
747 let expansions = variable_length_expand(
748 storage,
749 src_node_id,
750 op.direction,
751 &op.types,
752 min_hops,
753 max_hops,
754 bind_relationships,
755 );
756
757 for result in expansions {
758 if bound_dst.is_some_and(|dst| dst != result.dst_node_id) {
759 continue;
760 }
761 let mut new_row = row.clone();
762 new_row.insert(op.dst, LoraValue::Node(result.dst_node_id));
763
764 if let Some(rel_var) = op.rel {
765 let rel_list = LoraValue::List(
766 result
767 .rel_ids
768 .into_iter()
769 .map(LoraValue::Relationship)
770 .collect(),
771 );
772 new_row.insert(rel_var, rel_list);
773 }
774
775 out.push(new_row);
776 }
777 }
778
779 Ok(out)
780}
781
782pub(super) fn properties_to_value_map(props: &Properties) -> LoraValue {
783 let mut map = BTreeMap::new();
784 for (k, v) in props.iter() {
785 map.insert(k.to_string(), LoraValue::from(v));
786 }
787 LoraValue::Map(map)
788}
789
790pub(crate) fn dedup_rows_by_vars(rows: Vec<Row>) -> Vec<Row> {
794 let mut seen: BTreeSet<Vec<GroupValueKey>> = BTreeSet::new();
795 let mut out = Vec::new();
796
797 for row in rows {
798 let key: Vec<GroupValueKey> = row
799 .iter()
800 .map(|(_, val)| GroupValueKey::from_value(val))
801 .collect();
802 if seen.insert(key) {
803 out.push(row);
804 }
805 }
806
807 out
808}
809
810pub(crate) fn dedup_rows(rows: Vec<Row>) -> Vec<Row> {
814 let mut seen: BTreeSet<Vec<(String, GroupValueKey)>> = BTreeSet::new();
815 let mut out = Vec::new();
816
817 for row in rows {
818 let key: Vec<(String, GroupValueKey)> = row
819 .iter_named()
820 .map(|(_, name, val)| (name.into_owned(), GroupValueKey::from_value(val)))
821 .collect();
822 if seen.insert(key) {
823 out.push(row);
824 }
825 }
826
827 out
828}
829
830pub(super) fn eval_properties_expr<S: GraphStorage>(
831 expr: &ResolvedExpr,
832 row: &Row,
833 storage: &S,
834 params: &BTreeMap<String, LoraValue>,
835) -> ExecResult<Properties> {
836 let eval_ctx = EvalContext { storage, params };
837
838 if let ResolvedExpr::Map(items) = expr {
841 let mut out = Properties::new();
842 for (k, v) in items {
843 let value = eval_expr(v, row, &eval_ctx);
844 if matches!(value, LoraValue::Null) {
845 continue;
846 }
847 let prop = lora_value_to_property(value)
848 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
849 out.insert(lora_store::intern(k), prop);
850 }
851 return Ok(out);
852 }
853
854 match eval_expr(expr, row, &eval_ctx) {
855 LoraValue::Map(map) => {
856 let mut out = Properties::new();
857 for (k, v) in map {
858 if matches!(v, LoraValue::Null) {
859 continue;
860 }
861 let prop = lora_value_to_property(v)
862 .map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
863 out.insert(lora_store::intern_owned(k), prop);
870 }
871 Ok(out)
872 }
873 other => Err(ExecutorError::ExpectedPropertyMap {
874 found: value_kind(&other),
875 }),
876 }
877}
878
879pub(crate) fn compute_aggregate_expr<S: GraphStorage>(
880 expr: &ResolvedExpr,
881 rows: &[Row],
882 eval_ctx: &EvalContext<'_, S>,
883) -> ExecResult<LoraValue> {
884 match expr {
885 ResolvedExpr::Function {
886 function,
887 distinct,
888 args,
889 } => {
890 let func = function.as_aggregate();
891 match func {
892 Some(AggregateFunction::Count) => {
893 if args.is_empty() {
894 return Ok(LoraValue::Int(rows.len() as i64));
895 }
896
897 let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
898 values.retain(|v| !matches!(v, LoraValue::Null));
899
900 if *distinct {
901 values = dedup_values(values);
902 }
903
904 Ok(LoraValue::Int(values.len() as i64))
905 }
906
907 Some(AggregateFunction::Collect) => {
908 if args.is_empty() {
909 return Ok(LoraValue::List(Vec::new()));
910 }
911
912 let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
913
914 if *distinct {
915 values = dedup_values(values);
916 }
917
918 Ok(LoraValue::List(values))
919 }
920
921 Some(AggregateFunction::Sum) => {
922 if args.is_empty() {
923 return Ok(LoraValue::Null);
924 }
925
926 let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
927
928 if *distinct {
929 values = dedup_values(values);
930 }
931
932 let nums = values
933 .into_iter()
934 .filter_map(as_f64_lossy)
935 .collect::<Vec<_>>();
936
937 if nums.is_empty() {
938 Ok(LoraValue::Null)
939 } else if nums.iter().all(|n| n.fract() == 0.0) {
940 Ok(LoraValue::Int(nums.iter().sum::<f64>() as i64))
941 } else {
942 Ok(LoraValue::Float(nums.iter().sum::<f64>()))
943 }
944 }
945
946 Some(AggregateFunction::Avg) => {
947 if args.is_empty() {
948 return Ok(LoraValue::Null);
949 }
950
951 let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
952
953 if *distinct {
954 values = dedup_values(values);
955 }
956
957 let nums = values
958 .into_iter()
959 .filter_map(as_f64_lossy)
960 .collect::<Vec<_>>();
961
962 if nums.is_empty() {
963 Ok(LoraValue::Null)
964 } else {
965 Ok(LoraValue::Float(
966 nums.iter().sum::<f64>() / nums.len() as f64,
967 ))
968 }
969 }
970
971 Some(AggregateFunction::Min) => {
972 if args.is_empty() {
973 return Ok(LoraValue::Null);
974 }
975
976 let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
977 values.retain(|v| !matches!(v, LoraValue::Null));
978
979 if *distinct {
980 values = dedup_values(values);
981 }
982
983 Ok(values
984 .into_iter()
985 .min_by(compare_values_total)
986 .unwrap_or(LoraValue::Null))
987 }
988
989 Some(AggregateFunction::Max) => {
990 if args.is_empty() {
991 return Ok(LoraValue::Null);
992 }
993
994 let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
995 values.retain(|v| !matches!(v, LoraValue::Null));
996
997 if *distinct {
998 values = dedup_values(values);
999 }
1000
1001 Ok(values
1002 .into_iter()
1003 .max_by(compare_values_total)
1004 .unwrap_or(LoraValue::Null))
1005 }
1006
1007 Some(AggregateFunction::Stdev | AggregateFunction::Stdevp) => {
1008 if args.is_empty() {
1009 return Ok(LoraValue::Null);
1010 }
1011
1012 let nums: Vec<f64> = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?
1013 .into_iter()
1014 .filter_map(as_f64_lossy)
1015 .collect();
1016
1017 let is_population = matches!(func, Some(AggregateFunction::Stdevp));
1018
1019 if nums.is_empty() || (!is_population && nums.len() < 2) {
1020 return Ok(LoraValue::Float(0.0));
1021 }
1022
1023 let mean = nums.iter().sum::<f64>() / nums.len() as f64;
1024 let variance_sum: f64 = nums.iter().map(|x| (x - mean).powi(2)).sum();
1025 let denom = if is_population {
1026 nums.len() as f64
1027 } else {
1028 (nums.len() - 1) as f64
1029 };
1030 Ok(LoraValue::Float((variance_sum / denom).sqrt()))
1031 }
1032
1033 Some(AggregateFunction::PercentileCont) => {
1034 if args.len() < 2 {
1035 return Ok(LoraValue::Null);
1036 }
1037
1038 let Some(first) = rows.first() else {
1039 return Ok(LoraValue::Null);
1040 };
1041
1042 let percentile = eval_expr_result(&args[1], first, eval_ctx)
1043 .map_err(ExecutorError::RuntimeError)?
1044 .as_f64()
1045 .map(normalize_percentile)
1046 .unwrap_or(0.5);
1047 let mut nums: Vec<f64> = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?
1048 .into_iter()
1049 .filter_map(as_f64_lossy)
1050 .collect();
1051
1052 if nums.is_empty() {
1053 return Ok(LoraValue::Null);
1054 }
1055
1056 nums.sort_by(|a, b| a.partial_cmp(b).unwrap_or(Ordering::Equal));
1057
1058 let index = percentile * (nums.len() - 1) as f64;
1059 let lower = index.floor() as usize;
1060 let upper = index.ceil() as usize;
1061 let fraction = index - lower as f64;
1062
1063 if lower == upper || upper >= nums.len() {
1064 Ok(LoraValue::Float(nums[lower]))
1065 } else {
1066 Ok(LoraValue::Float(
1067 nums[lower] * (1.0 - fraction) + nums[upper] * fraction,
1068 ))
1069 }
1070 }
1071
1072 Some(AggregateFunction::PercentileDisc) => {
1073 if args.len() < 2 {
1074 return Ok(LoraValue::Null);
1075 }
1076
1077 let Some(first) = rows.first() else {
1078 return Ok(LoraValue::Null);
1079 };
1080
1081 let percentile = eval_expr_result(&args[1], first, eval_ctx)
1082 .map_err(ExecutorError::RuntimeError)?
1083 .as_f64()
1084 .map(normalize_percentile)
1085 .unwrap_or(0.5);
1086 let mut nums: Vec<f64> = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?
1087 .into_iter()
1088 .filter_map(as_f64_lossy)
1089 .collect();
1090
1091 if nums.is_empty() {
1092 return Ok(LoraValue::Null);
1093 }
1094
1095 nums.sort_by(|a, b| a.partial_cmp(b).unwrap_or(Ordering::Equal));
1096
1097 let index = (percentile * (nums.len() - 1) as f64).round() as usize;
1098 let index = index.min(nums.len() - 1);
1099 Ok(LoraValue::Float(nums[index]))
1100 }
1101
1102 _ => eval_first_or_null(expr, rows, eval_ctx),
1103 }
1104 }
1105
1106 _ => eval_first_or_null(expr, rows, eval_ctx),
1107 }
1108}
1109
1110fn eval_aggregate_arg_values<S: GraphStorage>(
1111 expr: &ResolvedExpr,
1112 rows: &[Row],
1113 eval_ctx: &EvalContext<'_, S>,
1114) -> ExecResult<Vec<LoraValue>> {
1115 rows.iter()
1116 .map(|row| eval_expr_result(expr, row, eval_ctx).map_err(ExecutorError::RuntimeError))
1117 .collect()
1118}
1119
1120fn normalize_percentile(value: f64) -> f64 {
1121 if value.is_finite() {
1122 value.clamp(0.0, 1.0)
1123 } else {
1124 0.5
1125 }
1126}
1127
1128fn eval_first_or_null<S: GraphStorage>(
1129 expr: &ResolvedExpr,
1130 rows: &[Row],
1131 eval_ctx: &EvalContext<'_, S>,
1132) -> ExecResult<LoraValue> {
1133 match rows.first() {
1134 Some(row) => eval_expr_result(expr, row, eval_ctx).map_err(ExecutorError::RuntimeError),
1135 None => Ok(LoraValue::Null),
1136 }
1137}
1138
1139fn dedup_values(values: Vec<LoraValue>) -> Vec<LoraValue> {
1140 let mut seen: BTreeSet<GroupValueKey> = BTreeSet::new();
1141 let mut out = Vec::new();
1142
1143 for value in values {
1144 let key = GroupValueKey::from_value(&value);
1145 if seen.insert(key) {
1146 out.push(value);
1147 }
1148 }
1149
1150 out
1151}
1152
1153fn as_f64_lossy(v: LoraValue) -> Option<f64> {
1154 match v {
1155 LoraValue::Int(i) => Some(i as f64),
1156 LoraValue::Float(f) => Some(f),
1157 _ => None,
1158 }
1159}
1160
1161pub(super) fn compare_values_total(a: &LoraValue, b: &LoraValue) -> Ordering {
1162 use LoraValue::*;
1163
1164 match (a, b) {
1165 (Bool(x), Bool(y)) => x.cmp(y),
1166 (Int(x), Int(y)) => x.cmp(y),
1167 (Float(x), Float(y)) => x.partial_cmp(y).unwrap_or(Ordering::Equal),
1168 (Int(x), Float(y)) => (*x as f64).partial_cmp(y).unwrap_or(Ordering::Equal),
1169 (Float(x), Int(y)) => x.partial_cmp(&(*y as f64)).unwrap_or(Ordering::Equal),
1170 (String(x), String(y)) => x.cmp(y),
1171 (Binary(x), Binary(y)) => x.segments().cmp(y.segments()),
1172 (Node(x), Node(y)) => x.cmp(y),
1173 (Relationship(x), Relationship(y)) => x.cmp(y),
1174 (Date(x), Date(y)) => x.cmp(y),
1175 (DateTime(x), DateTime(y)) => x.cmp(y),
1176 (LocalDateTime(x), LocalDateTime(y)) => x.cmp(y),
1177 (Time(x), Time(y)) => x.cmp(y),
1178 (LocalTime(x), LocalTime(y)) => x.cmp(y),
1179 (Duration(x), Duration(y)) => x.cmp(y),
1180 (Vector(x), Vector(y)) => x.to_key_string().cmp(&y.to_key_string()),
1181 _ => type_rank(a)
1182 .cmp(&type_rank(b))
1183 .then_with(|| format!("{a:?}").cmp(&format!("{b:?}"))),
1184 }
1185}
1186
1187pub fn value_matches_property_value(expected: &LoraValue, actual: &PropertyValue) -> bool {
1188 match (expected, actual) {
1189 (LoraValue::Null, PropertyValue::Null) => true,
1190 (LoraValue::Bool(a), PropertyValue::Bool(b)) => a == b,
1191 (LoraValue::Int(a), PropertyValue::Int(b)) => a == b,
1192 (LoraValue::Float(a), PropertyValue::Float(b)) => a == b,
1193 (LoraValue::Int(a), PropertyValue::Float(b)) => (*a as f64) == *b,
1194 (LoraValue::Float(a), PropertyValue::Int(b)) => *a == (*b as f64),
1195 (LoraValue::String(a), PropertyValue::String(b)) => a == b,
1196 (LoraValue::Binary(a), PropertyValue::Binary(b)) => a == b,
1197
1198 (LoraValue::List(xs), PropertyValue::List(ys)) => {
1199 xs.len() == ys.len()
1200 && xs
1201 .iter()
1202 .zip(ys.iter())
1203 .all(|(x, y)| value_matches_property_value(x, y))
1204 }
1205
1206 (LoraValue::Map(xm), PropertyValue::Map(ym)) => xm.iter().all(|(k, xv)| {
1207 ym.get(k)
1208 .map(|yv| value_matches_property_value(xv, yv))
1209 .unwrap_or(false)
1210 }),
1211
1212 (LoraValue::Date(a), PropertyValue::Date(b)) => a == b,
1213 (LoraValue::DateTime(a), PropertyValue::DateTime(b)) => a == b,
1214 (LoraValue::LocalDateTime(a), PropertyValue::LocalDateTime(b)) => a == b,
1215 (LoraValue::Time(a), PropertyValue::Time(b)) => a == b,
1216 (LoraValue::LocalTime(a), PropertyValue::LocalTime(b)) => a == b,
1217 (LoraValue::Duration(a), PropertyValue::Duration(b)) => a == b,
1218 (LoraValue::Point(a), PropertyValue::Point(b)) => a == b,
1219 (LoraValue::Vector(a), PropertyValue::Vector(b)) => a == b,
1220
1221 _ => false,
1222 }
1223}
1224
1225pub(crate) fn node_matches_property_filter<S: GraphStorage>(
1226 storage: &S,
1227 node_id: NodeId,
1228 labels: &[Vec<String>],
1229 key: &str,
1230 expected: &LoraValue,
1231) -> bool {
1232 storage
1233 .with_node(node_id, |node| {
1234 node_matches_label_groups(&node.labels, labels)
1235 && node
1236 .properties
1237 .get(key)
1238 .map(|actual| value_matches_property_value(expected, actual))
1239 .unwrap_or(false)
1240 })
1241 .unwrap_or(false)
1242}
1243
1244fn single_label_hint(labels: &[Vec<String>]) -> Option<&str> {
1245 if labels.len() == 1 && labels[0].len() == 1 {
1246 Some(labels[0][0].as_str())
1247 } else {
1248 None
1249 }
1250}
1251
1252fn property_lookup_values(expected: &LoraValue) -> Option<Vec<PropertyValue>> {
1253 let property = lora_value_to_property(expected.clone()).ok()?;
1254 let mut values = vec![property.clone()];
1255
1256 match property {
1257 PropertyValue::Int(i) => {
1258 values.push(PropertyValue::Float(i as f64));
1259 }
1260 PropertyValue::Float(f)
1261 if f.is_finite()
1262 && f.fract() == 0.0
1263 && f >= i64::MIN as f64
1264 && f <= i64::MAX as f64 =>
1265 {
1266 values.push(PropertyValue::Int(f as i64));
1267 }
1268 _ => {}
1269 }
1270
1271 Some(values)
1272}
1273
1274pub(crate) struct NodePropertyCandidates {
1275 pub(crate) ids: Vec<NodeId>,
1276 pub(crate) prefiltered: bool,
1277}
1278
1279pub(crate) fn node_by_property_range_scan_rows<S: GraphStorage>(
1280 storage: &S,
1281 params: &BTreeMap<String, LoraValue>,
1282 base_rows: Vec<Row>,
1283 op: &lora_compiler::NodeByPropertyRangeScanExec,
1284 deadline: Option<Instant>,
1285) -> ExecResult<Vec<Row>> {
1286 let eval_ctx = EvalContext { storage, params };
1287 let mut out = Vec::new();
1288 let mut other_kinds = OtherKindScan::default();
1289
1290 for row in base_rows {
1291 check_optional_deadline(deadline)?;
1292 let lo_value = op.lo.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
1293 let hi_value = op
1297 .hi
1298 .as_ref()
1299 .map(|expr| eval_expr(expr, &row, &eval_ctx))
1300 .filter(|v| !matches!(v, LoraValue::Null));
1301 let lo_prop = lo_value
1302 .clone()
1303 .and_then(|v| lora_value_to_property(v).ok());
1304 let hi_prop = hi_value
1305 .clone()
1306 .and_then(|v| lora_value_to_property(v).ok());
1307 let filter = NodeRangeFilter {
1308 labels: &op.labels,
1309 key: &op.key,
1310 lo: lo_value.as_ref(),
1311 lo_inclusive: op.lo_inclusive,
1312 hi: hi_value.as_ref(),
1313 hi_inclusive: op.hi_inclusive,
1314 };
1315 let bounds = [lo_value.as_ref(), hi_value.as_ref()];
1316 let bound_id = bound_node_id_for_expand(&row, op.var)?;
1317
1318 if let Some(existing_id) = bound_id {
1319 if node_matches_range_filter(storage, existing_id, &filter)
1320 || bound_node_has_other_kind(storage, existing_id, op, bounds)
1321 {
1322 out.push(row);
1323 }
1324 continue;
1325 }
1326
1327 for id in other_kind_node_ids(storage, op, bounds, &mut other_kinds) {
1329 let mut new_row = row.clone();
1330 new_row.insert(op.var, LoraValue::Node(id));
1331 out.push(new_row);
1332 }
1333
1334 if op.order.is_some() {
1335 let mut cursor =
1336 OrderedRangeCursor::new(storage, op, lo_value.clone(), hi_value.clone());
1337 while let Some(id) = cursor.next_id(storage) {
1338 check_optional_deadline(deadline)?;
1339 let mut new_row = row.clone();
1340 new_row.insert(op.var, LoraValue::Node(id));
1341 out.push(new_row);
1342 }
1343 continue;
1344 }
1345
1346 let candidate_ids = match single_label_hint(&op.labels) {
1347 Some(label) => storage
1348 .node_range_candidates(label, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
1349 .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1350 None => scan_node_ids_for_label_groups(storage, &op.labels),
1351 };
1352
1353 for id in candidate_ids {
1354 check_optional_deadline(deadline)?;
1355 if node_matches_range_filter(storage, id, &filter) {
1356 let mut new_row = row.clone();
1357 new_row.insert(op.var, LoraValue::Node(id));
1358 out.push(new_row);
1359 }
1360 }
1361 }
1362
1363 Ok(out)
1364}
1365
1366pub(crate) fn node_by_text_scan_rows<S: GraphStorage>(
1367 storage: &S,
1368 params: &BTreeMap<String, LoraValue>,
1369 base_rows: Vec<Row>,
1370 op: &lora_compiler::NodeByTextScanExec,
1371 deadline: Option<Instant>,
1372) -> ExecResult<Vec<Row>> {
1373 let eval_ctx = EvalContext { storage, params };
1374 let mut out = Vec::new();
1375
1376 for row in base_rows {
1377 check_optional_deadline(deadline)?;
1378 let query = eval_expr(&op.query, &row, &eval_ctx);
1379 let LoraValue::String(query_str) = &query else {
1380 continue;
1383 };
1384
1385 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
1386 if node_matches_text_filter(
1387 storage,
1388 existing_id,
1389 &op.labels,
1390 &op.key,
1391 op.predicate,
1392 query_str,
1393 ) {
1394 out.push(row);
1395 }
1396 continue;
1397 }
1398
1399 let candidate_ids = match single_label_hint(&op.labels) {
1400 Some(label) => storage
1401 .node_text_candidates(label, &op.key, query_str)
1402 .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1403 None => scan_node_ids_for_label_groups(storage, &op.labels),
1404 };
1405
1406 for id in candidate_ids {
1407 check_optional_deadline(deadline)?;
1408 if node_matches_text_filter(storage, id, &op.labels, &op.key, op.predicate, query_str) {
1409 let mut new_row = row.clone();
1410 new_row.insert(op.var, LoraValue::Node(id));
1411 out.push(new_row);
1412 }
1413 }
1414 }
1415
1416 Ok(out)
1417}
1418
1419fn temporal_bound_kinds(bounds: [Option<&LoraValue>; 2]) -> Vec<(&'static str, &LoraValue)> {
1421 let mut kinds: Vec<(&'static str, &LoraValue)> = Vec::with_capacity(2);
1422 for bound in bounds.into_iter().flatten() {
1423 if let Some(kind) = temporal_kind_name(bound) {
1424 if !kinds.iter().any(|(k, _)| *k == kind) {
1425 kinds.push((kind, bound));
1426 }
1427 }
1428 }
1429 kinds
1430}
1431
1432fn is_other_temporal_kind(value: &PropertyValue, kinds: &[(&'static str, &LoraValue)]) -> bool {
1434 temporal_kind_name(&LoraValue::from(value))
1435 .is_some_and(|kind| kinds.iter().any(|(k, _)| *k != kind))
1436}
1437
1438#[derive(Default)]
1442pub(crate) struct OtherKindScan<Id> {
1443 temporals: Option<Vec<(Id, &'static str)>>,
1444}
1445
1446impl<Id: Copy> OtherKindScan<Id> {
1447 fn ids(
1448 &mut self,
1449 kinds: &[(&'static str, &LoraValue)],
1450 scan: impl FnOnce() -> Vec<(Id, &'static str)>,
1451 ) -> Vec<Id> {
1452 self.temporals
1453 .get_or_insert_with(scan)
1454 .iter()
1455 .filter(|(_, kind)| kinds.iter().any(|(k, _)| k != kind))
1456 .map(|(id, _)| *id)
1457 .collect()
1458 }
1459}
1460
1461pub(crate) fn other_kind_node_ids<S: GraphStorage>(
1471 storage: &S,
1472 op: &lora_compiler::NodeByPropertyRangeScanExec,
1473 bounds: [Option<&LoraValue>; 2],
1474 fallback: &mut OtherKindScan<NodeId>,
1475) -> Vec<NodeId> {
1476 let kinds = temporal_bound_kinds(bounds);
1477 if kinds.is_empty() {
1478 return Vec::new();
1479 }
1480 let index_label = op
1482 .labels
1483 .iter()
1484 .find(|group| group.len() == 1)
1485 .map(|group| group[0].as_str());
1486 let mut ids = BTreeSet::new();
1487 let mut need_scan = index_label.is_none();
1488 if let Some(label) = index_label {
1489 for (_, bound) in &kinds {
1490 let listed = lora_value_to_property((*bound).clone())
1491 .ok()
1492 .and_then(|like| storage.node_range_other_temporal_kind_ids(label, &op.key, &like));
1493 match listed {
1494 Some(listed) => ids.extend(listed),
1495 None => need_scan = true,
1496 }
1497 }
1498 }
1501 if need_scan {
1502 ids.extend(fallback.ids(&kinds, || {
1503 scan_node_ids_for_label_groups(storage, &op.labels)
1504 .into_iter()
1505 .filter_map(|id| {
1506 storage
1507 .with_node(id, |n| {
1508 n.properties
1509 .get(op.key.as_str())
1510 .and_then(|v| temporal_kind_name(&LoraValue::from(v)))
1511 })
1512 .flatten()
1513 .map(|kind| (id, kind))
1514 })
1515 .collect()
1516 }));
1517 return ids.into_iter().collect();
1518 }
1519 let single = op.labels.len() == 1 && op.labels[0].len() == 1;
1520 ids.into_iter()
1521 .filter(|id| {
1522 single
1523 || storage
1524 .with_node(*id, |n| node_matches_label_groups(&n.labels, &op.labels))
1525 .unwrap_or(false)
1526 })
1527 .collect()
1528}
1529
1530pub(crate) fn bound_node_has_other_kind<S: GraphStorage>(
1534 storage: &S,
1535 id: NodeId,
1536 op: &lora_compiler::NodeByPropertyRangeScanExec,
1537 bounds: [Option<&LoraValue>; 2],
1538) -> bool {
1539 let kinds = temporal_bound_kinds(bounds);
1540 !kinds.is_empty()
1541 && storage
1542 .with_node(id, |n| {
1543 node_matches_label_groups(&n.labels, &op.labels)
1544 && n.properties
1545 .get(op.key.as_str())
1546 .is_some_and(|v| is_other_temporal_kind(v, &kinds))
1547 })
1548 .unwrap_or(false)
1549}
1550
1551fn other_kind_rel_ids<S: GraphStorage>(
1553 storage: &S,
1554 op: &lora_compiler::RelByPropertyRangeScanExec,
1555 bounds: [Option<&LoraValue>; 2],
1556 fallback: &mut OtherKindScan<RelationshipId>,
1557) -> Vec<RelationshipId> {
1558 let kinds = temporal_bound_kinds(bounds);
1559 if kinds.is_empty() {
1560 return Vec::new();
1561 }
1562 let mut ids = BTreeSet::new();
1563 let mut need_scan = op.types.is_empty();
1564 for ty in &op.types {
1565 for (_, bound) in &kinds {
1566 let listed = lora_value_to_property((*bound).clone())
1567 .ok()
1568 .and_then(|like| {
1569 storage.relationship_range_other_temporal_kind_ids(ty, &op.key, &like)
1570 });
1571 match listed {
1572 Some(listed) => ids.extend(listed),
1573 None => need_scan = true,
1574 }
1575 }
1576 }
1577 if need_scan {
1578 ids.extend(fallback.ids(&kinds, || {
1579 rel_candidate_ids(storage, &op.types, |_| None)
1580 .into_iter()
1581 .filter_map(|id| {
1582 storage
1583 .with_relationship(id, |r| {
1584 r.properties
1585 .get(op.key.as_str())
1586 .and_then(|v| temporal_kind_name(&LoraValue::from(v)))
1587 })
1588 .flatten()
1589 .map(|kind| (id, kind))
1590 })
1591 .collect()
1592 }));
1593 }
1594 ids.into_iter().collect()
1595}
1596
1597fn temporal_kind_name(value: &LoraValue) -> Option<&'static str> {
1598 Some(match value {
1599 LoraValue::Date(_) => "DATE",
1600 LoraValue::DateTime(_) => "DATETIME",
1601 LoraValue::LocalDateTime(_) => "LOCAL_DATETIME",
1602 LoraValue::Time(_) => "TIME",
1603 LoraValue::LocalTime(_) => "LOCAL_TIME",
1604 _ => return None,
1605 })
1606}
1607
1608pub(crate) struct NodeRangeFilter<'a> {
1609 labels: &'a [Vec<String>],
1610 key: &'a str,
1611 lo: Option<&'a LoraValue>,
1612 lo_inclusive: bool,
1613 hi: Option<&'a LoraValue>,
1614 hi_inclusive: bool,
1615}
1616
1617fn node_matches_range_filter<S: GraphStorage>(
1618 storage: &S,
1619 id: NodeId,
1620 filter: &NodeRangeFilter<'_>,
1621) -> bool {
1622 storage
1623 .with_node(id, |n| {
1624 if !node_matches_label_groups(&n.labels, filter.labels) {
1625 return false;
1626 }
1627 let Some(actual) = n.properties.get(filter.key) else {
1628 return false;
1629 };
1630 let actual_lv = lora_store_property_to_value(actual);
1631 range_predicate_holds(
1632 &actual_lv,
1633 filter.lo,
1634 filter.lo_inclusive,
1635 filter.hi,
1636 filter.hi_inclusive,
1637 )
1638 })
1639 .unwrap_or(false)
1640}
1641
1642fn node_matches_text_filter<S: GraphStorage>(
1643 storage: &S,
1644 id: NodeId,
1645 labels: &[Vec<String>],
1646 key: &str,
1647 predicate: lora_compiler::TextPredicate,
1648 query: &str,
1649) -> bool {
1650 storage
1651 .with_node(id, |n| {
1652 if !node_matches_label_groups(&n.labels, labels) {
1653 return false;
1654 }
1655 let Some(PropertyValue::String(actual)) = n.properties.get(key) else {
1656 return false;
1657 };
1658 text_predicate_holds(actual, predicate, query)
1659 })
1660 .unwrap_or(false)
1661}
1662
1663fn text_predicate_holds(
1664 actual: &str,
1665 predicate: lora_compiler::TextPredicate,
1666 query: &str,
1667) -> bool {
1668 match predicate {
1669 lora_compiler::TextPredicate::StartsWith => actual.starts_with(query),
1670 lora_compiler::TextPredicate::EndsWith => actual.ends_with(query),
1671 lora_compiler::TextPredicate::Contains => actual.contains(query),
1672 }
1673}
1674
1675fn range_predicate_holds(
1676 actual: &LoraValue,
1677 lo: Option<&LoraValue>,
1678 lo_inclusive: bool,
1679 hi: Option<&LoraValue>,
1680 hi_inclusive: bool,
1681) -> bool {
1682 if let Some(lo) = lo {
1683 match range_comparison(actual, lo) {
1684 None => return false,
1685 Some(Ordering::Less) => return false,
1686 Some(Ordering::Equal) if !lo_inclusive => return false,
1687 _ => {}
1688 }
1689 }
1690 if let Some(hi) = hi {
1691 match range_comparison(actual, hi) {
1692 None => return false,
1693 Some(Ordering::Greater) => return false,
1694 Some(Ordering::Equal) if !hi_inclusive => return false,
1695 _ => {}
1696 }
1697 }
1698 true
1699}
1700
1701fn range_comparison(actual: &LoraValue, bound: &LoraValue) -> Option<Ordering> {
1702 match (actual, bound) {
1703 (LoraValue::Null, _) | (_, LoraValue::Null) => None,
1704 (LoraValue::String(a), LoraValue::String(b)) => Some(a.cmp(b)),
1705 (
1706 LoraValue::Date(_)
1707 | LoraValue::DateTime(_)
1708 | LoraValue::LocalDateTime(_)
1709 | LoraValue::Time(_)
1710 | LoraValue::LocalTime(_),
1711 _,
1712 ) => actual.temporal_cmp(bound),
1713 (LoraValue::Duration(a), LoraValue::Duration(b)) => a
1714 .total_seconds_approx()
1715 .partial_cmp(&b.total_seconds_approx()),
1716 _ => actual.as_f64()?.partial_cmp(&bound.as_f64()?),
1717 }
1718}
1719
1720fn lora_store_property_to_value(value: &PropertyValue) -> LoraValue {
1721 LoraValue::from(value)
1722}
1723
1724pub(crate) fn node_by_point_scan_rows<S: GraphStorage>(
1725 storage: &S,
1726 params: &BTreeMap<String, LoraValue>,
1727 base_rows: Vec<Row>,
1728 op: &lora_compiler::NodeByPointScanExec,
1729 deadline: Option<Instant>,
1730) -> ExecResult<Vec<Row>> {
1731 let eval_ctx = EvalContext { storage, params };
1732 let mut out = Vec::new();
1733
1734 for row in base_rows {
1735 check_optional_deadline(deadline)?;
1736
1737 let probe = match &op.predicate {
1741 lora_compiler::PointPredicate::WithinBBox {
1742 lower_left,
1743 upper_right,
1744 } => {
1745 let ll = eval_expr(lower_left, &row, &eval_ctx);
1746 let ur = eval_expr(upper_right, &row, &eval_ctx);
1747 match (ll, ur) {
1748 (LoraValue::Point(a), LoraValue::Point(b)) => {
1749 Probe::WithinBBox { ll: a, ur: b }
1750 }
1751 _ => continue,
1752 }
1753 }
1754 lora_compiler::PointPredicate::WithinDistance {
1755 center,
1756 max_distance,
1757 inclusive,
1758 } => {
1759 let c = eval_expr(center, &row, &eval_ctx);
1760 let d = eval_expr(max_distance, &row, &eval_ctx);
1761 match (c, d) {
1762 (LoraValue::Point(c), LoraValue::Float(d)) => Probe::WithinDistance {
1763 center: c,
1764 max: d,
1765 inclusive: *inclusive,
1766 },
1767 (LoraValue::Point(c), LoraValue::Int(d)) => Probe::WithinDistance {
1768 center: c,
1769 max: d as f64,
1770 inclusive: *inclusive,
1771 },
1772 _ => continue,
1773 }
1774 }
1775 };
1776
1777 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
1778 if node_matches_point_filter(storage, existing_id, &op.labels, &op.key, &probe) {
1779 out.push(row);
1780 }
1781 continue;
1782 }
1783
1784 let candidate_ids = match single_label_hint(&op.labels) {
1785 Some(label) => match &probe {
1786 Probe::WithinBBox { ll, ur } => {
1787 let (ranges, n) = lora_store::bbox_x_ranges(ll, ur);
1789 let (lo_y, hi_y) = (ll.y.min(ur.y), ll.y.max(ur.y));
1790 let mut ids = Vec::new();
1791 let mut indexed = true;
1792 for (lo_x, hi_x) in &ranges[..n] {
1793 match storage.node_point_within_bbox(
1794 label,
1795 &op.key,
1796 (*lo_x, lo_y),
1797 (*hi_x, hi_y),
1798 ) {
1799 Some(found) => ids.extend(found),
1800 None => indexed = false,
1801 }
1802 }
1803 if !indexed {
1804 scan_node_ids_for_label_groups(storage, &op.labels)
1805 } else {
1806 if n > 1 {
1807 ids.sort_unstable();
1808 ids.dedup();
1809 }
1810 ids
1811 }
1812 }
1813 Probe::WithinDistance { center, max, .. } => storage
1814 .node_point_within_distance(label, &op.key, (center.x, center.y), *max)
1815 .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1816 },
1817 None => scan_node_ids_for_label_groups(storage, &op.labels),
1818 };
1819
1820 for id in candidate_ids {
1821 check_optional_deadline(deadline)?;
1822 if node_matches_point_filter(storage, id, &op.labels, &op.key, &probe) {
1823 let mut new_row = row.clone();
1824 new_row.insert(op.var, LoraValue::Node(id));
1825 out.push(new_row);
1826 }
1827 }
1828 }
1829
1830 Ok(out)
1831}
1832
1833fn node_matches_point_filter<S: GraphStorage>(
1834 storage: &S,
1835 id: NodeId,
1836 labels: &[Vec<String>],
1837 key: &str,
1838 probe: &Probe,
1839) -> bool {
1840 storage
1841 .with_node(id, |n| {
1842 if !node_matches_label_groups(&n.labels, labels) {
1843 return false;
1844 }
1845 let Some(PropertyValue::Point(point)) = n.properties.get(key) else {
1846 return false;
1847 };
1848 point_predicate_holds(point, probe)
1849 })
1850 .unwrap_or(false)
1851}
1852
1853enum Probe {
1854 WithinBBox {
1855 ll: lora_store::LoraPoint,
1856 ur: lora_store::LoraPoint,
1857 },
1858 WithinDistance {
1859 center: lora_store::LoraPoint,
1860 max: f64,
1861 inclusive: bool,
1862 },
1863}
1864
1865fn point_predicate_holds(actual: &lora_store::LoraPoint, probe: &Probe) -> bool {
1866 match probe {
1867 Probe::WithinBBox { ll, ur } => lora_store::bbox_contains(actual, ll, ur).unwrap_or(false),
1868 Probe::WithinDistance {
1869 center,
1870 max,
1871 inclusive,
1872 } => {
1873 let Some(d) = lora_store::point_distance(actual, center) else {
1874 return false;
1875 };
1876 if *inclusive {
1877 d <= *max
1878 } else {
1879 d < *max
1880 }
1881 }
1882 }
1883}
1884
1885pub(crate) fn indexed_node_property_candidates<S: GraphStorage>(
1886 storage: &S,
1887 labels: &[Vec<String>],
1888 key: &str,
1889 expected: &LoraValue,
1890) -> NodePropertyCandidates {
1891 let Some(values) = property_lookup_values(expected) else {
1892 return NodePropertyCandidates {
1893 ids: scan_node_ids_for_label_groups(storage, labels),
1894 prefiltered: false,
1895 };
1896 };
1897
1898 let label_hint = single_label_hint(labels);
1899 let mut seen = BTreeSet::new();
1900 let mut out = Vec::new();
1901 for value in values {
1902 for id in storage.find_node_ids_by_property(label_hint, key, &value) {
1903 if seen.insert(id) {
1904 out.push(id);
1905 }
1906 }
1907 }
1908 NodePropertyCandidates {
1909 ids: out,
1910 prefiltered: labels.is_empty() || label_hint.is_some(),
1911 }
1912}
1913
1914pub(crate) fn property_scan_candidates<S: GraphStorage>(
1920 storage: &S,
1921 labels: &[Vec<String>],
1922 key: &str,
1923 expected: &LoraValue,
1924 in_list: bool,
1925) -> NodePropertyCandidates {
1926 if !in_list {
1927 return indexed_node_property_candidates(storage, labels, key, expected);
1928 }
1929 match expected {
1930 LoraValue::Null => NodePropertyCandidates {
1931 ids: Vec::new(),
1932 prefiltered: true,
1933 },
1934 LoraValue::List(items) => {
1935 let mut seen = BTreeSet::new();
1936 let mut ids = Vec::new();
1937 for item in items {
1938 if matches!(item, LoraValue::Null) {
1939 continue;
1940 }
1941 let candidates = indexed_node_property_candidates(storage, labels, key, item);
1942 for id in candidates.ids {
1943 if seen.contains(&id) {
1944 continue;
1945 }
1946 if !candidates.prefiltered
1947 && !node_matches_property_filter(storage, id, labels, key, item)
1948 {
1949 continue;
1950 }
1951 seen.insert(id);
1952 ids.push(id);
1953 }
1954 }
1955 NodePropertyCandidates {
1956 ids,
1957 prefiltered: true,
1958 }
1959 }
1960 _ => NodePropertyCandidates {
1961 ids: scan_node_ids_for_label_groups(storage, labels)
1962 .into_iter()
1963 .filter(|&id| {
1964 storage
1965 .with_node(id, |n| node_matches_label_groups(&n.labels, labels))
1966 .unwrap_or(false)
1967 })
1968 .collect(),
1969 prefiltered: true,
1970 },
1971 }
1972}
1973
1974pub(crate) fn property_scan_matches<S: GraphStorage>(
1978 storage: &S,
1979 node_id: NodeId,
1980 labels: &[Vec<String>],
1981 key: &str,
1982 expected: &LoraValue,
1983 in_list: bool,
1984) -> bool {
1985 if !in_list {
1986 return node_matches_property_filter(storage, node_id, labels, key, expected);
1987 }
1988 match expected {
1989 LoraValue::Null => false,
1990 LoraValue::List(items) => items.iter().any(|item| {
1991 !matches!(item, LoraValue::Null)
1992 && node_matches_property_filter(storage, node_id, labels, key, item)
1993 }),
1994 _ => storage
1995 .with_node(node_id, |n| node_matches_label_groups(&n.labels, labels))
1996 .unwrap_or(false),
1997 }
1998}
1999
2000pub(crate) fn build_path_value<S: GraphStorage>(
2006 row: &Row,
2007 node_vars: &[VarId],
2008 rel_vars: &[VarId],
2009 storage: &S,
2010) -> LoraValue {
2011 let (raw_nodes, rels, has_var_len) = path_bindings(row, node_vars, rel_vars);
2012
2013 let nodes = if has_var_len && !rels.is_empty() && raw_nodes.len() == 2 {
2014 reconstruct_var_len_nodes(raw_nodes[0], &rels, storage)
2015 } else {
2016 raw_nodes
2017 };
2018
2019 LoraValue::Path(LoraPath { nodes, rels })
2020}
2021
2022#[inline]
2023fn path_bindings(
2024 row: &Row,
2025 node_vars: &[VarId],
2026 rel_vars: &[VarId],
2027) -> (Vec<NodeId>, Vec<RelationshipId>, bool) {
2028 let mut raw_nodes = Vec::new();
2029 let mut rels = Vec::new();
2030 let mut has_var_len = false;
2031
2032 for &nv in node_vars {
2033 match row.get(nv) {
2034 Some(LoraValue::Node(id)) => raw_nodes.push(*id),
2035 Some(LoraValue::List(items)) => {
2036 for item in items {
2037 if let LoraValue::Node(id) = item {
2038 raw_nodes.push(*id);
2039 }
2040 }
2041 }
2042 _ => {}
2043 }
2044 }
2045
2046 for &rv in rel_vars {
2047 match row.get(rv) {
2048 Some(LoraValue::Relationship(id)) => rels.push(*id),
2049 Some(LoraValue::List(items)) => {
2050 has_var_len = true;
2051 for item in items {
2052 if let LoraValue::Relationship(id) = item {
2053 rels.push(*id);
2054 }
2055 }
2056 }
2057 _ => {}
2058 }
2059 }
2060
2061 (raw_nodes, rels, has_var_len)
2062}
2063
2064#[inline]
2065fn reconstruct_var_len_nodes<S: GraphStorage>(
2066 start: NodeId,
2067 rels: &[RelationshipId],
2068 storage: &S,
2069) -> Vec<NodeId> {
2070 let mut ordered = Vec::with_capacity(rels.len() + 1);
2071 ordered.push(start);
2072 let mut current = start;
2073 for &rel_id in rels {
2074 if let Some((src, dst)) = storage.relationship_endpoints(rel_id) {
2075 let next = if src == current { dst } else { src };
2076 ordered.push(next);
2077 current = next;
2078 }
2079 }
2080 ordered
2081}
2082
2083fn type_rank(v: &LoraValue) -> u8 {
2084 match v {
2085 LoraValue::Null => 0,
2086 LoraValue::Bool(_) => 1,
2087 LoraValue::Int(_) | LoraValue::Float(_) => 2,
2088 LoraValue::String(_) => 3,
2089 LoraValue::Binary(_) => 4,
2090 LoraValue::Date(_) => 5,
2091 LoraValue::DateTime(_) => 6,
2092 LoraValue::LocalDateTime(_) => 7,
2093 LoraValue::Time(_) => 8,
2094 LoraValue::LocalTime(_) => 9,
2095 LoraValue::Duration(_) => 10,
2096 LoraValue::Point(_) => 11,
2097 LoraValue::Vector(_) => 12,
2098 LoraValue::List(_) => 13,
2099 LoraValue::Map(_) => 14,
2100 LoraValue::Node(_) => 15,
2101 LoraValue::Relationship(_) => 16,
2102 LoraValue::Path(_) => 17,
2103 }
2104}
2105
2106pub(crate) fn node_matches_label_groups(node_labels: &[String], groups: &[Vec<String>]) -> bool {
2110 groups
2111 .iter()
2112 .all(|group| group.iter().any(|l| node_labels.iter().any(|nl| nl == l)))
2113}
2114
2115pub(crate) fn scan_node_ids_for_label_groups<S: GraphStorage>(
2118 storage: &S,
2119 groups: &[Vec<String>],
2120) -> Vec<NodeId> {
2121 if groups.is_empty() {
2122 return storage.all_node_ids();
2123 }
2124 if groups.len() == 1 {
2125 return label_group_candidate_ids(storage, &groups[0]);
2126 }
2127
2128 let mut best: Option<Vec<NodeId>> = None;
2129 for group in groups {
2130 let ids = label_group_candidate_ids(storage, group);
2131 if ids.is_empty() {
2132 return Vec::new();
2133 }
2134 if best
2135 .as_ref()
2136 .map(|current| ids.len() < current.len())
2137 .unwrap_or(true)
2138 {
2139 best = Some(ids);
2140 }
2141 }
2142
2143 best.unwrap_or_default()
2144}
2145
2146pub(crate) fn label_group_candidates_prefiltered(groups: &[Vec<String>]) -> bool {
2147 groups.len() <= 1
2148}
2149
2150fn label_group_candidate_ids<S: GraphStorage>(storage: &S, group: &[String]) -> Vec<NodeId> {
2151 match group {
2152 [] => Vec::new(),
2153 [label] => storage.node_ids_by_label(label),
2154 labels => {
2155 let mut seen = BTreeSet::new();
2156 let mut out = Vec::new();
2157 for label in labels {
2158 for id in storage.node_ids_by_label(label) {
2159 if seen.insert(id) {
2160 out.push(id);
2161 }
2162 }
2163 }
2164 out
2165 }
2166 }
2167}
2168
2169pub(crate) fn hydrate_node_record(node: &lora_store::NodeRecord) -> LoraValue {
2170 let mut map = BTreeMap::new();
2171 map.insert("kind".to_string(), LoraValue::String("node".to_string()));
2172 map.insert("id".to_string(), LoraValue::Int(node.id as i64));
2173 map.insert(
2174 "labels".to_string(),
2175 LoraValue::List(
2176 node.labels
2177 .iter()
2178 .map(|s| LoraValue::String(s.clone()))
2179 .collect(),
2180 ),
2181 );
2182 map.insert(
2183 "properties".to_string(),
2184 properties_to_value_map(&node.properties),
2185 );
2186 LoraValue::Map(map)
2187}
2188
2189pub(crate) fn hydrate_relationship_record(rel: &lora_store::RelationshipRecord) -> LoraValue {
2190 let mut map = BTreeMap::new();
2191 map.insert(
2192 "kind".to_string(),
2193 LoraValue::String("relationship".to_string()),
2194 );
2195 map.insert("id".to_string(), LoraValue::Int(rel.id as i64));
2196 map.insert("startId".to_string(), LoraValue::Int(rel.src as i64));
2197 map.insert("endId".to_string(), LoraValue::Int(rel.dst as i64));
2198 map.insert("type".to_string(), LoraValue::String(rel.rel_type.clone()));
2199 map.insert(
2200 "properties".to_string(),
2201 properties_to_value_map(&rel.properties),
2202 );
2203 LoraValue::Map(map)
2204}
2205
2206pub(super) fn flatten_label_groups(groups: &[Vec<String>]) -> Vec<String> {
2209 groups.iter().flat_map(|g| g.iter().cloned()).collect()
2210}
2211
2212#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
2213pub(crate) enum GroupValueKey {
2214 Null,
2215 Bool(bool),
2216 Int(i64),
2217 Float(String),
2218 String(String),
2219 Binary(Vec<Vec<u8>>),
2220 List(Vec<GroupValueKey>),
2221 Map(Vec<(String, GroupValueKey)>),
2222 Node(u64),
2223 Relationship(u64),
2224}
2225
2226impl GroupValueKey {
2227 pub(crate) fn from_value(v: &LoraValue) -> Self {
2228 match v {
2229 LoraValue::Null => Self::Null,
2230 LoraValue::Bool(x) => Self::Bool(*x),
2231 LoraValue::Int(x) => Self::Int(*x),
2232 LoraValue::Float(x) => Self::Float(x.to_string()),
2233 LoraValue::String(x) => Self::String(x.clone()),
2234 LoraValue::Binary(x) => Self::Binary(x.segments().to_vec()),
2235 LoraValue::List(xs) => Self::List(xs.iter().map(Self::from_value).collect()),
2236 LoraValue::Map(m) => Self::Map(
2237 m.iter()
2238 .map(|(k, v)| (k.clone(), Self::from_value(v)))
2239 .collect(),
2240 ),
2241 LoraValue::Node(id) => Self::Node(*id),
2242 LoraValue::Relationship(id) => Self::Relationship(*id),
2243 LoraValue::Path(_) => Self::Null,
2244 LoraValue::Date(d) => Self::String(d.to_string()),
2246 LoraValue::DateTime(dt) => Self::String(dt.to_string()),
2247 LoraValue::LocalDateTime(dt) => Self::String(dt.to_string()),
2248 LoraValue::Time(t) => Self::String(t.to_string()),
2249 LoraValue::LocalTime(t) => Self::String(t.to_string()),
2250 LoraValue::Duration(dur) => Self::String(dur.to_string()),
2251 LoraValue::Point(p) => Self::String(p.to_string()),
2252 LoraValue::Vector(v) => Self::String(format!("vector:{}", v.to_key_string())),
2253 }
2254 }
2255}
2256
2257const MAX_VAR_LEN_HOPS: u64 = 100;
2269
2270pub(crate) fn resolve_range(range: &RangeLiteral) -> (u64, u64) {
2271 let min_hops = range.start.unwrap_or(1);
2272 let max_hops = range.end.unwrap_or(MAX_VAR_LEN_HOPS);
2273 (min_hops, max_hops)
2274}
2275
2276pub(crate) struct VarLenResult {
2278 pub(crate) dst_node_id: NodeId,
2280 pub(crate) rel_ids: Vec<u64>,
2282}
2283
2284pub(crate) fn variable_length_expand<S: GraphStorage>(
2291 storage: &S,
2292 start_node_id: NodeId,
2293 direction: Direction,
2294 types: &[String],
2295 min_hops: u64,
2296 max_hops: u64,
2297 bind_relationships: bool,
2298) -> Vec<VarLenResult> {
2299 let mut results = Vec::new();
2300
2301 let mut frontier: Vec<(NodeId, Vec<u64>)> = vec![(start_node_id, Vec::new())];
2303
2304 for depth in 1..=max_hops {
2305 let is_last_hop = depth == max_hops;
2309 let mut next_frontier: Vec<(NodeId, Vec<u64>)> = Vec::new();
2310
2311 for (current_node, rels_used) in &frontier {
2312 for (rel_id, neighbor_id) in storage.expand_ids(*current_node, direction, types) {
2315 if rels_used.contains(&rel_id) {
2318 continue;
2319 }
2320
2321 if is_last_hop {
2322 if depth >= min_hops {
2325 let mut rel_ids = Vec::with_capacity(rels_used.len() + 1);
2326 rel_ids.extend_from_slice(rels_used);
2327 rel_ids.push(rel_id);
2328 results.push(VarLenResult {
2329 dst_node_id: neighbor_id,
2330 rel_ids: if bind_relationships {
2331 rel_ids
2332 } else {
2333 Vec::new()
2334 },
2335 });
2336 }
2337 continue;
2338 }
2339
2340 let mut new_rels = Vec::with_capacity(rels_used.len() + 1);
2341 new_rels.extend_from_slice(rels_used);
2342 new_rels.push(rel_id);
2343
2344 if depth >= min_hops {
2345 results.push(VarLenResult {
2346 dst_node_id: neighbor_id,
2347 rel_ids: if bind_relationships {
2348 new_rels.clone()
2349 } else {
2350 Vec::new()
2351 },
2352 });
2353 }
2354
2355 next_frontier.push((neighbor_id, new_rels));
2356 }
2357 }
2358
2359 if is_last_hop || next_frontier.is_empty() {
2360 break;
2361 }
2362
2363 frontier = next_frontier;
2364 }
2365
2366 if min_hops == 0 {
2368 results.insert(
2369 0,
2370 VarLenResult {
2371 dst_node_id: start_node_id,
2372 rel_ids: Vec::new(),
2373 },
2374 );
2375 }
2376
2377 results
2378}
2379
2380pub(crate) fn filter_shortest_paths(rows: Vec<Row>, path_var: VarId, all: bool) -> Vec<Row> {
2383 if rows.is_empty() {
2384 return rows;
2385 }
2386
2387 let lengths: Vec<usize> = rows
2389 .iter()
2390 .map(|row| match row.get(path_var) {
2391 Some(LoraValue::Path(p)) => p.rels.len(),
2392 _ => usize::MAX,
2393 })
2394 .collect();
2395
2396 let min_len = lengths.iter().copied().min().unwrap_or(usize::MAX);
2397
2398 let mut result: Vec<Row> = rows
2399 .into_iter()
2400 .zip(lengths.iter())
2401 .filter(|(_, len)| **len == min_len)
2402 .map(|(row, _)| row)
2403 .collect();
2404
2405 if !all && result.len() > 1 {
2406 result.truncate(1);
2407 }
2408
2409 result
2410}
2411
2412fn emit_rel_rows(
2429 direction: Direction,
2430 src_var: VarId,
2431 rel_var: VarId,
2432 dst_var: VarId,
2433 rel: &lora_store::RelationshipRecord,
2434 base: &Row,
2435 out: &mut Vec<Row>,
2436) -> ExecResult<()> {
2437 match direction {
2438 Direction::Right => {
2439 emit_one_rel_row(
2440 src_var, rel_var, dst_var, rel.src, rel.id, rel.dst, base, out,
2441 )?;
2442 }
2443 Direction::Left => {
2444 emit_one_rel_row(
2445 src_var, rel_var, dst_var, rel.dst, rel.id, rel.src, base, out,
2446 )?;
2447 }
2448 Direction::Undirected => {
2449 emit_one_rel_row(
2450 src_var, rel_var, dst_var, rel.src, rel.id, rel.dst, base, out,
2451 )?;
2452 if rel.src != rel.dst {
2456 emit_one_rel_row(
2457 src_var, rel_var, dst_var, rel.dst, rel.id, rel.src, base, out,
2458 )?;
2459 }
2460 }
2461 }
2462 Ok(())
2463}
2464
2465#[allow(clippy::too_many_arguments)]
2466fn emit_one_rel_row(
2467 src_var: VarId,
2468 rel_var: VarId,
2469 dst_var: VarId,
2470 src_id: NodeId,
2471 rel_id: RelationshipId,
2472 dst_id: NodeId,
2473 base: &Row,
2474 out: &mut Vec<Row>,
2475) -> ExecResult<()> {
2476 let mut row = base.clone();
2477 if bind_node_value(&mut row, src_var, src_id)?
2478 && bind_relationship_value(&mut row, rel_var, rel_id)?
2479 && bind_node_value(&mut row, dst_var, dst_id)?
2480 {
2481 out.push(row);
2482 }
2483 Ok(())
2484}
2485
2486fn bind_node_value(row: &mut Row, var: VarId, id: NodeId) -> ExecResult<bool> {
2487 match row.get(var) {
2488 Some(LoraValue::Node(existing)) => Ok(*existing == id),
2489 Some(other) => Err(ExecutorError::ExpectedNodeForExpand {
2490 var: format!("{var:?}"),
2491 found: value_kind(other),
2492 }),
2493 None => {
2494 row.insert(var, LoraValue::Node(id));
2495 Ok(true)
2496 }
2497 }
2498}
2499
2500fn bind_relationship_value(row: &mut Row, var: VarId, id: RelationshipId) -> ExecResult<bool> {
2501 match row.get(var) {
2502 Some(LoraValue::Relationship(existing)) => Ok(*existing == id),
2503 Some(other) => Err(ExecutorError::ExpectedRelationshipForExpand {
2504 var: format!("{var:?}"),
2505 found: value_kind(other),
2506 }),
2507 None => {
2508 row.insert(var, LoraValue::Relationship(id));
2509 Ok(true)
2510 }
2511 }
2512}
2513
2514fn rel_candidate_ids<S, F>(storage: &S, types: &[String], indexed: F) -> Vec<RelationshipId>
2515where
2516 S: GraphStorage,
2517 F: Fn(&str) -> Option<Vec<RelationshipId>>,
2518{
2519 if types.is_empty() {
2520 return storage.all_rel_ids();
2523 }
2524 let mut all = Vec::new();
2525 let mut seen = BTreeSet::new();
2526 for ty in types {
2527 let ids = indexed(ty).unwrap_or_else(|| storage.rel_ids_by_type(ty));
2528 for id in ids {
2529 if seen.insert(id) {
2530 all.push(id);
2531 }
2532 }
2533 }
2534 all
2535}
2536
2537pub(crate) fn rel_by_property_range_scan_rows<S: GraphStorage>(
2538 storage: &S,
2539 params: &BTreeMap<String, LoraValue>,
2540 base_rows: Vec<Row>,
2541 op: &lora_compiler::RelByPropertyRangeScanExec,
2542 deadline: Option<Instant>,
2543) -> ExecResult<Vec<Row>> {
2544 let eval_ctx = EvalContext { storage, params };
2545 let mut out = Vec::new();
2546 let mut other_kinds = OtherKindScan::default();
2547
2548 for row in base_rows {
2549 check_optional_deadline(deadline)?;
2550 let lo_value = op.lo.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
2551 let hi_value = op
2555 .hi
2556 .as_ref()
2557 .map(|expr| eval_expr(expr, &row, &eval_ctx))
2558 .filter(|v| !matches!(v, LoraValue::Null));
2559 let lo_prop = lo_value
2560 .clone()
2561 .and_then(|v| lora_value_to_property(v).ok());
2562 let hi_prop = hi_value
2563 .clone()
2564 .and_then(|v| lora_value_to_property(v).ok());
2565
2566 let bounds = [lo_value.as_ref(), hi_value.as_ref()];
2569 for rel_id in other_kind_rel_ids(storage, op, bounds, &mut other_kinds) {
2570 if let Some(result) = storage.with_relationship(rel_id, |rel| {
2571 if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
2572 return Ok(());
2573 }
2574 emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2575 }) {
2576 result?;
2577 }
2578 }
2579 let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| {
2580 storage.relationship_range_candidates(ty, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
2581 });
2582
2583 for rel_id in candidate_ids {
2584 check_optional_deadline(deadline)?;
2585 if let Some(result) = storage.with_relationship(rel_id, |rel| {
2586 if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
2587 return Ok(());
2588 }
2589 let Some(actual) = rel.properties.get(op.key.as_str()) else {
2590 return Ok(());
2591 };
2592 let actual_lv = LoraValue::from(actual);
2593 if !range_predicate_holds(
2594 &actual_lv,
2595 lo_value.as_ref(),
2596 op.lo_inclusive,
2597 hi_value.as_ref(),
2598 op.hi_inclusive,
2599 ) {
2600 return Ok(());
2601 }
2602 emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2603 }) {
2604 result?;
2605 }
2606 }
2607 }
2608
2609 Ok(out)
2610}
2611
2612pub(crate) fn rel_by_text_scan_rows<S: GraphStorage>(
2613 storage: &S,
2614 params: &BTreeMap<String, LoraValue>,
2615 base_rows: Vec<Row>,
2616 op: &lora_compiler::RelByTextScanExec,
2617 deadline: Option<Instant>,
2618) -> ExecResult<Vec<Row>> {
2619 let eval_ctx = EvalContext { storage, params };
2620 let mut out = Vec::new();
2621
2622 for row in base_rows {
2623 check_optional_deadline(deadline)?;
2624 let query = eval_expr(&op.query, &row, &eval_ctx);
2625 let LoraValue::String(query_str) = &query else {
2626 continue;
2627 };
2628
2629 let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| {
2630 storage.relationship_text_candidates(ty, &op.key, query_str)
2631 });
2632
2633 for rel_id in candidate_ids {
2634 check_optional_deadline(deadline)?;
2635 if let Some(result) = storage.with_relationship(rel_id, |rel| {
2636 if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
2637 return Ok(());
2638 }
2639 let Some(PropertyValue::String(actual)) = rel.properties.get(op.key.as_str())
2640 else {
2641 return Ok(());
2642 };
2643 if !text_predicate_holds(actual, op.predicate, query_str) {
2644 return Ok(());
2645 }
2646 emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2647 }) {
2648 result?;
2649 }
2650 }
2651 }
2652
2653 Ok(out)
2654}
2655
2656pub(crate) fn rel_by_point_scan_rows<S: GraphStorage>(
2657 storage: &S,
2658 params: &BTreeMap<String, LoraValue>,
2659 base_rows: Vec<Row>,
2660 op: &lora_compiler::RelByPointScanExec,
2661 deadline: Option<Instant>,
2662) -> ExecResult<Vec<Row>> {
2663 let eval_ctx = EvalContext { storage, params };
2664 let mut out = Vec::new();
2665
2666 for row in base_rows {
2667 check_optional_deadline(deadline)?;
2668
2669 let probe = match &op.predicate {
2670 lora_compiler::PointPredicate::WithinBBox {
2671 lower_left,
2672 upper_right,
2673 } => {
2674 let ll = eval_expr(lower_left, &row, &eval_ctx);
2675 let ur = eval_expr(upper_right, &row, &eval_ctx);
2676 match (ll, ur) {
2677 (LoraValue::Point(a), LoraValue::Point(b)) => {
2678 Probe::WithinBBox { ll: a, ur: b }
2679 }
2680 _ => continue,
2681 }
2682 }
2683 lora_compiler::PointPredicate::WithinDistance {
2684 center,
2685 max_distance,
2686 inclusive,
2687 } => {
2688 let c = eval_expr(center, &row, &eval_ctx);
2689 let d = eval_expr(max_distance, &row, &eval_ctx);
2690 match (c, d) {
2691 (LoraValue::Point(c), LoraValue::Float(d)) => Probe::WithinDistance {
2692 center: c,
2693 max: d,
2694 inclusive: *inclusive,
2695 },
2696 (LoraValue::Point(c), LoraValue::Int(d)) => Probe::WithinDistance {
2697 center: c,
2698 max: d as f64,
2699 inclusive: *inclusive,
2700 },
2701 _ => continue,
2702 }
2703 }
2704 };
2705
2706 let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| match &probe {
2707 Probe::WithinBBox { ll, ur } => {
2708 let (ranges, n) = lora_store::bbox_x_ranges(ll, ur);
2710 let (lo_y, hi_y) = (ll.y.min(ur.y), ll.y.max(ur.y));
2711 let mut ids = Vec::new();
2712 for (lo_x, hi_x) in &ranges[..n] {
2713 ids.extend(storage.relationship_point_within_bbox(
2714 ty,
2715 &op.key,
2716 (*lo_x, lo_y),
2717 (*hi_x, hi_y),
2718 )?);
2719 }
2720 if n > 1 {
2721 ids.sort_unstable();
2722 ids.dedup();
2723 }
2724 Some(ids)
2725 }
2726 Probe::WithinDistance { center, max, .. } => {
2727 storage.relationship_point_within_distance(ty, &op.key, (center.x, center.y), *max)
2728 }
2729 });
2730
2731 for rel_id in candidate_ids {
2732 check_optional_deadline(deadline)?;
2733 if let Some(result) = storage.with_relationship(rel_id, |rel| {
2734 if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
2735 return Ok(());
2736 }
2737 let Some(PropertyValue::Point(actual)) = rel.properties.get(op.key.as_str()) else {
2738 return Ok(());
2739 };
2740 if !point_predicate_holds(actual, &probe) {
2741 return Ok(());
2742 }
2743 emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2744 }) {
2745 result?;
2746 }
2747 }
2748 }
2749
2750 Ok(out)
2751}
2752
2753pub(crate) struct OrderedRangeCursor {
2764 lo: Option<LoraValue>,
2765 hi: Option<LoraValue>,
2766 mode: OrderedMode,
2767 pending: std::vec::IntoIter<NodeId>,
2768 labels: Vec<Vec<String>>,
2769 key: String,
2770 lo_inclusive: bool,
2771 hi_inclusive: bool,
2772}
2773
2774enum OrderedMode {
2775 Index {
2776 label: String,
2777 descending: bool,
2778 after: Option<(PropertyValue, NodeId)>,
2779 chunk: usize,
2780 exhausted: bool,
2781 },
2782 Buffered,
2783}
2784
2785impl OrderedRangeCursor {
2786 pub(crate) fn new<S: GraphStorage>(
2787 storage: &S,
2788 op: &lora_compiler::NodeByPropertyRangeScanExec,
2789 lo: Option<LoraValue>,
2790 hi: Option<LoraValue>,
2791 ) -> Self {
2792 let descending = matches!(op.order, Some(SortDirection::Desc));
2793 let mut bounds = lo.iter().chain(hi.iter());
2797 let string_bounds = match bounds.next() {
2798 Some(LoraValue::String(_)) => bounds.all(|v| matches!(v, LoraValue::String(_))),
2799 Some(
2800 first @ (LoraValue::Date(_)
2801 | LoraValue::DateTime(_)
2802 | LoraValue::LocalDateTime(_)
2803 | LoraValue::Time(_)
2804 | LoraValue::LocalTime(_)),
2805 ) => bounds.all(|v| std::mem::discriminant(v) == std::mem::discriminant(first)),
2806 _ => false,
2807 };
2808 let index_label = single_label_hint(&op.labels).filter(|label| {
2809 string_bounds
2810 && storage
2811 .node_range_ordered_chunk(label, &op.key, None, None, descending, None, 0)
2812 .is_some()
2813 });
2814
2815 let mut cursor = Self {
2816 lo,
2817 hi,
2818 mode: OrderedMode::Buffered,
2819 pending: Vec::new().into_iter(),
2820 labels: op.labels.clone(),
2821 key: op.key.clone(),
2822 lo_inclusive: op.lo_inclusive,
2823 hi_inclusive: op.hi_inclusive,
2824 };
2825 match index_label {
2826 Some(label) => {
2827 cursor.mode = OrderedMode::Index {
2828 label: label.to_string(),
2829 descending,
2830 after: None,
2831 chunk: 64,
2832 exhausted: false,
2833 }
2834 }
2835 None => {
2836 cursor.pending = cursor
2837 .sorted_candidates(storage, op, descending)
2838 .into_iter()
2839 }
2840 }
2841 cursor
2842 }
2843
2844 pub(crate) fn filter(&self) -> NodeRangeFilter<'_> {
2845 NodeRangeFilter {
2846 labels: &self.labels,
2847 key: &self.key,
2848 lo: self.lo.as_ref(),
2849 lo_inclusive: self.lo_inclusive,
2850 hi: self.hi.as_ref(),
2851 hi_inclusive: self.hi_inclusive,
2852 }
2853 }
2854
2855 fn sorted_candidates<S: GraphStorage>(
2858 &self,
2859 storage: &S,
2860 op: &lora_compiler::NodeByPropertyRangeScanExec,
2861 descending: bool,
2862 ) -> Vec<NodeId> {
2863 let lo_prop = self.lo.clone().and_then(|v| lora_value_to_property(v).ok());
2864 let hi_prop = self.hi.clone().and_then(|v| lora_value_to_property(v).ok());
2865 let candidates = match single_label_hint(&op.labels) {
2866 Some(label) => storage
2867 .node_range_candidates(label, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
2868 .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
2869 None => scan_node_ids_for_label_groups(storage, &op.labels),
2870 };
2871 let filter = self.filter();
2872 let mut keyed: Vec<(LoraValue, NodeId)> = candidates
2873 .into_iter()
2874 .filter(|&id| node_matches_range_filter(storage, id, &filter))
2875 .map(|id| {
2876 let value = storage
2877 .with_node(id, |n| n.properties.get(op.key.as_str()).cloned())
2878 .flatten()
2879 .map(LoraValue::from)
2880 .unwrap_or(LoraValue::Null);
2881 (value, id)
2882 })
2883 .collect();
2884 keyed.sort_by(|(a, ai), (b, bi)| {
2885 let ord = compare_values_total(a, b).then(ai.cmp(bi));
2886 if descending {
2887 ord.reverse()
2888 } else {
2889 ord
2890 }
2891 });
2892 keyed.into_iter().map(|(_, id)| id).collect()
2893 }
2894
2895 pub(crate) fn next_id<S: GraphStorage>(&mut self, storage: &S) -> Option<NodeId> {
2896 loop {
2897 if let Some(id) = self.pending.next() {
2898 if matches!(self.mode, OrderedMode::Buffered)
2899 || node_matches_range_filter(storage, id, &self.filter())
2900 {
2901 return Some(id);
2902 }
2903 continue;
2904 }
2905 let lo_prop = self.lo.clone().and_then(|v| lora_value_to_property(v).ok());
2906 let hi_prop = self.hi.clone().and_then(|v| lora_value_to_property(v).ok());
2907 let OrderedMode::Index {
2908 label,
2909 descending,
2910 after,
2911 chunk,
2912 exhausted,
2913 } = &mut self.mode
2914 else {
2915 return None;
2916 };
2917 if *exhausted {
2918 return None;
2919 }
2920 let ids = storage.node_range_ordered_chunk(
2921 label,
2922 &self.key,
2923 lo_prop.as_ref(),
2924 hi_prop.as_ref(),
2925 *descending,
2926 after.as_ref().map(|(v, id)| (v, *id)),
2927 *chunk,
2928 )?;
2929 if ids.len() < *chunk {
2930 *exhausted = true;
2931 }
2932 let last = *ids.last()?;
2933 let last_value = storage
2934 .with_node(last, |n| n.properties.get(self.key.as_str()).cloned())
2935 .flatten()?;
2936 *after = Some((last_value, last));
2937 *chunk = (*chunk * 2).min(4096);
2938 self.pending = ids.into_iter();
2939 }
2940 }
2941}