use std::borrow::Cow;
use std::collections::HashSet;
use std::sync::Arc;
use async_channel::Sender;
#[allow(unused_imports)]
pub(crate) use surrealdb_engine_api::{
Command, EngineContext, MlExportConfig, RequestData, Route, RouteChannelEngine, SurrealEngine,
single_result,
};
use surrealdb_rpc::QueryResult;
use uuid::Uuid;
use super::opt::Config;
use crate::method::BoxFuture;
use crate::opt::Endpoint;
use crate::types::{SurrealValue, Value, Variables};
use crate::{Error, ExtraFeatures, Result, Surreal};
#[derive(Debug, Clone)]
pub struct Router {
pub(crate) engine: Arc<dyn SurrealEngine>,
#[allow(dead_code)]
pub(crate) config: Config,
pub(crate) features: HashSet<ExtraFeatures>,
}
#[derive(Debug)]
pub(crate) struct QueryRequest {
pub(crate) txn: Option<Uuid>,
pub(crate) query: Cow<'static, str>,
pub(crate) variables: Variables,
}
impl Router {
pub(crate) fn run_query_opt<R>(
&self,
session: Uuid,
request: QueryRequest,
) -> BoxFuture<'_, Result<Option<R>>>
where
R: SurrealValue,
{
self.query_opt(ctx_txn(session, request.txn), request.query, request.variables)
}
pub(crate) fn run_query_vec<R>(
&self,
session: Uuid,
request: QueryRequest,
) -> BoxFuture<'_, Result<Vec<R>>>
where
R: SurrealValue,
{
self.query_vec(ctx_txn(session, request.txn), request.query, request.variables)
}
pub(crate) fn run_query_value(
&self,
session: Uuid,
request: QueryRequest,
) -> BoxFuture<'_, Result<Value>> {
self.query_value(ctx_txn(session, request.txn), request.query, request.variables)
}
pub(crate) fn from_route_sender(
sender: Sender<Route>,
features: HashSet<ExtraFeatures>,
config: Config,
) -> Self {
Self {
engine: Arc::new(RouteChannelEngine::new(sender)),
config,
features,
}
}
#[cfg_attr(
not(any(
feature = "protocol-grpc",
feature = "kv-mem",
feature = "kv-tikv",
feature = "kv-rocksdb",
feature = "kv-indxdb",
feature = "kv-surrealkv",
)),
allow(dead_code)
)]
pub(crate) fn from_engine(
engine: Arc<dyn SurrealEngine>,
features: HashSet<ExtraFeatures>,
config: Config,
) -> Self {
Self {
engine,
config,
features,
}
}
pub(crate) fn query_value(
&self,
ctx: EngineContext,
query: Cow<'static, str>,
variables: Variables,
) -> BoxFuture<'_, Result<Value>> {
Box::pin(async move { single_result(self.engine.query(ctx, query, variables).await?) })
}
pub(crate) fn query_opt<R>(
&self,
ctx: EngineContext,
query: Cow<'static, str>,
variables: Variables,
) -> BoxFuture<'_, Result<Option<R>>>
where
R: SurrealValue,
{
Box::pin(async move {
match self.query_value(ctx, query, variables).await? {
Value::None | Value::Null => Ok(None),
Value::Array(array) => match array.len() {
0 => Ok(None),
1 => {
let value =
array.into_iter().next().expect("array has exactly one element");
Ok(Some(R::from_value(value).map_err(deserialization_error)?))
}
_ => {
Ok(Some(R::from_value(Value::Array(array)).map_err(deserialization_error)?))
}
},
value => Ok(Some(R::from_value(value).map_err(deserialization_error)?)),
}
})
}
pub(crate) fn query_vec<R>(
&self,
ctx: EngineContext,
query: Cow<'static, str>,
variables: Variables,
) -> BoxFuture<'_, Result<Vec<R>>>
where
R: SurrealValue,
{
Box::pin(async move {
match self.query_value(ctx, query, variables).await? {
Value::None | Value::Null => Ok(Vec::new()),
Value::Array(array) => array
.into_iter()
.map(|value| R::from_value(value).map_err(deserialization_error))
.collect(),
value => Ok(vec![R::from_value(value).map_err(deserialization_error)?]),
}
})
}
pub(crate) fn query_results(
&self,
ctx: EngineContext,
query: Cow<'static, str>,
variables: Variables,
) -> BoxFuture<'_, Result<Vec<QueryResult>>> {
Box::pin(async move { self.engine.query(ctx, query, variables).await })
}
}
fn deserialization_error(error: impl std::fmt::Display) -> Error {
Error::serialization(error.to_string(), crate::types::SerializationError::Deserialization)
}
pub(crate) fn ctx(session: Uuid) -> EngineContext {
EngineContext::new(session)
}
pub(crate) fn ctx_txn(session: Uuid, txn: Option<Uuid>) -> EngineContext {
EngineContext::with_transaction(session, txn)
}
pub trait Sealed: Sized + Send + Sync + 'static {
#[allow(private_interfaces)]
fn connect(
address: Endpoint,
capacity: usize,
session_clone: Option<crate::SessionClone>,
) -> BoxFuture<'static, Result<Surreal<Self>>>
where
Self: crate::Connection;
}