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