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<'a>(props: impl Into<lora_store::PropsRef<'a>>) -> LoraValue {
840 let mut map = BTreeMap::new();
841 for (k, v) in props.into() {
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<'a>(
1238 expected: &LoraValue,
1239 actual: impl Into<lora_store::ValueRef<'a>>,
1240) -> bool {
1241 use lora_store::ValueRef;
1242 match (expected, actual.into()) {
1243 (LoraValue::Null, ValueRef::Null) => true,
1244 (LoraValue::Bool(a), ValueRef::Bool(b)) => *a == b,
1245 (LoraValue::Int(a), ValueRef::Int(b)) => *a == b,
1246 (LoraValue::Float(a), ValueRef::Float(b)) => *a == b,
1247 (LoraValue::Int(a), ValueRef::Float(b)) => (*a as f64) == b,
1248 (LoraValue::Float(a), ValueRef::Int(b)) => *a == (b as f64),
1249 (LoraValue::String(a), ValueRef::String(b)) => a == b,
1250 (expected, ValueRef::Other(actual)) => owned_value_matches(expected, &actual.get()),
1251 _ => false,
1252 }
1253}
1254
1255fn owned_value_matches(expected: &LoraValue, actual: &PropertyValue) -> bool {
1259 match (expected, actual) {
1260 (LoraValue::Binary(a), PropertyValue::Binary(b)) => a == b,
1261
1262 (LoraValue::List(xs), PropertyValue::List(ys)) => {
1263 xs.len() == ys.len()
1264 && xs
1265 .iter()
1266 .zip(ys.iter())
1267 .all(|(x, y)| value_matches_property_value(x, y))
1268 }
1269
1270 (LoraValue::Map(xm), PropertyValue::Map(ym)) => xm.iter().all(|(k, xv)| {
1271 ym.get(k)
1272 .map(|yv| value_matches_property_value(xv, yv))
1273 .unwrap_or(false)
1274 }),
1275
1276 (LoraValue::Date(a), PropertyValue::Date(b)) => a == b,
1277 (LoraValue::DateTime(a), PropertyValue::DateTime(b)) => a == b,
1278 (LoraValue::LocalDateTime(a), PropertyValue::LocalDateTime(b)) => a == b,
1279 (LoraValue::Time(a), PropertyValue::Time(b)) => a == b,
1280 (LoraValue::LocalTime(a), PropertyValue::LocalTime(b)) => a == b,
1281 (LoraValue::Duration(a), PropertyValue::Duration(b)) => a == b,
1282 (LoraValue::Point(a), PropertyValue::Point(b)) => a == b,
1283 (LoraValue::Vector(a), PropertyValue::Vector(b)) => a == b,
1284
1285 _ => false,
1286 }
1287}
1288
1289pub(crate) fn node_matches_property_filter<S: GraphStorage>(
1290 storage: &S,
1291 node_id: NodeId,
1292 labels: &[Vec<String>],
1293 key: &str,
1294 expected: &LoraValue,
1295) -> bool {
1296 storage
1297 .with_node(node_id, |node| {
1298 node_matches_label_groups(node.labels(), labels)
1299 && node
1300 .properties()
1301 .get(key)
1302 .map(|actual| value_matches_property_value(expected, actual))
1303 .unwrap_or(false)
1304 })
1305 .unwrap_or(false)
1306}
1307
1308fn single_label_hint(labels: &[Vec<String>]) -> Option<&str> {
1309 if labels.len() == 1 && labels[0].len() == 1 {
1310 Some(labels[0][0].as_str())
1311 } else {
1312 None
1313 }
1314}
1315
1316fn property_lookup_values(expected: &LoraValue) -> Option<Vec<PropertyValue>> {
1317 let property = lora_value_to_property(expected.clone()).ok()?;
1318 let mut values = vec![property.clone()];
1319
1320 match property {
1321 PropertyValue::Int(i) => {
1322 values.push(PropertyValue::Float(i as f64));
1323 }
1324 PropertyValue::Float(f)
1325 if f.is_finite()
1326 && f.fract() == 0.0
1327 && f >= i64::MIN as f64
1328 && f <= i64::MAX as f64 =>
1329 {
1330 values.push(PropertyValue::Int(f as i64));
1331 }
1332 _ => {}
1333 }
1334
1335 Some(values)
1336}
1337
1338pub(crate) struct NodePropertyCandidates {
1339 pub(crate) ids: Vec<NodeId>,
1340 pub(crate) prefiltered: bool,
1341}
1342
1343pub(crate) fn node_by_property_range_scan_rows<S: GraphStorage>(
1344 storage: &S,
1345 params: &BTreeMap<String, LoraValue>,
1346 base_rows: Vec<Row>,
1347 op: &lora_compiler::NodeByPropertyRangeScanExec,
1348 deadline: Option<Instant>,
1349) -> ExecResult<Vec<Row>> {
1350 let eval_ctx = EvalContext { storage, params };
1351 let mut out = Vec::new();
1352 let mut other_kinds = OtherKindScan::default();
1353
1354 for row in base_rows {
1355 check_optional_deadline(deadline)?;
1356 let lo_value = op.lo.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
1357 let hi_value = op
1361 .hi
1362 .as_ref()
1363 .map(|expr| eval_expr(expr, &row, &eval_ctx))
1364 .filter(|v| !matches!(v, LoraValue::Null));
1365 let lo_prop = lo_value
1366 .clone()
1367 .and_then(|v| lora_value_to_property(v).ok());
1368 let hi_prop = hi_value
1369 .clone()
1370 .and_then(|v| lora_value_to_property(v).ok());
1371 let filter = NodeRangeFilter {
1372 labels: &op.labels,
1373 key: &op.key,
1374 lo: lo_value.as_ref(),
1375 lo_inclusive: op.lo_inclusive,
1376 hi: hi_value.as_ref(),
1377 hi_inclusive: op.hi_inclusive,
1378 };
1379 let bounds = [lo_value.as_ref(), hi_value.as_ref()];
1380 let bound_id = bound_node_id_for_expand(&row, op.var)?;
1381
1382 if let Some(existing_id) = bound_id {
1383 if node_matches_range_filter(storage, existing_id, &filter)
1384 || bound_node_has_other_kind(storage, existing_id, op, bounds)
1385 {
1386 out.push(row);
1387 }
1388 continue;
1389 }
1390
1391 for id in other_kind_node_ids(storage, op, bounds, &mut other_kinds) {
1393 let mut new_row = row.clone();
1394 new_row.insert(op.var, LoraValue::Node(id));
1395 out.push(new_row);
1396 }
1397
1398 if op.order.is_some() {
1399 let mut cursor =
1400 OrderedRangeCursor::new(storage, op, lo_value.clone(), hi_value.clone());
1401 while let Some(id) = cursor.next_id(storage) {
1402 check_optional_deadline(deadline)?;
1403 let mut new_row = row.clone();
1404 new_row.insert(op.var, LoraValue::Node(id));
1405 out.push(new_row);
1406 }
1407 continue;
1408 }
1409
1410 let candidate_ids = match single_label_hint(&op.labels) {
1411 Some(label) => storage
1412 .node_range_candidates(label, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
1413 .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1414 None => scan_node_ids_for_label_groups(storage, &op.labels),
1415 };
1416
1417 for id in candidate_ids {
1418 check_optional_deadline(deadline)?;
1419 if node_matches_range_filter(storage, id, &filter) {
1420 let mut new_row = row.clone();
1421 new_row.insert(op.var, LoraValue::Node(id));
1422 out.push(new_row);
1423 }
1424 }
1425 }
1426
1427 Ok(out)
1428}
1429
1430pub(crate) fn node_by_text_scan_rows<S: GraphStorage>(
1431 storage: &S,
1432 params: &BTreeMap<String, LoraValue>,
1433 base_rows: Vec<Row>,
1434 op: &lora_compiler::NodeByTextScanExec,
1435 deadline: Option<Instant>,
1436) -> ExecResult<Vec<Row>> {
1437 let eval_ctx = EvalContext { storage, params };
1438 let mut out = Vec::new();
1439
1440 for row in base_rows {
1441 check_optional_deadline(deadline)?;
1442 let query = eval_expr(&op.query, &row, &eval_ctx);
1443 let LoraValue::String(query_str) = &query else {
1444 continue;
1447 };
1448
1449 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
1450 if node_matches_text_filter(
1451 storage,
1452 existing_id,
1453 &op.labels,
1454 &op.key,
1455 op.predicate,
1456 query_str,
1457 ) {
1458 out.push(row);
1459 }
1460 continue;
1461 }
1462
1463 let candidate_ids = match single_label_hint(&op.labels) {
1464 Some(label) => storage
1465 .node_text_candidates(label, &op.key, query_str)
1466 .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1467 None => scan_node_ids_for_label_groups(storage, &op.labels),
1468 };
1469
1470 for id in candidate_ids {
1471 check_optional_deadline(deadline)?;
1472 if node_matches_text_filter(storage, id, &op.labels, &op.key, op.predicate, query_str) {
1473 let mut new_row = row.clone();
1474 new_row.insert(op.var, LoraValue::Node(id));
1475 out.push(new_row);
1476 }
1477 }
1478 }
1479
1480 Ok(out)
1481}
1482
1483fn temporal_bound_kinds(bounds: [Option<&LoraValue>; 2]) -> Vec<(&'static str, &LoraValue)> {
1485 let mut kinds: Vec<(&'static str, &LoraValue)> = Vec::with_capacity(2);
1486 for bound in bounds.into_iter().flatten() {
1487 if let Some(kind) = temporal_kind_name(bound) {
1488 if !kinds.iter().any(|(k, _)| *k == kind) {
1489 kinds.push((kind, bound));
1490 }
1491 }
1492 }
1493 kinds
1494}
1495
1496fn is_other_temporal_kind<'a>(
1498 value: impl Into<lora_store::ValueRef<'a>>,
1499 kinds: &[(&'static str, &LoraValue)],
1500) -> bool {
1501 let lora_store::ValueRef::Other(value) = value.into() else {
1503 return false;
1504 };
1505 temporal_kind_name(&LoraValue::from(value.to_owned()))
1506 .is_some_and(|kind| kinds.iter().any(|(k, _)| *k != kind))
1507}
1508
1509#[derive(Default)]
1513pub(crate) struct OtherKindScan<Id> {
1514 temporals: Option<Vec<(Id, &'static str)>>,
1515}
1516
1517impl<Id: Copy> OtherKindScan<Id> {
1518 fn ids(
1519 &mut self,
1520 kinds: &[(&'static str, &LoraValue)],
1521 scan: impl FnOnce() -> Vec<(Id, &'static str)>,
1522 ) -> Vec<Id> {
1523 self.temporals
1524 .get_or_insert_with(scan)
1525 .iter()
1526 .filter(|(_, kind)| kinds.iter().any(|(k, _)| k != kind))
1527 .map(|(id, _)| *id)
1528 .collect()
1529 }
1530}
1531
1532pub(crate) fn other_kind_node_ids<S: GraphStorage>(
1542 storage: &S,
1543 op: &lora_compiler::NodeByPropertyRangeScanExec,
1544 bounds: [Option<&LoraValue>; 2],
1545 fallback: &mut OtherKindScan<NodeId>,
1546) -> Vec<NodeId> {
1547 let kinds = temporal_bound_kinds(bounds);
1548 if kinds.is_empty() {
1549 return Vec::new();
1550 }
1551 let index_label = op
1553 .labels
1554 .iter()
1555 .find(|group| group.len() == 1)
1556 .map(|group| group[0].as_str());
1557 let mut ids = BTreeSet::new();
1558 let mut need_scan = index_label.is_none();
1559 if let Some(label) = index_label {
1560 for (_, bound) in &kinds {
1561 let listed = lora_value_to_property((*bound).clone())
1562 .ok()
1563 .and_then(|like| storage.node_range_other_temporal_kind_ids(label, &op.key, &like));
1564 match listed {
1565 Some(listed) => ids.extend(listed),
1566 None => need_scan = true,
1567 }
1568 }
1569 }
1572 if need_scan {
1573 ids.extend(fallback.ids(&kinds, || {
1574 scan_node_ids_for_label_groups(storage, &op.labels)
1575 .into_iter()
1576 .filter_map(|id| {
1577 storage
1578 .with_node(id, |n| {
1579 n.properties()
1580 .get(op.key.as_str())
1581 .and_then(|v| temporal_kind_name(&LoraValue::from(v)))
1582 })
1583 .flatten()
1584 .map(|kind| (id, kind))
1585 })
1586 .collect()
1587 }));
1588 return ids.into_iter().collect();
1589 }
1590 let single = op.labels.len() == 1 && op.labels[0].len() == 1;
1591 ids.into_iter()
1592 .filter(|id| {
1593 single
1594 || storage
1595 .with_node(*id, |n| node_matches_label_groups(n.labels(), &op.labels))
1596 .unwrap_or(false)
1597 })
1598 .collect()
1599}
1600
1601pub(crate) fn bound_node_has_other_kind<S: GraphStorage>(
1605 storage: &S,
1606 id: NodeId,
1607 op: &lora_compiler::NodeByPropertyRangeScanExec,
1608 bounds: [Option<&LoraValue>; 2],
1609) -> bool {
1610 let kinds = temporal_bound_kinds(bounds);
1611 !kinds.is_empty()
1612 && storage
1613 .with_node(id, |n| {
1614 node_matches_label_groups(n.labels(), &op.labels)
1615 && n.properties()
1616 .get(op.key.as_str())
1617 .is_some_and(|v| is_other_temporal_kind(v, &kinds))
1618 })
1619 .unwrap_or(false)
1620}
1621
1622fn other_kind_rel_ids<S: GraphStorage>(
1624 storage: &S,
1625 op: &lora_compiler::RelByPropertyRangeScanExec,
1626 bounds: [Option<&LoraValue>; 2],
1627 fallback: &mut OtherKindScan<RelationshipId>,
1628) -> Vec<RelationshipId> {
1629 let kinds = temporal_bound_kinds(bounds);
1630 if kinds.is_empty() {
1631 return Vec::new();
1632 }
1633 let mut ids = BTreeSet::new();
1634 let mut need_scan = op.types.is_empty();
1635 for ty in &op.types {
1636 for (_, bound) in &kinds {
1637 let listed = lora_value_to_property((*bound).clone())
1638 .ok()
1639 .and_then(|like| {
1640 storage.relationship_range_other_temporal_kind_ids(ty, &op.key, &like)
1641 });
1642 match listed {
1643 Some(listed) => ids.extend(listed),
1644 None => need_scan = true,
1645 }
1646 }
1647 }
1648 if need_scan {
1649 ids.extend(fallback.ids(&kinds, || {
1650 rel_candidate_ids(storage, &op.types, |_| None)
1651 .into_iter()
1652 .filter_map(|id| {
1653 storage
1654 .with_relationship(id, |r| {
1655 r.properties()
1656 .get(op.key.as_str())
1657 .and_then(|v| temporal_kind_name(&LoraValue::from(v)))
1658 })
1659 .flatten()
1660 .map(|kind| (id, kind))
1661 })
1662 .collect()
1663 }));
1664 }
1665 ids.into_iter().collect()
1666}
1667
1668fn temporal_kind_name(value: &LoraValue) -> Option<&'static str> {
1669 Some(match value {
1670 LoraValue::Date(_) => "DATE",
1671 LoraValue::DateTime(_) => "DATETIME",
1672 LoraValue::LocalDateTime(_) => "LOCAL_DATETIME",
1673 LoraValue::Time(_) => "TIME",
1674 LoraValue::LocalTime(_) => "LOCAL_TIME",
1675 _ => return None,
1676 })
1677}
1678
1679pub(crate) struct NodeRangeFilter<'a> {
1680 labels: &'a [Vec<String>],
1681 key: &'a str,
1682 lo: Option<&'a LoraValue>,
1683 lo_inclusive: bool,
1684 hi: Option<&'a LoraValue>,
1685 hi_inclusive: bool,
1686}
1687
1688fn node_matches_range_filter<S: GraphStorage>(
1689 storage: &S,
1690 id: NodeId,
1691 filter: &NodeRangeFilter<'_>,
1692) -> bool {
1693 storage
1694 .with_node(id, |n| {
1695 if !node_matches_label_groups(n.labels(), filter.labels) {
1696 return false;
1697 }
1698 let Some(actual) = n.properties().get(filter.key) else {
1699 return false;
1700 };
1701 let actual_lv = LoraValue::from(actual);
1702 range_predicate_holds(
1703 &actual_lv,
1704 filter.lo,
1705 filter.lo_inclusive,
1706 filter.hi,
1707 filter.hi_inclusive,
1708 )
1709 })
1710 .unwrap_or(false)
1711}
1712
1713fn node_matches_text_filter<S: GraphStorage>(
1714 storage: &S,
1715 id: NodeId,
1716 labels: &[Vec<String>],
1717 key: &str,
1718 predicate: lora_compiler::TextPredicate,
1719 query: &str,
1720) -> bool {
1721 storage
1722 .with_node(id, |n| {
1723 if !node_matches_label_groups(n.labels(), labels) {
1724 return false;
1725 }
1726 let Some(lora_store::ValueRef::String(actual)) = n.properties().get(key) else {
1727 return false;
1728 };
1729 text_predicate_holds(actual, predicate, query)
1730 })
1731 .unwrap_or(false)
1732}
1733
1734fn text_predicate_holds(
1735 actual: &str,
1736 predicate: lora_compiler::TextPredicate,
1737 query: &str,
1738) -> bool {
1739 match predicate {
1740 lora_compiler::TextPredicate::StartsWith => actual.starts_with(query),
1741 lora_compiler::TextPredicate::EndsWith => actual.ends_with(query),
1742 lora_compiler::TextPredicate::Contains => actual.contains(query),
1743 }
1744}
1745
1746fn range_predicate_holds(
1747 actual: &LoraValue,
1748 lo: Option<&LoraValue>,
1749 lo_inclusive: bool,
1750 hi: Option<&LoraValue>,
1751 hi_inclusive: bool,
1752) -> bool {
1753 if let Some(lo) = lo {
1754 match range_comparison(actual, lo) {
1755 None => return false,
1756 Some(Ordering::Less) => return false,
1757 Some(Ordering::Equal) if !lo_inclusive => return false,
1758 _ => {}
1759 }
1760 }
1761 if let Some(hi) = hi {
1762 match range_comparison(actual, hi) {
1763 None => return false,
1764 Some(Ordering::Greater) => return false,
1765 Some(Ordering::Equal) if !hi_inclusive => return false,
1766 _ => {}
1767 }
1768 }
1769 true
1770}
1771
1772fn range_comparison(actual: &LoraValue, bound: &LoraValue) -> Option<Ordering> {
1773 match (actual, bound) {
1774 (LoraValue::Null, _) | (_, LoraValue::Null) => None,
1775 (LoraValue::String(a), LoraValue::String(b)) => Some(a.cmp(b)),
1776 (
1777 LoraValue::Date(_)
1778 | LoraValue::DateTime(_)
1779 | LoraValue::LocalDateTime(_)
1780 | LoraValue::Time(_)
1781 | LoraValue::LocalTime(_),
1782 _,
1783 ) => actual.temporal_cmp(bound),
1784 (LoraValue::Duration(a), LoraValue::Duration(b)) => a
1785 .total_seconds_approx()
1786 .partial_cmp(&b.total_seconds_approx()),
1787 _ => actual.as_f64()?.partial_cmp(&bound.as_f64()?),
1788 }
1789}
1790
1791pub(crate) fn node_by_point_scan_rows<S: GraphStorage>(
1792 storage: &S,
1793 params: &BTreeMap<String, LoraValue>,
1794 base_rows: Vec<Row>,
1795 op: &lora_compiler::NodeByPointScanExec,
1796 deadline: Option<Instant>,
1797) -> ExecResult<Vec<Row>> {
1798 let eval_ctx = EvalContext { storage, params };
1799 let mut out = Vec::new();
1800
1801 for row in base_rows {
1802 check_optional_deadline(deadline)?;
1803
1804 let probe = match &op.predicate {
1808 lora_compiler::PointPredicate::WithinBBox {
1809 lower_left,
1810 upper_right,
1811 } => {
1812 let ll = eval_expr(lower_left, &row, &eval_ctx);
1813 let ur = eval_expr(upper_right, &row, &eval_ctx);
1814 match (ll, ur) {
1815 (LoraValue::Point(a), LoraValue::Point(b)) => {
1816 Probe::WithinBBox { ll: a, ur: b }
1817 }
1818 _ => continue,
1819 }
1820 }
1821 lora_compiler::PointPredicate::WithinDistance {
1822 center,
1823 max_distance,
1824 inclusive,
1825 } => {
1826 let c = eval_expr(center, &row, &eval_ctx);
1827 let d = eval_expr(max_distance, &row, &eval_ctx);
1828 match (c, d) {
1829 (LoraValue::Point(c), LoraValue::Float(d)) => Probe::WithinDistance {
1830 center: c,
1831 max: d,
1832 inclusive: *inclusive,
1833 },
1834 (LoraValue::Point(c), LoraValue::Int(d)) => Probe::WithinDistance {
1835 center: c,
1836 max: d as f64,
1837 inclusive: *inclusive,
1838 },
1839 _ => continue,
1840 }
1841 }
1842 };
1843
1844 if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
1845 if node_matches_point_filter(storage, existing_id, &op.labels, &op.key, &probe) {
1846 out.push(row);
1847 }
1848 continue;
1849 }
1850
1851 let candidate_ids = match single_label_hint(&op.labels) {
1852 Some(label) => match &probe {
1853 Probe::WithinBBox { ll, ur } => {
1854 let (ranges, n) = lora_store::bbox_x_ranges(ll, ur);
1856 let (lo_y, hi_y) = (ll.y.min(ur.y), ll.y.max(ur.y));
1857 let mut ids = Vec::new();
1858 let mut indexed = true;
1859 for (lo_x, hi_x) in &ranges[..n] {
1860 match storage.node_point_within_bbox(
1861 label,
1862 &op.key,
1863 (*lo_x, lo_y),
1864 (*hi_x, hi_y),
1865 ) {
1866 Some(found) => ids.extend(found),
1867 None => indexed = false,
1868 }
1869 }
1870 if !indexed {
1871 scan_node_ids_for_label_groups(storage, &op.labels)
1872 } else {
1873 if n > 1 {
1874 ids.sort_unstable();
1875 ids.dedup();
1876 }
1877 ids
1878 }
1879 }
1880 Probe::WithinDistance { center, max, .. } => storage
1881 .node_point_within_distance(label, &op.key, (center.x, center.y), *max)
1882 .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
1883 },
1884 None => scan_node_ids_for_label_groups(storage, &op.labels),
1885 };
1886
1887 for id in candidate_ids {
1888 check_optional_deadline(deadline)?;
1889 if node_matches_point_filter(storage, id, &op.labels, &op.key, &probe) {
1890 let mut new_row = row.clone();
1891 new_row.insert(op.var, LoraValue::Node(id));
1892 out.push(new_row);
1893 }
1894 }
1895 }
1896
1897 Ok(out)
1898}
1899
1900fn node_matches_point_filter<S: GraphStorage>(
1901 storage: &S,
1902 id: NodeId,
1903 labels: &[Vec<String>],
1904 key: &str,
1905 probe: &Probe,
1906) -> bool {
1907 storage
1908 .with_node(id, |n| {
1909 if !node_matches_label_groups(n.labels(), labels) {
1910 return false;
1911 }
1912 let Some(lora_store::ValueRef::Other(value)) = n.properties().get(key) else {
1913 return false;
1914 };
1915 match &*value.get() {
1916 PropertyValue::Point(point) => point_predicate_holds(point, probe),
1917 _ => false,
1918 }
1919 })
1920 .unwrap_or(false)
1921}
1922
1923enum Probe {
1924 WithinBBox {
1925 ll: lora_store::LoraPoint,
1926 ur: lora_store::LoraPoint,
1927 },
1928 WithinDistance {
1929 center: lora_store::LoraPoint,
1930 max: f64,
1931 inclusive: bool,
1932 },
1933}
1934
1935fn point_predicate_holds(actual: &lora_store::LoraPoint, probe: &Probe) -> bool {
1936 match probe {
1937 Probe::WithinBBox { ll, ur } => lora_store::bbox_contains(actual, ll, ur).unwrap_or(false),
1938 Probe::WithinDistance {
1939 center,
1940 max,
1941 inclusive,
1942 } => {
1943 let Some(d) = lora_store::point_distance(actual, center) else {
1944 return false;
1945 };
1946 if *inclusive {
1947 d <= *max
1948 } else {
1949 d < *max
1950 }
1951 }
1952 }
1953}
1954
1955pub(crate) fn indexed_node_property_candidates<S: GraphStorage>(
1956 storage: &S,
1957 labels: &[Vec<String>],
1958 key: &str,
1959 expected: &LoraValue,
1960) -> NodePropertyCandidates {
1961 let Some(values) = property_lookup_values(expected) else {
1962 return NodePropertyCandidates {
1963 ids: scan_node_ids_for_label_groups(storage, labels),
1964 prefiltered: false,
1965 };
1966 };
1967
1968 let label_hint = single_label_hint(labels);
1969 let mut seen = BTreeSet::new();
1970 let mut out = Vec::new();
1971 for value in values {
1972 for id in storage.find_node_ids_by_property(label_hint, key, &value) {
1973 if seen.insert(id) {
1974 out.push(id);
1975 }
1976 }
1977 }
1978 NodePropertyCandidates {
1979 ids: out,
1980 prefiltered: labels.is_empty() || label_hint.is_some(),
1981 }
1982}
1983
1984pub(crate) fn property_scan_candidates<S: GraphStorage>(
1990 storage: &S,
1991 labels: &[Vec<String>],
1992 key: &str,
1993 expected: &LoraValue,
1994 in_list: bool,
1995) -> NodePropertyCandidates {
1996 if !in_list {
1997 return indexed_node_property_candidates(storage, labels, key, expected);
1998 }
1999 match expected {
2000 LoraValue::Null => NodePropertyCandidates {
2001 ids: Vec::new(),
2002 prefiltered: true,
2003 },
2004 LoraValue::List(items) => {
2005 let mut seen = BTreeSet::new();
2006 let mut ids = Vec::new();
2007 for item in items {
2008 if matches!(item, LoraValue::Null) {
2009 continue;
2010 }
2011 let candidates = indexed_node_property_candidates(storage, labels, key, item);
2012 for id in candidates.ids {
2013 if seen.contains(&id) {
2014 continue;
2015 }
2016 if !candidates.prefiltered
2017 && !node_matches_property_filter(storage, id, labels, key, item)
2018 {
2019 continue;
2020 }
2021 seen.insert(id);
2022 ids.push(id);
2023 }
2024 }
2025 NodePropertyCandidates {
2026 ids,
2027 prefiltered: true,
2028 }
2029 }
2030 _ => NodePropertyCandidates {
2031 ids: scan_node_ids_for_label_groups(storage, labels)
2032 .into_iter()
2033 .filter(|&id| {
2034 storage
2035 .with_node(id, |n| node_matches_label_groups(n.labels(), labels))
2036 .unwrap_or(false)
2037 })
2038 .collect(),
2039 prefiltered: true,
2040 },
2041 }
2042}
2043
2044pub(crate) fn property_scan_matches<S: GraphStorage>(
2048 storage: &S,
2049 node_id: NodeId,
2050 labels: &[Vec<String>],
2051 key: &str,
2052 expected: &LoraValue,
2053 in_list: bool,
2054) -> bool {
2055 if !in_list {
2056 return node_matches_property_filter(storage, node_id, labels, key, expected);
2057 }
2058 match expected {
2059 LoraValue::Null => false,
2060 LoraValue::List(items) => items.iter().any(|item| {
2061 !matches!(item, LoraValue::Null)
2062 && node_matches_property_filter(storage, node_id, labels, key, item)
2063 }),
2064 _ => storage
2065 .with_node(node_id, |n| node_matches_label_groups(n.labels(), labels))
2066 .unwrap_or(false),
2067 }
2068}
2069
2070pub(crate) fn build_path_value<S: GraphStorage>(
2076 row: &Row,
2077 node_vars: &[VarId],
2078 rel_vars: &[VarId],
2079 storage: &S,
2080) -> LoraValue {
2081 let (raw_nodes, rels, has_var_len) = path_bindings(row, node_vars, rel_vars);
2082
2083 let nodes = if has_var_len && !rels.is_empty() && raw_nodes.len() == 2 {
2084 reconstruct_var_len_nodes(raw_nodes[0], &rels, storage)
2085 } else {
2086 raw_nodes
2087 };
2088
2089 LoraValue::Path(LoraPath { nodes, rels })
2090}
2091
2092#[inline]
2093fn path_bindings(
2094 row: &Row,
2095 node_vars: &[VarId],
2096 rel_vars: &[VarId],
2097) -> (Vec<NodeId>, Vec<RelationshipId>, bool) {
2098 let mut raw_nodes = Vec::new();
2099 let mut rels = Vec::new();
2100 let mut has_var_len = false;
2101
2102 for &nv in node_vars {
2103 match row.get(nv) {
2104 Some(LoraValue::Node(id)) => raw_nodes.push(*id),
2105 Some(LoraValue::List(items)) => {
2106 for item in items {
2107 if let LoraValue::Node(id) = item {
2108 raw_nodes.push(*id);
2109 }
2110 }
2111 }
2112 _ => {}
2113 }
2114 }
2115
2116 for &rv in rel_vars {
2117 match row.get(rv) {
2118 Some(LoraValue::Relationship(id)) => rels.push(*id),
2119 Some(LoraValue::List(items)) => {
2120 has_var_len = true;
2121 for item in items {
2122 if let LoraValue::Relationship(id) = item {
2123 rels.push(*id);
2124 }
2125 }
2126 }
2127 _ => {}
2128 }
2129 }
2130
2131 (raw_nodes, rels, has_var_len)
2132}
2133
2134#[inline]
2135fn reconstruct_var_len_nodes<S: GraphStorage>(
2136 start: NodeId,
2137 rels: &[RelationshipId],
2138 storage: &S,
2139) -> Vec<NodeId> {
2140 let mut ordered = Vec::with_capacity(rels.len() + 1);
2141 ordered.push(start);
2142 let mut current = start;
2143 for &rel_id in rels {
2144 if let Some((src, dst)) = storage.relationship_endpoints(rel_id) {
2145 let next = if src == current { dst } else { src };
2146 ordered.push(next);
2147 current = next;
2148 }
2149 }
2150 ordered
2151}
2152
2153fn type_rank(v: &LoraValue) -> u8 {
2154 match v {
2155 LoraValue::Null => 0,
2156 LoraValue::Bool(_) => 1,
2157 LoraValue::Int(_) | LoraValue::Float(_) => 2,
2158 LoraValue::String(_) => 3,
2159 LoraValue::Binary(_) => 4,
2160 LoraValue::Date(_) => 5,
2161 LoraValue::DateTime(_) => 6,
2162 LoraValue::LocalDateTime(_) => 7,
2163 LoraValue::Time(_) => 8,
2164 LoraValue::LocalTime(_) => 9,
2165 LoraValue::Duration(_) => 10,
2166 LoraValue::Point(_) => 11,
2167 LoraValue::Vector(_) => 12,
2168 LoraValue::List(_) => 13,
2169 LoraValue::Map(_) => 14,
2170 LoraValue::Node(_) => 15,
2171 LoraValue::Relationship(_) => 16,
2172 LoraValue::Path(_) => 17,
2173 }
2174}
2175
2176pub(crate) fn node_matches_label_groups<'a>(
2180 node_labels: impl Into<lora_store::LabelsRef<'a>>,
2181 groups: &[Vec<String>],
2182) -> bool {
2183 let node_labels = node_labels.into();
2184 groups
2185 .iter()
2186 .all(|group| group.iter().any(|l| node_labels.has(l)))
2187}
2188
2189pub(crate) fn scan_node_ids_for_label_groups<S: GraphStorage>(
2192 storage: &S,
2193 groups: &[Vec<String>],
2194) -> Vec<NodeId> {
2195 if groups.is_empty() {
2196 return storage.all_node_ids();
2197 }
2198 if groups.len() == 1 {
2199 return label_group_candidate_ids(storage, &groups[0]);
2200 }
2201
2202 let mut best: Option<Vec<NodeId>> = None;
2203 for group in groups {
2204 let ids = label_group_candidate_ids(storage, group);
2205 if ids.is_empty() {
2206 return Vec::new();
2207 }
2208 if best
2209 .as_ref()
2210 .map(|current| ids.len() < current.len())
2211 .unwrap_or(true)
2212 {
2213 best = Some(ids);
2214 }
2215 }
2216
2217 best.unwrap_or_default()
2218}
2219
2220pub(crate) fn label_group_candidates_prefiltered(groups: &[Vec<String>]) -> bool {
2221 groups.len() <= 1
2222}
2223
2224fn label_group_candidate_ids<S: GraphStorage>(storage: &S, group: &[String]) -> Vec<NodeId> {
2225 match group {
2226 [] => Vec::new(),
2227 [label] => storage.node_ids_by_label(label),
2228 labels => {
2229 let mut seen = BTreeSet::new();
2230 let mut out = Vec::new();
2231 for label in labels {
2232 for id in storage.node_ids_by_label(label) {
2233 if seen.insert(id) {
2234 out.push(id);
2235 }
2236 }
2237 }
2238 out
2239 }
2240 }
2241}
2242
2243pub(crate) fn hydrate_node_record<'a>(node: impl Into<lora_store::NodeRef<'a>>) -> LoraValue {
2244 let node = node.into();
2245 let mut map = BTreeMap::new();
2246 map.insert("kind".to_string(), LoraValue::String("node".to_string()));
2247 map.insert("id".to_string(), LoraValue::Int(node.id() as i64));
2248 map.insert(
2249 "labels".to_string(),
2250 LoraValue::List(
2251 node.labels()
2252 .strs()
2253 .map(|s| LoraValue::String(s.to_string()))
2254 .collect(),
2255 ),
2256 );
2257 map.insert(
2258 "properties".to_string(),
2259 properties_to_value_map(node.properties()),
2260 );
2261 LoraValue::Map(map)
2262}
2263
2264pub(crate) fn hydrate_relationship_record<'a>(rel: impl Into<lora_store::RelRef<'a>>) -> LoraValue {
2265 let rel = rel.into();
2266 let mut map = BTreeMap::new();
2267 map.insert(
2268 "kind".to_string(),
2269 LoraValue::String("relationship".to_string()),
2270 );
2271 map.insert("id".to_string(), LoraValue::Int(rel.id() as i64));
2272 map.insert("startId".to_string(), LoraValue::Int(rel.src() as i64));
2273 map.insert("endId".to_string(), LoraValue::Int(rel.dst() as i64));
2274 map.insert(
2275 "type".to_string(),
2276 LoraValue::String(rel.rel_type().to_string()),
2277 );
2278 map.insert(
2279 "properties".to_string(),
2280 properties_to_value_map(rel.properties()),
2281 );
2282 LoraValue::Map(map)
2283}
2284
2285pub(super) fn flatten_label_groups(groups: &[Vec<String>]) -> Vec<String> {
2288 groups.iter().flat_map(|g| g.iter().cloned()).collect()
2289}
2290
2291#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
2292pub(crate) enum GroupValueKey {
2293 Null,
2294 Bool(bool),
2295 Int(i64),
2296 Float(String),
2297 String(String),
2298 Binary(Vec<Vec<u8>>),
2299 List(Vec<GroupValueKey>),
2300 Map(Vec<(String, GroupValueKey)>),
2301 Node(u64),
2302 Relationship(u64),
2303}
2304
2305impl GroupValueKey {
2306 pub(crate) fn from_value(v: &LoraValue) -> Self {
2307 match v {
2308 LoraValue::Null => Self::Null,
2309 LoraValue::Bool(x) => Self::Bool(*x),
2310 LoraValue::Int(x) => Self::Int(*x),
2311 LoraValue::Float(x) => Self::Float(x.to_string()),
2312 LoraValue::String(x) => Self::String(x.clone()),
2313 LoraValue::Binary(x) => Self::Binary(x.segments().to_vec()),
2314 LoraValue::List(xs) => Self::List(xs.iter().map(Self::from_value).collect()),
2315 LoraValue::Map(m) => Self::Map(
2316 m.iter()
2317 .map(|(k, v)| (k.clone(), Self::from_value(v)))
2318 .collect(),
2319 ),
2320 LoraValue::Node(id) => Self::Node(*id),
2321 LoraValue::Relationship(id) => Self::Relationship(*id),
2322 LoraValue::Path(_) => Self::Null,
2323 LoraValue::Date(d) => Self::String(d.to_string()),
2325 LoraValue::DateTime(dt) => Self::String(dt.to_string()),
2326 LoraValue::LocalDateTime(dt) => Self::String(dt.to_string()),
2327 LoraValue::Time(t) => Self::String(t.to_string()),
2328 LoraValue::LocalTime(t) => Self::String(t.to_string()),
2329 LoraValue::Duration(dur) => Self::String(dur.to_string()),
2330 LoraValue::Point(p) => Self::String(p.to_string()),
2331 LoraValue::Vector(v) => Self::String(format!("vector:{}", v.to_key_string())),
2332 }
2333 }
2334}
2335
2336const MAX_VAR_LEN_HOPS: u64 = 100;
2348
2349pub(crate) fn resolve_range(range: &RangeLiteral) -> (u64, u64) {
2350 let min_hops = range.start.unwrap_or(1);
2351 let max_hops = range.end.unwrap_or(MAX_VAR_LEN_HOPS);
2352 (min_hops, max_hops)
2353}
2354
2355pub(crate) struct VarLenResult {
2357 pub(crate) dst_node_id: NodeId,
2359 pub(crate) rel_ids: Vec<u64>,
2361}
2362
2363pub(crate) fn variable_length_expand<S: GraphStorage>(
2370 storage: &S,
2371 start_node_id: NodeId,
2372 direction: Direction,
2373 types: &[String],
2374 min_hops: u64,
2375 max_hops: u64,
2376 bind_relationships: bool,
2377) -> Vec<VarLenResult> {
2378 let mut results = Vec::new();
2379
2380 let mut frontier: Vec<(NodeId, Vec<u64>)> = vec![(start_node_id, Vec::new())];
2382
2383 for depth in 1..=max_hops {
2384 let is_last_hop = depth == max_hops;
2388 let mut next_frontier: Vec<(NodeId, Vec<u64>)> = Vec::new();
2389
2390 for (current_node, rels_used) in &frontier {
2391 for (rel_id, neighbor_id) in storage.expand_ids(*current_node, direction, types) {
2394 if rels_used.contains(&rel_id) {
2397 continue;
2398 }
2399
2400 if is_last_hop {
2401 if depth >= min_hops {
2404 let mut rel_ids = Vec::with_capacity(rels_used.len() + 1);
2405 rel_ids.extend_from_slice(rels_used);
2406 rel_ids.push(rel_id);
2407 results.push(VarLenResult {
2408 dst_node_id: neighbor_id,
2409 rel_ids: if bind_relationships {
2410 rel_ids
2411 } else {
2412 Vec::new()
2413 },
2414 });
2415 }
2416 continue;
2417 }
2418
2419 let mut new_rels = Vec::with_capacity(rels_used.len() + 1);
2420 new_rels.extend_from_slice(rels_used);
2421 new_rels.push(rel_id);
2422
2423 if depth >= min_hops {
2424 results.push(VarLenResult {
2425 dst_node_id: neighbor_id,
2426 rel_ids: if bind_relationships {
2427 new_rels.clone()
2428 } else {
2429 Vec::new()
2430 },
2431 });
2432 }
2433
2434 next_frontier.push((neighbor_id, new_rels));
2435 }
2436 }
2437
2438 if is_last_hop || next_frontier.is_empty() {
2439 break;
2440 }
2441
2442 frontier = next_frontier;
2443 }
2444
2445 if min_hops == 0 {
2447 results.insert(
2448 0,
2449 VarLenResult {
2450 dst_node_id: start_node_id,
2451 rel_ids: Vec::new(),
2452 },
2453 );
2454 }
2455
2456 results
2457}
2458
2459pub(crate) fn filter_shortest_paths(rows: Vec<Row>, path_var: VarId, all: bool) -> Vec<Row> {
2462 if rows.is_empty() {
2463 return rows;
2464 }
2465
2466 let lengths: Vec<usize> = rows
2468 .iter()
2469 .map(|row| match row.get(path_var) {
2470 Some(LoraValue::Path(p)) => p.rels.len(),
2471 _ => usize::MAX,
2472 })
2473 .collect();
2474
2475 let min_len = lengths.iter().copied().min().unwrap_or(usize::MAX);
2476
2477 let mut result: Vec<Row> = rows
2478 .into_iter()
2479 .zip(lengths.iter())
2480 .filter(|(_, len)| **len == min_len)
2481 .map(|(row, _)| row)
2482 .collect();
2483
2484 if !all && result.len() > 1 {
2485 result.truncate(1);
2486 }
2487
2488 result
2489}
2490
2491struct RelEnds {
2509 id: lora_store::RelationshipId,
2510 src: NodeId,
2511 dst: NodeId,
2512}
2513
2514fn emit_rel_rows(
2515 direction: Direction,
2516 src_var: VarId,
2517 rel_var: VarId,
2518 dst_var: VarId,
2519 rel: lora_store::RelRef<'_>,
2520 base: &Row,
2521 out: &mut Vec<Row>,
2522) -> ExecResult<()> {
2523 let rel = RelEnds {
2524 id: rel.id(),
2525 src: rel.src(),
2526 dst: rel.dst(),
2527 };
2528 match direction {
2529 Direction::Right => {
2530 emit_one_rel_row(
2531 src_var, rel_var, dst_var, rel.src, rel.id, rel.dst, base, out,
2532 )?;
2533 }
2534 Direction::Left => {
2535 emit_one_rel_row(
2536 src_var, rel_var, dst_var, rel.dst, rel.id, rel.src, base, out,
2537 )?;
2538 }
2539 Direction::Undirected => {
2540 emit_one_rel_row(
2541 src_var, rel_var, dst_var, rel.src, rel.id, rel.dst, base, out,
2542 )?;
2543 if rel.src != rel.dst {
2547 emit_one_rel_row(
2548 src_var, rel_var, dst_var, rel.dst, rel.id, rel.src, base, out,
2549 )?;
2550 }
2551 }
2552 }
2553 Ok(())
2554}
2555
2556#[allow(clippy::too_many_arguments)]
2557fn emit_one_rel_row(
2558 src_var: VarId,
2559 rel_var: VarId,
2560 dst_var: VarId,
2561 src_id: NodeId,
2562 rel_id: RelationshipId,
2563 dst_id: NodeId,
2564 base: &Row,
2565 out: &mut Vec<Row>,
2566) -> ExecResult<()> {
2567 let mut row = base.clone();
2568 if bind_node_value(&mut row, src_var, src_id)?
2569 && bind_relationship_value(&mut row, rel_var, rel_id)?
2570 && bind_node_value(&mut row, dst_var, dst_id)?
2571 {
2572 out.push(row);
2573 }
2574 Ok(())
2575}
2576
2577fn bind_node_value(row: &mut Row, var: VarId, id: NodeId) -> ExecResult<bool> {
2578 match row.get(var) {
2579 Some(LoraValue::Node(existing)) => Ok(*existing == id),
2580 Some(other) => Err(ExecutorError::ExpectedNodeForExpand {
2581 var: format!("{var:?}"),
2582 found: value_kind(other),
2583 }),
2584 None => {
2585 row.insert(var, LoraValue::Node(id));
2586 Ok(true)
2587 }
2588 }
2589}
2590
2591fn bind_relationship_value(row: &mut Row, var: VarId, id: RelationshipId) -> ExecResult<bool> {
2592 match row.get(var) {
2593 Some(LoraValue::Relationship(existing)) => Ok(*existing == id),
2594 Some(other) => Err(ExecutorError::ExpectedRelationshipForExpand {
2595 var: format!("{var:?}"),
2596 found: value_kind(other),
2597 }),
2598 None => {
2599 row.insert(var, LoraValue::Relationship(id));
2600 Ok(true)
2601 }
2602 }
2603}
2604
2605fn rel_candidate_ids<S, F>(storage: &S, types: &[String], indexed: F) -> Vec<RelationshipId>
2606where
2607 S: GraphStorage,
2608 F: Fn(&str) -> Option<Vec<RelationshipId>>,
2609{
2610 if types.is_empty() {
2611 return storage.all_rel_ids();
2614 }
2615 let mut all = Vec::new();
2616 let mut seen = BTreeSet::new();
2617 for ty in types {
2618 let ids = indexed(ty).unwrap_or_else(|| storage.rel_ids_by_type(ty));
2619 for id in ids {
2620 if seen.insert(id) {
2621 all.push(id);
2622 }
2623 }
2624 }
2625 all
2626}
2627
2628pub(crate) fn rel_by_property_range_scan_rows<S: GraphStorage>(
2629 storage: &S,
2630 params: &BTreeMap<String, LoraValue>,
2631 base_rows: Vec<Row>,
2632 op: &lora_compiler::RelByPropertyRangeScanExec,
2633 deadline: Option<Instant>,
2634) -> ExecResult<Vec<Row>> {
2635 let eval_ctx = EvalContext { storage, params };
2636 let mut out = Vec::new();
2637 let mut other_kinds = OtherKindScan::default();
2638
2639 for row in base_rows {
2640 check_optional_deadline(deadline)?;
2641 let lo_value = op.lo.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
2642 let hi_value = op
2646 .hi
2647 .as_ref()
2648 .map(|expr| eval_expr(expr, &row, &eval_ctx))
2649 .filter(|v| !matches!(v, LoraValue::Null));
2650 let lo_prop = lo_value
2651 .clone()
2652 .and_then(|v| lora_value_to_property(v).ok());
2653 let hi_prop = hi_value
2654 .clone()
2655 .and_then(|v| lora_value_to_property(v).ok());
2656
2657 let bounds = [lo_value.as_ref(), hi_value.as_ref()];
2660 for rel_id in other_kind_rel_ids(storage, op, bounds, &mut other_kinds) {
2661 if let Some(result) = storage.with_relationship(rel_id, |rel| {
2662 if !op.types.is_empty() && !op.types.iter().any(|t| t == rel.rel_type()) {
2663 return Ok(());
2664 }
2665 emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2666 }) {
2667 result?;
2668 }
2669 }
2670 let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| {
2671 storage.relationship_range_candidates(ty, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
2672 });
2673
2674 for rel_id in candidate_ids {
2675 check_optional_deadline(deadline)?;
2676 if let Some(result) = storage.with_relationship(rel_id, |rel| {
2677 if !op.types.is_empty() && !op.types.iter().any(|t| t == rel.rel_type()) {
2678 return Ok(());
2679 }
2680 let Some(actual) = rel.properties().get(op.key.as_str()) else {
2681 return Ok(());
2682 };
2683 let actual_lv = LoraValue::from(actual);
2684 if !range_predicate_holds(
2685 &actual_lv,
2686 lo_value.as_ref(),
2687 op.lo_inclusive,
2688 hi_value.as_ref(),
2689 op.hi_inclusive,
2690 ) {
2691 return Ok(());
2692 }
2693 emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2694 }) {
2695 result?;
2696 }
2697 }
2698 }
2699
2700 Ok(out)
2701}
2702
2703pub(crate) fn rel_by_text_scan_rows<S: GraphStorage>(
2704 storage: &S,
2705 params: &BTreeMap<String, LoraValue>,
2706 base_rows: Vec<Row>,
2707 op: &lora_compiler::RelByTextScanExec,
2708 deadline: Option<Instant>,
2709) -> ExecResult<Vec<Row>> {
2710 let eval_ctx = EvalContext { storage, params };
2711 let mut out = Vec::new();
2712
2713 for row in base_rows {
2714 check_optional_deadline(deadline)?;
2715 let query = eval_expr(&op.query, &row, &eval_ctx);
2716 let LoraValue::String(query_str) = &query else {
2717 continue;
2718 };
2719
2720 let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| {
2721 storage.relationship_text_candidates(ty, &op.key, query_str)
2722 });
2723
2724 for rel_id in candidate_ids {
2725 check_optional_deadline(deadline)?;
2726 if let Some(result) = storage.with_relationship(rel_id, |rel| {
2727 if !op.types.is_empty() && !op.types.iter().any(|t| t == rel.rel_type()) {
2728 return Ok(());
2729 }
2730 let Some(lora_store::ValueRef::String(actual)) =
2731 rel.properties().get(op.key.as_str())
2732 else {
2733 return Ok(());
2734 };
2735 if !text_predicate_holds(actual, op.predicate, query_str) {
2736 return Ok(());
2737 }
2738 emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2739 }) {
2740 result?;
2741 }
2742 }
2743 }
2744
2745 Ok(out)
2746}
2747
2748pub(crate) fn rel_by_point_scan_rows<S: GraphStorage>(
2749 storage: &S,
2750 params: &BTreeMap<String, LoraValue>,
2751 base_rows: Vec<Row>,
2752 op: &lora_compiler::RelByPointScanExec,
2753 deadline: Option<Instant>,
2754) -> ExecResult<Vec<Row>> {
2755 let eval_ctx = EvalContext { storage, params };
2756 let mut out = Vec::new();
2757
2758 for row in base_rows {
2759 check_optional_deadline(deadline)?;
2760
2761 let probe = match &op.predicate {
2762 lora_compiler::PointPredicate::WithinBBox {
2763 lower_left,
2764 upper_right,
2765 } => {
2766 let ll = eval_expr(lower_left, &row, &eval_ctx);
2767 let ur = eval_expr(upper_right, &row, &eval_ctx);
2768 match (ll, ur) {
2769 (LoraValue::Point(a), LoraValue::Point(b)) => {
2770 Probe::WithinBBox { ll: a, ur: b }
2771 }
2772 _ => continue,
2773 }
2774 }
2775 lora_compiler::PointPredicate::WithinDistance {
2776 center,
2777 max_distance,
2778 inclusive,
2779 } => {
2780 let c = eval_expr(center, &row, &eval_ctx);
2781 let d = eval_expr(max_distance, &row, &eval_ctx);
2782 match (c, d) {
2783 (LoraValue::Point(c), LoraValue::Float(d)) => Probe::WithinDistance {
2784 center: c,
2785 max: d,
2786 inclusive: *inclusive,
2787 },
2788 (LoraValue::Point(c), LoraValue::Int(d)) => Probe::WithinDistance {
2789 center: c,
2790 max: d as f64,
2791 inclusive: *inclusive,
2792 },
2793 _ => continue,
2794 }
2795 }
2796 };
2797
2798 let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| match &probe {
2799 Probe::WithinBBox { ll, ur } => {
2800 let (ranges, n) = lora_store::bbox_x_ranges(ll, ur);
2802 let (lo_y, hi_y) = (ll.y.min(ur.y), ll.y.max(ur.y));
2803 let mut ids = Vec::new();
2804 for (lo_x, hi_x) in &ranges[..n] {
2805 ids.extend(storage.relationship_point_within_bbox(
2806 ty,
2807 &op.key,
2808 (*lo_x, lo_y),
2809 (*hi_x, hi_y),
2810 )?);
2811 }
2812 if n > 1 {
2813 ids.sort_unstable();
2814 ids.dedup();
2815 }
2816 Some(ids)
2817 }
2818 Probe::WithinDistance { center, max, .. } => {
2819 storage.relationship_point_within_distance(ty, &op.key, (center.x, center.y), *max)
2820 }
2821 });
2822
2823 for rel_id in candidate_ids {
2824 check_optional_deadline(deadline)?;
2825 if let Some(result) = storage.with_relationship(rel_id, |rel| {
2826 if !op.types.is_empty() && !op.types.iter().any(|t| t == rel.rel_type()) {
2827 return Ok(());
2828 }
2829 let Some(lora_store::ValueRef::Other(value)) =
2830 rel.properties().get(op.key.as_str())
2831 else {
2832 return Ok(());
2833 };
2834 let holds = match &*value.get() {
2835 PropertyValue::Point(actual) => point_predicate_holds(actual, &probe),
2836 _ => false,
2837 };
2838 if !holds {
2839 return Ok(());
2840 }
2841 emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
2842 }) {
2843 result?;
2844 }
2845 }
2846 }
2847
2848 Ok(out)
2849}
2850
2851pub(crate) struct OrderedRangeCursor {
2862 lo: Option<LoraValue>,
2863 hi: Option<LoraValue>,
2864 mode: OrderedMode,
2865 pending: std::vec::IntoIter<NodeId>,
2866 labels: Vec<Vec<String>>,
2867 key: String,
2868 lo_inclusive: bool,
2869 hi_inclusive: bool,
2870}
2871
2872enum OrderedMode {
2873 Index {
2874 label: String,
2875 descending: bool,
2876 after: Option<(PropertyValue, NodeId)>,
2877 chunk: usize,
2878 exhausted: bool,
2879 },
2880 Buffered,
2881}
2882
2883impl OrderedRangeCursor {
2884 pub(crate) fn new<S: GraphStorage>(
2885 storage: &S,
2886 op: &lora_compiler::NodeByPropertyRangeScanExec,
2887 lo: Option<LoraValue>,
2888 hi: Option<LoraValue>,
2889 ) -> Self {
2890 let descending = matches!(op.order, Some(SortDirection::Desc));
2891 let mut bounds = lo.iter().chain(hi.iter());
2895 let string_bounds = match bounds.next() {
2896 Some(LoraValue::String(_)) => bounds.all(|v| matches!(v, LoraValue::String(_))),
2897 Some(
2898 first @ (LoraValue::Date(_)
2899 | LoraValue::DateTime(_)
2900 | LoraValue::LocalDateTime(_)
2901 | LoraValue::Time(_)
2902 | LoraValue::LocalTime(_)),
2903 ) => bounds.all(|v| std::mem::discriminant(v) == std::mem::discriminant(first)),
2904 _ => false,
2905 };
2906 let index_label = single_label_hint(&op.labels).filter(|label| {
2907 string_bounds
2908 && storage
2909 .node_range_ordered_chunk(label, &op.key, None, None, descending, None, 0)
2910 .is_some()
2911 });
2912
2913 let mut cursor = Self {
2914 lo,
2915 hi,
2916 mode: OrderedMode::Buffered,
2917 pending: Vec::new().into_iter(),
2918 labels: op.labels.clone(),
2919 key: op.key.clone(),
2920 lo_inclusive: op.lo_inclusive,
2921 hi_inclusive: op.hi_inclusive,
2922 };
2923 match index_label {
2924 Some(label) => {
2925 cursor.mode = OrderedMode::Index {
2926 label: label.to_string(),
2927 descending,
2928 after: None,
2929 chunk: 64,
2930 exhausted: false,
2931 }
2932 }
2933 None => {
2934 cursor.pending = cursor
2935 .sorted_candidates(storage, op, descending)
2936 .into_iter()
2937 }
2938 }
2939 cursor
2940 }
2941
2942 pub(crate) fn filter(&self) -> NodeRangeFilter<'_> {
2943 NodeRangeFilter {
2944 labels: &self.labels,
2945 key: &self.key,
2946 lo: self.lo.as_ref(),
2947 lo_inclusive: self.lo_inclusive,
2948 hi: self.hi.as_ref(),
2949 hi_inclusive: self.hi_inclusive,
2950 }
2951 }
2952
2953 fn sorted_candidates<S: GraphStorage>(
2956 &self,
2957 storage: &S,
2958 op: &lora_compiler::NodeByPropertyRangeScanExec,
2959 descending: bool,
2960 ) -> Vec<NodeId> {
2961 let lo_prop = self.lo.clone().and_then(|v| lora_value_to_property(v).ok());
2962 let hi_prop = self.hi.clone().and_then(|v| lora_value_to_property(v).ok());
2963 let candidates = match single_label_hint(&op.labels) {
2964 Some(label) => storage
2965 .node_range_candidates(label, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
2966 .unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
2967 None => scan_node_ids_for_label_groups(storage, &op.labels),
2968 };
2969 let filter = self.filter();
2970 let mut keyed: Vec<(LoraValue, NodeId)> = candidates
2971 .into_iter()
2972 .filter(|&id| node_matches_range_filter(storage, id, &filter))
2973 .map(|id| {
2974 let value = storage
2975 .with_node(id, |n| {
2976 n.properties().get(op.key.as_str()).map(|v| v.to_owned())
2977 })
2978 .flatten()
2979 .map(LoraValue::from)
2980 .unwrap_or(LoraValue::Null);
2981 (value, id)
2982 })
2983 .collect();
2984 keyed.sort_by(|(a, ai), (b, bi)| {
2985 let ord = compare_values_total(a, b).then(ai.cmp(bi));
2986 if descending {
2987 ord.reverse()
2988 } else {
2989 ord
2990 }
2991 });
2992 keyed.into_iter().map(|(_, id)| id).collect()
2993 }
2994
2995 pub(crate) fn next_id<S: GraphStorage>(&mut self, storage: &S) -> Option<NodeId> {
2996 loop {
2997 if let Some(id) = self.pending.next() {
2998 if matches!(self.mode, OrderedMode::Buffered)
2999 || node_matches_range_filter(storage, id, &self.filter())
3000 {
3001 return Some(id);
3002 }
3003 continue;
3004 }
3005 let lo_prop = self.lo.clone().and_then(|v| lora_value_to_property(v).ok());
3006 let hi_prop = self.hi.clone().and_then(|v| lora_value_to_property(v).ok());
3007 let OrderedMode::Index {
3008 label,
3009 descending,
3010 after,
3011 chunk,
3012 exhausted,
3013 } = &mut self.mode
3014 else {
3015 return None;
3016 };
3017 if *exhausted {
3018 return None;
3019 }
3020 let ids = storage.node_range_ordered_chunk(
3021 label,
3022 &self.key,
3023 lo_prop.as_ref(),
3024 hi_prop.as_ref(),
3025 *descending,
3026 after.as_ref().map(|(v, id)| (v, *id)),
3027 *chunk,
3028 )?;
3029 if ids.len() < *chunk {
3030 *exhausted = true;
3031 }
3032 let last = *ids.last()?;
3033 let last_value = storage
3034 .with_node(last, |n| {
3035 n.properties().get(self.key.as_str()).map(|v| v.to_owned())
3036 })
3037 .flatten()?;
3038 *after = Some((last_value, last));
3039 *chunk = (*chunk * 2).min(4096);
3040 self.pending = ids.into_iter();
3041 }
3042 }
3043}
3044
3045pub(crate) fn id_seek_candidates(
3059 value: &LoraValue,
3060 in_list: bool,
3061 all_ids: impl Fn() -> Vec<u64>,
3062) -> Vec<u64> {
3063 fn push_candidates(value: &LoraValue, all_ids: &dyn Fn() -> Vec<u64>, out: &mut Vec<u64>) {
3064 const EXACT_F64_LIMIT: f64 = 9_007_199_254_740_992.0;
3066 match value {
3067 LoraValue::Int(i) if *i >= 0 => out.push(*i as u64),
3068 LoraValue::Float(f) if f.is_finite() && *f >= 0.0 && f.fract() == 0.0 => {
3069 if *f < EXACT_F64_LIMIT {
3070 out.push(*f as u64);
3071 } else {
3072 out.extend(all_ids().into_iter().filter(|id| (*id as i64) as f64 == *f));
3073 }
3074 }
3075 _ => {}
3076 }
3077 }
3078
3079 let mut out = Vec::new();
3080 if in_list {
3081 if let LoraValue::List(items) = value {
3082 for item in items {
3083 push_candidates(item, &all_ids, &mut out);
3084 }
3085 }
3086 } else {
3087 push_candidates(value, &all_ids, &mut out);
3088 }
3089 out.sort_unstable();
3090 out.dedup();
3091 out
3092}
3093
3094pub(crate) fn node_by_id_seek_rows<S: GraphStorage>(
3098 storage: &S,
3099 params: &BTreeMap<String, LoraValue>,
3100 base_rows: Vec<Row>,
3101 op: &lora_compiler::NodeByIdSeekExec,
3102 deadline: Option<Instant>,
3103) -> ExecResult<Vec<Row>> {
3104 let eval_ctx = EvalContext { storage, params };
3105 let node_ok = |id: NodeId| {
3106 if op.labels.is_empty() {
3107 storage.contains_node(id)
3108 } else {
3109 storage
3110 .with_node(id, |n| node_matches_label_groups(n.labels(), &op.labels))
3111 .unwrap_or(false)
3112 }
3113 };
3114 let mut out = Vec::new();
3115 for row in base_rows {
3116 check_optional_deadline(deadline)?;
3117 let value = eval_expr(&op.ids, &row, &eval_ctx);
3118 let ids = id_seek_candidates(&value, op.in_list, || storage.all_node_ids());
3119 if let Some(bound) = bound_node_id_for_expand(&row, op.var)? {
3120 if ids.binary_search(&bound).is_ok() && node_ok(bound) {
3121 out.push(row);
3122 }
3123 continue;
3124 }
3125 for id in ids {
3126 if node_ok(id) {
3127 let mut new_row = row.clone();
3128 new_row.insert(op.var, LoraValue::Node(id));
3129 out.push(new_row);
3130 }
3131 }
3132 }
3133 Ok(out)
3134}
3135
3136pub(crate) fn rel_by_id_seek_rows<S: GraphStorage>(
3143 storage: &S,
3144 params: &BTreeMap<String, LoraValue>,
3145 base_rows: Vec<Row>,
3146 op: &lora_compiler::RelByIdSeekExec,
3147 deadline: Option<Instant>,
3148) -> ExecResult<Vec<Row>> {
3149 let eval_ctx = EvalContext { storage, params };
3150 let mut out = Vec::new();
3151 let mut emitted = Vec::new();
3152 for row in base_rows {
3153 check_optional_deadline(deadline)?;
3154 let value = eval_expr(&op.ids, &row, &eval_ctx);
3155 let ids = id_seek_candidates(&value, op.in_list, || storage.all_rel_ids());
3156 for rel_id in ids {
3157 emitted.clear();
3158 if let Some(result) = storage.with_relationship(rel_id, |rel| {
3159 if !op.types.is_empty() && !op.types.iter().any(|t| t == rel.rel_type()) {
3160 return Ok(());
3161 }
3162 emit_rel_rows(
3163 op.direction,
3164 op.src,
3165 op.rel,
3166 op.dst,
3167 rel,
3168 &row,
3169 &mut emitted,
3170 )
3171 }) {
3172 result?;
3173 }
3174 for new_row in emitted.drain(..) {
3175 if !op.src_labels.is_empty() {
3176 let Some(src) = bound_node_id_for_expand(&new_row, op.src)? else {
3177 continue;
3178 };
3179 let labels_ok = storage
3180 .with_node(src, |n| {
3181 node_matches_label_groups(n.labels(), &op.src_labels)
3182 })
3183 .unwrap_or(false);
3184 if !labels_ok {
3185 continue;
3186 }
3187 }
3188 out.push(new_row);
3189 }
3190 }
3191 }
3192 Ok(out)
3193}