use std::sync::Arc;
use super::Planner;
use crate::err::Error;
use crate::exec::ExecOperator;
use crate::exec::operators::{
Aggregate, AggregateField, Bind, DeleteBinding, Distinct, DistinctEdges, DrainSink,
EdgeBinding, EndpointBind, EndpointField, Expand, ExpandDir, FieldSelection, Filter, HashJoin,
InsertEdgeOp, InsertGraph, InsertNodeOp, JoinType, Limit, OrderByField, PathExpand, PathMode,
Project, ShortestPathExpand, ShortestSelector, SingleRowScan, Sort, SortDirection,
UpdateBinding,
};
use crate::expr::match_plan::{
BindingId, BindingKind, EdgeStep, ExpandDirection, MatchClausePlan, MatchOutput, MatchPlan,
MatchPredicate, MatchStage, MutationStage, NodeStep, PathMode as IrPathMode, PathPrefixPlan,
PathSearch, PatternPlan,
};
use crate::expr::{Cond, Expr, Function, Idiom, Literal, Part};
use crate::val::TableName;
enum SearchRouting {
Every,
Shortest {
selector: ShortestSelector,
},
}
fn resolve_path_search(search: Option<PathPrefixPlan>) -> (PathMode, SearchRouting) {
let Some(prefix) = search else {
return (PathMode::Walk, SearchRouting::Every);
};
let mode = match prefix.mode {
IrPathMode::Walk => PathMode::Walk,
IrPathMode::Trail => PathMode::Trail,
IrPathMode::Simple => PathMode::Simple,
IrPathMode::Acyclic => PathMode::Acyclic,
};
let routing = match prefix.search {
PathSearch::All => SearchRouting::Every,
PathSearch::Any {
count,
} => SearchRouting::Shortest {
selector: ShortestSelector::Any {
count,
},
},
PathSearch::AllShortest => SearchRouting::Shortest {
selector: ShortestSelector::AllShortest,
},
PathSearch::AnyShortest => SearchRouting::Shortest {
selector: ShortestSelector::AnyShortest,
},
PathSearch::ShortestCounted {
count,
} => SearchRouting::Shortest {
selector: ShortestSelector::Counted(count),
},
PathSearch::ShortestGroups {
count,
} => SearchRouting::Shortest {
selector: ShortestSelector::CountedGroups(count),
},
};
(mode, routing)
}
struct ChainStage {
operator: Arc<dyn ExecOperator>,
bound: Vec<BindingId>,
tip: BindingId,
}
struct PlannedClause {
operator: Arc<dyn ExecOperator>,
bound: Vec<BindingId>,
deferred: Vec<MatchPredicate>,
}
impl<'ctx> Planner<'ctx> {
pub(crate) async fn plan_match(&self, plan: MatchPlan) -> Result<Arc<dyn ExecOperator>, Error> {
let units = stage_units(&plan);
let mut units = units.into_iter();
let Some(first) = units.next() else {
return Err(Error::Internal(
"GQL planning: a MatchPlan must carry at least one stage".to_string(),
));
};
let (mut acc_op, mut acc_bound) = match first {
StageUnit::Mandatory(clause) => {
let planned = self.plan_clause(&plan, clause, &[]).await?;
debug_assert!(
planned.deferred.is_empty(),
"GQL planning: the first clause cannot defer predicates (nothing precedes it)",
);
(planned.operator, planned.bound)
}
StageUnit::Optional(_) => {
return Err(Error::Internal(
"GQL planning: a query cannot start with OPTIONAL (the lowering must reject it)"
.to_string(),
));
}
StageUnit::Mutation(stage) => {
let seed = Arc::new(SingleRowScan::new()) as Arc<dyn ExecOperator>;
self.fold_mutation(&plan, seed, Vec::new(), stage)
}
};
for unit in units {
let (op, bound) = match unit {
StageUnit::Mandatory(clause) => {
let planned = self.plan_clause(&plan, clause, &acc_bound).await?;
self.fold_mandatory(&plan, acc_op, &acc_bound, planned).await?
}
StageUnit::Optional(clauses) => {
self.fold_optional(&plan, acc_op, &acc_bound, &clauses).await?
}
StageUnit::Mutation(stage) => self.fold_mutation(&plan, acc_op, acc_bound, stage),
};
acc_op = op;
acc_bound = bound;
}
self.plan_match_tail(&plan, acc_op).await
}
fn fold_mutation(
&self,
plan: &MatchPlan,
op: Arc<dyn ExecOperator>,
mut bound: Vec<BindingId>,
stage: &MutationStage,
) -> (Arc<dyn ExecOperator>, Vec<BindingId>) {
let op = match stage {
MutationStage::Update {
target,
data,
} => {
let name = binding_name(plan, *target).to_string();
Arc::new(UpdateBinding::new(op, name, data.clone())) as Arc<dyn ExecOperator>
}
MutationStage::Delete {
target,
detach,
} => {
let name = binding_name(plan, *target).to_string();
let is_edge = matches!(
plan.binding(*target).kind,
BindingKind::Edge | BindingKind::EdgeGroup
);
Arc::new(DeleteBinding::new(op, name, *detach, is_edge)) as Arc<dyn ExecOperator>
}
MutationStage::Insert(insert) => {
let nodes: Vec<InsertNodeOp> = insert
.nodes
.iter()
.map(|n| InsertNodeOp {
name: binding_name(plan, n.binding).to_string(),
table: n.label.clone(),
props: n.props.clone(),
})
.collect();
let edges: Vec<InsertEdgeOp> = insert
.edges
.iter()
.map(|e| InsertEdgeOp {
name: binding_name(plan, e.binding).to_string(),
table: e.label.clone(),
from: binding_name(plan, e.from).to_string(),
to: binding_name(plan, e.to).to_string(),
props: e.props.clone(),
})
.collect();
bound.extend(insert.nodes.iter().map(|n| n.binding));
bound.extend(insert.edges.iter().map(|e| e.binding));
Arc::new(InsertGraph::new(op, nodes, edges)) as Arc<dyn ExecOperator>
}
};
(op, bound)
}
async fn fold_mandatory(
&self,
plan: &MatchPlan,
acc_op: Arc<dyn ExecOperator>,
acc_bound: &[BindingId],
clause: PlannedClause,
) -> Result<(Arc<dyn ExecOperator>, Vec<BindingId>), Error> {
let shared = shared_bindings(acc_bound, &clause.bound);
let keys = binding_names(plan, &shared);
let join_type = if shared.is_empty() {
JoinType::Cross
} else {
JoinType::Inner
};
let join =
Arc::new(HashJoin::new(acc_op, clause.operator, keys, join_type, Vec::new(), None))
as Arc<dyn ExecOperator>;
let bound = union_bindings(acc_bound, &clause.bound);
let op = self.place_above(join, &bound, clause.deferred).await?;
Ok((op, bound))
}
async fn fold_optional(
&self,
plan: &MatchPlan,
acc_op: Arc<dyn ExecOperator>,
acc_bound: &[BindingId],
clauses: &[&MatchClausePlan],
) -> Result<(Arc<dyn ExecOperator>, Vec<BindingId>), Error> {
if let Some(stage) = self
.try_optional_expand_fast_path(plan, Arc::clone(&acc_op), acc_bound, clauses)
.await?
{
let bound = union_bindings(acc_bound, &stage.bound);
return Ok((stage.operator, bound));
}
let subplan = self.plan_optional_block(plan, clauses).await?;
let shared = shared_bindings(acc_bound, &subplan.bound);
let keys = binding_names(plan, &shared);
let introduced = difference_bindings(&subplan.bound, acc_bound);
let null_template = binding_names(plan, &introduced);
let merged_bound = union_bindings(acc_bound, &subplan.bound);
let mut residual_exprs = Vec::with_capacity(subplan.deferred.len());
for predicate in subplan.deferred {
if !deps_subset(&predicate.deps, &merged_bound) {
return Err(Error::Internal(
"GQL MATCH planning: an OPTIONAL block's deferred predicate has deps \
unsatisfied by the accumulator and block bindings combined"
.to_string(),
));
}
residual_exprs.push(predicate.expr);
}
let residual = match conjoin(residual_exprs) {
Some(joined) => Some(self.physical_expr(joined).await?),
None => None,
};
let join = Arc::new(HashJoin::new(
subplan.operator,
acc_op,
keys,
JoinType::Left,
null_template,
residual,
)) as Arc<dyn ExecOperator>;
Ok((join, merged_bound))
}
async fn try_optional_expand_fast_path(
&self,
plan: &MatchPlan,
acc_op: Arc<dyn ExecOperator>,
acc_bound: &[BindingId],
clauses: &[&MatchClausePlan],
) -> Result<Option<ChainStage>, Error> {
let [clause] = clauses else {
return Ok(None);
};
let [pattern] = clause.patterns.as_slice() else {
return Ok(None);
};
let [(edge, far_node)] = pattern.steps.as_slice() else {
return Ok(None);
};
if edge.quantifier.is_some() || pattern.path_var.is_some() {
return Ok(None);
}
let forward = acc_bound.contains(&pattern.start.binding);
let reverse = acc_bound.contains(&far_node.binding);
let (source_node, target_node, direction) = if forward && !reverse {
(&pattern.start, far_node, edge.direction)
} else if reverse && !forward {
(far_node, &pattern.start, edge.direction.reverse())
} else {
return Ok(None);
};
let mut pending: Vec<MatchPredicate> = clause.predicates.clone();
let source = binding_name(plan, source_node.binding).to_string();
let dir = expand_dir(direction);
let edge_tables = edge.label.iter().cloned().collect::<Vec<_>>();
let target_binding = binding_name(plan, target_node.binding).to_string();
let target_label = target_node.label.clone();
let step_bound = {
let mut b = acc_bound.to_vec();
b.push(edge.binding);
b.push(target_node.binding);
b
};
let predicates = drain_satisfiable(&mut pending, &step_bound);
if !pending.is_empty() {
return Err(Error::Internal(
"GQL MATCH planning: an OPTIONAL single-hop clause owns a predicate its hop cannot \
satisfy; the lowering must scope inside-optional predicates to the clause"
.to_string(),
));
}
let predicate = match conjoin(predicates) {
Some(joined) => Some(self.physical_expr(joined).await?),
None => None,
};
let edge_def = plan.binding(edge.binding);
let edge_name = edge_def.name.clone();
let edge_binding = if edge_def.user_named {
EdgeBinding::Full(edge_name)
} else {
EdgeBinding::IdOnly(edge_name)
};
let operator = Arc::new(Expand::new(
acc_op,
source,
dir,
edge_tables,
edge_binding,
target_binding,
target_label,
predicate,
true,
)) as Arc<dyn ExecOperator>;
Ok(Some(ChainStage {
operator,
bound: vec![edge.binding, target_node.binding],
tip: target_node.binding,
}))
}
async fn plan_optional_block(
&self,
plan: &MatchPlan,
clauses: &[&MatchClausePlan],
) -> Result<PlannedClause, Error> {
let Some((first, rest)) = clauses.split_first() else {
return Err(Error::Internal(
"GQL MATCH planning: an OPTIONAL block must carry at least one clause".to_string(),
));
};
let seed = self.plan_clause(plan, first, &[]).await?;
let mut block_op = seed.operator;
let mut block_bound = seed.bound;
let mut block_deferred = seed.deferred;
for clause in rest {
let planned = self
.plan_clause_onto(
plan,
clause,
&block_bound,
Some((Arc::clone(&block_op), block_bound.clone())),
)
.await?;
block_op = planned.operator;
block_bound = planned.bound;
block_deferred.extend(planned.deferred);
}
Ok(PlannedClause {
operator: block_op,
bound: block_bound,
deferred: block_deferred,
})
}
async fn plan_clause(
&self,
plan: &MatchPlan,
clause: &MatchClausePlan,
bound_acc: &[BindingId],
) -> Result<PlannedClause, Error> {
self.plan_clause_onto(plan, clause, bound_acc, None).await
}
async fn plan_clause_onto(
&self,
plan: &MatchPlan,
clause: &MatchClausePlan,
bound_acc: &[BindingId],
seed: Option<(Arc<dyn ExecOperator>, Vec<BindingId>)>,
) -> Result<PlannedClause, Error> {
let Some((first_pattern, rest_patterns)) = clause.patterns.split_first() else {
return Err(Error::Internal(
"GQL MATCH planning: a clause must carry at least one pattern".to_string(),
));
};
let mut pending: Vec<MatchPredicate> = clause.predicates.clone();
let mut clause_op: Arc<dyn ExecOperator>;
let mut clause_bound: Vec<BindingId>;
match self.plan_pattern_selfrooted(plan, first_pattern, &mut pending).await? {
Some(stage) => match &seed {
None => {
clause_bound = stage.bound.clone();
clause_op = stage.operator;
}
Some((seed_op, seed_bound)) => {
let shared = shared_bindings(seed_bound, &stage.bound);
let keys = binding_names(plan, &shared);
let join_type = if shared.is_empty() {
JoinType::Cross
} else {
JoinType::Inner
};
clause_op = Arc::new(HashJoin::new(
Arc::clone(seed_op),
stage.operator,
keys,
join_type,
Vec::new(),
None,
)) as Arc<dyn ExecOperator>;
clause_bound = union_bindings(seed_bound, &stage.bound);
}
},
None => match &seed {
Some((seed_op, seed_bound))
if shared_node_anchor(plan, first_pattern, seed_bound).is_some() =>
{
let stage = self
.expand_bound_pattern(
plan,
first_pattern,
seed_bound,
Arc::clone(seed_op),
&mut pending,
)
.await?;
clause_bound = union_bindings(seed_bound, &stage.bound);
clause_op = stage.operator;
}
_ => {
return Err(Error::Internal(
"GQL MATCH planning: a clause's leading pattern is unanchorable (no labeled \
element) and no expandable accumulator was supplied; the lowering must \
reject it"
.to_string(),
));
}
},
}
clause_op = self.place_satisfiable(clause_op, &clause_bound, &mut pending).await?;
for pattern in rest_patterns {
let visible = union_bindings(bound_acc, &clause_bound);
if let Some(stage) = self.plan_pattern_selfrooted(plan, pattern, &mut pending).await? {
let shared = shared_bindings(&clause_bound, &stage.bound);
let keys = binding_names(plan, &shared);
let join_type = if shared.is_empty() {
JoinType::Cross
} else {
JoinType::Inner
};
clause_op = Arc::new(HashJoin::new(
clause_op,
stage.operator,
keys,
join_type,
Vec::new(),
None,
)) as Arc<dyn ExecOperator>;
clause_bound = union_bindings(&clause_bound, &stage.bound);
} else if shared_node_anchor(plan, pattern, &visible).is_some() {
let stage = self
.expand_bound_pattern(
plan,
pattern,
&visible,
Arc::clone(&clause_op),
&mut pending,
)
.await?;
clause_op = stage.operator;
clause_bound = union_bindings(&clause_bound, &stage.bound);
} else {
return Err(Error::Internal(
"GQL MATCH planning: pattern is neither self-anchorable nor reuses a bound \
variable (unanchorable); the lowering must reject it"
.to_string(),
));
}
clause_op = self.place_satisfiable(clause_op, &clause_bound, &mut pending).await?;
}
clause_op = self.wrap_distinct_edges(plan, clause, clause_op);
let mut deferred = Vec::new();
for predicate in pending {
if deps_subset(&predicate.deps, &clause_bound) {
return Err(Error::Internal(
"GQL MATCH planning: a clause predicate was not placed despite all its \
dependencies being bound within the clause"
.to_string(),
));
}
deferred.push(predicate);
}
Ok(PlannedClause {
operator: clause_op,
bound: clause_bound,
deferred,
})
}
async fn plan_pattern_selfrooted(
&self,
plan: &MatchPlan,
pattern: &PatternPlan,
pending: &mut Vec<MatchPredicate>,
) -> Result<Option<ChainStage>, Error> {
if pattern.start.label.is_some() {
return Ok(Some(self.plan_node_anchored(plan, pattern, pending).await?));
}
if let Some(stage) = self.plan_edge_anchored(plan, pattern).await? {
return Ok(Some(stage));
}
Ok(None)
}
async fn plan_node_anchored(
&self,
plan: &MatchPlan,
pattern: &PatternPlan,
pending: &mut Vec<MatchPredicate>,
) -> Result<ChainStage, Error> {
let Some(anchor_label) = pattern.start.label.clone() else {
return Err(Error::Internal(
"GQL MATCH planning: plan_node_anchored called on an unlabeled start node"
.to_string(),
));
};
let anchor_binding = pattern.start.binding;
let mut anchor_pushed: Vec<Expr> = Vec::new();
let mut anchor_filtered: Vec<Expr> = Vec::new();
pending.retain(|predicate| {
if predicate.deps.as_slice() == [anchor_binding] {
match prefix_strip(&predicate.expr, binding_name(plan, anchor_binding)) {
Some(stripped) => anchor_pushed.push(stripped),
None => anchor_filtered.push(predicate.expr.clone()),
}
false
} else {
true
}
});
let mut stage = self.plan_anchor(anchor_binding, anchor_label, plan, anchor_pushed).await?;
stage = self.maybe_filter(stage, anchor_filtered).await?;
for (edge, node) in pattern.steps.iter() {
let step_bound = step_bindings(&stage.bound, pattern, edge, node);
let mut step_predicates: Vec<Expr> = Vec::new();
pending.retain(|predicate| {
if deps_subset(&predicate.deps, &step_bound) {
step_predicates.push(predicate.expr.clone());
false
} else {
true
}
});
stage = self
.plan_step(plan, pattern, stage, edge, node, step_predicates, edge.direction, false)
.await?;
}
Ok(stage)
}
async fn plan_edge_anchored(
&self,
plan: &MatchPlan,
pattern: &PatternPlan,
) -> Result<Option<ChainStage>, Error> {
let [(edge, far_node)] = pattern.steps.as_slice() else {
return Ok(None);
};
let Some(edge_label) = edge.label.clone() else {
return Ok(None);
};
if edge.quantifier.is_some() {
return Ok(None);
}
let edge_binding = edge.binding;
let edge_name = binding_name(plan, edge_binding).to_string();
let scan = self.plan_table_scan(edge_label, Vec::new()).await?;
let bind = Arc::new(Bind::new(scan, edge_name.clone())) as Arc<dyn ExecOperator>;
let (in_node, out_node) = match edge.direction {
ExpandDirection::Out => (&pattern.start, far_node),
ExpandDirection::In => (far_node, &pattern.start),
};
let in_bind = Arc::new(EndpointBind::new(
bind,
edge_name.clone(),
EndpointField::In,
binding_name(plan, in_node.binding).to_string(),
in_node.label.clone(),
)) as Arc<dyn ExecOperator>;
let out_bind = Arc::new(EndpointBind::new(
in_bind,
edge_name,
EndpointField::Out,
binding_name(plan, out_node.binding).to_string(),
out_node.label.clone(),
)) as Arc<dyn ExecOperator>;
Ok(Some(ChainStage {
operator: out_bind,
bound: vec![edge_binding, in_node.binding, out_node.binding],
tip: far_node.binding,
}))
}
async fn expand_bound_pattern(
&self,
plan: &MatchPlan,
pattern: &PatternPlan,
visible: &[BindingId],
input: Arc<dyn ExecOperator>,
pending: &mut Vec<MatchPredicate>,
) -> Result<ChainStage, Error> {
if visible.contains(&pattern.start.binding) {
let mut stage = ChainStage {
operator: input,
bound: visible.to_vec(),
tip: pattern.start.binding,
};
for (edge, node) in pattern.steps.iter() {
let step_bound = step_bindings(&stage.bound, pattern, edge, node);
let step_predicates = drain_satisfiable(pending, &step_bound);
stage = self
.plan_step(
plan,
pattern,
stage,
edge,
node,
step_predicates,
edge.direction,
false,
)
.await?;
}
return Ok(stage);
}
if let [(edge, far_node)] = pattern.steps.as_slice()
&& visible.contains(&far_node.binding)
{
let stage = ChainStage {
operator: input,
bound: visible.to_vec(),
tip: far_node.binding,
};
let step_bound = step_bindings(&stage.bound, pattern, edge, &pattern.start);
let step_predicates = drain_satisfiable(pending, &step_bound);
return self
.plan_step(
plan,
pattern,
stage,
edge,
&pattern.start,
step_predicates,
edge.direction.reverse(),
true,
)
.await;
}
Err(Error::Internal(
"GQL MATCH planning: bound-variable expansion found no shared anchor node \
(multi-hop reverse anchoring is out of PR-B scope)"
.to_string(),
))
}
async fn place_satisfiable(
&self,
op: Arc<dyn ExecOperator>,
bound: &[BindingId],
pending: &mut Vec<MatchPredicate>,
) -> Result<Arc<dyn ExecOperator>, Error> {
let mut ready: Vec<Expr> = Vec::new();
pending.retain(|predicate| {
if deps_subset(&predicate.deps, bound) {
ready.push(predicate.expr.clone());
false
} else {
true
}
});
self.filter_above(op, ready).await
}
async fn place_above(
&self,
op: Arc<dyn ExecOperator>,
bound: &[BindingId],
predicates: Vec<MatchPredicate>,
) -> Result<Arc<dyn ExecOperator>, Error> {
let mut exprs: Vec<Expr> = Vec::with_capacity(predicates.len());
for predicate in predicates {
if !deps_subset(&predicate.deps, bound) {
return Err(Error::Internal(
"GQL MATCH planning: a deferred clause predicate's deps are unsatisfied at \
the clause-combine join"
.to_string(),
));
}
exprs.push(predicate.expr);
}
self.filter_above(op, exprs).await
}
fn wrap_distinct_edges(
&self,
plan: &MatchPlan,
clause: &MatchClausePlan,
op: Arc<dyn ExecOperator>,
) -> Arc<dyn ExecOperator> {
let edges = clause_edge_bindings(plan, clause);
if edges.len() < 2 {
return op;
}
if edge_tables_pairwise_disjoint(&edges) {
return op;
}
let names = edges.iter().map(|(name, _)| name.clone()).collect::<Vec<_>>();
Arc::new(DistinctEdges::new(op, names)) as Arc<dyn ExecOperator>
}
async fn plan_anchor(
&self,
binding: BindingId,
label: TableName,
plan: &MatchPlan,
pushed: Vec<Expr>,
) -> Result<ChainStage, Error> {
let scan = self.plan_table_scan(label, pushed).await?;
let name = binding_name(plan, binding).to_string();
let operator = Arc::new(Bind::new(scan, name)) as Arc<dyn ExecOperator>;
Ok(ChainStage {
operator,
bound: vec![binding],
tip: binding,
})
}
async fn plan_table_scan(
&self,
label: TableName,
pushed: Vec<Expr>,
) -> Result<Arc<dyn ExecOperator>, Error> {
let cond = conjoin(pushed).map(Cond);
let scan_predicate = match cond.as_ref() {
Some(c) => Some(self.physical_expr(c.0.clone()).await?),
None => None,
};
let planned = self
.plan_source(
Expr::Table(label),
None,
cond.as_ref(),
None,
None,
None,
scan_predicate,
None,
None,
false,
&super::select::TopKPushdownRequest::NotApplicable,
)
.await?;
self.apply_source_filter(planned).await
}
#[allow(clippy::too_many_arguments)]
async fn plan_step(
&self,
plan: &MatchPlan,
pattern: &PatternPlan,
prev: ChainStage,
edge: &EdgeStep,
node: &NodeStep,
predicates: Vec<Expr>,
direction: ExpandDirection,
reversed: bool,
) -> Result<ChainStage, Error> {
let source = binding_name(plan, prev.tip).to_string();
let dir = expand_dir(direction);
let edge_tables = edge.label.iter().cloned().collect::<Vec<_>>();
let target_binding = binding_name(plan, node.binding).to_string();
let target_label = node.label.clone();
let mut bound = prev.bound.clone();
bound.push(edge.binding);
bound.push(node.binding);
match edge.quantifier {
None => {
let edge_def = plan.binding(edge.binding);
let edge_name = edge_def.name.clone();
let edge_binding = if edge_def.user_named {
EdgeBinding::Full(edge_name)
} else {
EdgeBinding::IdOnly(edge_name)
};
let predicate = match conjoin(predicates) {
Some(joined) => Some(self.physical_expr(joined).await?),
None => None,
};
let operator = Arc::new(Expand::new(
prev.operator,
source,
dir,
edge_tables,
edge_binding,
target_binding,
target_label,
predicate,
false,
)) as Arc<dyn ExecOperator>;
Ok(ChainStage {
operator,
bound,
tip: node.binding,
})
}
Some(quantifier) => {
let group_binding = Some(binding_name(plan, edge.binding).to_string());
let path_binding = pattern.path_var.map(|id| binding_name(plan, id).to_string());
let (exec_mode, routing) = resolve_path_search(pattern.search);
debug_assert!(
reversed || source == binding_name(plan, pattern.start.binding),
"a forward-anchored quantified step must expand from the pattern's start",
);
let operator: Arc<dyn ExecOperator> = match routing {
SearchRouting::Every => Arc::new(PathExpand::new(
prev.operator,
source,
dir,
edge_tables,
quantifier.min,
quantifier.max,
target_binding,
target_label,
group_binding,
path_binding,
exec_mode,
reversed,
)) as Arc<dyn ExecOperator>,
SearchRouting::Shortest {
selector,
} => Arc::new(ShortestPathExpand::new(
prev.operator,
source,
dir,
edge_tables,
quantifier.min,
quantifier.max,
target_binding,
target_label,
group_binding,
path_binding,
exec_mode,
selector,
reversed,
)) as Arc<dyn ExecOperator>,
};
if let Some(id) = pattern.path_var {
bound.push(id);
}
self.maybe_filter(
ChainStage {
operator,
bound,
tip: node.binding,
},
predicates,
)
.await
}
}
}
async fn maybe_filter(&self, stage: ChainStage, exprs: Vec<Expr>) -> Result<ChainStage, Error> {
let Some(joined) = conjoin(exprs) else {
return Ok(stage);
};
let predicate = self.physical_expr(joined).await?;
let operator = Arc::new(Filter::new(stage.operator, predicate)) as Arc<dyn ExecOperator>;
Ok(ChainStage {
operator,
bound: stage.bound,
tip: stage.tip,
})
}
async fn filter_above(
&self,
op: Arc<dyn ExecOperator>,
exprs: Vec<Expr>,
) -> Result<Arc<dyn ExecOperator>, Error> {
let Some(joined) = conjoin(exprs) else {
return Ok(op);
};
let predicate = self.physical_expr(joined).await?;
Ok(Arc::new(Filter::new(op, predicate)) as Arc<dyn ExecOperator>)
}
async fn plan_match_tail(
&self,
plan: &MatchPlan,
body: Arc<dyn ExecOperator>,
) -> Result<Arc<dyn ExecOperator>, Error> {
let Some(output) = plan.output.as_ref() else {
return Ok(Arc::new(DrainSink::new(body)) as Arc<dyn ExecOperator>);
};
if let Some(group_keys) = output.group_by.as_ref() {
let mut op = self.plan_match_aggregate(body, output, group_keys).await?;
if output.distinct {
op = Arc::new(Distinct::new(op)) as Arc<dyn ExecOperator>;
}
op = self.plan_match_sort(op, output).await?;
op = self.plan_match_limit(op, output).await?;
if output.columns.iter().any(|c| c.hidden) {
op = self.plan_match_drop_hidden(op, output).await?;
}
Ok(op)
} else if output.distinct {
let mut op = self.plan_match_project(body, output).await?;
op = Arc::new(Distinct::new(op)) as Arc<dyn ExecOperator>;
op = self.plan_match_sort(op, output).await?;
op = self.plan_match_limit(op, output).await?;
Ok(op)
} else {
let mut op = self.plan_match_sort(body, output).await?;
op = self.plan_match_limit(op, output).await?;
op = self.plan_match_project(op, output).await?;
Ok(op)
}
}
async fn plan_match_aggregate(
&self,
input: Arc<dyn ExecOperator>,
output: &MatchOutput,
group_keys: &[Expr],
) -> Result<Arc<dyn ExecOperator>, Error> {
use surrealdb_types::ToSql;
let mut group_by_exprs = Vec::with_capacity(group_keys.len());
for key in group_keys {
group_by_exprs.push(self.physical_expr(key.clone()).await?);
}
let group_by_idioms: Vec<Idiom> = group_keys
.iter()
.map(|key| match key {
Expr::Idiom(idiom) => idiom.clone(),
other => Idiom::field(other.to_sql()),
})
.collect();
let mut aggregates = Vec::with_capacity(output.columns.len());
for column in output.columns.iter() {
if let Some(idx) = group_keys.iter().position(|k| *k == column.expr) {
aggregates.push(AggregateField::new(
column.name.clone(),
true,
Some(idx),
None,
None,
));
} else if expr_has_aggregate(self.function_registry(), &column.expr) {
let (info, fallback) = self.extract_aggregate_info(column.expr.clone()).await?;
aggregates.push(AggregateField::new(
column.name.clone(),
false,
None,
info,
fallback,
));
} else {
let fallback = self.physical_expr(column.expr.clone()).await?;
aggregates.push(AggregateField::new(
column.name.clone(),
false,
None,
None,
Some(fallback),
));
}
}
Ok(Arc::new(Aggregate::new(input, group_by_idioms, group_by_exprs, aggregates))
as Arc<dyn ExecOperator>)
}
async fn plan_match_drop_hidden(
&self,
input: Arc<dyn ExecOperator>,
output: &MatchOutput,
) -> Result<Arc<dyn ExecOperator>, Error> {
let mut fields = Vec::new();
for column in output.columns.iter().filter(|c| !c.hidden) {
let expr = self.physical_expr(Expr::Idiom(Idiom::field(column.name.clone()))).await?;
fields.push(FieldSelection::new(&column.name, expr));
}
Ok(Arc::new(Project::new(input, fields, Vec::new(), false)) as Arc<dyn ExecOperator>)
}
async fn plan_match_sort(
&self,
input: Arc<dyn ExecOperator>,
output: &MatchOutput,
) -> Result<Arc<dyn ExecOperator>, Error> {
if output.order.is_empty() {
return Ok(input);
}
let mut order_by = Vec::with_capacity(output.order.len());
for order in output.order.iter() {
let expr = self.physical_expr(order.expr.clone()).await?;
order_by.push(OrderByField {
expr,
direction: if order.ascending {
SortDirection::Asc
} else {
SortDirection::Desc
},
collate: false,
numeric: false,
});
}
Ok(Arc::new(Sort::new(input, order_by)) as Arc<dyn ExecOperator>)
}
async fn plan_match_limit(
&self,
input: Arc<dyn ExecOperator>,
output: &MatchOutput,
) -> Result<Arc<dyn ExecOperator>, Error> {
let limit = match output.limit.as_ref() {
Some(e) => Some(self.physical_expr(e.clone()).await?),
None => None,
};
let offset = match output.skip.as_ref() {
Some(e) => Some(self.physical_expr(e.clone()).await?),
None => None,
};
if limit.is_none() && offset.is_none() {
return Ok(input);
}
Ok(Arc::new(Limit::new(input, limit, offset)) as Arc<dyn ExecOperator>)
}
async fn plan_match_project(
&self,
input: Arc<dyn ExecOperator>,
output: &MatchOutput,
) -> Result<Arc<dyn ExecOperator>, Error> {
let mut fields = Vec::with_capacity(output.columns.len());
for column in output.columns.iter() {
let expr = self.physical_expr(column.expr.clone()).await?;
fields.push(FieldSelection::new(&column.name, expr));
}
Ok(Arc::new(Project::new(input, fields, Vec::new(), false)) as Arc<dyn ExecOperator>)
}
async fn apply_source_filter(
&self,
planned: super::select::PlannedSource,
) -> Result<Arc<dyn ExecOperator>, Error> {
use super::select::FilterAction;
match planned.filter_action {
FilterAction::FullyConsumed => Ok(planned.operator),
FilterAction::UseOriginal => Ok(planned.operator),
FilterAction::Residual(residual) => {
let predicate = self.physical_expr(residual.0).await?;
Ok(Arc::new(Filter::new(planned.operator, predicate)) as Arc<dyn ExecOperator>)
}
}
}
}
fn shared_bindings(a: &[BindingId], b: &[BindingId]) -> Vec<BindingId> {
a.iter().copied().filter(|id| b.contains(id)).collect()
}
fn drain_satisfiable(pending: &mut Vec<MatchPredicate>, bound: &[BindingId]) -> Vec<Expr> {
let mut ready = Vec::new();
pending.retain(|predicate| {
if deps_subset(&predicate.deps, bound) {
ready.push(predicate.expr.clone());
false
} else {
true
}
});
ready
}
fn step_bindings(
prior: &[BindingId],
pattern: &PatternPlan,
edge: &EdgeStep,
node: &NodeStep,
) -> Vec<BindingId> {
let mut bound = prior.to_vec();
bound.push(edge.binding);
bound.push(node.binding);
if edge.quantifier.is_some()
&& let Some(id) = pattern.path_var
{
bound.push(id);
}
bound
}
fn expr_has_aggregate(registry: &crate::exec::function::FunctionRegistry, expr: &Expr) -> bool {
let mut stack = vec![expr];
while let Some(e) = stack.pop() {
match e {
Expr::FunctionCall(call) => {
if let Function::Normal(name) = &call.receiver
&& registry.get_aggregate(name.as_str()).is_some()
{
return true;
}
stack.extend(call.arguments.iter());
}
Expr::Binary {
left,
right,
..
} => {
stack.push(left);
stack.push(right);
}
Expr::Prefix {
expr,
..
}
| Expr::Postfix {
expr,
..
} => stack.push(expr),
Expr::Literal(Literal::Array(items)) => stack.extend(items.iter()),
Expr::Literal(Literal::Object(entries)) => {
stack.extend(entries.iter().map(|entry| &entry.value));
}
_ => {}
}
}
false
}
fn union_bindings(a: &[BindingId], b: &[BindingId]) -> Vec<BindingId> {
let mut out = a.to_vec();
for id in b {
if !out.contains(id) {
out.push(*id);
}
}
out
}
fn difference_bindings(a: &[BindingId], b: &[BindingId]) -> Vec<BindingId> {
a.iter().copied().filter(|id| !b.contains(id)).collect()
}
enum StageUnit<'p> {
Mandatory(&'p MatchClausePlan),
Optional(Vec<&'p MatchClausePlan>),
Mutation(&'p MutationStage),
}
fn stage_units(plan: &MatchPlan) -> Vec<StageUnit<'_>> {
let mut units: Vec<StageUnit<'_>> = Vec::new();
let mut current_group: Option<u32> = None;
for stage in plan.stages.iter() {
match stage {
MatchStage::Read(clause) => match clause.optional_group {
Some(group) => {
if current_group == Some(group)
&& let Some(StageUnit::Optional(block)) = units.last_mut()
{
block.push(clause);
} else {
units.push(StageUnit::Optional(vec![clause]));
current_group = Some(group);
}
}
None => {
units.push(StageUnit::Mandatory(clause));
current_group = None;
}
},
MatchStage::Mutate(mutation) => {
units.push(StageUnit::Mutation(mutation));
current_group = None;
}
}
}
units
}
fn shared_node_anchor(
_plan: &MatchPlan,
pattern: &PatternPlan,
visible: &[BindingId],
) -> Option<BindingId> {
if visible.contains(&pattern.start.binding) {
return Some(pattern.start.binding);
}
if let [(_, far)] = pattern.steps.as_slice()
&& visible.contains(&far.binding)
{
return Some(far.binding);
}
None
}
fn clause_edge_bindings(plan: &MatchPlan, clause: &MatchClausePlan) -> Vec<(String, EdgeTables)> {
let mut edges = Vec::new();
for pattern in clause.patterns.iter() {
for (edge, _) in pattern.steps.iter() {
let def = plan.binding(edge.binding);
if matches!(def.kind, BindingKind::Edge | BindingKind::EdgeGroup) {
let tables = match edge.label.as_ref() {
Some(label) => EdgeTables::Known(label.clone()),
None => EdgeTables::Any,
};
edges.push((def.name.clone(), tables));
}
}
}
edges
}
#[derive(Clone)]
enum EdgeTables {
Known(TableName),
Any,
}
impl EdgeTables {
fn disjoint_from(&self, other: &EdgeTables) -> bool {
match (self, other) {
(EdgeTables::Known(a), EdgeTables::Known(b)) => a != b,
_ => false,
}
}
}
fn edge_tables_pairwise_disjoint(edges: &[(String, EdgeTables)]) -> bool {
for i in 0..edges.len() {
for j in (i + 1)..edges.len() {
if !edges[i].1.disjoint_from(&edges[j].1) {
return false;
}
}
}
true
}
fn binding_names(plan: &MatchPlan, ids: &[BindingId]) -> Vec<String> {
ids.iter().map(|id| binding_name(plan, *id).to_string()).collect()
}
fn binding_name(plan: &MatchPlan, id: BindingId) -> &str {
plan.binding(id).name.as_str()
}
fn expand_dir(direction: ExpandDirection) -> ExpandDir {
match direction {
ExpandDirection::Out => ExpandDir::Out,
ExpandDirection::In => ExpandDir::In,
}
}
fn deps_subset(deps: &[BindingId], bound: &[BindingId]) -> bool {
deps.iter().all(|d| bound.contains(d))
}
fn conjoin(mut exprs: Vec<Expr>) -> Option<Expr> {
let mut acc = exprs.pop()?;
while let Some(next) = exprs.pop() {
acc = Expr::Binary {
left: Box::new(next),
op: crate::expr::BinaryOperator::And,
right: Box::new(acc),
};
}
Some(acc)
}
fn prefix_strip(expr: &Expr, binding: &str) -> Option<Expr> {
match expr {
Expr::Idiom(idiom) => {
let parts = idiom.0.as_slice();
match parts {
[Part::Field(field)] if field.as_str() == binding => None,
[Part::Field(field), rest @ ..] if field.as_str() == binding => {
Some(Expr::Idiom(Idiom(rest.to_vec())))
}
_ => Some(expr.clone()),
}
}
Expr::Binary {
left,
op,
right,
} => Some(Expr::Binary {
left: Box::new(prefix_strip(left, binding)?),
op: op.clone(),
right: Box::new(prefix_strip(right, binding)?),
}),
Expr::Prefix {
op,
expr,
} => Some(Expr::Prefix {
op: op.clone(),
expr: Box::new(prefix_strip(expr, binding)?),
}),
Expr::Postfix {
op,
expr,
} => Some(Expr::Postfix {
op: op.clone(),
expr: Box::new(prefix_strip(expr, binding)?),
}),
other => Some(other.clone()),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ctx::Context;
use crate::expr::match_plan::{
BindingDef, BindingKind, EdgeQuantifier, MatchColumn, MatchOrder, MatchOutput,
};
use crate::expr::{BinaryOperator, Literal};
use crate::kvs::{Datastore, LockType, TransactionType};
use crate::val::TableName;
fn stages_of(clauses: Vec<MatchClausePlan>) -> Vec<MatchStage> {
clauses.into_iter().map(MatchStage::Read).collect()
}
fn clause_mut(plan: &mut MatchPlan, i: usize) -> &mut MatchClausePlan {
match &mut plan.stages[i] {
MatchStage::Read(c) => c,
MatchStage::Mutate(_) => panic!("stage {i} is not a read clause"),
}
}
fn binding(name: &str, kind: BindingKind, user_named: bool) -> BindingDef {
BindingDef {
name: name.to_string(),
kind,
user_named,
}
}
fn node(name: &str) -> BindingDef {
binding(name, BindingKind::Node, true)
}
fn edge(name: &str) -> BindingDef {
binding(name, BindingKind::Edge, true)
}
fn hidden_edge(name: &str) -> BindingDef {
binding(name, BindingKind::Edge, false)
}
fn field_path(b: &str, field: &str) -> Expr {
Expr::Idiom(Idiom(vec![
Part::Field(b.to_string().into()),
Part::Field(field.to_string().into()),
]))
}
fn var(b: &str) -> Expr {
Expr::Idiom(Idiom(vec![Part::Field(b.to_string().into())]))
}
fn col(name: &str, expr: Expr) -> MatchColumn {
MatchColumn {
name: name.to_string(),
expr,
hidden: false,
}
}
fn tbl(name: &str) -> TableName {
TableName::new(name.to_string())
}
fn nodestep(b: BindingId, label: Option<&str>) -> NodeStep {
NodeStep {
binding: b,
label: label.map(tbl),
}
}
fn edgestep(b: BindingId, label: Option<&str>, dir: ExpandDirection) -> EdgeStep {
EdgeStep {
binding: b,
label: label.map(tbl),
direction: dir,
quantifier: None,
}
}
fn out_columns(cols: Vec<MatchColumn>) -> Option<MatchOutput> {
Some(MatchOutput {
columns: cols,
distinct: false,
group_by: None,
order: Vec::new(),
skip: None,
limit: None,
})
}
fn plan_tree_i() -> MatchPlan {
MatchPlan {
bindings: vec![node("a"), edge("k"), node("b")],
stages: stages_of(vec![MatchClausePlan {
optional_group: None,
patterns: vec![PatternPlan {
path_var: None,
search: None,
start: nodestep(0, Some("person")),
steps: vec![(
edgestep(1, Some("knows"), ExpandDirection::Out),
nodestep(2, Some("person")),
)],
}],
predicates: vec![MatchPredicate {
expr: Expr::Binary {
left: Box::new(field_path("k", "since")),
op: BinaryOperator::MoreThan,
right: Box::new(Expr::Literal(Literal::Integer(2020))),
},
deps: vec![1],
}],
}]),
output: out_columns(vec![
col("a_name", field_path("a", "name")),
col("b_name", field_path("b", "name")),
]),
}
}
fn plan_tree_iv() -> MatchPlan {
MatchPlan {
bindings: vec![
node("a"),
binding("__e0", BindingKind::EdgeGroup, false),
node("b"),
node("p"),
],
stages: stages_of(vec![MatchClausePlan {
optional_group: None,
patterns: vec![PatternPlan {
path_var: Some(3),
search: None,
start: nodestep(0, Some("person")),
steps: vec![(
EdgeStep {
binding: 1,
label: Some(tbl("knows")),
direction: ExpandDirection::Out,
quantifier: Some(EdgeQuantifier {
min: 1,
max: Some(3),
}),
},
nodestep(2, Some("person")),
)],
}],
predicates: Vec::new(),
}]),
output: Some(MatchOutput {
columns: vec![col("p", var("p")), col("b", var("b"))],
distinct: false,
group_by: None,
order: vec![MatchOrder {
expr: field_path("a", "age"),
ascending: true,
}],
skip: None,
limit: None,
}),
}
}
fn plan_tree_ii() -> MatchPlan {
MatchPlan {
bindings: vec![
node("a"),
hidden_edge("__e0"),
node("b"),
node("c"),
hidden_edge("__e1"),
],
stages: stages_of(vec![MatchClausePlan {
optional_group: None,
patterns: vec![
PatternPlan {
path_var: None,
search: None,
start: nodestep(0, None),
steps: vec![(
edgestep(1, Some("x"), ExpandDirection::Out),
nodestep(2, None),
)],
},
PatternPlan {
path_var: None,
search: None,
start: nodestep(3, None),
steps: vec![(
edgestep(4, Some("y"), ExpandDirection::Out),
nodestep(2, None),
)],
},
],
predicates: Vec::new(),
}]),
output: out_columns(vec![col("a", var("a")), col("c", var("c"))]),
}
}
fn plan_shared_node_same_edge() -> MatchPlan {
MatchPlan {
bindings: vec![node("a"), edge("k"), node("b"), node("c"), edge("k2")],
stages: stages_of(vec![MatchClausePlan {
optional_group: None,
patterns: vec![
PatternPlan {
path_var: None,
search: None,
start: nodestep(0, Some("person")),
steps: vec![(
edgestep(1, Some("knows"), ExpandDirection::Out),
nodestep(2, Some("person")),
)],
},
PatternPlan {
path_var: None,
search: None,
start: nodestep(3, Some("person")),
steps: vec![(
edgestep(4, Some("knows"), ExpandDirection::Out),
nodestep(2, None),
)],
},
],
predicates: Vec::new(),
}]),
output: out_columns(vec![col("a", var("a")), col("c", var("c"))]),
}
}
fn plan_anon_group_second_edge() -> MatchPlan {
MatchPlan {
bindings: vec![
node("a"),
binding("__e0", BindingKind::EdgeGroup, false),
node("b"),
edge("k2"),
node("c"),
],
stages: stages_of(vec![MatchClausePlan {
optional_group: None,
patterns: vec![
PatternPlan {
path_var: None,
search: None,
start: nodestep(0, Some("person")),
steps: vec![(
EdgeStep {
binding: 1,
label: Some(tbl("knows")),
direction: ExpandDirection::Out,
quantifier: Some(EdgeQuantifier {
min: 1,
max: Some(1),
}),
},
nodestep(2, Some("person")),
)],
},
PatternPlan {
path_var: None,
search: None,
start: nodestep(0, None),
steps: vec![(
edgestep(3, Some("knows"), ExpandDirection::Out),
nodestep(4, Some("person")),
)],
},
],
predicates: Vec::new(),
}]),
output: out_columns(vec![col("b", var("b")), col("c", var("c"))]),
}
}
fn plan_sequential_clauses() -> MatchPlan {
MatchPlan {
bindings: vec![node("a"), edge("k"), node("b"), edge("k2"), node("c")],
stages: stages_of(vec![
MatchClausePlan {
optional_group: None,
patterns: vec![PatternPlan {
path_var: None,
search: None,
start: nodestep(0, Some("person")),
steps: vec![(
edgestep(1, Some("knows"), ExpandDirection::Out),
nodestep(2, Some("person")),
)],
}],
predicates: Vec::new(),
},
MatchClausePlan {
optional_group: None,
patterns: vec![PatternPlan {
path_var: None,
search: None,
start: nodestep(2, Some("person")),
steps: vec![(
edgestep(3, Some("likes"), ExpandDirection::Out),
nodestep(4, Some("person")),
)],
}],
predicates: Vec::new(),
},
]),
output: out_columns(vec![col("a", var("a")), col("c", var("c"))]),
}
}
fn render_plan(plan: &dyn ExecOperator, out: &mut String, prefix: &str) {
use std::fmt::Write;
let _ = write!(out, "{} [ctx: {}]", plan.name(), plan.required_context().short_name());
let attrs = plan.attrs();
if !attrs.is_empty() {
let _ = write!(out, " [");
for (i, (k, v)) in attrs.iter().enumerate() {
if i > 0 {
let _ = write!(out, ", ");
}
let _ = write!(out, "{k}: {v}");
}
let _ = write!(out, "]");
}
let _ = writeln!(out);
let children = plan.children();
if !children.is_empty() {
let child_prefix = format!("{prefix} ");
for child in children.iter() {
let _ = write!(out, "{child_prefix}");
render_plan(child.as_ref(), out, &child_prefix);
}
}
}
async fn explain(plan: MatchPlan) -> String {
let ds = Datastore::new("memory").await.expect("datastore");
let session = crate::dbs::Session::owner().with_ns("test").with_db("test");
ds.execute(
"DEFINE NAMESPACE test; DEFINE DATABASE test; \
DEFINE TABLE person; DEFINE TABLE knows; DEFINE TABLE likes; \
DEFINE TABLE x; DEFINE TABLE y;",
&session,
None,
)
.await
.expect("define schema");
let base = ds.setup_ctx().expect("setup_ctx").freeze();
let txn = Arc::new(
ds.transaction(TransactionType::Read, LockType::Optimistic).await.expect("txn"),
);
let mut ctx = Context::new_child(&base);
ctx.set_transaction(Arc::clone(&txn));
let ctx = ctx.freeze();
let planner =
Planner::with_txn(&ctx, txn, Some("test".to_string()), Some("test".to_string()));
let operator = planner.plan_match(plan).await.expect("plan_match");
let mut out = String::new();
render_plan(operator.as_ref(), &mut out, "");
out
}
#[tokio::test]
async fn explain_tree_i_single_pattern_edge_predicate() {
let rendered = explain(plan_tree_i()).await;
let expected = "\
Project [ctx: Db]
Expand [ctx: Db] [source: a, direction: ->, tables: knows, edge: k, target: b, target_label: person, predicate: k.since > 2020]
Bind [ctx: Db] [binding: a]
TableScan [ctx: Db] [table: person, direction: Forward]
";
assert_eq!(rendered, expected, "\n--- got ---\n{rendered}");
}
#[tokio::test]
async fn explain_tree_iv_quantified_path_order() {
let rendered = explain(plan_tree_iv()).await;
let expected = "\
Project [ctx: Db]
Sort [ctx: Db] [order_by: a.age ASC]
PathExpand [ctx: Db] [source: a, direction: ->, tables: knows, min: 1, max: 3, target_binding: b, target_label: person, group: __e0, path: p]
Bind [ctx: Db] [binding: a]
TableScan [ctx: Db] [table: person, direction: Forward]
";
assert_eq!(rendered, expected, "\n--- got ---\n{rendered}");
}
#[tokio::test]
async fn explain_shortest_routes_to_shortest_path_expand() {
let mut plan = plan_tree_iv();
clause_mut(&mut plan, 0).patterns[0].search = Some(PathPrefixPlan {
search: PathSearch::AllShortest,
mode: IrPathMode::Acyclic,
});
let rendered = explain(plan).await;
assert!(rendered.contains("ShortestPathExpand"), "\n--- got ---\n{rendered}");
assert!(rendered.contains("search: all shortest"), "\n--- got ---\n{rendered}");
assert!(rendered.contains("mode: acyclic"), "\n--- got ---\n{rendered}");
}
#[tokio::test]
async fn explain_any_routes_to_shortest_path_expand() {
let mut plan = plan_tree_iv();
clause_mut(&mut plan, 0).patterns[0].search = Some(PathPrefixPlan {
search: PathSearch::Any {
count: 3,
},
mode: IrPathMode::Walk,
});
let rendered = explain(plan).await;
assert!(rendered.contains("ShortestPathExpand"), "\n--- got ---\n{rendered}");
assert!(rendered.contains("search: any 3"), "\n--- got ---\n{rendered}");
assert!(!rendered.contains("mode:"), "\n--- got ---\n{rendered}");
}
#[tokio::test]
async fn explain_tree_ii_edge_anchored_hashjoin_distinctedges_elided() {
let rendered = explain(plan_tree_ii()).await;
let expected = "\
Project [ctx: Db]
HashJoin [ctx: Db] [type: Inner, keys: b]
EndpointBind [ctx: Db] [edge: __e0, field: out, node: b]
EndpointBind [ctx: Db] [edge: __e0, field: in, node: a]
Bind [ctx: Db] [binding: __e0]
TableScan [ctx: Db] [table: x, direction: Forward]
EndpointBind [ctx: Db] [edge: __e1, field: out, node: b]
EndpointBind [ctx: Db] [edge: __e1, field: in, node: c]
Bind [ctx: Db] [binding: __e1]
TableScan [ctx: Db] [table: y, direction: Forward]
";
assert_eq!(rendered, expected, "\n--- got ---\n{rendered}");
}
#[tokio::test]
async fn explain_shared_node_same_edge_keeps_distinctedges() {
let rendered = explain(plan_shared_node_same_edge()).await;
let expected = "\
Project [ctx: Db]
DistinctEdges [ctx: Db] [edges: k, k2]
HashJoin [ctx: Db] [type: Inner, keys: b]
Expand [ctx: Db] [source: a, direction: ->, tables: knows, edge: k, target: b, target_label: person]
Bind [ctx: Db] [binding: a]
TableScan [ctx: Db] [table: person, direction: Forward]
Expand [ctx: Db] [source: c, direction: ->, tables: knows, edge: k2, target: b]
Bind [ctx: Db] [binding: c]
TableScan [ctx: Db] [table: person, direction: Forward]
";
assert_eq!(rendered, expected, "\n--- got ---\n{rendered}");
}
#[tokio::test]
async fn explain_sequential_clause_hashjoin_on_shared_node() {
let rendered = explain(plan_sequential_clauses()).await;
let expected = "\
Project [ctx: Db]
HashJoin [ctx: Db] [type: Inner, keys: b]
Expand [ctx: Db] [source: a, direction: ->, tables: knows, edge: k, target: b, target_label: person]
Bind [ctx: Db] [binding: a]
TableScan [ctx: Db] [table: person, direction: Forward]
Expand [ctx: Db] [source: b, direction: ->, tables: likes, edge: k2, target: c, target_label: person]
Bind [ctx: Db] [binding: b]
TableScan [ctx: Db] [table: person, direction: Forward]
";
assert_eq!(rendered, expected, "\n--- got ---\n{rendered}");
}
#[tokio::test]
async fn explain_anon_group_second_edge_binds_group_and_keeps_distinctedges() {
let rendered = explain(plan_anon_group_second_edge()).await;
let expected = "\
Project [ctx: Db]
DistinctEdges [ctx: Db] [edges: __e0, k2]
HashJoin [ctx: Db] [type: Inner, keys: a]
PathExpand [ctx: Db] [source: a, direction: ->, tables: knows, min: 1, max: 1, target_binding: b, target_label: person, group: __e0]
Bind [ctx: Db] [binding: a]
TableScan [ctx: Db] [table: person, direction: Forward]
EndpointBind [ctx: Db] [edge: k2, field: out, node: c, target_label: person]
EndpointBind [ctx: Db] [edge: k2, field: in, node: a]
Bind [ctx: Db] [binding: k2]
TableScan [ctx: Db] [table: knows, direction: Forward]
";
assert_eq!(rendered, expected, "\n--- got ---\n{rendered}");
}
fn plan_optional_fast_path() -> MatchPlan {
MatchPlan {
bindings: vec![node("a"), edge("k"), node("b")],
stages: stages_of(vec![
MatchClausePlan {
optional_group: None,
patterns: vec![PatternPlan {
path_var: None,
search: None,
start: nodestep(0, Some("person")),
steps: Vec::new(),
}],
predicates: Vec::new(),
},
MatchClausePlan {
optional_group: Some(0),
patterns: vec![PatternPlan {
path_var: None,
search: None,
start: nodestep(0, None),
steps: vec![(
edgestep(1, Some("knows"), ExpandDirection::Out),
nodestep(2, None),
)],
}],
predicates: Vec::new(),
},
]),
output: out_columns(vec![
col("a_name", field_path("a", "name")),
col("b_name", field_path("b", "name")),
]),
}
}
fn plan_optional_block_unit() -> MatchPlan {
MatchPlan {
bindings: vec![node("a"), hidden_edge("__e0"), node("b"), edge("k2"), node("c")],
stages: stages_of(vec![
MatchClausePlan {
optional_group: None,
patterns: vec![PatternPlan {
path_var: None,
search: None,
start: nodestep(0, Some("person")),
steps: Vec::new(),
}],
predicates: Vec::new(),
},
MatchClausePlan {
optional_group: Some(0),
patterns: vec![PatternPlan {
path_var: None,
search: None,
start: nodestep(0, None),
steps: vec![(
edgestep(1, Some("knows"), ExpandDirection::Out),
nodestep(2, Some("person")),
)],
}],
predicates: Vec::new(),
},
MatchClausePlan {
optional_group: Some(0),
patterns: vec![PatternPlan {
path_var: None,
search: None,
start: nodestep(2, None),
steps: vec![(
edgestep(3, None, ExpandDirection::Out),
nodestep(4, Some("person")),
)],
}],
predicates: Vec::new(),
},
]),
output: out_columns(vec![col("a", var("a"))]),
}
}
fn plan_chained_optionals() -> MatchPlan {
MatchPlan {
bindings: vec![
node("a"),
hidden_edge("__e0"),
node("b"),
hidden_edge("__e1"),
node("c"),
],
stages: stages_of(vec![
MatchClausePlan {
optional_group: None,
patterns: vec![PatternPlan {
path_var: None,
search: None,
start: nodestep(0, Some("person")),
steps: Vec::new(),
}],
predicates: Vec::new(),
},
MatchClausePlan {
optional_group: Some(0),
patterns: vec![PatternPlan {
path_var: None,
search: None,
start: nodestep(0, None),
steps: vec![(
edgestep(1, Some("knows"), ExpandDirection::Out),
nodestep(2, Some("person")),
)],
}],
predicates: Vec::new(),
},
MatchClausePlan {
optional_group: Some(1),
patterns: vec![PatternPlan {
path_var: None,
search: None,
start: nodestep(0, None),
steps: vec![(
edgestep(3, Some("likes"), ExpandDirection::Out),
nodestep(4, Some("person")),
)],
}],
predicates: Vec::new(),
},
]),
output: out_columns(vec![col("a", var("a"))]),
}
}
fn plan_optional_multipattern() -> MatchPlan {
MatchPlan {
bindings: vec![
node("a"),
node("z"),
hidden_edge("__e0"),
node("b"),
hidden_edge("__e1"),
node("c"),
],
stages: stages_of(vec![
MatchClausePlan {
optional_group: None,
patterns: vec![
PatternPlan {
path_var: None,
search: None,
start: nodestep(0, Some("person")),
steps: Vec::new(),
},
PatternPlan {
path_var: None,
search: None,
start: nodestep(1, Some("person")),
steps: Vec::new(),
},
],
predicates: Vec::new(),
},
MatchClausePlan {
optional_group: Some(0),
patterns: vec![
PatternPlan {
path_var: None,
search: None,
start: nodestep(0, None),
steps: vec![(
edgestep(2, Some("knows"), ExpandDirection::Out),
nodestep(3, Some("person")),
)],
},
PatternPlan {
path_var: None,
search: None,
start: nodestep(1, None),
steps: vec![(
edgestep(4, Some("likes"), ExpandDirection::Out),
nodestep(5, Some("person")),
)],
},
],
predicates: Vec::new(),
},
]),
output: out_columns(vec![col("a", var("a"))]),
}
}
#[tokio::test]
async fn explain_optional_fast_path_is_optionalexpand() {
let rendered = explain(plan_optional_fast_path()).await;
let expected = "\
Project [ctx: Db]
OptionalExpand [ctx: Db] [source: a, direction: ->, tables: knows, edge: k, target: b]
Bind [ctx: Db] [binding: a]
TableScan [ctx: Db] [table: person, direction: Forward]
";
assert_eq!(rendered, expected, "\n--- got ---\n{rendered}");
}
#[tokio::test]
async fn explain_optional_block_is_one_leftjoin_unit() {
let rendered = explain(plan_optional_block_unit()).await;
let expected = "\
Project [ctx: Db]
HashJoin [ctx: Db] [type: Left, keys: a, null_template: __e0, b, k2, c]
Expand [ctx: Db] [source: b, direction: ->, tables: *, edge: k2, target: c, target_label: person]
EndpointBind [ctx: Db] [edge: __e0, field: out, node: b, target_label: person]
EndpointBind [ctx: Db] [edge: __e0, field: in, node: a]
Bind [ctx: Db] [binding: __e0]
TableScan [ctx: Db] [table: knows, direction: Forward]
Bind [ctx: Db] [binding: a]
TableScan [ctx: Db] [table: person, direction: Forward]
";
assert_eq!(rendered, expected, "\n--- got ---\n{rendered}");
}
#[tokio::test]
async fn explain_chained_optionals_stack_optionalexpands() {
let rendered = explain(plan_chained_optionals()).await;
let expected = "\
Project [ctx: Db]
OptionalExpand [ctx: Db] [source: a, direction: ->, tables: likes, edge: __e1, target: c, target_label: person]
OptionalExpand [ctx: Db] [source: a, direction: ->, tables: knows, edge: __e0, target: b, target_label: person]
Bind [ctx: Db] [binding: a]
TableScan [ctx: Db] [table: person, direction: Forward]
";
assert_eq!(rendered, expected, "\n--- got ---\n{rendered}");
}
#[tokio::test]
async fn explain_optional_multipattern_is_leftjoin_over_cross() {
let rendered = explain(plan_optional_multipattern()).await;
let expected = "\
Project [ctx: Db]
HashJoin [ctx: Db] [type: Left, keys: a, z, null_template: __e0, b, __e1, c]
HashJoin [ctx: Db] [type: Cross]
EndpointBind [ctx: Db] [edge: __e0, field: out, node: b, target_label: person]
EndpointBind [ctx: Db] [edge: __e0, field: in, node: a]
Bind [ctx: Db] [binding: __e0]
TableScan [ctx: Db] [table: knows, direction: Forward]
EndpointBind [ctx: Db] [edge: __e1, field: out, node: c, target_label: person]
EndpointBind [ctx: Db] [edge: __e1, field: in, node: z]
Bind [ctx: Db] [binding: __e1]
TableScan [ctx: Db] [table: likes, direction: Forward]
HashJoin [ctx: Db] [type: Cross]
Bind [ctx: Db] [binding: a]
TableScan [ctx: Db] [table: person, direction: Forward]
Bind [ctx: Db] [binding: z]
TableScan [ctx: Db] [table: person, direction: Forward]
";
assert_eq!(rendered, expected, "\n--- got ---\n{rendered}");
}
#[test]
fn shared_bindings_intersects_in_a_order() {
assert_eq!(shared_bindings(&[0, 2, 4], &[4, 2, 9]), vec![2, 4]);
assert!(shared_bindings(&[0, 1], &[2, 3]).is_empty());
}
#[test]
fn union_bindings_appends_new_ids() {
assert_eq!(union_bindings(&[0, 1], &[1, 2, 3]), vec![0, 1, 2, 3]);
}
#[test]
fn edge_tables_disjoint_only_for_distinct_known_tables() {
assert!(EdgeTables::Known(tbl("x")).disjoint_from(&EdgeTables::Known(tbl("y"))));
assert!(!EdgeTables::Known(tbl("x")).disjoint_from(&EdgeTables::Known(tbl("x"))));
assert!(!EdgeTables::Any.disjoint_from(&EdgeTables::Known(tbl("x"))));
assert!(!EdgeTables::Known(tbl("x")).disjoint_from(&EdgeTables::Any));
assert!(!EdgeTables::Any.disjoint_from(&EdgeTables::Any));
}
#[test]
fn pairwise_disjoint_skips_distinctedges_for_distinct_tables() {
let edges = vec![
("e0".to_string(), EdgeTables::Known(tbl("x"))),
("e1".to_string(), EdgeTables::Known(tbl("y"))),
];
assert!(edge_tables_pairwise_disjoint(&edges));
let same = vec![
("e0".to_string(), EdgeTables::Known(tbl("knows"))),
("e1".to_string(), EdgeTables::Known(tbl("knows"))),
];
assert!(!edge_tables_pairwise_disjoint(&same));
let with_any = vec![
("e0".to_string(), EdgeTables::Known(tbl("x"))),
("e1".to_string(), EdgeTables::Any),
];
assert!(!edge_tables_pairwise_disjoint(&with_any));
}
#[test]
fn shared_node_anchor_picks_start_then_far() {
let pattern = PatternPlan {
path_var: None,
search: None,
start: nodestep(0, None),
steps: vec![(edgestep(1, Some("e"), ExpandDirection::Out), nodestep(2, None))],
};
let plan = MatchPlan {
bindings: vec![node("a"), hidden_edge("e"), node("b")],
stages: Vec::new(),
output: out_columns(Vec::new()),
};
assert_eq!(shared_node_anchor(&plan, &pattern, &[0]), Some(0));
assert_eq!(shared_node_anchor(&plan, &pattern, &[2]), Some(2));
assert_eq!(shared_node_anchor(&plan, &pattern, &[9]), None);
}
#[test]
fn prefix_strip_rewrites_binding_field_path() {
let expr = Expr::Binary {
left: Box::new(field_path("a", "since")),
op: BinaryOperator::MoreThan,
right: Box::new(Expr::Literal(Literal::Integer(2020))),
};
let stripped = prefix_strip(&expr, "a").expect("strippable");
let expected = Expr::Binary {
left: Box::new(Expr::Idiom(Idiom(vec![Part::Field("since".to_string().into())]))),
op: BinaryOperator::MoreThan,
right: Box::new(Expr::Literal(Literal::Integer(2020))),
};
assert_eq!(stripped, expected);
}
#[test]
fn prefix_strip_bails_on_whole_record_reference() {
let expr = var("a");
assert!(prefix_strip(&expr, "a").is_none());
}
#[test]
fn deps_subset_checks_membership() {
assert!(deps_subset(&[1, 2], &[0, 1, 2, 3]));
assert!(!deps_subset(&[1, 4], &[0, 1, 2, 3]));
assert!(deps_subset(&[], &[0]));
}
#[test]
fn conjoin_builds_left_leaning_and_chain() {
assert!(conjoin(Vec::new()).is_none());
let one = conjoin(vec![Expr::Literal(Literal::Bool(true))]).expect("one");
assert_eq!(one, Expr::Literal(Literal::Bool(true)));
}
#[test]
fn drain_satisfiable_drains_only_ready_predicates() {
let mut pending = vec![
MatchPredicate {
expr: var("a"),
deps: vec![0],
},
MatchPredicate {
expr: var("b"),
deps: vec![5],
},
];
let ready = drain_satisfiable(&mut pending, &[0, 1]);
assert_eq!(ready, vec![var("a")]);
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].deps, vec![5]);
}
}