mod aggregate;
mod cycle_guard;
mod idiom;
#[cfg(feature = "gql")]
mod match_plan;
mod row_scope;
mod select;
mod source;
pub(crate) mod util;
use std::sync::Arc;
pub(crate) use cycle_guard::CycleGuard;
use self::util::literal_to_value;
use crate::ctx::FrozenContext;
use crate::dbs::NewPlannerStrategy;
use crate::err::Error;
use crate::exec::ExecOperator;
use crate::exec::function::FunctionRegistry;
use crate::exec::operators::{
AnalyzePlan, DatabaseInfoPlan, ExplainPlan, ExprPlan, Fetch, ForeachPlan, IfElsePlan,
IndexInfoPlan, NamespaceInfoPlan, ReturnPlan, RootInfoPlan, SequencePlan, SleepPlan,
TableInfoPlan, UserInfoPlan,
};
use crate::exec::physical_expr::{
ArrayLiteral, BinaryOp, BlockPhysicalExpr, BuiltinFunctionExec, ClosureCallExec, ClosureExec,
ControlFlowExpr, ControlFlowKind, IfElseExpr, JsFunctionExec, Literal as PhysicalLiteral,
MockExpr, ModelFunctionExec, ObjectLiteral, Param, PostfixOp, ProjectionFunctionExec,
RecordIdExpr, ScalarSubquery, SetLiteral, SiloModuleExec, SurrealismModuleExec, UnaryOp,
UserDefinedFunctionExec,
};
use crate::expr::statements::IfelseStatement;
use crate::expr::{Expr, Function, FunctionCall};
pub struct Planner<'ctx> {
ctx: &'ctx FrozenContext,
function_registry: &'ctx FunctionRegistry,
pub(crate) txn: Option<Arc<crate::kvs::Transaction>>,
pub(crate) ns: Option<String>,
pub(crate) db: Option<String>,
pub(crate) version: Option<Arc<dyn crate::exec::PhysicalExpr>>,
pub(crate) auth: Option<Arc<crate::iam::Auth>>,
pub(crate) cycle_guard: CycleGuard,
ns_db_ids_cache:
tokio::sync::OnceCell<Option<(crate::catalog::NamespaceId, crate::catalog::DatabaseId)>>,
planner_strategy: NewPlannerStrategy,
depth: u32,
}
impl<'ctx> Planner<'ctx> {
pub fn new(ctx: &'ctx FrozenContext) -> Self {
Self {
ctx,
function_registry: ctx.function_registry(),
txn: None,
ns: None,
db: None,
version: None,
auth: None,
cycle_guard: CycleGuard::default(),
ns_db_ids_cache: tokio::sync::OnceCell::new(),
planner_strategy: *ctx.new_planner_strategy(),
depth: 0,
}
}
pub fn with_txn(
ctx: &'ctx FrozenContext,
txn: Arc<crate::kvs::Transaction>,
ns: Option<String>,
db: Option<String>,
) -> Self {
Self {
ctx,
function_registry: ctx.function_registry(),
txn: Some(txn),
ns,
db,
version: None,
auth: None,
cycle_guard: CycleGuard::default(),
ns_db_ids_cache: tokio::sync::OnceCell::new(),
planner_strategy: *ctx.new_planner_strategy(),
depth: 0,
}
}
#[inline]
pub(crate) fn for_database(
ctx: &'ctx FrozenContext,
txn: Arc<crate::kvs::Transaction>,
db_ctx: &crate::exec::DatabaseContext,
) -> Self {
Self::with_txn(
ctx,
txn,
Some(db_ctx.ns_name().to_owned()),
Some(db_ctx.db_name().to_owned()),
)
}
pub fn with_version(mut self, version: Option<Arc<dyn crate::exec::PhysicalExpr>>) -> Self {
self.version = version;
self
}
#[must_use]
pub(crate) fn with_auth(mut self, auth: Arc<crate::iam::Auth>) -> Self {
self.auth = Some(auth);
self
}
pub(crate) fn should_check_perms_for_view(&self, ns: &str, db: &str) -> bool {
let Some(ref auth) = self.auth else {
return true;
};
if !self.ctx.auth_enabled() && auth.is_anon() {
return false;
}
let allowed = auth.has_viewer_role();
let db_in_actor_level = auth.is_root() || auth.is_ns_check(ns) || auth.is_db_check(ns, db);
!allowed || !db_in_actor_level
}
#[must_use]
pub(crate) fn with_cycle_guard(mut self, guard: CycleGuard) -> Self {
self.cycle_guard = guard;
self
}
#[inline]
pub(crate) fn cycle_guard(&self) -> CycleGuard {
self.cycle_guard.clone()
}
#[must_use]
pub(crate) fn with_depth(mut self, depth: u32) -> Self {
self.depth = depth;
self
}
#[inline]
pub(crate) fn current_depth(&self) -> u32 {
self.depth
}
#[inline]
fn check_depth(&self) -> Result<(), Error> {
if self.depth > self.ctx.config.max_computation_depth {
return Err(Error::ComputationDepthExceeded);
}
Ok(())
}
#[inline]
pub(crate) fn txn(&self) -> Option<&Arc<crate::kvs::Transaction>> {
self.txn.as_ref()
}
#[inline]
pub(crate) fn ns(&self) -> Option<&str> {
self.ns.as_deref()
}
#[inline]
pub(crate) fn db(&self) -> Option<&str> {
self.db.as_deref()
}
#[inline]
pub fn function_registry(&self) -> &'ctx FunctionRegistry {
self.function_registry
}
#[cfg(feature = "surrealism")]
async fn resolve_module_writeable(
&self,
module: &str,
sub: Option<&str>,
) -> Result<bool, Error> {
use crate::catalog::providers::DatabaseProvider;
use crate::ctx::Context;
use crate::expr::module::ModuleExecutable;
let Some(txn) = &self.txn else {
return Ok(false);
};
let (Some(ns), Some(db)) = (&self.ns, &self.db) else {
return Ok(false);
};
let Some(db_def) =
txn.get_db_by_name(ns, db, None).await.map_err(|e| Error::Internal(e.to_string()))?
else {
return Ok(false);
};
let mod_name = format!("mod::{module}");
let val =
match txn.get_db_module(db_def.namespace_id, db_def.database_id, &mod_name, None).await
{
Ok(v) => v,
Err(e) => {
if let Some(Error::MdNotFound {
..
}) = e.downcast_ref::<Error>()
{
return Ok(false);
}
return Err(Error::Internal(e.to_string()));
}
};
let executable: ModuleExecutable = val.executable.clone().into();
let mut plan_ctx = Context::new_child(self.ctx);
plan_ctx.set_transaction(Arc::clone(txn));
let frozen = plan_ctx.freeze();
let sig = executable
.signature(&frozen, &db_def.namespace_id, &db_def.database_id, sub)
.await
.map_err(|e| Error::Internal(e.to_string()))?;
Ok(sig.writeable)
}
#[cfg(not(feature = "surrealism"))]
async fn resolve_module_writeable(
&self,
_module: &str,
_sub: Option<&str>,
) -> Result<bool, Error> {
Ok(false)
}
#[cfg(feature = "surrealism")]
async fn resolve_silo_writeable(
&self,
org: &str,
pkg: &str,
major: u32,
minor: u32,
patch: u32,
sub: Option<&str>,
) -> Result<bool, Error> {
use crate::ctx::Context;
use crate::expr::module::SiloExecutable;
let executable = SiloExecutable {
organisation: org.to_string(),
package: pkg.to_string(),
major,
minor,
patch,
};
let ctx = if let Some(txn) = &self.txn {
let mut plan_ctx = Context::new_child(self.ctx);
plan_ctx.set_transaction(Arc::clone(txn));
plan_ctx.freeze()
} else {
Arc::clone(self.ctx)
};
let sig =
executable.signature(&ctx, sub).await.map_err(|e| Error::Internal(e.to_string()))?;
Ok(sig.writeable)
}
#[cfg(not(feature = "surrealism"))]
async fn resolve_silo_writeable(
&self,
_org: &str,
_pkg: &str,
_major: u32,
_minor: u32,
_patch: u32,
_sub: Option<&str>,
) -> Result<bool, Error> {
Ok(false)
}
pub async fn plan(&self, expr: &Expr) -> Result<Arc<dyn ExecOperator>, Error> {
let result = self.plan_expr(expr.clone()).await;
self.require_planned(result)
}
pub async fn physical_expr(
&self,
expr: Expr,
) -> Result<Arc<dyn crate::exec::PhysicalExpr>, Error> {
self.check_depth()?;
match expr {
Expr::Literal(lit) => Box::pin(self.physical_literal(lit)).await,
Expr::Constant(c) => Ok(Arc::new(PhysicalLiteral(c.compute()))),
Expr::Table(t) => Ok(Arc::new(PhysicalLiteral(crate::val::Value::Table(t)))),
Expr::Param(p) => Ok(Arc::new(Param(p.into_strand()))),
Expr::Idiom(idiom) => Box::pin(self.convert_idiom(idiom)).await,
Expr::Binary {
left,
op,
right,
} => Box::pin(self.physical_binary_expr(*left, op, *right)).await,
Expr::Prefix {
op,
expr,
} => Box::pin(self.physical_prefix_expr(op, *expr)).await,
Expr::Postfix {
op,
expr,
} => Box::pin(self.physical_postfix_expr(op, *expr)).await,
Expr::FunctionCall(fc) => Box::pin(self.physical_function_call(*fc)).await,
Expr::Closure(c) => Ok(Arc::new(ClosureExec {
closure: *c,
})),
Expr::IfElse(stmt) => Box::pin(self.physical_if_else(*stmt)).await,
Expr::Mock(m) => Ok(Arc::new(MockExpr(m))),
Expr::Block(b) => Ok(Arc::new(BlockPhysicalExpr {
block: *b,
})),
Expr::Break => Ok(Arc::new(ControlFlowExpr {
kind: ControlFlowKind::Break,
inner: None,
})),
Expr::Continue => Ok(Arc::new(ControlFlowExpr {
kind: ControlFlowKind::Continue,
inner: None,
})),
Expr::Return(s) => {
let inner = Box::pin(self.physical_expr(s.what)).await?;
Ok(Arc::new(ControlFlowExpr {
kind: ControlFlowKind::Return,
inner: Some(inner),
}))
}
Expr::Throw(e) => {
let inner = Box::pin(self.physical_expr(*e)).await?;
Ok(Arc::new(ControlFlowExpr {
kind: ControlFlowKind::Throw,
inner: Some(inner),
}))
}
Expr::Select(_)
| Expr::Info(_)
| Expr::Foreach(_)
| Expr::Sleep(_)
| Expr::Explain {
..
} => Box::pin(self.physical_statement_subquery(expr)).await,
Expr::Let(_) => Err(Error::InvalidStatement(
"LET statements can only appear at the top level of a query or inside a block \
expression"
.to_string(),
)),
Expr::Define(_) | Expr::Remove(_) | Expr::Rebuild(_) | Expr::Alter(_) => {
Err(Error::PlannerUnsupported(
"DDL statements cannot be used in expression context".to_string(),
))
}
Expr::Create(_)
| Expr::Update(_)
| Expr::Upsert(_)
| Expr::Delete(_)
| Expr::Relate(_)
| Expr::Insert(_) => Err(Error::PlannerUnsupported(
"DML subqueries not yet supported in execution plans".to_string(),
)),
#[cfg(feature = "gql")]
Expr::Match(_) => Err(Error::PlannerUnsupported(
"GQL MATCH cannot be used as a sub-expression".to_string(),
)),
}
}
async fn physical_literal(
&self,
lit: crate::expr::literal::Literal,
) -> Result<Arc<dyn crate::exec::PhysicalExpr>, Error> {
use crate::expr::literal::Literal;
match lit {
Literal::Array(elements) => {
let elements = self.physical_args(elements).await?;
Ok(Arc::new(ArrayLiteral {
elements,
}))
}
Literal::Object(entries) => {
let mut phys_entries = Vec::with_capacity(entries.len());
for entry in entries {
let value = Box::pin(self.physical_expr(entry.value)).await?;
phys_entries.push((entry.key, value));
}
Ok(Arc::new(ObjectLiteral {
entries: phys_entries,
}))
}
Literal::Set(elements) => {
let elements = self.physical_args(elements).await?;
Ok(Arc::new(SetLiteral {
elements,
}))
}
Literal::RecordId(rid_lit) => {
let key = self.convert_record_key_to_physical(&rid_lit.key).await?;
Ok(Arc::new(RecordIdExpr {
table: rid_lit.table,
key,
}))
}
other => {
let value = literal_to_value(other)?;
Ok(Arc::new(PhysicalLiteral(value)))
}
}
}
async fn physical_binary_expr(
&self,
left: Expr,
op: crate::expr::operator::BinaryOperator,
right: Expr,
) -> Result<Arc<dyn crate::exec::PhysicalExpr>, Error> {
if let crate::expr::operator::BinaryOperator::Matches(ref matches_op) = op
&& let Expr::Idiom(idiom) = left
{
let resolved_query = match &right {
Expr::Literal(crate::expr::literal::Literal::String(s)) => {
Some(s.as_str().to_owned())
}
Expr::Param(param) => self.ctx.value(param.as_str()).and_then(|v| {
if let crate::val::Value::String(s) = v {
Some(s.as_str().to_owned())
} else {
None
}
}),
_ => None,
};
if let Some(query) = resolved_query {
if idiom.0.len() > 1 {
return Err(Error::PlannerUnimplemented(
"MATCHES with multi-part field path not yet supported \
in streaming executor"
.to_string(),
));
}
let idiom_clone = idiom.clone();
let query_clone = query.clone();
let left_phys = Box::pin(self.physical_expr(Expr::Idiom(idiom))).await?;
let right_phys = Box::pin(self.physical_expr(Expr::Literal(
crate::expr::literal::Literal::String(query.into()),
)))
.await?;
return Ok(Arc::new(crate::exec::physical_expr::MatchesOp::new(
left_phys,
right_phys,
matches_op.clone(),
idiom_clone,
query_clone,
)));
}
let left_phys = Box::pin(self.physical_expr(Expr::Idiom(idiom))).await?;
let right_phys = Box::pin(self.physical_expr(right)).await?;
return Ok(Arc::new(BinaryOp {
left: left_phys,
op,
right: right_phys,
}));
}
let left_phys = Box::pin(self.physical_expr(left)).await?;
let right_phys = Box::pin(self.physical_expr(right)).await?;
if is_simple_binary_eligible(&op) {
if let Some(field) = left_phys.try_simple_field()
&& let Some(lit) = right_phys.try_literal()
{
return Ok(Arc::new(crate::exec::physical_expr::SimpleBinaryOp {
field_name: field.to_string(),
op,
literal: lit.clone(),
reversed: false,
}));
} else if let Some(field) = right_phys.try_simple_field()
&& let Some(lit) = left_phys.try_literal()
{
return Ok(Arc::new(crate::exec::physical_expr::SimpleBinaryOp {
field_name: field.to_string(),
op,
literal: lit.clone(),
reversed: true,
}));
}
}
Ok(Arc::new(BinaryOp {
left: left_phys,
op,
right: right_phys,
}))
}
async fn physical_prefix_expr(
&self,
op: crate::expr::operator::PrefixOperator,
expr: Expr,
) -> Result<Arc<dyn crate::exec::PhysicalExpr>, Error> {
{
let mut d = self.depth;
let mut cur = &expr;
while let Expr::Prefix {
expr: inner,
..
} = cur
{
d += 1;
if d > self.ctx.config.max_computation_depth {
return Err(Error::ComputationDepthExceeded);
}
cur = inner;
}
}
let inner = Box::pin(self.physical_expr(expr)).await?;
Ok(Arc::new(UnaryOp {
op,
expr: inner,
}))
}
async fn physical_postfix_expr(
&self,
op: crate::expr::operator::PostfixOperator,
expr: Expr,
) -> Result<Arc<dyn crate::exec::PhysicalExpr>, Error> {
use crate::expr::operator::PostfixOperator;
match op {
PostfixOperator::Call(args) => {
let target = Box::pin(self.physical_expr(expr)).await?;
let arguments = self.physical_args(args).await?;
Ok(Arc::new(ClosureCallExec {
target,
arguments,
}))
}
_ => {
let inner = Box::pin(self.physical_expr(expr)).await?;
Ok(Arc::new(PostfixOp {
op,
expr: inner,
}))
}
}
}
async fn physical_function_call(
&self,
func_call: FunctionCall,
) -> Result<Arc<dyn crate::exec::PhysicalExpr>, Error> {
let FunctionCall {
receiver,
arguments,
} = func_call;
match receiver {
Function::Normal(name) => {
let registry = self.function_registry();
if registry.is_index_function(&name) {
return Box::pin(self.plan_index_function(&name, arguments)).await;
}
let arguments = self.physical_args(arguments).await?;
if registry.is_projection(&name) {
let func_ctx = registry
.get_projection(&name)
.map(|f| f.required_context())
.unwrap_or(crate::exec::ContextLevel::Database);
Ok(Arc::new(ProjectionFunctionExec {
name,
arguments,
func_required_context: func_ctx,
}))
} else {
let func_ctx = registry
.get(&name)
.map(|f| f.required_context())
.unwrap_or(crate::exec::ContextLevel::Root);
Ok(Arc::new(BuiltinFunctionExec {
name,
arguments,
func_required_context: func_ctx,
plan_depth: self.current_depth(),
}))
}
}
Function::Custom(name) => {
let arguments = self.physical_args(arguments).await?;
Ok(Arc::new(UserDefinedFunctionExec {
name,
arguments,
plan_depth: self.current_depth(),
}))
}
Function::Script(script) => {
let arguments = self.physical_args(arguments).await?;
Ok(Arc::new(JsFunctionExec {
script,
arguments,
plan_depth: self.current_depth(),
}))
}
Function::Model(model) => {
let arguments = self.physical_args(arguments).await?;
Ok(Arc::new(ModelFunctionExec {
model,
arguments,
}))
}
Function::Module(module, sub) => {
let arguments = self.physical_args(arguments).await?;
let writeable = self.resolve_module_writeable(&module, sub.as_deref()).await?;
Ok(Arc::new(SurrealismModuleExec {
module,
sub,
arguments,
writeable,
}))
}
Function::Silo {
org,
pkg,
major,
minor,
patch,
sub,
} => {
let arguments = self.physical_args(arguments).await?;
let writeable = self
.resolve_silo_writeable(&org, &pkg, major, minor, patch, sub.as_deref())
.await?;
Ok(Arc::new(SiloModuleExec {
org,
pkg,
major,
minor,
patch,
sub,
arguments,
writeable,
}))
}
}
}
async fn physical_args(
&self,
args: Vec<Expr>,
) -> Result<Vec<Arc<dyn crate::exec::PhysicalExpr>>, Error> {
let mut phys = Vec::with_capacity(args.len());
for arg in args {
phys.push(Box::pin(self.physical_expr(arg)).await?);
}
Ok(phys)
}
async fn physical_if_else(
&self,
stmt: IfelseStatement,
) -> Result<Arc<dyn crate::exec::PhysicalExpr>, Error> {
let IfelseStatement {
exprs,
close,
} = stmt;
let mut branches = Vec::with_capacity(exprs.len());
for (condition, body) in exprs {
let cond_phys = Box::pin(self.physical_expr(condition)).await?;
let body_phys = Box::pin(self.physical_expr(body)).await?;
branches.push((cond_phys, body_phys));
}
let otherwise = if let Some(else_expr) = close {
Some(Box::pin(self.physical_expr(else_expr)).await?)
} else {
None
};
Ok(Arc::new(IfElseExpr {
branches,
otherwise,
}))
}
async fn physical_statement_subquery(
&self,
expr: Expr,
) -> Result<Arc<dyn crate::exec::PhysicalExpr>, Error> {
let plan: Arc<dyn ExecOperator> = match expr {
Expr::Select(select) => Box::pin(self.plan_select_statement(*select)).await?,
Expr::Info(info) => self.plan_info_statement(*info).await?,
Expr::Foreach(stmt) => self.plan_foreach_statement(*stmt)?,
Expr::Sleep(stmt) => self.plan_sleep_statement(*stmt)?,
Expr::Explain {
format,
analyze,
statement,
} => {
let inner_plan = self.plan_expr(*statement).await?;
if analyze {
Arc::new(AnalyzePlan {
plan: inner_plan,
format,
redact_volatile_explain_attrs: self.ctx.redact_volatile_explain_attrs(),
})
} else {
Arc::new(ExplainPlan {
plan: inner_plan,
format,
})
}
}
other => {
tracing::error!(
expr = ?other,
"physical_statement_subquery dispatched with non-statement expr"
);
return Err(Error::Internal(
"physical_statement_subquery dispatched with non-statement expr; \
only Select/Info/Foreach/Sleep/Explain are valid here"
.into(),
));
}
};
Ok(Arc::new(ScalarSubquery {
plan,
}))
}
pub async fn physical_expr_as_name(
&self,
expr: Expr,
) -> Result<Arc<dyn crate::exec::PhysicalExpr>, Error> {
use crate::exec::physical_expr::Literal as PhysicalLiteral;
use crate::expr::part::Part;
if let Expr::Idiom(ref idiom) = expr
&& idiom.0.len() == 1
&& let Part::Field(name) = &idiom.0[0]
{
return Ok(Arc::new(PhysicalLiteral(crate::val::Value::String(name.as_str().into()))));
}
if let Expr::Table(name) = expr {
return Ok(Arc::new(PhysicalLiteral(crate::val::Value::String(name.as_str().into()))));
}
Box::pin(self.physical_expr(expr)).await
}
fn convert_record_key_to_physical<'a>(
&'a self,
key: &'a crate::expr::RecordIdKeyLit,
) -> crate::exec::BoxFut<
'a,
Result<crate::exec::physical_expr::record_id::PhysicalRecordIdKey, Error>,
> {
Box::pin(async move {
use crate::exec::physical_expr::record_id::PhysicalRecordIdKey;
use crate::expr::RecordIdKeyLit;
match key {
RecordIdKeyLit::Number(n) => Ok(PhysicalRecordIdKey::Number(*n)),
RecordIdKeyLit::String(s) => Ok(PhysicalRecordIdKey::String(s.clone())),
RecordIdKeyLit::Uuid(u) => Ok(PhysicalRecordIdKey::Uuid(*u)),
RecordIdKeyLit::Generate(generator) => {
Ok(PhysicalRecordIdKey::Generate(generator.clone()))
}
RecordIdKeyLit::Array(exprs) => {
let mut phys = Vec::with_capacity(exprs.len());
for expr in exprs {
phys.push(Box::pin(self.physical_expr(expr.clone())).await?);
}
Ok(PhysicalRecordIdKey::Array(phys))
}
RecordIdKeyLit::Object(entries) => {
let mut phys = Vec::with_capacity(entries.len());
for entry in entries {
let value = Box::pin(self.physical_expr(entry.value.clone())).await?;
phys.push((entry.key.clone(), value));
}
Ok(PhysicalRecordIdKey::Object(phys))
}
RecordIdKeyLit::Range(range) => {
let start = self.convert_bound_to_physical(&range.start).await?;
let end = self.convert_bound_to_physical(&range.end).await?;
Ok(PhysicalRecordIdKey::Range {
start,
end,
})
}
}
})
}
async fn convert_bound_to_physical(
&self,
bound: &std::ops::Bound<crate::expr::RecordIdKeyLit>,
) -> Result<
std::ops::Bound<Box<crate::exec::physical_expr::record_id::PhysicalRecordIdKey>>,
Error,
> {
match bound {
std::ops::Bound::Unbounded => Ok(std::ops::Bound::Unbounded),
std::ops::Bound::Included(key) => Ok(std::ops::Bound::Included(Box::new(
self.convert_record_key_to_physical(key).await?,
))),
std::ops::Bound::Excluded(key) => Ok(std::ops::Bound::Excluded(Box::new(
self.convert_record_key_to_physical(key).await?,
))),
}
}
fn require_planned<T>(&self, result: Result<T, Error>) -> Result<T, Error> {
match result {
Err(Error::PlannerUnimplemented(msg))
if self.planner_strategy == NewPlannerStrategy::AllReadOnlyStatements =>
{
Err(Error::Query {
message: format!("New executor does not support: {msg}"),
})
}
other => other,
}
}
fn plan_expr(
&self,
expr: Expr,
) -> crate::exec::BoxFut<'_, Result<Arc<dyn ExecOperator>, Error>> {
Box::pin(async move {
self.check_depth()?;
match expr {
Expr::Select(select) => self.plan_select_statement(*select).await,
Expr::Block(block) => self.plan_block(*block).await,
Expr::Return(output_stmt) => self.plan_return_statement(*output_stmt).await,
Expr::Let(let_stmt) => self.plan_let_statement(*let_stmt).await,
Expr::Explain {
format,
analyze,
statement,
} => self.plan_explain_statement(format, analyze, *statement).await,
Expr::Info(info) => self.plan_info_statement(*info).await,
Expr::Foreach(stmt) => self.plan_foreach_statement(*stmt),
Expr::IfElse(stmt) => self.plan_if_else_statement(*stmt),
Expr::Sleep(sleep_stmt) => self.plan_sleep_statement(*sleep_stmt),
expr @ (Expr::FunctionCall(_)
| Expr::Closure(_)
| Expr::Literal(_)
| Expr::Param(_)
| Expr::Constant(_)
| Expr::Prefix {
..
}
| Expr::Binary {
..
}
| Expr::Postfix {
..
}
| Expr::Table(_)
| Expr::Idiom(_)
| Expr::Mock(_)
| Expr::Throw(_)
| Expr::Break
| Expr::Continue) => self.plan_expr_as_operator(expr).await,
Expr::Create(_)
| Expr::Update(_)
| Expr::Upsert(_)
| Expr::Delete(_)
| Expr::Insert(_)
| Expr::Relate(_) => Err(Error::PlannerUnsupported(
"DML statements not yet supported in execution plans".to_string(),
)),
Expr::Define(_) | Expr::Remove(_) | Expr::Rebuild(_) | Expr::Alter(_) => {
Err(Error::PlannerUnsupported(
"DDL statements not yet supported in execution plans".to_string(),
))
}
#[cfg(feature = "gql")]
Expr::Match(m) => self.plan_match(*m).await,
}
})
}
async fn plan_block(&self, block: crate::expr::Block) -> Result<Arc<dyn ExecOperator>, Error> {
if block.0.is_empty() {
use crate::exec::physical_expr::Literal as PhysicalLiteral;
Ok(Arc::new(ExprPlan::new(Arc::new(PhysicalLiteral(crate::val::Value::None))))
as Arc<dyn ExecOperator>)
} else if block.0.len() == 1 {
self.plan_expr(block.0.into_iter().next().expect("block verified non-empty")).await
} else {
Ok(Arc::new(SequencePlan::new(block, self.current_depth())) as Arc<dyn ExecOperator>)
}
}
async fn plan_return_statement(
&self,
output_stmt: crate::expr::statements::OutputStatement,
) -> Result<Arc<dyn ExecOperator>, Error> {
let inner = self.plan_expr(output_stmt.what).await?;
let inner = if let Some(fetchs) = output_stmt.fetch {
let mut fields = Vec::with_capacity(fetchs.len());
for fetch_item in fetchs {
let mut idioms = self.resolve_field_idioms(fetch_item.0).await?;
fields.append(&mut idioms);
}
if fields.is_empty() {
inner
} else {
Arc::new(Fetch::new(inner, fields)) as Arc<dyn ExecOperator>
}
} else {
inner
};
Ok(Arc::new(ReturnPlan::new(inner)))
}
async fn plan_explain_statement(
&self,
format: crate::expr::ExplainFormat,
analyze: bool,
statement: Expr,
) -> Result<Arc<dyn ExecOperator>, Error> {
let inner_plan = self.plan_expr(statement).await?;
if analyze {
Ok(Arc::new(AnalyzePlan {
plan: inner_plan,
format,
redact_volatile_explain_attrs: self.ctx.redact_volatile_explain_attrs(),
}))
} else {
Ok(Arc::new(ExplainPlan {
plan: inner_plan,
format,
}))
}
}
async fn plan_let_statement(
&self,
let_stmt: crate::expr::statements::SetStatement,
) -> Result<Arc<dyn ExecOperator>, Error> {
let crate::expr::statements::SetStatement {
name,
what,
kind,
} = let_stmt;
if crate::cnf::PROTECTED_PARAM_NAMES.contains(&name.as_str()) {
return Err(Error::InvalidParam {
name: name.to_string(),
});
}
let value: Arc<dyn ExecOperator> = match what {
Expr::Select(select) => self.plan_select_statement(*select).await?,
Expr::Create(_) => {
return Err(Error::PlannerUnsupported(
"CREATE statements in LET not yet supported in execution plans".to_string(),
));
}
Expr::Update(_) => {
return Err(Error::PlannerUnsupported(
"UPDATE statements in LET not yet supported in execution plans".to_string(),
));
}
Expr::Upsert(_) => {
return Err(Error::PlannerUnsupported(
"UPSERT statements in LET not yet supported in execution plans".to_string(),
));
}
Expr::Delete(_) => {
return Err(Error::PlannerUnsupported(
"DELETE statements in LET not yet supported in execution plans".to_string(),
));
}
Expr::Insert(_) => {
return Err(Error::PlannerUnsupported(
"INSERT statements in LET not yet supported in execution plans".to_string(),
));
}
Expr::Relate(_) => {
return Err(Error::PlannerUnsupported(
"RELATE statements in LET not yet supported in execution plans".to_string(),
));
}
other => {
let expr = Box::pin(self.physical_expr(other)).await?;
Arc::new(ExprPlan::new(expr))
}
};
Ok(Arc::new(crate::exec::operators::LetPlan::new(name, kind, value)))
}
async fn plan_info_statement(
&self,
info: crate::expr::statements::info::InfoStatement,
) -> Result<Arc<dyn ExecOperator>, Error> {
use crate::expr::statements::info::InfoStatement;
match info {
InfoStatement::Root(structured, version) => {
let version = match version {
Some(v) => Some(Box::pin(self.physical_expr(v)).await?),
None => None,
};
Ok(Arc::new(RootInfoPlan::new(structured, version)) as Arc<dyn ExecOperator>)
}
InfoStatement::Ns(structured, version) => {
let version = match version {
Some(v) => Some(Box::pin(self.physical_expr(v)).await?),
None => None,
};
Ok(Arc::new(NamespaceInfoPlan::new(structured, version)) as Arc<dyn ExecOperator>)
}
InfoStatement::Db(structured, version) => {
let version = match version {
Some(v) => Some(Box::pin(self.physical_expr(v)).await?),
None => None,
};
Ok(Arc::new(DatabaseInfoPlan::new(structured, version)) as Arc<dyn ExecOperator>)
}
InfoStatement::Tb(table, structured, version) => {
let table = self.physical_expr_as_name(table).await?;
let version = match version {
Some(v) => Some(Box::pin(self.physical_expr(v)).await?),
None => None,
};
Ok(Arc::new(TableInfoPlan::new(table, structured, version))
as Arc<dyn ExecOperator>)
}
InfoStatement::User(user, base, structured) => {
let user = self.physical_expr_as_name(user).await?;
Ok(Arc::new(UserInfoPlan::new(user, base, structured)) as Arc<dyn ExecOperator>)
}
InfoStatement::Index(index, table, structured) => {
let index = self.physical_expr_as_name(index).await?;
let table = self.physical_expr_as_name(table).await?;
Ok(Arc::new(IndexInfoPlan::new(index, table, structured)) as Arc<dyn ExecOperator>)
}
}
}
fn plan_foreach_statement(
&self,
stmt: crate::expr::statements::ForeachStatement,
) -> Result<Arc<dyn ExecOperator>, Error> {
let crate::expr::statements::ForeachStatement {
param,
range,
block,
} = stmt;
Ok(Arc::new(ForeachPlan::new(param, range, block, self.current_depth()))
as Arc<dyn ExecOperator>)
}
fn plan_if_else_statement(
&self,
stmt: IfelseStatement,
) -> Result<Arc<dyn ExecOperator>, Error> {
let IfelseStatement {
exprs,
close,
} = stmt;
Ok(Arc::new(IfElsePlan::new(exprs, close, self.current_depth())) as Arc<dyn ExecOperator>)
}
fn plan_sleep_statement(
&self,
sleep_stmt: crate::expr::statements::SleepStatement,
) -> Result<Arc<dyn ExecOperator>, Error> {
Ok(Arc::new(SleepPlan::new(sleep_stmt.duration)))
}
async fn plan_expr_as_operator(&self, expr: Expr) -> Result<Arc<dyn ExecOperator>, Error> {
let phys_expr = Box::pin(self.physical_expr(expr)).await?;
Ok(Arc::new(ExprPlan::new(phys_expr)) as Arc<dyn ExecOperator>)
}
}
macro_rules! try_plan_expr {
($expr:expr, $ctx:expr, $txn:expr) => {{ $crate::exec::planner::try_plan_expr!($expr, $ctx, $txn, None) }};
($expr:expr, $ctx:expr, $txn:expr, $auth:expr) => {{ $crate::exec::planner::try_plan_expr!($expr, $ctx, $txn, $auth, 0u32) }};
($expr:expr, $ctx:expr, $txn:expr, $auth:expr, $depth:expr) => {{
let __expr: &$crate::expr::Expr = $expr;
if matches!(
__expr,
$crate::expr::Expr::Create(_)
| $crate::expr::Expr::Update(_)
| $crate::expr::Expr::Upsert(_)
| $crate::expr::Expr::Delete(_)
| $crate::expr::Expr::Insert(_)
| $crate::expr::Expr::Relate(_)
| $crate::expr::Expr::Define(_)
| $crate::expr::Expr::Remove(_)
| $crate::expr::Expr::Rebuild(_)
| $crate::expr::Expr::Alter(_)
) {
Err($crate::err::Error::PlannerUnsupported(String::new()))
} else if *$ctx.new_planner_strategy() == $crate::dbs::NewPlannerStrategy::ComputeOnly {
Err($crate::err::Error::PlannerUnsupported(String::new()))
} else {
$crate::exec::planner::plan_expr_inner(__expr, $ctx, $txn, $auth, $depth).await
}
}};
}
pub(crate) use try_plan_expr;
pub(crate) async fn plan_expr_inner(
expr: &Expr,
ctx: &FrozenContext,
txn: Arc<crate::kvs::Transaction>,
auth: Option<Arc<crate::iam::Auth>>,
depth: u32,
) -> Result<Arc<dyn ExecOperator>, Error> {
let ns =
ctx.value("session").and_then(|v| v.as_object()).and_then(|o| o.get("ns")).and_then(|v| {
match v {
crate::val::Value::String(s) => Some(s.as_str().to_owned()),
_ => None,
}
});
let db =
ctx.value("session").and_then(|v| v.as_object()).and_then(|o| o.get("db")).and_then(|v| {
match v {
crate::val::Value::String(s) => Some(s.as_str().to_owned()),
_ => None,
}
});
let mut planner = Planner::with_txn(ctx, txn, ns, db).with_depth(depth);
if let Some(auth) = auth {
planner = planner.with_auth(auth);
}
planner.plan(expr).await
}
pub(crate) async fn expr_to_physical_expr(
expr: Expr,
ctx: &FrozenContext,
) -> Result<Arc<dyn crate::exec::PhysicalExpr>, Error> {
Planner::new(ctx).physical_expr(expr).await
}
pub(crate) async fn expr_to_physical_expr_at_depth(
expr: Expr,
ctx: &FrozenContext,
depth: u32,
) -> Result<Arc<dyn crate::exec::PhysicalExpr>, Error> {
Planner::new(ctx).with_depth(depth).physical_expr(expr).await
}
pub(crate) fn is_simple_binary_eligible(op: &crate::expr::operator::BinaryOperator) -> bool {
use crate::expr::operator::BinaryOperator;
matches!(
op,
BinaryOperator::Equal
| BinaryOperator::ExactEqual
| BinaryOperator::NotEqual
| BinaryOperator::AllEqual
| BinaryOperator::AnyEqual
| BinaryOperator::LessThan
| BinaryOperator::LessThanEqual
| BinaryOperator::MoreThan
| BinaryOperator::MoreThanEqual
| BinaryOperator::Contain
| BinaryOperator::NotContain
| BinaryOperator::ContainAll
| BinaryOperator::ContainAny
| BinaryOperator::ContainNone
| BinaryOperator::Inside
| BinaryOperator::NotInside
| BinaryOperator::AllInside
| BinaryOperator::AnyInside
| BinaryOperator::NoneInside
| BinaryOperator::Outside
| BinaryOperator::Intersects
)
}
#[cfg(test)]
mod planner_tests {
use surrealdb_strand::Strand;
use super::*;
use crate::ctx::Context;
#[tokio::test]
async fn test_planner_creates_let_operator() {
let expr = Expr::Let(Box::new(crate::expr::statements::SetStatement {
name: Strand::new_static("x"),
what: Expr::Literal(crate::expr::literal::Literal::Integer(42)),
kind: None,
}));
let ctx = Arc::new(Context::new_test());
let plan = Planner::new(&ctx).plan(&expr).await.expect("Planning failed");
assert_eq!(plan.name(), "Let");
assert!(plan.mutates_context());
}
#[tokio::test]
async fn test_planner_creates_scalar_plan() {
let expr = Expr::Literal(crate::expr::literal::Literal::Integer(42));
let ctx = Arc::new(Context::new_test());
let plan = Planner::new(&ctx).plan(&expr).await.expect("Planning failed");
assert_eq!(plan.name(), "Expr");
assert!(plan.is_scalar());
}
}