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