use anyhow::Result;
use crate::catalog::providers::{DatabaseProvider, NamespaceProvider};
use crate::ctx::Context;
use crate::dbs::Variables;
use crate::exec::function::{FunctionRegistry, ScalarFunction, Signature};
use crate::exec::physical_expr::EvalContext;
use crate::exec::plan_or_compute::evaluate_expr_at_depth;
use crate::expr::{FlowResultExt as _, Kind};
use crate::fnc::args::{FromArgs, Optional};
use crate::fnc::eval::{Dialect, prepare};
use crate::val::{Object, Value};
async fn evaluate_streaming(
ctx: &EvalContext<'_>,
dialect: Dialect,
query: String,
bindings: Option<Object>,
) -> Result<Value> {
let exec_ctx = ctx.exec_ctx;
let caps = exec_ctx.capabilities();
let depth = ctx.plan_depth + 1;
let block = prepare(&caps, exec_ctx.auth(), dialect, &query)?;
let mut isolated = Context::new_isolated(exec_ctx.ctx());
if let Some(bindings) = bindings {
isolated.attach_variables(Variables::from(bindings))?;
}
let isolated = isolated.freeze();
let mut eval_ctx = exec_ctx.with_new_ctx(isolated);
if let Some(session) = exec_ctx.session()
&& let (Some(ns), Some(db)) = (session.ns.clone(), session.db.clone())
{
let txn = exec_ctx.txn();
let ns_def = txn.expect_ns_by_name(ns.as_str()).await?;
let db_def = txn.expect_db_by_name(ns.as_str(), db.as_str()).await?;
eval_ctx = eval_ctx.with_database(ns_def, db_def);
}
let mut result = Value::None;
for stmt in block.iter() {
result = evaluate_expr_at_depth(stmt, &eval_ctx, depth).await.catch_return()?;
}
Ok(result)
}
fn eval_signature() -> Signature {
Signature::new()
.arg("query", Kind::String)
.optional("bindings", Kind::Object)
.returns(Kind::Any)
}
#[derive(Debug, Clone, Copy, Default)]
pub struct EvalSurql;
impl ScalarFunction for EvalSurql {
fn name(&self) -> &'static str {
"eval::surql"
}
fn signature(&self) -> Signature {
eval_signature()
}
fn is_pure(&self) -> bool {
false
}
fn is_async(&self) -> bool {
true
}
fn invoke(&self, _args: Vec<Value>) -> Result<Value> {
Err(anyhow::anyhow!("Function '{}' requires async execution", self.name()))
}
fn invoke_async<'a>(
&'a self,
ctx: &'a EvalContext<'_>,
args: Vec<Value>,
) -> crate::exec::BoxFut<'a, Result<Value>> {
Box::pin(async move {
let (query, Optional(bindings)) = FromArgs::from_args("eval::surql", args)?;
evaluate_streaming(ctx, Dialect::Surql, query, bindings).await
})
}
}
#[derive(Debug, Clone, Copy, Default)]
pub struct EvalGql;
impl ScalarFunction for EvalGql {
fn name(&self) -> &'static str {
"eval::gql"
}
fn signature(&self) -> Signature {
eval_signature()
}
fn is_pure(&self) -> bool {
false
}
fn is_async(&self) -> bool {
true
}
fn invoke(&self, _args: Vec<Value>) -> Result<Value> {
Err(anyhow::anyhow!("Function '{}' requires async execution", self.name()))
}
fn invoke_async<'a>(
&'a self,
ctx: &'a EvalContext<'_>,
args: Vec<Value>,
) -> crate::exec::BoxFut<'a, Result<Value>> {
Box::pin(async move {
let (query, Optional(bindings)) = FromArgs::from_args("eval::gql", args)?;
evaluate_streaming(ctx, Dialect::Gql, query, bindings).await
})
}
}
pub fn register(registry: &mut FunctionRegistry) {
registry.register(EvalSurql);
registry.register(EvalGql);
}