use std::collections::HashMap;
use std::sync::Arc;
use tokio_util::sync::CancellationToken;
use crate::catalog::providers::{CatalogProvider, DatabaseProvider, NamespaceProvider};
use crate::ctx::Context;
use crate::dbs::Session;
use crate::exec::context::{DatabaseContext, NamespaceContext, RootContext, SessionInfo};
use crate::exec::function::FunctionRegistry;
use crate::exec::{
AccessMode, CardinalityHint, ContextLevel, EvalContext, ExecOperator, ExecutionContext,
FlowResult, OutputOrdering, PhysicalExpr, ValueBatch, ValueBatchStream,
};
use crate::iam::Auth;
use crate::kvs::{Datastore, TransactionType};
use crate::val::Value;
#[derive(Debug)]
pub(crate) struct ValuesOperator {
values: Vec<Value>,
cardinality: CardinalityHint,
ordering: OutputOrdering,
}
impl ValuesOperator {
#[allow(clippy::new_ret_no_self)]
pub(crate) fn new(values: Vec<Value>) -> Arc<dyn ExecOperator> {
Arc::new(Self {
values,
cardinality: CardinalityHint::Unbounded,
ordering: OutputOrdering::Unordered,
})
}
}
impl ExecOperator for ValuesOperator {
fn name(&self) -> &'static str {
"Values"
}
fn required_context(&self) -> ContextLevel {
ContextLevel::Root
}
fn access_mode(&self) -> AccessMode {
AccessMode::ReadOnly
}
fn cardinality_hint(&self) -> CardinalityHint {
self.cardinality
}
fn output_ordering(&self) -> OutputOrdering {
self.ordering.clone()
}
fn execute(&self, _ctx: &ExecutionContext) -> FlowResult<ValueBatchStream> {
let values = self.values.clone();
Ok(Box::pin(futures::stream::once(std::future::ready(Ok(ValueBatch::new(values))))))
}
}
pub(crate) fn root_ctx() -> ExecutionContext {
ExecutionContext::Root(RootContext {
ctx: Context::new_test().freeze(),
function_registry: Arc::new(FunctionRegistry::with_builtins()),
options: None,
datastore: None,
cancellation: CancellationToken::new(),
auth: Arc::new(Auth::default()),
session: None,
current_value: None,
skip_fetch_perms: false,
computing_field: false,
version_stamp: None,
})
}
pub(crate) fn root_ctx_with_auth(auth: Auth) -> ExecutionContext {
let ExecutionContext::Root(mut root) = root_ctx() else {
unreachable!("root_ctx builds a Root context")
};
root.auth = Arc::new(auth);
ExecutionContext::Root(root)
}
pub(crate) async fn collect(op: &Arc<dyn ExecOperator>, ctx: &ExecutionContext) -> Vec<Value> {
use futures::StreamExt;
let mut stream = op.execute(ctx).expect("execute should succeed");
let mut out = Vec::new();
while let Some(batch) = stream.next().await {
out.extend(batch.expect("batch should be Ok").into_values());
}
out
}
pub(crate) async fn try_collect(
op: &Arc<dyn ExecOperator>,
ctx: &ExecutionContext,
) -> FlowResult<Vec<Value>> {
use futures::StreamExt;
let mut stream = op.execute(ctx)?;
let mut out = Vec::new();
while let Some(batch) = stream.next().await {
out.extend(batch?.into_values());
}
Ok(out)
}
pub(crate) fn parse_expr(src: &str) -> crate::expr::Expr {
let ast = crate::syn::parse(&format!("RETURN {src};")).expect("fragment should parse");
let mut exprs = ast.expressions;
assert_eq!(exprs.len(), 1, "expected exactly one statement in {src:?}");
let top: crate::expr::TopLevelExpr = exprs.remove(0).into();
match top {
crate::expr::TopLevelExpr::Expr(crate::expr::Expr::Return(ret)) => ret.what.clone(),
other => panic!("unexpected statement shape for {src:?}: {other:?}"),
}
}
pub(crate) fn parse_idiom(src: &str) -> crate::expr::Idiom {
match parse_expr(src) {
crate::expr::Expr::Idiom(idiom) => idiom,
other => panic!("expected an idiom for {src:?}, got {other:?}"),
}
}
pub(crate) async fn physical_expr(src: &str, ctx: &ExecutionContext) -> Arc<dyn PhysicalExpr> {
crate::exec::planner::expr_to_physical_expr(parse_expr(src), ctx.ctx(), ctx.function_registry())
.await
.expect("fragment should compile to a physical expression")
}
pub(crate) async fn eval_on(src: &str, value: &Value, ctx: &ExecutionContext) -> FlowResult<Value> {
let expr = physical_expr(src, ctx).await;
expr.evaluate(EvalContext::from_exec_ctx(ctx).with_value_and_doc(value)).await
}
pub(crate) async fn eval(src: &str, ctx: &ExecutionContext) -> FlowResult<Value> {
let expr = physical_expr(src, ctx).await;
expr.evaluate(EvalContext::from_exec_ctx(ctx)).await
}
pub(crate) async fn val(src: &str) -> Value {
let ctx = root_ctx();
eval(src, &ctx).await.expect("literal should evaluate")
}
pub(crate) struct TestDb {
ds: Arc<Datastore>,
}
impl TestDb {
pub(crate) async fn new(setup: &str) -> Self {
Self::build(setup, false, surrealdb_cnf::ConfigMap::empty()).await
}
pub(crate) async fn new_with_auth(setup: &str) -> Self {
Self::build(setup, true, surrealdb_cnf::ConfigMap::empty()).await
}
#[cfg(feature = "kv-mem")]
pub(crate) async fn new_with_config(setup: &str, config: surrealdb_cnf::ConfigMap) -> Self {
Self::build(setup, false, config).await
}
async fn build(setup: &str, auth_enabled: bool, config: surrealdb_cnf::ConfigMap) -> Self {
let ds = Datastore::builder()
.with_capabilities(crate::dbs::Capabilities::all())
.with_auth(auth_enabled)
.with_config(config)
.build_with_path("memory")
.await
.expect("in-memory datastore");
{
let txn = ds.transaction(TransactionType::Write).await.expect("write transaction");
txn.ensure_ns_db(None, "test", "test").await.expect("ensure test/test");
txn.commit().await.expect("commit ns/db");
}
let db = Self {
ds,
};
if !setup.trim().is_empty() {
db.run(setup).await;
}
db
}
pub(crate) fn owner() -> Session {
Session::owner().with_ns("test").with_db("test")
}
pub(crate) async fn run(&self, sql: &str) {
self.run_as(&Self::owner(), sql).await
}
pub(crate) async fn run_as(&self, session: &Session, sql: &str) {
for response in self.ds.execute(sql, session, None).await.expect("query should execute") {
response.result.expect("statement should succeed");
}
}
pub(crate) async fn exec_ctx(&self) -> ExecutionContext {
self.exec_ctx_as(&Self::owner(), TransactionType::Read).await
}
pub(crate) async fn exec_ctx_as(
&self,
session: &Session,
mode: TransactionType,
) -> ExecutionContext {
let txn = Arc::new(self.ds.transaction(mode).await.expect("transaction"));
let mut ctx = self.ds.setup_ctx().expect("context from datastore");
ctx.attach_session(session).expect("attach session");
ctx.set_transaction(Arc::clone(&txn));
let ns = txn.expect_ns_by_name("test").await.expect("namespace test");
let db = txn.expect_db_by_name("test", "test").await.expect("database test");
let root = RootContext {
ctx: ctx.freeze(),
function_registry: Arc::new(FunctionRegistry::with_builtins()),
options: Some(self.ds.setup_options(session)),
datastore: Some(Arc::clone(&self.ds)),
cancellation: CancellationToken::new(),
auth: Arc::clone(&session.au),
session: Some(Arc::new(session_info(session))),
current_value: None,
skip_fetch_perms: false,
computing_field: false,
version_stamp: None,
};
ExecutionContext::Database(DatabaseContext {
ns_ctx: NamespaceContext {
root,
ns,
},
db,
field_state_cache: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
table_def_cache: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
index_def_cache: Arc::new(tokio::sync::RwLock::new(HashMap::new())),
})
}
}
fn session_info(session: &Session) -> SessionInfo {
use crate::val::convert_public::convert_public_value_to_internal;
SessionInfo {
ns: session.ns.as_deref().map(Into::into),
db: session.db.as_deref().map(Into::into),
id: session.id,
ip: session.ip.as_deref().map(Into::into),
origin: session.or.as_deref().map(Into::into),
ac: session.ac.as_deref().map(Into::into),
rd: session.rd.clone().map(convert_public_value_to_internal),
token: session.tk.clone().map(convert_public_value_to_internal),
exp: None,
}
}
pub(crate) async fn drain_err(op: &impl ExecOperator, ctx: &ExecutionContext) -> anyhow::Error {
use futures::StreamExt;
let unwrap = |ctrl| match ctrl {
crate::expr::ControlFlow::Err(e) => e,
other => panic!("expected an error, got the control-flow signal: {other:?}"),
};
let mut stream = match op.execute(ctx) {
Ok(stream) => stream,
Err(ctrl) => return unwrap(ctrl),
};
match stream.next().await {
Some(Ok(batch)) => panic!("expected an error, got rows: {:?}", batch.values()),
Some(Err(ctrl)) => unwrap(ctrl),
None => panic!("expected an error, got an empty stream"),
}
}