use smol_str::SmolStr;
use crate::{
engine::{volcano::builder::PhysicalPlanBuilder, GraphCtx},
gremlin::{
type_bridge,
type_bridge::{push_has_step, value_to_primitive},
value::{Predicate, Value},
},
planner::{
apply_rules,
logical_step::{
AddEStep, AddVStep, AndStep, AsStep, BothEStep, BothStep, ChooseStep, CoalesceStep, ConstantStep,
CountStep, CyclicPathStep, DedupStep, DropStep, EStep, EmitSpec, FoldStep, FromStep, GroupCountStep,
GroupStep, HasIdStep, HasLabelStep, HasRankStep, IdStep, IdentityStep, InEStep, InStep, InVStep, LabelStep,
LimitStep, LocalStep, LogicalPlan, LogicalStep, MaxStep, MeanStep, MinStep, NotStep, OrStep, Order,
OrderKey, OrderKeySpec, OrderStep, OtherVStep, OutEStep, OutStep, OutVStep, PathStep, PropertiesStep,
PropertyStep, RangeStep, RankStep, RepeatStep, ScalarFilterStep, SelectStep, SimplePathStep, SkipStep,
SumStep, TailStep, ToStep, UnfoldStep, UnionStep, ValuesStep, WhereStep,
},
},
types::{prop_key::LABEL, StoreError},
};
pub(crate) mod built;
pub use built::BuiltTraversal;
#[derive(Clone)]
pub(crate) struct RepeatBuilder {
body: LogicalPlan,
until: Option<LogicalPlan>,
times: Option<i64>,
emit: EmitSpec,
}
fn follows_group_step(plan: &LogicalPlan) -> bool {
matches!(plan.steps.last(), Some(LogicalStep::Group(_)) | Some(LogicalStep::GroupCount(_)))
}
fn by_after_group_error(caller: &str) -> StoreError {
StoreError::TraversalError(format!(
"{caller} is not supported immediately after group()/group_count() — they have no by() \
modulator yet; see docs/design_group_step.md"
))
}
#[allow(private_interfaces)]
pub trait PlanAppender: Sized {
fn plan_mut(&mut self) -> &mut LogicalPlan;
fn record_error(&mut self, err: StoreError);
fn pending_repeat_mut(&mut self) -> &mut Option<RepeatBuilder>;
fn flush_pending_repeat(&mut self) {
if let Some(rb) = self.pending_repeat_mut().take() {
if rb.until.is_none() && rb.times.is_none() {
self.record_error(StoreError::TraversalError(
"repeat() must have at least one stop condition: .times(n) or .until(cond).".to_string(),
));
return;
}
self.plan_mut().steps.push(LogicalStep::Repeat(RepeatStep {
body: rb.body,
until: rb.until,
times: rb.times,
emit: rb.emit,
}));
}
}
fn push_step(&mut self, step: LogicalStep) {
self.flush_pending_repeat();
self.plan_mut().steps.push(step);
}
}
pub struct GraphTraversal {
plan: LogicalPlan,
pub(crate) error: Option<StoreError>,
pending_repeat: Option<RepeatBuilder>,
}
impl Clone for GraphTraversal {
fn clone(&self) -> Self {
Self { plan: self.plan.clone(), error: None, pending_repeat: None }
}
}
pub fn __() -> GraphTraversal {
GraphTraversal { plan: LogicalPlan { steps: vec![] }, error: None, pending_repeat: None }
}
#[allow(non_snake_case)]
impl GraphTraversal {
pub(crate) fn build(
self,
graph: &mut dyn GraphCtx,
prop_keys: Option<Vec<SmolStr>>,
) -> Result<BuiltTraversal<'_>, StoreError> {
if let Some(err) = self.error {
return Err(err);
}
let mut logical = self.plan;
if self.pending_repeat.is_some() {
return Err(StoreError::TraversalError(
"repeat() requires at least one stop condition — call .times(n) or .until(cond).".to_string(),
));
}
apply_rules(&mut logical)?;
let schema_lock = graph.schema();
let plan = PhysicalPlanBuilder::default().build(&logical, &schema_lock)?;
let schema = graph.schema();
let cache = built::SchemaCache::from_schema(&schema.read().unwrap());
Ok(BuiltTraversal { graph, plan, cache, prop_keys })
}
pub(crate) fn into_plan(self) -> LogicalPlan {
if self.pending_repeat.is_some() {
self.plan
} else {
self.plan
}
}
pub fn addV(mut self, label: impl Into<SmolStr>) -> Self {
self.push_step(LogicalStep::AddV(AddVStep {
label: label.into(),
vertex_id: None,
properties: smallvec::smallvec![],
}));
self
}
pub fn addE(mut self, label: impl Into<SmolStr>) -> Self {
self.push_step(LogicalStep::AddE(AddEStep {
label: label.into(),
out_v_id: None,
in_v_id: None,
properties: smallvec::smallvec![],
rank: None,
}));
self
}
pub fn from(mut self, vertex_id: i64) -> Self {
self.push_step(LogicalStep::From(FromStep { vertex_id }));
self
}
pub fn to(mut self, vertex_id: i64) -> Self {
self.push_step(LogicalStep::To(ToStep { vertex_id }));
self
}
pub fn property(mut self, key: impl Into<SmolStr>, value: impl Into<Value>) -> Self {
let key_smol = key.into();
if key_smol == LABEL {
self.record_error(StoreError::SchemaViolation(
"Cannot manually set or update the reserved property 'label'. Vertex and edge labels must be specified when creating elements via addV()/addE().".to_string()
));
return self;
}
let val = value.into();
if let Some(prim) = value_to_primitive(val.clone()) {
self.push_step(LogicalStep::Property(PropertyStep { prop_key: key_smol, prop_value: prim }));
} else {
self.record_error(StoreError::UnexpectedDataType(format!(
"property() expects a scalar primitive value, got complex type: {:?}",
val
)));
}
self
}
pub fn drop(mut self) -> Self {
self.push_step(LogicalStep::Drop(DropStep {}));
self
}
}
#[allow(private_interfaces)]
impl PlanAppender for GraphTraversal {
fn plan_mut(&mut self) -> &mut LogicalPlan {
&mut self.plan
}
fn record_error(&mut self, err: StoreError) {
if self.error.is_none() {
self.error = Some(err);
}
}
fn pending_repeat_mut(&mut self) -> &mut Option<RepeatBuilder> {
&mut self.pending_repeat
}
}
pub trait TraversalBuilder: PlanAppender {
#[allow(non_snake_case)]
fn V(mut self, ids: impl IntoIterator<Item = i64>) -> Self {
use crate::planner::logical_step::VStep;
self.push_step(LogicalStep::V(VStep { ids: ids.into_iter().collect() }));
self
}
#[allow(non_snake_case)]
fn E(mut self, keys: impl IntoIterator<Item = String>) -> Self {
self.push_step(LogicalStep::E(EStep { keys: keys.into_iter().collect() }));
self
}
fn out<'a>(mut self, labels: impl IntoIterator<Item = &'a str>) -> Self {
self.push_step(LogicalStep::Out(OutStep {
labels: labels.into_iter().map(SmolStr::from).collect(),
end_vertex_ids: None,
}));
self
}
fn r#in<'a>(mut self, labels: impl IntoIterator<Item = &'a str>) -> Self {
self.push_step(LogicalStep::In(InStep {
labels: labels.into_iter().map(SmolStr::from).collect(),
end_vertex_ids: None,
}));
self
}
fn both<'a>(mut self, labels: impl IntoIterator<Item = &'a str>) -> Self {
self.push_step(LogicalStep::Both(BothStep {
labels: labels.into_iter().map(SmolStr::from).collect(),
end_vertex_ids: None,
}));
self
}
#[allow(non_snake_case)]
fn outE<'a>(mut self, labels: impl IntoIterator<Item = &'a str>) -> Self {
self.push_step(LogicalStep::OutE(OutEStep {
labels: labels.into_iter().map(SmolStr::from).collect(),
end_vertex_ids: None,
rank: None,
}));
self
}
#[allow(non_snake_case)]
fn inE<'a>(mut self, labels: impl IntoIterator<Item = &'a str>) -> Self {
self.push_step(LogicalStep::InE(InEStep {
labels: labels.into_iter().map(SmolStr::from).collect(),
end_vertex_ids: None,
rank: None,
}));
self
}
#[allow(non_snake_case)]
fn bothE<'a>(mut self, labels: impl IntoIterator<Item = &'a str>) -> Self {
self.push_step(LogicalStep::BothE(BothEStep {
labels: labels.into_iter().map(SmolStr::from).collect(),
end_vertex_ids: None,
rank: None,
}));
self
}
#[allow(non_snake_case)]
fn inV(mut self) -> Self {
self.push_step(LogicalStep::InV(InVStep {}));
self
}
#[allow(non_snake_case)]
fn outV(mut self) -> Self {
self.push_step(LogicalStep::OutV(OutVStep {}));
self
}
#[allow(non_snake_case)]
fn otherV(mut self) -> Self {
self.push_step(LogicalStep::OtherV(OtherVStep {}));
self
}
fn has(mut self, key: impl Into<SmolStr>, pred: impl Into<Predicate>) -> Self {
if let Err(err) = push_has_step(self.plan_mut().steps.as_mut(), key.into(), pred.into()) {
self.record_error(err);
}
self
}
#[allow(non_snake_case)]
fn hasLabel(mut self, pred: impl Into<Predicate>) -> Self {
let pred = pred.into();
if let Err(err) = type_bridge::validate_label_predicate(&pred) {
self.record_error(err);
return self;
}
match type_bridge::predicate_to_primitive_predicate(pred) {
Ok(prim_pred) => self.push_step(LogicalStep::HasLabel(HasLabelStep { pred: prim_pred })),
Err(err) => self.record_error(err),
}
self
}
#[allow(non_snake_case)]
fn hasId(mut self, pred: impl Into<Predicate>) -> Self {
match type_bridge::predicate_to_primitive_predicate(pred.into()) {
Ok(prim_pred) => self.push_step(LogicalStep::HasId(HasIdStep { pred: prim_pred })),
Err(err) => self.record_error(err),
}
self
}
#[allow(non_snake_case)]
fn hasRank(mut self, pred: impl Into<Predicate>) -> Self {
match type_bridge::predicate_to_primitive_predicate(pred.into()) {
Ok(prim_pred) => self.push_step(LogicalStep::HasRank(HasRankStep { pred: prim_pred })),
Err(err) => self.record_error(err),
}
self
}
fn is(mut self, pred: impl Into<Predicate>) -> Self {
let p = pred.into();
match &p {
Predicate::Eq(v)
| Predicate::Ne(v)
| Predicate::Gt(v)
| Predicate::Gte(v)
| Predicate::Lt(v)
| Predicate::Lte(v) => {
if value_to_primitive(v.clone()).is_none() {
self.record_error(StoreError::UnexpectedDataType(format!(
"is() expects scalar values, got: {:?}",
v
)));
return self;
}
}
Predicate::Between(lo, hi) => {
if value_to_primitive(lo.clone()).is_none() || value_to_primitive(hi.clone()).is_none() {
self.record_error(StoreError::UnexpectedDataType(format!(
"is() expects scalar values, got: {:?}, {:?}",
lo, hi
)));
return self;
}
}
Predicate::Within(vs) | Predicate::Without(vs) => {
for v in vs {
if value_to_primitive(v.clone()).is_none() {
self.record_error(StoreError::UnexpectedDataType(format!(
"is() expects scalar values, got: {:?}",
v
)));
return self;
}
}
}
}
match type_bridge::predicate_to_primitive_predicate(p) {
Ok(prim_pred) => self.push_step(LogicalStep::ScalarFilter(ScalarFilterStep { pred: prim_pred })),
Err(err) => self.record_error(err),
}
self
}
fn values<'a>(mut self, keys: impl IntoIterator<Item = &'a str>) -> Self {
self.push_step(LogicalStep::Values(ValuesStep {
property_keys: keys.into_iter().map(SmolStr::from).collect(),
}));
self
}
fn properties<'a>(mut self, keys: impl IntoIterator<Item = &'a str>) -> Self {
self.push_step(LogicalStep::Properties(PropertiesStep {
property_keys: keys.into_iter().map(SmolStr::from).collect(),
}));
self
}
fn count(mut self) -> Self {
self.push_step(LogicalStep::Count(CountStep {}));
self
}
fn limit(mut self, n: i64) -> Self {
self.push_step(LogicalStep::Limit(LimitStep { limit: n }));
self
}
fn path(mut self) -> Self {
self.push_step(LogicalStep::Path(PathStep {}));
self
}
fn as_(mut self, label: impl Into<SmolStr>) -> Self {
self.push_step(LogicalStep::As(AsStep { labels: smallvec::smallvec![label.into()] }));
self
}
fn range(mut self, lo: i64, hi: i64) -> Self {
self.push_step(LogicalStep::Range(RangeStep { lo, hi }));
self
}
fn skip(mut self, n: i64) -> Self {
self.push_step(LogicalStep::Skip(SkipStep { n }));
self
}
fn tail(mut self, n: i64) -> Self {
self.push_step(LogicalStep::Tail(TailStep { n }));
self
}
fn order(mut self) -> Self {
let keys = smallvec::smallvec![OrderKey { spec: OrderKeySpec::Value, order: Order::Asc }];
self.push_step(LogicalStep::Order(OrderStep { keys }));
self
}
fn by(mut self, key: impl Into<SmolStr>) -> Self {
if follows_group_step(self.plan_mut()) {
self.record_error(by_after_group_error("by()"));
return self;
}
let key: SmolStr = key.into();
let key2 = key.clone();
let needs_order = {
let plan = self.plan_mut();
match plan.steps.last_mut() {
Some(LogicalStep::Order(order_step)) => {
let is_default = matches!(order_step.keys.as_slice(), [OrderKey { spec: OrderKeySpec::Value, .. }]);
if is_default {
order_step.keys =
smallvec::smallvec![OrderKey { spec: OrderKeySpec::Property(key), order: Order::Asc }];
} else {
order_step.keys.push(OrderKey { spec: OrderKeySpec::Property(key), order: Order::Asc });
}
false
}
_ => true,
}
};
if needs_order {
self = self.order();
let plan = self.plan_mut();
match plan.steps.last_mut() {
Some(LogicalStep::Order(order_step)) => {
let is_default = matches!(order_step.keys.as_slice(), [OrderKey { spec: OrderKeySpec::Value, .. }]);
if is_default {
order_step.keys =
smallvec::smallvec![OrderKey { spec: OrderKeySpec::Property(key2), order: Order::Asc }];
} else {
order_step.keys.push(OrderKey { spec: OrderKeySpec::Property(key2), order: Order::Asc });
}
}
_ => unreachable!("order() just pushed an Order step"),
}
}
self
}
fn order_by(mut self, key: impl Into<SmolStr>, order: Order) -> Self {
if follows_group_step(self.plan_mut()) {
self.record_error(by_after_group_error("order_by()"));
return self;
}
let key: SmolStr = key.into();
let needs_order = {
let plan = self.plan_mut();
match plan.steps.last_mut() {
Some(LogicalStep::Order(order_step)) => {
let is_default = matches!(order_step.keys.as_slice(), [OrderKey { spec: OrderKeySpec::Value, .. }]);
if is_default {
order_step.keys =
smallvec::smallvec![OrderKey { spec: OrderKeySpec::Property(key.clone()), order }];
} else {
order_step.keys.push(OrderKey { spec: OrderKeySpec::Property(key.clone()), order });
}
false
}
_ => true,
}
};
if needs_order {
self = self.order();
let plan = self.plan_mut();
match plan.steps.last_mut() {
Some(LogicalStep::Order(order_step)) => {
let is_default = matches!(order_step.keys.as_slice(), [OrderKey { spec: OrderKeySpec::Value, .. }]);
if is_default {
order_step.keys =
smallvec::smallvec![OrderKey { spec: OrderKeySpec::Property(key.clone()), order }];
} else {
order_step.keys.push(OrderKey { spec: OrderKeySpec::Property(key.clone()), order });
}
}
_ => unreachable!("order() just pushed an Order step"),
}
}
self
}
fn simple_path(mut self) -> Self {
self.push_step(LogicalStep::SimplePath(SimplePathStep {}));
self
}
fn cyclic_path(mut self) -> Self {
self.push_step(LogicalStep::CyclicPath(CyclicPathStep {}));
self
}
fn choose(
mut self,
mut predicate: GraphTraversal,
mut true_choice: GraphTraversal,
false_choice: Option<GraphTraversal>,
) -> Self {
if let Some(err) = predicate.error.take() {
self.record_error(err);
}
if let Some(err) = true_choice.error.take() {
self.record_error(err);
}
let fc = false_choice.map(|mut f| {
if let Some(err) = f.error.take() {
self.record_error(err);
}
f.into_plan()
});
self.push_step(LogicalStep::Choose(ChooseStep {
predicate: predicate.into_plan(),
true_choice: true_choice.into_plan(),
false_choice: fc,
}));
self
}
fn select(mut self, label: impl Into<SmolStr>) -> Self {
self.push_step(LogicalStep::Select(SelectStep { labels: smallvec::smallvec![label.into()] }));
self
}
fn id(mut self) -> Self {
self.push_step(LogicalStep::Id(IdStep {}));
self
}
fn label(mut self) -> Self {
self.push_step(LogicalStep::Label(LabelStep {}));
self
}
fn rank(mut self) -> Self {
self.push_step(LogicalStep::Rank(RankStep {}));
self
}
fn identity(mut self) -> Self {
self.push_step(LogicalStep::Identity(IdentityStep {}));
self
}
fn constant(mut self, value: impl Into<crate::types::gvalue::Primitive>) -> Self {
self.push_step(LogicalStep::Constant(ConstantStep { value: value.into() }));
self
}
fn local(mut self, mut traversal: GraphTraversal) -> Self {
if let Some(err) = traversal.error.take() {
self.record_error(err);
}
self.push_step(LogicalStep::Local(LocalStep { plan: traversal.into_plan() }));
self
}
fn dedup(mut self) -> Self {
self.push_step(LogicalStep::Dedup(DedupStep {}));
self
}
fn group(mut self) -> Self {
self.push_step(LogicalStep::Group(GroupStep { key: None }));
self
}
fn group_count(mut self) -> Self {
self.push_step(LogicalStep::GroupCount(GroupCountStep { key: None }));
self
}
fn fold(mut self) -> Self {
self.push_step(LogicalStep::Fold(FoldStep {}));
self
}
fn sum(mut self) -> Self {
self.push_step(LogicalStep::Sum(SumStep {}));
self
}
fn mean(mut self) -> Self {
self.push_step(LogicalStep::Mean(MeanStep {}));
self
}
fn max(mut self) -> Self {
self.push_step(LogicalStep::Max(MaxStep {}));
self
}
fn min(mut self) -> Self {
self.push_step(LogicalStep::Min(MinStep {}));
self
}
fn unfold(mut self) -> Self {
self.push_step(LogicalStep::Unfold(UnfoldStep {}));
self
}
fn r#where(mut self, mut sub: GraphTraversal) -> Self {
if let Some(err) = sub.error.take() {
self.record_error(err);
}
self.push_step(LogicalStep::Where(WhereStep { plan: sub.into_plan() }));
self
}
fn not(mut self, mut sub: GraphTraversal) -> Self {
if let Some(err) = sub.error.take() {
self.record_error(err);
}
self.push_step(LogicalStep::Not(NotStep { plan: sub.into_plan() }));
self
}
fn and(mut self, subs: impl IntoIterator<Item = GraphTraversal>) -> Self {
let mut plans = Vec::new();
for mut sub in subs {
if let Some(err) = sub.error.take() {
self.record_error(err);
}
plans.push(sub.into_plan());
}
self.push_step(LogicalStep::And(AndStep { plans }));
self
}
fn or(mut self, subs: impl IntoIterator<Item = GraphTraversal>) -> Self {
let mut plans = Vec::new();
for mut sub in subs {
if let Some(err) = sub.error.take() {
self.record_error(err);
}
plans.push(sub.into_plan());
}
self.push_step(LogicalStep::Or(OrStep { plans }));
self
}
fn coalesce(mut self, subs: impl IntoIterator<Item = GraphTraversal>) -> Self {
let mut plans = Vec::new();
for mut sub in subs {
if let Some(err) = sub.error.take() {
self.record_error(err);
}
plans.push(sub.into_plan());
}
self.push_step(LogicalStep::Coalesce(CoalesceStep { plans }));
self
}
fn union(mut self, subs: impl IntoIterator<Item = GraphTraversal>) -> Self {
let mut plans = Vec::new();
for mut sub in subs {
if let Some(err) = sub.error.take() {
self.record_error(err);
}
plans.push(sub.into_plan());
}
self.push_step(LogicalStep::Union(UnionStep { plans: plans.into_iter().collect() }));
self
}
fn repeat(mut self, body: GraphTraversal) -> Self {
self.flush_pending_repeat();
let mut body = body;
if let Some(err) = body.error.take() {
self.record_error(err);
}
*self.pending_repeat_mut() =
Some(RepeatBuilder { body: body.into_plan(), until: None, times: None, emit: EmitSpec::Never });
self
}
fn times(mut self, n: i64) -> Self {
if n == 0 {
self.record_error(StoreError::TraversalError(
"times(0) is invalid: a repeat body must run at least once.".to_string(),
));
return self;
}
match self.pending_repeat_mut() {
Some(ref mut rb) => rb.times = Some(n),
None => {
self.record_error(StoreError::TraversalError("times() must immediately follow repeat().".to_string()))
}
}
self
}
fn until(mut self, cond: GraphTraversal) -> Self {
let mut cond = cond;
if let Some(err) = cond.error.take() {
self.record_error(err);
}
match self.pending_repeat_mut() {
Some(ref mut rb) => rb.until = Some(cond.into_plan()),
None => {
self.record_error(StoreError::TraversalError("until() must immediately follow repeat().".to_string()))
}
}
self
}
fn emit(mut self) -> Self {
match self.pending_repeat_mut() {
Some(ref mut rb) => rb.emit = EmitSpec::Always,
None => {
self.record_error(StoreError::TraversalError("emit() must immediately follow repeat().".to_string()))
}
}
self
}
fn emit_if(mut self, cond: GraphTraversal) -> Self {
let mut cond = cond;
if let Some(err) = cond.error.take() {
self.record_error(err);
}
match self.pending_repeat_mut() {
Some(ref mut rb) => rb.emit = EmitSpec::If(cond.into_plan()),
None => {
self.record_error(StoreError::TraversalError("emit_if() must immediately follow repeat().".to_string()))
}
}
self
}
}
impl<T: PlanAppender> TraversalBuilder for T {}
mod terminals;
pub use terminals::{ReadTraversal, WriteTraversal};