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