use std::cmp::Ordering;
use std::collections::{BTreeMap, BTreeSet};
use std::time::Instant;
use lora_analyzer::symbols::VarId;
use lora_analyzer::{AggregateFunction, FunctionId, ResolvedExpr, ResolvedMapSelector};
use lora_ast::{Direction, RangeLiteral};
use lora_compiler::physical::{
ExpandExec, HashAggregationExec, LimitExec, NodeByLabelScanExec, NodeByPropertyScanExec,
NodeScanExec, PhysicalNodeId, PhysicalOp, PhysicalPlan, ProjectionExec, UnwindExec,
};
use lora_store::{GraphStorage, NodeId, Properties, PropertyValue, RelationshipId};
use crate::errors::{value_kind, ExecResult, ExecutorError};
use crate::eval::{eval_expr, eval_expr_result, eval_truthy_result, EvalContext};
use crate::value::{lora_value_to_property, LoraPath, LoraValue, Row};
#[inline]
pub(super) fn check_deadline_at(deadline: Instant) -> ExecResult<()> {
if Instant::now() >= deadline {
Err(ExecutorError::QueryTimeout)
} else {
Ok(())
}
}
pub(super) fn filter_rows_checked<S: GraphStorage>(
input_rows: Vec<Row>,
predicate: &ResolvedExpr,
eval_ctx: &EvalContext<'_, S>,
) -> ExecResult<Vec<Row>> {
let mut out = Vec::with_capacity(input_rows.len());
for row in input_rows {
if eval_truthy_result(predicate, &row, eval_ctx).map_err(ExecutorError::RuntimeError)? {
out.push(row);
}
}
Ok(out)
}
pub(super) fn project_rows_checked<S: GraphStorage>(
input_rows: Vec<Row>,
op: &ProjectionExec,
eval_ctx: &EvalContext<'_, S>,
) -> ExecResult<Vec<Row>> {
let mut out = Vec::with_capacity(input_rows.len());
for row in input_rows {
if op.include_existing {
let mut projected = row;
for item in &op.items {
let value = eval_expr_result(&item.expr, &projected, eval_ctx)
.map_err(ExecutorError::RuntimeError)?;
projected.insert_named(item.output, item.name.clone(), value);
}
out.push(projected);
} else {
let mut projected = Row::new();
for item in &op.items {
let value = eval_expr_result(&item.expr, &row, eval_ctx)
.map_err(ExecutorError::RuntimeError)?;
projected.insert_named(item.output, item.name.clone(), value);
}
out.push(projected);
}
}
Ok(if op.distinct {
dedup_rows_by_vars(out)
} else {
out
})
}
pub(super) fn unwind_rows<S: GraphStorage>(
input_rows: Vec<Row>,
op: &UnwindExec,
eval_ctx: &EvalContext<'_, S>,
) -> Vec<Row> {
let mut out = Vec::new();
for row in input_rows {
match eval_expr(&op.expr, &row, eval_ctx) {
LoraValue::List(values) => {
for value in values {
let mut new_row = row.clone();
new_row.insert(op.alias, value);
out.push(new_row);
}
}
LoraValue::Null => {}
other => {
let mut new_row = row;
new_row.insert(op.alias, other);
out.push(new_row);
}
}
}
out
}
pub(super) fn limit_rows<S: GraphStorage>(
mut rows: Vec<Row>,
op: &LimitExec,
eval_ctx: &EvalContext<'_, S>,
) -> Vec<Row> {
let limit = op
.limit
.as_ref()
.and_then(|e| eval_expr(e, &Row::new(), eval_ctx).as_i64())
.unwrap_or(rows.len() as i64)
.max(0) as usize;
let skip = op
.skip
.as_ref()
.and_then(|e| eval_expr(e, &Row::new(), eval_ctx).as_i64())
.unwrap_or(0)
.max(0) as usize;
if skip >= rows.len() {
return Vec::new();
}
rows.drain(0..skip);
rows.truncate(limit);
rows
}
#[inline]
pub(crate) fn bound_node_id_for_expand(row: &Row, var: VarId) -> ExecResult<Option<NodeId>> {
match row.get(var) {
Some(LoraValue::Node(id)) => Ok(Some(*id)),
Some(other) => Err(ExecutorError::ExpectedNodeForExpand {
var: format!("{var:?}"),
found: value_kind(other),
}),
None => Ok(None),
}
}
#[inline]
pub(crate) fn bound_relationship_id_for_expand(
row: &Row,
var: VarId,
) -> ExecResult<Option<RelationshipId>> {
match row.get(var) {
Some(LoraValue::Relationship(id)) => Ok(Some(*id)),
Some(other) => Err(ExecutorError::ExpectedRelationshipForExpand {
var: format!("{var:?}"),
found: value_kind(other),
}),
None => Ok(None),
}
}
pub(super) fn node_scan_rows<S: GraphStorage>(
storage: &S,
base_rows: Vec<Row>,
op: &NodeScanExec,
deadline: Option<Instant>,
) -> ExecResult<Vec<Row>> {
let node_ids = storage.all_node_ids();
let mut out = Vec::with_capacity(base_rows.len().saturating_mul(node_ids.len()));
if deadline.is_none() {
for row in base_rows {
if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
if storage.has_node(existing_id) {
out.push(row);
}
continue;
}
for &id in &node_ids {
let mut new_row = row.clone();
new_row.insert(op.var, LoraValue::Node(id));
out.push(new_row);
}
}
return Ok(out);
}
for row in base_rows {
check_optional_deadline(deadline)?;
if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
if storage.has_node(existing_id) {
out.push(row);
}
continue;
}
for &id in &node_ids {
check_optional_deadline(deadline)?;
let mut new_row = row.clone();
new_row.insert(op.var, LoraValue::Node(id));
out.push(new_row);
}
}
Ok(out)
}
pub(super) fn node_by_label_scan_rows<S: GraphStorage>(
storage: &S,
base_rows: Vec<Row>,
op: &NodeByLabelScanExec,
deadline: Option<Instant>,
) -> ExecResult<Vec<Row>> {
let candidate_ids = scan_node_ids_for_label_groups(storage, &op.labels);
let candidates_prefiltered = label_group_candidates_prefiltered(&op.labels);
let mut out = Vec::with_capacity(base_rows.len().saturating_mul(candidate_ids.len()));
if deadline.is_none() {
for row in base_rows {
if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
let labels_ok = storage
.with_node(existing_id, |n| {
node_matches_label_groups(&n.labels, &op.labels)
})
.unwrap_or(false);
if labels_ok {
out.push(row);
}
continue;
}
for &id in &candidate_ids {
if !candidates_prefiltered {
let labels_ok = storage
.with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
.unwrap_or(false);
if !labels_ok {
continue;
}
}
let mut new_row = row.clone();
new_row.insert(op.var, LoraValue::Node(id));
out.push(new_row);
}
}
return Ok(out);
}
for row in base_rows {
check_optional_deadline(deadline)?;
if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
let labels_ok = storage
.with_node(existing_id, |n| {
node_matches_label_groups(&n.labels, &op.labels)
})
.unwrap_or(false);
if labels_ok {
out.push(row);
}
continue;
}
for &id in &candidate_ids {
check_optional_deadline(deadline)?;
if !candidates_prefiltered {
let labels_ok = storage
.with_node(id, |n| node_matches_label_groups(&n.labels, &op.labels))
.unwrap_or(false);
if !labels_ok {
continue;
}
}
let mut new_row = row.clone();
new_row.insert(op.var, LoraValue::Node(id));
out.push(new_row);
}
}
Ok(out)
}
pub(super) fn node_by_property_scan_rows<S: GraphStorage>(
storage: &S,
params: &BTreeMap<String, LoraValue>,
base_rows: Vec<Row>,
op: &NodeByPropertyScanExec,
deadline: Option<Instant>,
) -> ExecResult<Vec<Row>> {
let eval_ctx = EvalContext { storage, params };
let mut out = Vec::new();
if deadline.is_none() {
for row in base_rows {
let expected = eval_expr(&op.value, &row, &eval_ctx);
if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
if node_matches_property_filter(
storage,
existing_id,
&op.labels,
&op.key,
&expected,
) {
out.push(row);
}
continue;
}
let candidates =
indexed_node_property_candidates(storage, &op.labels, &op.key, &expected);
for id in candidates.ids {
if !candidates.prefiltered
&& !node_matches_property_filter(storage, id, &op.labels, &op.key, &expected)
{
continue;
}
let mut new_row = row.clone();
new_row.insert(op.var, LoraValue::Node(id));
out.push(new_row);
}
}
return Ok(out);
}
for row in base_rows {
check_optional_deadline(deadline)?;
let expected = eval_expr(&op.value, &row, &eval_ctx);
if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
if node_matches_property_filter(storage, existing_id, &op.labels, &op.key, &expected) {
out.push(row);
}
continue;
}
let candidates = indexed_node_property_candidates(storage, &op.labels, &op.key, &expected);
for id in candidates.ids {
check_optional_deadline(deadline)?;
if !candidates.prefiltered
&& !node_matches_property_filter(storage, id, &op.labels, &op.key, &expected)
{
continue;
}
let mut new_row = row.clone();
new_row.insert(op.var, LoraValue::Node(id));
out.push(new_row);
}
}
Ok(out)
}
#[inline]
fn check_optional_deadline(deadline: Option<Instant>) -> ExecResult<()> {
match deadline {
Some(deadline) => check_deadline_at(deadline),
None => Ok(()),
}
}
pub(crate) fn plan_may_need_hydration(plan: &PhysicalPlan) -> bool {
op_may_need_hydration(plan, plan.root)
}
pub(crate) fn count_all_scan_aggregation_rows<S: GraphStorage>(
storage: &S,
plan: &PhysicalPlan,
op: &HashAggregationExec,
) -> Option<Vec<Row>> {
if !op.group_by.is_empty() {
return None;
}
let specs = crate::pull::classify_streamable_aggregates(&op.aggregates)?;
if !specs
.iter()
.all(|spec| matches!(spec.kind, crate::pull::StreamableAggKind::CountAll))
{
return None;
}
let count = count_rows_for_scan_subtree(storage, plan, op.input)? as i64;
let value = LoraValue::Int(count);
let mut row = Row::new();
for proj in &op.aggregates {
row.insert_named(proj.output, proj.name.clone(), value.clone());
}
Some(vec![row])
}
fn count_rows_for_scan_subtree<S: GraphStorage>(
storage: &S,
plan: &PhysicalPlan,
node_id: PhysicalNodeId,
) -> Option<usize> {
match &plan.nodes[node_id] {
PhysicalOp::NodeScan(op) if scan_input_is_argument(plan, op.input) => {
Some(storage.node_count())
}
PhysicalOp::NodeByLabelScan(op) if scan_input_is_argument(plan, op.input) => {
let ids = scan_node_ids_for_label_groups(storage, &op.labels);
if label_group_candidates_prefiltered(&op.labels) {
return Some(ids.len());
}
Some(
ids.into_iter()
.filter(|&id| {
storage
.with_node(id, |node| {
node_matches_label_groups(&node.labels, &op.labels)
})
.unwrap_or(false)
})
.count(),
)
}
_ => None,
}
}
fn scan_input_is_argument(plan: &PhysicalPlan, input: Option<PhysicalNodeId>) -> bool {
match input {
None => true,
Some(id) => matches!(plan.nodes.get(id), Some(PhysicalOp::Argument(_))),
}
}
fn op_may_need_hydration(plan: &PhysicalPlan, node_id: PhysicalNodeId) -> bool {
match &plan.nodes[node_id] {
PhysicalOp::Projection(op) if !op.include_existing => op
.items
.iter()
.any(|item| expr_may_produce_hydratable_value(&item.expr)),
PhysicalOp::HashAggregation(op) => {
op.group_by
.iter()
.any(|item| expr_may_produce_hydratable_value(&item.expr))
|| op
.aggregates
.iter()
.any(|item| expr_may_produce_hydratable_value(&item.expr))
}
PhysicalOp::Sort(op) => op_may_need_hydration(plan, op.input),
PhysicalOp::Limit(op) => op_may_need_hydration(plan, op.input),
PhysicalOp::Filter(op) => op_may_need_hydration(plan, op.input),
PhysicalOp::Unwind(op) => op_may_need_hydration(plan, op.input),
PhysicalOp::PathBuild(op) => op_may_need_hydration(plan, op.input),
_ => true,
}
}
fn expr_may_produce_hydratable_value(expr: &ResolvedExpr) -> bool {
match expr {
ResolvedExpr::Variable(_) | ResolvedExpr::Parameter(_) => true,
ResolvedExpr::Literal(_)
| ResolvedExpr::Property { .. }
| ResolvedExpr::ExistsSubquery { .. }
| ResolvedExpr::Binary { .. }
| ResolvedExpr::Unary { .. }
| ResolvedExpr::ListPredicate { .. } => false,
ResolvedExpr::Function { function, args, .. } => {
function_may_produce_hydratable_value(*function, args)
}
ResolvedExpr::List(items) => items.iter().any(expr_may_produce_hydratable_value),
ResolvedExpr::Map(items) => items
.iter()
.any(|(_, value)| expr_may_produce_hydratable_value(value)),
ResolvedExpr::Case {
alternatives,
else_expr,
..
} => {
alternatives
.iter()
.any(|(_, value)| expr_may_produce_hydratable_value(value))
|| else_expr
.as_deref()
.is_some_and(expr_may_produce_hydratable_value)
}
ResolvedExpr::ListComprehension { map_expr, .. } => map_expr
.as_deref()
.map(expr_may_produce_hydratable_value)
.unwrap_or(true),
ResolvedExpr::Reduce { expr, .. } => expr_may_produce_hydratable_value(expr),
ResolvedExpr::MapProjection { selectors, .. } => selectors.iter().any(|selector| {
matches!(selector, ResolvedMapSelector::Literal(_, expr) if expr_may_produce_hydratable_value(expr))
}),
ResolvedExpr::Index { expr, .. } | ResolvedExpr::Slice { expr, .. } => {
expr_may_produce_hydratable_value(expr)
}
ResolvedExpr::PatternComprehension { map_expr, .. } => {
expr_may_produce_hydratable_value(map_expr)
}
}
}
fn function_may_produce_hydratable_value(function: FunctionId, args: &[ResolvedExpr]) -> bool {
match function.name() {
"path.nodes" | "path.edges" | "path.first" | "path.last" | "list.first" | "list.last" => {
true
}
"value.coalesce" | "value.first_non_null" | "collect" => {
args.iter().any(expr_may_produce_hydratable_value)
}
"list.rest" | "value.reverse" | "list.reverse" => {
args.first().is_some_and(expr_may_produce_hydratable_value)
}
_ => false,
}
}
pub(super) fn expand_rows<S: GraphStorage>(
storage: &S,
params: &BTreeMap<String, LoraValue>,
input_rows: Vec<Row>,
op: &ExpandExec,
) -> ExecResult<Vec<Row>> {
let eval_ctx = EvalContext { storage, params };
let mut out = Vec::new();
for row in input_rows {
let Some(src_node_id) = bound_node_id_for_expand(&row, op.src)? else {
continue;
};
let mut rel_property_filter = None;
storage.try_for_each_expand_id(
src_node_id,
op.direction,
&op.types,
|rel_id, dst_id| {
if let Some(expr) = op.rel_properties.as_ref() {
if rel_property_filter.is_none() {
let expected = eval_expr(expr, &row, &eval_ctx);
let LoraValue::Map(map) = expected else {
return Err(ExecutorError::ExpectedPropertyMap {
found: value_kind(&expected),
});
};
rel_property_filter = Some(map);
}
let Some(map) = rel_property_filter.as_ref() else {
return Ok(());
};
let matches = storage
.with_relationship(rel_id, |rel| {
map.iter().all(|(key, expected)| {
rel.properties
.get(key)
.map(|actual| value_matches_property_value(expected, actual))
.unwrap_or(false)
})
})
.unwrap_or(false);
if !matches {
return Ok(());
}
}
if let Some(existing_id) = bound_node_id_for_expand(&row, op.dst)? {
if existing_id != dst_id {
return Ok(());
}
}
if let Some(rel_var) = op.rel {
if let Some(existing_id) = bound_relationship_id_for_expand(&row, rel_var)? {
if existing_id != rel_id {
return Ok(());
}
}
}
let mut new_row = row.clone();
if !new_row.contains_key(op.dst) {
new_row.insert(op.dst, LoraValue::Node(dst_id));
}
if let Some(rel_var) = op.rel {
if !new_row.contains_key(rel_var) {
new_row.insert(rel_var, LoraValue::Relationship(rel_id));
}
}
out.push(new_row);
Ok(())
},
)?;
}
Ok(out)
}
pub(super) fn expand_var_len_rows<S: GraphStorage>(
storage: &S,
input_rows: Vec<Row>,
op: &ExpandExec,
range: &RangeLiteral,
) -> ExecResult<Vec<Row>> {
let (min_hops, max_hops) = resolve_range(range);
let bind_relationships = op.rel.is_some();
let mut out = Vec::new();
for row in input_rows {
let Some(src_node_id) = bound_node_id_for_expand(&row, op.src)? else {
continue;
};
let expansions = variable_length_expand(
storage,
src_node_id,
op.direction,
&op.types,
min_hops,
max_hops,
bind_relationships,
);
for result in expansions {
let mut new_row = row.clone();
new_row.insert(op.dst, LoraValue::Node(result.dst_node_id));
if let Some(rel_var) = op.rel {
let rel_list = LoraValue::List(
result
.rel_ids
.into_iter()
.map(LoraValue::Relationship)
.collect(),
);
new_row.insert(rel_var, rel_list);
}
out.push(new_row);
}
}
Ok(out)
}
pub(super) fn properties_to_value_map(props: &Properties) -> LoraValue {
let mut map = BTreeMap::new();
for (k, v) in props.iter() {
map.insert(k.clone(), LoraValue::from(v));
}
LoraValue::Map(map)
}
pub(crate) fn dedup_rows_by_vars(rows: Vec<Row>) -> Vec<Row> {
let mut seen: BTreeSet<Vec<GroupValueKey>> = BTreeSet::new();
let mut out = Vec::new();
for row in rows {
let key: Vec<GroupValueKey> = row
.iter()
.map(|(_, val)| GroupValueKey::from_value(val))
.collect();
if seen.insert(key) {
out.push(row);
}
}
out
}
pub(crate) fn dedup_rows(rows: Vec<Row>) -> Vec<Row> {
let mut seen: BTreeSet<Vec<(String, GroupValueKey)>> = BTreeSet::new();
let mut out = Vec::new();
for row in rows {
let key: Vec<(String, GroupValueKey)> = row
.iter_named()
.map(|(_, name, val)| (name.into_owned(), GroupValueKey::from_value(val)))
.collect();
if seen.insert(key) {
out.push(row);
}
}
out
}
pub(super) fn eval_properties_expr<S: GraphStorage>(
expr: &ResolvedExpr,
row: &Row,
storage: &S,
params: &BTreeMap<String, LoraValue>,
) -> ExecResult<Properties> {
let eval_ctx = EvalContext { storage, params };
match eval_expr(expr, row, &eval_ctx) {
LoraValue::Map(map) => {
let mut out = Properties::new();
for (k, v) in map {
let prop = lora_value_to_property(v)
.map_err(|e| ExecutorError::RuntimeError(e.to_string()))?;
out.insert(k, prop);
}
Ok(out)
}
other => Err(ExecutorError::ExpectedPropertyMap {
found: value_kind(&other),
}),
}
}
pub(crate) fn compute_aggregate_expr<S: GraphStorage>(
expr: &ResolvedExpr,
rows: &[Row],
eval_ctx: &EvalContext<'_, S>,
) -> ExecResult<LoraValue> {
match expr {
ResolvedExpr::Function {
function,
distinct,
args,
} => {
let func = function.as_aggregate();
match func {
Some(AggregateFunction::Count) => {
if args.is_empty() {
return Ok(LoraValue::Int(rows.len() as i64));
}
let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
values.retain(|v| !matches!(v, LoraValue::Null));
if *distinct {
values = dedup_values(values);
}
Ok(LoraValue::Int(values.len() as i64))
}
Some(AggregateFunction::Collect) => {
if args.is_empty() {
return Ok(LoraValue::List(Vec::new()));
}
let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
if *distinct {
values = dedup_values(values);
}
Ok(LoraValue::List(values))
}
Some(AggregateFunction::Sum) => {
if args.is_empty() {
return Ok(LoraValue::Null);
}
let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
if *distinct {
values = dedup_values(values);
}
let nums = values
.into_iter()
.filter_map(as_f64_lossy)
.collect::<Vec<_>>();
if nums.is_empty() {
Ok(LoraValue::Null)
} else if nums.iter().all(|n| n.fract() == 0.0) {
Ok(LoraValue::Int(nums.iter().sum::<f64>() as i64))
} else {
Ok(LoraValue::Float(nums.iter().sum::<f64>()))
}
}
Some(AggregateFunction::Avg) => {
if args.is_empty() {
return Ok(LoraValue::Null);
}
let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
if *distinct {
values = dedup_values(values);
}
let nums = values
.into_iter()
.filter_map(as_f64_lossy)
.collect::<Vec<_>>();
if nums.is_empty() {
Ok(LoraValue::Null)
} else {
Ok(LoraValue::Float(
nums.iter().sum::<f64>() / nums.len() as f64,
))
}
}
Some(AggregateFunction::Min) => {
if args.is_empty() {
return Ok(LoraValue::Null);
}
let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
values.retain(|v| !matches!(v, LoraValue::Null));
if *distinct {
values = dedup_values(values);
}
Ok(values
.into_iter()
.min_by(compare_values_total)
.unwrap_or(LoraValue::Null))
}
Some(AggregateFunction::Max) => {
if args.is_empty() {
return Ok(LoraValue::Null);
}
let mut values = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?;
values.retain(|v| !matches!(v, LoraValue::Null));
if *distinct {
values = dedup_values(values);
}
Ok(values
.into_iter()
.max_by(compare_values_total)
.unwrap_or(LoraValue::Null))
}
Some(AggregateFunction::Stdev | AggregateFunction::Stdevp) => {
if args.is_empty() {
return Ok(LoraValue::Null);
}
let nums: Vec<f64> = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?
.into_iter()
.filter_map(as_f64_lossy)
.collect();
let is_population = matches!(func, Some(AggregateFunction::Stdevp));
if nums.is_empty() || (!is_population && nums.len() < 2) {
return Ok(LoraValue::Float(0.0));
}
let mean = nums.iter().sum::<f64>() / nums.len() as f64;
let variance_sum: f64 = nums.iter().map(|x| (x - mean).powi(2)).sum();
let denom = if is_population {
nums.len() as f64
} else {
(nums.len() - 1) as f64
};
Ok(LoraValue::Float((variance_sum / denom).sqrt()))
}
Some(AggregateFunction::PercentileCont) => {
if args.len() < 2 {
return Ok(LoraValue::Null);
}
let Some(first) = rows.first() else {
return Ok(LoraValue::Null);
};
let percentile = eval_expr_result(&args[1], first, eval_ctx)
.map_err(ExecutorError::RuntimeError)?
.as_f64()
.map(normalize_percentile)
.unwrap_or(0.5);
let mut nums: Vec<f64> = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?
.into_iter()
.filter_map(as_f64_lossy)
.collect();
if nums.is_empty() {
return Ok(LoraValue::Null);
}
nums.sort_by(|a, b| a.partial_cmp(b).unwrap_or(Ordering::Equal));
let index = percentile * (nums.len() - 1) as f64;
let lower = index.floor() as usize;
let upper = index.ceil() as usize;
let fraction = index - lower as f64;
if lower == upper || upper >= nums.len() {
Ok(LoraValue::Float(nums[lower]))
} else {
Ok(LoraValue::Float(
nums[lower] * (1.0 - fraction) + nums[upper] * fraction,
))
}
}
Some(AggregateFunction::PercentileDisc) => {
if args.len() < 2 {
return Ok(LoraValue::Null);
}
let Some(first) = rows.first() else {
return Ok(LoraValue::Null);
};
let percentile = eval_expr_result(&args[1], first, eval_ctx)
.map_err(ExecutorError::RuntimeError)?
.as_f64()
.map(normalize_percentile)
.unwrap_or(0.5);
let mut nums: Vec<f64> = eval_aggregate_arg_values(&args[0], rows, eval_ctx)?
.into_iter()
.filter_map(as_f64_lossy)
.collect();
if nums.is_empty() {
return Ok(LoraValue::Null);
}
nums.sort_by(|a, b| a.partial_cmp(b).unwrap_or(Ordering::Equal));
let index = (percentile * (nums.len() - 1) as f64).round() as usize;
let index = index.min(nums.len() - 1);
Ok(LoraValue::Float(nums[index]))
}
_ => eval_first_or_null(expr, rows, eval_ctx),
}
}
_ => eval_first_or_null(expr, rows, eval_ctx),
}
}
fn eval_aggregate_arg_values<S: GraphStorage>(
expr: &ResolvedExpr,
rows: &[Row],
eval_ctx: &EvalContext<'_, S>,
) -> ExecResult<Vec<LoraValue>> {
rows.iter()
.map(|row| eval_expr_result(expr, row, eval_ctx).map_err(ExecutorError::RuntimeError))
.collect()
}
fn normalize_percentile(value: f64) -> f64 {
if value.is_finite() {
value.clamp(0.0, 1.0)
} else {
0.5
}
}
fn eval_first_or_null<S: GraphStorage>(
expr: &ResolvedExpr,
rows: &[Row],
eval_ctx: &EvalContext<'_, S>,
) -> ExecResult<LoraValue> {
match rows.first() {
Some(row) => eval_expr_result(expr, row, eval_ctx).map_err(ExecutorError::RuntimeError),
None => Ok(LoraValue::Null),
}
}
fn dedup_values(values: Vec<LoraValue>) -> Vec<LoraValue> {
let mut seen: BTreeSet<GroupValueKey> = BTreeSet::new();
let mut out = Vec::new();
for value in values {
let key = GroupValueKey::from_value(&value);
if seen.insert(key) {
out.push(value);
}
}
out
}
fn as_f64_lossy(v: LoraValue) -> Option<f64> {
match v {
LoraValue::Int(i) => Some(i as f64),
LoraValue::Float(f) => Some(f),
_ => None,
}
}
pub(super) fn compare_values_total(a: &LoraValue, b: &LoraValue) -> Ordering {
use LoraValue::*;
match (a, b) {
(Bool(x), Bool(y)) => x.cmp(y),
(Int(x), Int(y)) => x.cmp(y),
(Float(x), Float(y)) => x.partial_cmp(y).unwrap_or(Ordering::Equal),
(Int(x), Float(y)) => (*x as f64).partial_cmp(y).unwrap_or(Ordering::Equal),
(Float(x), Int(y)) => x.partial_cmp(&(*y as f64)).unwrap_or(Ordering::Equal),
(String(x), String(y)) => x.cmp(y),
(Binary(x), Binary(y)) => x.segments().cmp(y.segments()),
(Node(x), Node(y)) => x.cmp(y),
(Relationship(x), Relationship(y)) => x.cmp(y),
(Date(x), Date(y)) => x.cmp(y),
(DateTime(x), DateTime(y)) => x.cmp(y),
(Duration(x), Duration(y)) => x.cmp(y),
(Vector(x), Vector(y)) => x.to_key_string().cmp(&y.to_key_string()),
_ => type_rank(a)
.cmp(&type_rank(b))
.then_with(|| format!("{a:?}").cmp(&format!("{b:?}"))),
}
}
pub fn value_matches_property_value(expected: &LoraValue, actual: &PropertyValue) -> bool {
match (expected, actual) {
(LoraValue::Null, PropertyValue::Null) => true,
(LoraValue::Bool(a), PropertyValue::Bool(b)) => a == b,
(LoraValue::Int(a), PropertyValue::Int(b)) => a == b,
(LoraValue::Float(a), PropertyValue::Float(b)) => a == b,
(LoraValue::Int(a), PropertyValue::Float(b)) => (*a as f64) == *b,
(LoraValue::Float(a), PropertyValue::Int(b)) => *a == (*b as f64),
(LoraValue::String(a), PropertyValue::String(b)) => a == b,
(LoraValue::Binary(a), PropertyValue::Binary(b)) => a == b,
(LoraValue::List(xs), PropertyValue::List(ys)) => {
xs.len() == ys.len()
&& xs
.iter()
.zip(ys.iter())
.all(|(x, y)| value_matches_property_value(x, y))
}
(LoraValue::Map(xm), PropertyValue::Map(ym)) => xm.iter().all(|(k, xv)| {
ym.get(k)
.map(|yv| value_matches_property_value(xv, yv))
.unwrap_or(false)
}),
(LoraValue::Date(a), PropertyValue::Date(b)) => a == b,
(LoraValue::DateTime(a), PropertyValue::DateTime(b)) => a == b,
(LoraValue::LocalDateTime(a), PropertyValue::LocalDateTime(b)) => a == b,
(LoraValue::Time(a), PropertyValue::Time(b)) => a == b,
(LoraValue::LocalTime(a), PropertyValue::LocalTime(b)) => a == b,
(LoraValue::Duration(a), PropertyValue::Duration(b)) => a == b,
(LoraValue::Point(a), PropertyValue::Point(b)) => a == b,
(LoraValue::Vector(a), PropertyValue::Vector(b)) => a == b,
_ => false,
}
}
pub(crate) fn node_matches_property_filter<S: GraphStorage>(
storage: &S,
node_id: NodeId,
labels: &[Vec<String>],
key: &str,
expected: &LoraValue,
) -> bool {
storage
.with_node(node_id, |node| {
node_matches_label_groups(&node.labels, labels)
&& node
.properties
.get(key)
.map(|actual| value_matches_property_value(expected, actual))
.unwrap_or(false)
})
.unwrap_or(false)
}
fn single_label_hint(labels: &[Vec<String>]) -> Option<&str> {
if labels.len() == 1 && labels[0].len() == 1 {
Some(labels[0][0].as_str())
} else {
None
}
}
fn property_lookup_values(expected: &LoraValue) -> Option<Vec<PropertyValue>> {
let property = lora_value_to_property(expected.clone()).ok()?;
let mut values = vec![property.clone()];
match property {
PropertyValue::Int(i) => {
values.push(PropertyValue::Float(i as f64));
}
PropertyValue::Float(f)
if f.is_finite()
&& f.fract() == 0.0
&& f >= i64::MIN as f64
&& f <= i64::MAX as f64 =>
{
values.push(PropertyValue::Int(f as i64));
}
_ => {}
}
Some(values)
}
pub(crate) struct NodePropertyCandidates {
pub(crate) ids: Vec<NodeId>,
pub(crate) prefiltered: bool,
}
pub(crate) fn node_by_property_range_scan_rows<S: GraphStorage>(
storage: &S,
params: &BTreeMap<String, LoraValue>,
base_rows: Vec<Row>,
op: &lora_compiler::NodeByPropertyRangeScanExec,
deadline: Option<Instant>,
) -> ExecResult<Vec<Row>> {
let eval_ctx = EvalContext { storage, params };
let mut out = Vec::new();
for row in base_rows {
check_optional_deadline(deadline)?;
let lo_value = op.lo.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
let hi_value = op.hi.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
let lo_prop = lo_value
.clone()
.and_then(|v| lora_value_to_property(v).ok());
let hi_prop = hi_value
.clone()
.and_then(|v| lora_value_to_property(v).ok());
let filter = NodeRangeFilter {
labels: &op.labels,
key: &op.key,
lo: lo_value.as_ref(),
lo_inclusive: op.lo_inclusive,
hi: hi_value.as_ref(),
hi_inclusive: op.hi_inclusive,
};
if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
if node_matches_range_filter(storage, existing_id, &filter) {
out.push(row);
}
continue;
}
let candidate_ids = match single_label_hint(&op.labels) {
Some(label) => storage
.node_range_candidates(label, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
.unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
None => scan_node_ids_for_label_groups(storage, &op.labels),
};
for id in candidate_ids {
check_optional_deadline(deadline)?;
if node_matches_range_filter(storage, id, &filter) {
let mut new_row = row.clone();
new_row.insert(op.var, LoraValue::Node(id));
out.push(new_row);
}
}
}
Ok(out)
}
pub(crate) fn node_by_text_scan_rows<S: GraphStorage>(
storage: &S,
params: &BTreeMap<String, LoraValue>,
base_rows: Vec<Row>,
op: &lora_compiler::NodeByTextScanExec,
deadline: Option<Instant>,
) -> ExecResult<Vec<Row>> {
let eval_ctx = EvalContext { storage, params };
let mut out = Vec::new();
for row in base_rows {
check_optional_deadline(deadline)?;
let query = eval_expr(&op.query, &row, &eval_ctx);
let LoraValue::String(query_str) = &query else {
continue;
};
if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
if node_matches_text_filter(
storage,
existing_id,
&op.labels,
&op.key,
op.predicate,
query_str,
) {
out.push(row);
}
continue;
}
let candidate_ids = match single_label_hint(&op.labels) {
Some(label) => storage
.node_text_candidates(label, &op.key, query_str)
.unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
None => scan_node_ids_for_label_groups(storage, &op.labels),
};
for id in candidate_ids {
check_optional_deadline(deadline)?;
if node_matches_text_filter(storage, id, &op.labels, &op.key, op.predicate, query_str) {
let mut new_row = row.clone();
new_row.insert(op.var, LoraValue::Node(id));
out.push(new_row);
}
}
}
Ok(out)
}
struct NodeRangeFilter<'a> {
labels: &'a [Vec<String>],
key: &'a str,
lo: Option<&'a LoraValue>,
lo_inclusive: bool,
hi: Option<&'a LoraValue>,
hi_inclusive: bool,
}
fn node_matches_range_filter<S: GraphStorage>(
storage: &S,
id: NodeId,
filter: &NodeRangeFilter<'_>,
) -> bool {
storage
.with_node(id, |n| {
if !node_matches_label_groups(&n.labels, filter.labels) {
return false;
}
let Some(actual) = n.properties.get(filter.key) else {
return false;
};
let actual_lv = lora_store_property_to_value(actual);
range_predicate_holds(
&actual_lv,
filter.lo,
filter.lo_inclusive,
filter.hi,
filter.hi_inclusive,
)
})
.unwrap_or(false)
}
fn node_matches_text_filter<S: GraphStorage>(
storage: &S,
id: NodeId,
labels: &[Vec<String>],
key: &str,
predicate: lora_compiler::TextPredicate,
query: &str,
) -> bool {
storage
.with_node(id, |n| {
if !node_matches_label_groups(&n.labels, labels) {
return false;
}
let Some(PropertyValue::String(actual)) = n.properties.get(key) else {
return false;
};
text_predicate_holds(actual, predicate, query)
})
.unwrap_or(false)
}
fn text_predicate_holds(
actual: &str,
predicate: lora_compiler::TextPredicate,
query: &str,
) -> bool {
match predicate {
lora_compiler::TextPredicate::StartsWith => actual.starts_with(query),
lora_compiler::TextPredicate::EndsWith => actual.ends_with(query),
lora_compiler::TextPredicate::Contains => actual.contains(query),
}
}
fn range_predicate_holds(
actual: &LoraValue,
lo: Option<&LoraValue>,
lo_inclusive: bool,
hi: Option<&LoraValue>,
hi_inclusive: bool,
) -> bool {
if let Some(lo) = lo {
match range_comparison(actual, lo) {
None => return false,
Some(Ordering::Less) => return false,
Some(Ordering::Equal) if !lo_inclusive => return false,
_ => {}
}
}
if let Some(hi) = hi {
match range_comparison(actual, hi) {
None => return false,
Some(Ordering::Greater) => return false,
Some(Ordering::Equal) if !hi_inclusive => return false,
_ => {}
}
}
true
}
fn range_comparison(actual: &LoraValue, bound: &LoraValue) -> Option<Ordering> {
match (actual, bound) {
(LoraValue::Null, _) | (_, LoraValue::Null) => None,
(LoraValue::String(a), LoraValue::String(b)) => Some(a.cmp(b)),
(LoraValue::Date(a), LoraValue::Date(b)) => Some(a.to_epoch_days().cmp(&b.to_epoch_days())),
(LoraValue::DateTime(a), LoraValue::DateTime(b)) => {
Some(a.to_epoch_millis().cmp(&b.to_epoch_millis()))
}
(LoraValue::Duration(a), LoraValue::Duration(b)) => a
.total_seconds_approx()
.partial_cmp(&b.total_seconds_approx()),
_ => actual.as_f64()?.partial_cmp(&bound.as_f64()?),
}
}
fn lora_store_property_to_value(value: &PropertyValue) -> LoraValue {
LoraValue::from(value)
}
pub(crate) fn node_by_point_scan_rows<S: GraphStorage>(
storage: &S,
params: &BTreeMap<String, LoraValue>,
base_rows: Vec<Row>,
op: &lora_compiler::NodeByPointScanExec,
deadline: Option<Instant>,
) -> ExecResult<Vec<Row>> {
let eval_ctx = EvalContext { storage, params };
let mut out = Vec::new();
for row in base_rows {
check_optional_deadline(deadline)?;
let probe = match &op.predicate {
lora_compiler::PointPredicate::WithinBBox {
lower_left,
upper_right,
} => {
let ll = eval_expr(lower_left, &row, &eval_ctx);
let ur = eval_expr(upper_right, &row, &eval_ctx);
match (ll, ur) {
(LoraValue::Point(a), LoraValue::Point(b)) => {
Probe::WithinBBox { ll: a, ur: b }
}
_ => continue,
}
}
lora_compiler::PointPredicate::WithinDistance {
center,
max_distance,
inclusive,
} => {
let c = eval_expr(center, &row, &eval_ctx);
let d = eval_expr(max_distance, &row, &eval_ctx);
match (c, d) {
(LoraValue::Point(c), LoraValue::Float(d)) => Probe::WithinDistance {
center: c,
max: d,
inclusive: *inclusive,
},
(LoraValue::Point(c), LoraValue::Int(d)) => Probe::WithinDistance {
center: c,
max: d as f64,
inclusive: *inclusive,
},
_ => continue,
}
}
};
if let Some(existing_id) = bound_node_id_for_expand(&row, op.var)? {
if node_matches_point_filter(storage, existing_id, &op.labels, &op.key, &probe) {
out.push(row);
}
continue;
}
let candidate_ids = match single_label_hint(&op.labels) {
Some(label) => match &probe {
Probe::WithinBBox { ll, ur } => storage
.node_point_within_bbox(label, &op.key, (ll.x, ll.y), (ur.x, ur.y))
.unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
Probe::WithinDistance { center, max, .. } => storage
.node_point_within_distance(label, &op.key, (center.x, center.y), *max)
.unwrap_or_else(|| scan_node_ids_for_label_groups(storage, &op.labels)),
},
None => scan_node_ids_for_label_groups(storage, &op.labels),
};
for id in candidate_ids {
check_optional_deadline(deadline)?;
if node_matches_point_filter(storage, id, &op.labels, &op.key, &probe) {
let mut new_row = row.clone();
new_row.insert(op.var, LoraValue::Node(id));
out.push(new_row);
}
}
}
Ok(out)
}
fn node_matches_point_filter<S: GraphStorage>(
storage: &S,
id: NodeId,
labels: &[Vec<String>],
key: &str,
probe: &Probe,
) -> bool {
storage
.with_node(id, |n| {
if !node_matches_label_groups(&n.labels, labels) {
return false;
}
let Some(PropertyValue::Point(point)) = n.properties.get(key) else {
return false;
};
point_predicate_holds(point, probe)
})
.unwrap_or(false)
}
enum Probe {
WithinBBox {
ll: lora_store::LoraPoint,
ur: lora_store::LoraPoint,
},
WithinDistance {
center: lora_store::LoraPoint,
max: f64,
inclusive: bool,
},
}
fn point_predicate_holds(actual: &lora_store::LoraPoint, probe: &Probe) -> bool {
match probe {
Probe::WithinBBox { ll, ur } => {
if actual.srid != ll.srid || actual.srid != ur.srid {
return false;
}
let in_x = actual.x >= ll.x.min(ur.x) && actual.x <= ll.x.max(ur.x);
let in_y = actual.y >= ll.y.min(ur.y) && actual.y <= ll.y.max(ur.y);
let in_z = match (actual.z, ll.z, ur.z) {
(Some(pz), Some(lz), Some(uz)) => pz >= lz.min(uz) && pz <= lz.max(uz),
(None, None, None) => true,
_ => return false,
};
in_x && in_y && in_z
}
Probe::WithinDistance {
center,
max,
inclusive,
} => {
let Some(d) = lora_store::point_distance(actual, center) else {
return false;
};
if *inclusive {
d <= *max
} else {
d < *max
}
}
}
}
pub(crate) fn indexed_node_property_candidates<S: GraphStorage>(
storage: &S,
labels: &[Vec<String>],
key: &str,
expected: &LoraValue,
) -> NodePropertyCandidates {
let Some(values) = property_lookup_values(expected) else {
return NodePropertyCandidates {
ids: scan_node_ids_for_label_groups(storage, labels),
prefiltered: false,
};
};
let label_hint = single_label_hint(labels);
let mut seen = BTreeSet::new();
let mut out = Vec::new();
for value in values {
for id in storage.find_node_ids_by_property(label_hint, key, &value) {
if seen.insert(id) {
out.push(id);
}
}
}
NodePropertyCandidates {
ids: out,
prefiltered: labels.is_empty() || label_hint.is_some(),
}
}
pub(crate) fn build_path_value<S: GraphStorage>(
row: &Row,
node_vars: &[VarId],
rel_vars: &[VarId],
storage: &S,
) -> LoraValue {
let (raw_nodes, rels, has_var_len) = path_bindings(row, node_vars, rel_vars);
let nodes = if has_var_len && !rels.is_empty() && raw_nodes.len() == 2 {
reconstruct_var_len_nodes(raw_nodes[0], &rels, storage)
} else {
raw_nodes
};
LoraValue::Path(LoraPath { nodes, rels })
}
#[inline]
fn path_bindings(
row: &Row,
node_vars: &[VarId],
rel_vars: &[VarId],
) -> (Vec<NodeId>, Vec<RelationshipId>, bool) {
let mut raw_nodes = Vec::new();
let mut rels = Vec::new();
let mut has_var_len = false;
for &nv in node_vars {
match row.get(nv) {
Some(LoraValue::Node(id)) => raw_nodes.push(*id),
Some(LoraValue::List(items)) => {
for item in items {
if let LoraValue::Node(id) = item {
raw_nodes.push(*id);
}
}
}
_ => {}
}
}
for &rv in rel_vars {
match row.get(rv) {
Some(LoraValue::Relationship(id)) => rels.push(*id),
Some(LoraValue::List(items)) => {
has_var_len = true;
for item in items {
if let LoraValue::Relationship(id) = item {
rels.push(*id);
}
}
}
_ => {}
}
}
(raw_nodes, rels, has_var_len)
}
#[inline]
fn reconstruct_var_len_nodes<S: GraphStorage>(
start: NodeId,
rels: &[RelationshipId],
storage: &S,
) -> Vec<NodeId> {
let mut ordered = Vec::with_capacity(rels.len() + 1);
ordered.push(start);
let mut current = start;
for &rel_id in rels {
if let Some((src, dst)) = storage.relationship_endpoints(rel_id) {
let next = if src == current { dst } else { src };
ordered.push(next);
current = next;
}
}
ordered
}
fn type_rank(v: &LoraValue) -> u8 {
match v {
LoraValue::Null => 0,
LoraValue::Bool(_) => 1,
LoraValue::Int(_) | LoraValue::Float(_) => 2,
LoraValue::String(_) => 3,
LoraValue::Binary(_) => 4,
LoraValue::Date(_) => 5,
LoraValue::DateTime(_) => 6,
LoraValue::LocalDateTime(_) => 7,
LoraValue::Time(_) => 8,
LoraValue::LocalTime(_) => 9,
LoraValue::Duration(_) => 10,
LoraValue::Point(_) => 11,
LoraValue::Vector(_) => 12,
LoraValue::List(_) => 13,
LoraValue::Map(_) => 14,
LoraValue::Node(_) => 15,
LoraValue::Relationship(_) => 16,
LoraValue::Path(_) => 17,
}
}
pub(crate) fn node_matches_label_groups(node_labels: &[String], groups: &[Vec<String>]) -> bool {
groups
.iter()
.all(|group| group.iter().any(|l| node_labels.iter().any(|nl| nl == l)))
}
pub(crate) fn scan_node_ids_for_label_groups<S: GraphStorage>(
storage: &S,
groups: &[Vec<String>],
) -> Vec<NodeId> {
if groups.is_empty() {
return storage.all_node_ids();
}
if groups.len() == 1 {
return label_group_candidate_ids(storage, &groups[0]);
}
let mut best: Option<Vec<NodeId>> = None;
for group in groups {
let ids = label_group_candidate_ids(storage, group);
if ids.is_empty() {
return Vec::new();
}
if best
.as_ref()
.map(|current| ids.len() < current.len())
.unwrap_or(true)
{
best = Some(ids);
}
}
best.unwrap_or_default()
}
pub(crate) fn label_group_candidates_prefiltered(groups: &[Vec<String>]) -> bool {
groups.len() <= 1
}
fn label_group_candidate_ids<S: GraphStorage>(storage: &S, group: &[String]) -> Vec<NodeId> {
match group {
[] => Vec::new(),
[label] => storage.node_ids_by_label(label),
labels => {
let mut seen = BTreeSet::new();
let mut out = Vec::new();
for label in labels {
for id in storage.node_ids_by_label(label) {
if seen.insert(id) {
out.push(id);
}
}
}
out
}
}
}
pub(crate) fn hydrate_node_record(node: &lora_store::NodeRecord) -> LoraValue {
let mut map = BTreeMap::new();
map.insert("kind".to_string(), LoraValue::String("node".to_string()));
map.insert("id".to_string(), LoraValue::Int(node.id as i64));
map.insert(
"labels".to_string(),
LoraValue::List(
node.labels
.iter()
.map(|s| LoraValue::String(s.clone()))
.collect(),
),
);
map.insert(
"properties".to_string(),
properties_to_value_map(&node.properties),
);
LoraValue::Map(map)
}
pub(crate) fn hydrate_relationship_record(rel: &lora_store::RelationshipRecord) -> LoraValue {
let mut map = BTreeMap::new();
map.insert(
"kind".to_string(),
LoraValue::String("relationship".to_string()),
);
map.insert("id".to_string(), LoraValue::Int(rel.id as i64));
map.insert("startId".to_string(), LoraValue::Int(rel.src as i64));
map.insert("endId".to_string(), LoraValue::Int(rel.dst as i64));
map.insert("type".to_string(), LoraValue::String(rel.rel_type.clone()));
map.insert(
"properties".to_string(),
properties_to_value_map(&rel.properties),
);
LoraValue::Map(map)
}
pub(super) fn flatten_label_groups(groups: &[Vec<String>]) -> Vec<String> {
groups.iter().flat_map(|g| g.iter().cloned()).collect()
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
pub(crate) enum GroupValueKey {
Null,
Bool(bool),
Int(i64),
Float(String),
String(String),
Binary(Vec<Vec<u8>>),
List(Vec<GroupValueKey>),
Map(Vec<(String, GroupValueKey)>),
Node(u64),
Relationship(u64),
}
impl GroupValueKey {
pub(crate) fn from_value(v: &LoraValue) -> Self {
match v {
LoraValue::Null => Self::Null,
LoraValue::Bool(x) => Self::Bool(*x),
LoraValue::Int(x) => Self::Int(*x),
LoraValue::Float(x) => Self::Float(x.to_string()),
LoraValue::String(x) => Self::String(x.clone()),
LoraValue::Binary(x) => Self::Binary(x.segments().to_vec()),
LoraValue::List(xs) => Self::List(xs.iter().map(Self::from_value).collect()),
LoraValue::Map(m) => Self::Map(
m.iter()
.map(|(k, v)| (k.clone(), Self::from_value(v)))
.collect(),
),
LoraValue::Node(id) => Self::Node(*id),
LoraValue::Relationship(id) => Self::Relationship(*id),
LoraValue::Path(_) => Self::Null,
LoraValue::Date(d) => Self::String(d.to_string()),
LoraValue::DateTime(dt) => Self::String(dt.to_string()),
LoraValue::LocalDateTime(dt) => Self::String(dt.to_string()),
LoraValue::Time(t) => Self::String(t.to_string()),
LoraValue::LocalTime(t) => Self::String(t.to_string()),
LoraValue::Duration(dur) => Self::String(dur.to_string()),
LoraValue::Point(p) => Self::String(p.to_string()),
LoraValue::Vector(v) => Self::String(format!("vector:{}", v.to_key_string())),
}
}
}
const MAX_VAR_LEN_HOPS: u64 = 100;
pub(crate) fn resolve_range(range: &RangeLiteral) -> (u64, u64) {
let min_hops = range.start.unwrap_or(1);
let max_hops = range.end.unwrap_or(MAX_VAR_LEN_HOPS);
(min_hops, max_hops)
}
pub(crate) struct VarLenResult {
pub(crate) dst_node_id: NodeId,
pub(crate) rel_ids: Vec<u64>,
}
pub(crate) fn variable_length_expand<S: GraphStorage>(
storage: &S,
start_node_id: NodeId,
direction: Direction,
types: &[String],
min_hops: u64,
max_hops: u64,
bind_relationships: bool,
) -> Vec<VarLenResult> {
let mut results = Vec::new();
let mut frontier: Vec<(NodeId, Vec<u64>)> = vec![(start_node_id, Vec::new())];
for depth in 1..=max_hops {
let is_last_hop = depth == max_hops;
let mut next_frontier: Vec<(NodeId, Vec<u64>)> = Vec::new();
for (current_node, rels_used) in &frontier {
for (rel_id, neighbor_id) in storage.expand_ids(*current_node, direction, types) {
if rels_used.contains(&rel_id) {
continue;
}
if is_last_hop {
if depth >= min_hops {
let mut rel_ids = Vec::with_capacity(rels_used.len() + 1);
rel_ids.extend_from_slice(rels_used);
rel_ids.push(rel_id);
results.push(VarLenResult {
dst_node_id: neighbor_id,
rel_ids: if bind_relationships {
rel_ids
} else {
Vec::new()
},
});
}
continue;
}
let mut new_rels = Vec::with_capacity(rels_used.len() + 1);
new_rels.extend_from_slice(rels_used);
new_rels.push(rel_id);
if depth >= min_hops {
results.push(VarLenResult {
dst_node_id: neighbor_id,
rel_ids: if bind_relationships {
new_rels.clone()
} else {
Vec::new()
},
});
}
next_frontier.push((neighbor_id, new_rels));
}
}
if is_last_hop || next_frontier.is_empty() {
break;
}
frontier = next_frontier;
}
if min_hops == 0 {
results.insert(
0,
VarLenResult {
dst_node_id: start_node_id,
rel_ids: Vec::new(),
},
);
}
results
}
pub(crate) fn filter_shortest_paths(rows: Vec<Row>, path_var: VarId, all: bool) -> Vec<Row> {
if rows.is_empty() {
return rows;
}
let lengths: Vec<usize> = rows
.iter()
.map(|row| match row.get(path_var) {
Some(LoraValue::Path(p)) => p.rels.len(),
_ => usize::MAX,
})
.collect();
let min_len = lengths.iter().copied().min().unwrap_or(usize::MAX);
let mut result: Vec<Row> = rows
.into_iter()
.zip(lengths.iter())
.filter(|(_, len)| **len == min_len)
.map(|(row, _)| row)
.collect();
if !all && result.len() > 1 {
result.truncate(1);
}
result
}
fn emit_rel_rows(
direction: Direction,
src_var: VarId,
rel_var: VarId,
dst_var: VarId,
rel: &lora_store::RelationshipRecord,
base: &Row,
out: &mut Vec<Row>,
) -> ExecResult<()> {
match direction {
Direction::Right => {
emit_one_rel_row(
src_var, rel_var, dst_var, rel.src, rel.id, rel.dst, base, out,
)?;
}
Direction::Left => {
emit_one_rel_row(
src_var, rel_var, dst_var, rel.dst, rel.id, rel.src, base, out,
)?;
}
Direction::Undirected => {
emit_one_rel_row(
src_var, rel_var, dst_var, rel.src, rel.id, rel.dst, base, out,
)?;
if rel.src != rel.dst {
emit_one_rel_row(
src_var, rel_var, dst_var, rel.dst, rel.id, rel.src, base, out,
)?;
}
}
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
fn emit_one_rel_row(
src_var: VarId,
rel_var: VarId,
dst_var: VarId,
src_id: NodeId,
rel_id: RelationshipId,
dst_id: NodeId,
base: &Row,
out: &mut Vec<Row>,
) -> ExecResult<()> {
let mut row = base.clone();
if bind_node_value(&mut row, src_var, src_id)?
&& bind_relationship_value(&mut row, rel_var, rel_id)?
&& bind_node_value(&mut row, dst_var, dst_id)?
{
out.push(row);
}
Ok(())
}
fn bind_node_value(row: &mut Row, var: VarId, id: NodeId) -> ExecResult<bool> {
match row.get(var) {
Some(LoraValue::Node(existing)) => Ok(*existing == id),
Some(other) => Err(ExecutorError::ExpectedNodeForExpand {
var: format!("{var:?}"),
found: value_kind(other),
}),
None => {
row.insert(var, LoraValue::Node(id));
Ok(true)
}
}
}
fn bind_relationship_value(row: &mut Row, var: VarId, id: RelationshipId) -> ExecResult<bool> {
match row.get(var) {
Some(LoraValue::Relationship(existing)) => Ok(*existing == id),
Some(other) => Err(ExecutorError::ExpectedRelationshipForExpand {
var: format!("{var:?}"),
found: value_kind(other),
}),
None => {
row.insert(var, LoraValue::Relationship(id));
Ok(true)
}
}
}
fn rel_candidate_ids<S, F>(storage: &S, types: &[String], indexed: F) -> Vec<RelationshipId>
where
S: GraphStorage,
F: Fn(&str) -> Option<Vec<RelationshipId>>,
{
if types.is_empty() {
return storage.all_rel_ids();
}
let mut all = Vec::new();
let mut seen = BTreeSet::new();
for ty in types {
let ids = indexed(ty).unwrap_or_else(|| storage.rel_ids_by_type(ty));
for id in ids {
if seen.insert(id) {
all.push(id);
}
}
}
all
}
pub(crate) fn rel_by_property_range_scan_rows<S: GraphStorage>(
storage: &S,
params: &BTreeMap<String, LoraValue>,
base_rows: Vec<Row>,
op: &lora_compiler::RelByPropertyRangeScanExec,
deadline: Option<Instant>,
) -> ExecResult<Vec<Row>> {
let eval_ctx = EvalContext { storage, params };
let mut out = Vec::new();
for row in base_rows {
check_optional_deadline(deadline)?;
let lo_value = op.lo.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
let hi_value = op.hi.as_ref().map(|expr| eval_expr(expr, &row, &eval_ctx));
let lo_prop = lo_value
.clone()
.and_then(|v| lora_value_to_property(v).ok());
let hi_prop = hi_value
.clone()
.and_then(|v| lora_value_to_property(v).ok());
let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| {
storage.relationship_range_candidates(ty, &op.key, lo_prop.as_ref(), hi_prop.as_ref())
});
for rel_id in candidate_ids {
check_optional_deadline(deadline)?;
if let Some(result) = storage.with_relationship(rel_id, |rel| {
if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
return Ok(());
}
let Some(actual) = rel.properties.get(&op.key) else {
return Ok(());
};
let actual_lv = LoraValue::from(actual);
if !range_predicate_holds(
&actual_lv,
lo_value.as_ref(),
op.lo_inclusive,
hi_value.as_ref(),
op.hi_inclusive,
) {
return Ok(());
}
emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
}) {
result?;
}
}
}
Ok(out)
}
pub(crate) fn rel_by_text_scan_rows<S: GraphStorage>(
storage: &S,
params: &BTreeMap<String, LoraValue>,
base_rows: Vec<Row>,
op: &lora_compiler::RelByTextScanExec,
deadline: Option<Instant>,
) -> ExecResult<Vec<Row>> {
let eval_ctx = EvalContext { storage, params };
let mut out = Vec::new();
for row in base_rows {
check_optional_deadline(deadline)?;
let query = eval_expr(&op.query, &row, &eval_ctx);
let LoraValue::String(query_str) = &query else {
continue;
};
let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| {
storage.relationship_text_candidates(ty, &op.key, query_str)
});
for rel_id in candidate_ids {
check_optional_deadline(deadline)?;
if let Some(result) = storage.with_relationship(rel_id, |rel| {
if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
return Ok(());
}
let Some(PropertyValue::String(actual)) = rel.properties.get(&op.key) else {
return Ok(());
};
if !text_predicate_holds(actual, op.predicate, query_str) {
return Ok(());
}
emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
}) {
result?;
}
}
}
Ok(out)
}
pub(crate) fn rel_by_point_scan_rows<S: GraphStorage>(
storage: &S,
params: &BTreeMap<String, LoraValue>,
base_rows: Vec<Row>,
op: &lora_compiler::RelByPointScanExec,
deadline: Option<Instant>,
) -> ExecResult<Vec<Row>> {
let eval_ctx = EvalContext { storage, params };
let mut out = Vec::new();
for row in base_rows {
check_optional_deadline(deadline)?;
let probe = match &op.predicate {
lora_compiler::PointPredicate::WithinBBox {
lower_left,
upper_right,
} => {
let ll = eval_expr(lower_left, &row, &eval_ctx);
let ur = eval_expr(upper_right, &row, &eval_ctx);
match (ll, ur) {
(LoraValue::Point(a), LoraValue::Point(b)) => {
Probe::WithinBBox { ll: a, ur: b }
}
_ => continue,
}
}
lora_compiler::PointPredicate::WithinDistance {
center,
max_distance,
inclusive,
} => {
let c = eval_expr(center, &row, &eval_ctx);
let d = eval_expr(max_distance, &row, &eval_ctx);
match (c, d) {
(LoraValue::Point(c), LoraValue::Float(d)) => Probe::WithinDistance {
center: c,
max: d,
inclusive: *inclusive,
},
(LoraValue::Point(c), LoraValue::Int(d)) => Probe::WithinDistance {
center: c,
max: d as f64,
inclusive: *inclusive,
},
_ => continue,
}
}
};
let candidate_ids = rel_candidate_ids(storage, &op.types, |ty| match &probe {
Probe::WithinBBox { ll, ur } => {
storage.relationship_point_within_bbox(ty, &op.key, (ll.x, ll.y), (ur.x, ur.y))
}
Probe::WithinDistance { center, max, .. } => {
storage.relationship_point_within_distance(ty, &op.key, (center.x, center.y), *max)
}
});
for rel_id in candidate_ids {
check_optional_deadline(deadline)?;
if let Some(result) = storage.with_relationship(rel_id, |rel| {
if !op.types.is_empty() && !op.types.iter().any(|t| t == &rel.rel_type) {
return Ok(());
}
let Some(PropertyValue::Point(actual)) = rel.properties.get(&op.key) else {
return Ok(());
};
if !point_predicate_holds(actual, &probe) {
return Ok(());
}
emit_rel_rows(op.direction, op.src, op.rel, op.dst, rel, &row, &mut out)
}) {
result?;
}
}
}
Ok(out)
}