use std::fmt::Write;
use std::sync::Arc;
use std::{collections::BTreeMap, future::Future};
use teaql_core::{
CompactRow, DeleteCommand, Entity, InsertCommand, RecoverCommand, SelectQuery, SmartList,
UpdateCommand,
};
use crate::{
ContextError, DataServiceError, GraphMutationPlan, GraphNode, RuntimeError, UserContext,
};
use super::{
AggregationCacheBackend, ContextDataService, EntityDataService, InMemoryAggregationCache,
RuntimeDataService, UserContextMetadata, helpers::invalidate_aggregation_cache_namespace,
};
impl UserContext {
pub(crate) fn data_service_internal<E>(&self) -> Result<ContextDataService<'_, E>, ContextError>
where
E: teaql_data_service::QueryExecutor
+ teaql_data_service::MutationExecutor
+ Send
+ Sync
+ 'static,
{
if self.metadata.is_none() {
return Err(ContextError::MissingResource("metadata".to_owned()));
}
let executor = self.require_resource::<E>()?;
Ok(ContextDataService {
metadata: UserContextMetadata { context: self },
executor,
})
}
pub fn entity_data_service<E>(
&self,
entity: impl Into<String>,
) -> Result<EntityDataService<'_, E>, ContextError>
where
E: teaql_data_service::QueryExecutor
+ teaql_data_service::MutationExecutor
+ Send
+ Sync
+ 'static,
{
let entity = entity.into();
if !self.has_entity_data_service(&entity) {
return Err(ContextError::MissingEntityDataService(entity));
}
Ok(EntityDataService {
entity,
data_service: self.data_service_internal::<E>()?,
trace_context: Vec::new(),
})
}
pub fn register_executor<E>(&mut self, executor: E)
where
E: teaql_data_service::QueryExecutor
+ teaql_data_service::MutationExecutor
+ teaql_data_service::TransactionExecutor
+ Send
+ Sync
+ 'static,
for<'tx> <E as teaql_data_service::TransactionExecutor>::Tx<'tx>: Send + Sync,
{
use std::sync::Arc;
self.insert_resource::<Arc<dyn crate::entity_save::DynGraphSaver>>(Arc::new(
crate::entity_save::GraphSaverFor::<E>::new(),
));
self.insert_resource(executor);
}
}
impl<'a, E> ContextDataService<'a, E>
where
E: teaql_data_service::QueryExecutor + teaql_data_service::MutationExecutor + Send + Sync,
{
async fn observe<T, F, N>(
&self,
family: &str,
name: N,
entity: &str,
work: F,
) -> Result<T, DataServiceError<E::Error>>
where
F: Future<Output = Result<T, DataServiceError<E::Error>>>,
N: FnOnce() -> String,
{
if self.metadata.context.runtime_telemetry_is_noop() {
return work.await;
}
let operation = crate::RuntimeOperation::new(family, name())
.attribute("teaql.entity.type", entity.to_owned());
let scope = self.metadata.context.start_runtime_operation(operation);
let provider_kind = std::any::type_name::<E>().to_owned();
let provider_operation = family.to_owned();
let result = scope
.run(async {
let provider_scope = self.metadata.context.start_runtime_operation(
crate::RuntimeOperation::new(
"provider",
format!("{provider_kind}.{provider_operation}"),
)
.attribute("teaql.provider.kind", provider_kind)
.attribute("teaql.provider.operation", provider_operation),
);
let result = provider_scope.run(work).await;
match &result {
Ok(_) => provider_scope.success(BTreeMap::new()),
Err(_) => provider_scope.failure("data_service_error"),
}
result
})
.await;
match result {
Ok(value) => {
scope.success(BTreeMap::new());
Ok(value)
}
Err(error) => {
scope.failure("data_service_error");
Err(error)
}
}
}
fn data_service(&self) -> RuntimeDataService<'_, UserContextMetadata<'_>, E> {
RuntimeDataService::new(&self.metadata, self.executor)
}
pub(crate) async fn fetch_all(
&self,
mut query: SelectQuery,
) -> Result<Vec<CompactRow>, DataServiceError<E::Error>> {
let final_comment = self.resolve_final_comment(&query.trace_chain, query.comment.clone());
query.comment = final_comment;
self.observe(
"query",
|| format!("{}.list", query.entity),
&query.entity,
self.data_service().fetch_all(&query),
)
.await
}
pub(crate) async fn fetch_smart_list(
&self,
query: &SelectQuery,
) -> Result<SmartList<CompactRow>, DataServiceError<E::Error>> {
self.observe(
"query",
|| format!("{}.list", query.entity),
&query.entity,
self.data_service().fetch_smart_list(query),
)
.await
}
pub(crate) async fn fetch_entities<T>(
&self,
query: &SelectQuery,
) -> Result<SmartList<T>, DataServiceError<E::Error>>
where
T: Entity,
{
self.observe(
"query",
|| format!("{}.list", query.entity),
&query.entity,
self.data_service().fetch_entities(query),
)
.await
}
pub(crate) async fn fetch_enhanced_entities<T>(
&self,
query: &SelectQuery,
) -> Result<SmartList<T>, DataServiceError<E::Error>>
where
T: Entity,
{
self.observe(
"query",
|| format!("{}.list", query.entity),
&query.entity,
self.data_service().fetch_enhanced_entities(query),
)
.await
}
pub(crate) async fn insert(
&self,
command: &InsertCommand,
) -> Result<u64, DataServiceError<E::Error>> {
let affected = self
.observe(
"mutation",
|| format!("{}.insert", command.entity),
&command.entity,
self.data_service().insert(command),
)
.await?;
self.invalidate_aggregation_cache_for(&command.entity);
Ok(affected)
}
pub(crate) async fn update(
&self,
command: &UpdateCommand,
) -> Result<u64, DataServiceError<E::Error>> {
let affected = self
.observe(
"mutation",
|| format!("{}.update", command.entity),
&command.entity,
self.data_service().update(command),
)
.await?;
self.invalidate_aggregation_cache_for(&command.entity);
Ok(affected)
}
pub(crate) async fn batch_insert(
&self,
command: &teaql_core::BatchInsertCommand,
) -> Result<u64, DataServiceError<E::Error>> {
let affected = self
.observe(
"mutation",
|| format!("{}.batch_insert", command.entity),
&command.entity,
self.data_service().batch_insert(command),
)
.await?;
self.invalidate_aggregation_cache_for(&command.entity);
Ok(affected)
}
pub(crate) async fn batch_update(
&self,
command: &teaql_core::BatchUpdateCommand,
) -> Result<u64, DataServiceError<E::Error>> {
let affected = self
.observe(
"mutation",
|| format!("{}.batch_update", command.entity),
&command.entity,
self.data_service().batch_update(command),
)
.await?;
self.invalidate_aggregation_cache_for(&command.entity);
Ok(affected)
}
pub(crate) async fn delete(
&self,
command: &DeleteCommand,
) -> Result<u64, DataServiceError<E::Error>> {
let affected = self
.observe(
"mutation",
|| format!("{}.delete", command.entity),
&command.entity,
self.data_service().delete(command),
)
.await?;
self.invalidate_aggregation_cache_for(&command.entity);
Ok(affected)
}
pub(crate) async fn recover(
&self,
command: &RecoverCommand,
) -> Result<u64, DataServiceError<E::Error>> {
let affected = self
.observe(
"mutation",
|| format!("{}.recover", command.entity),
&command.entity,
self.data_service().recover(command),
)
.await?;
self.invalidate_aggregation_cache_for(&command.entity);
Ok(affected)
}
pub(super) fn invalidate_aggregation_cache_for(&self, entity: &str) {
if let Some(cache) = self
.metadata
.context
.get_resource::<Arc<dyn AggregationCacheBackend>>()
{
invalidate_aggregation_cache_namespace(cache.as_ref(), entity);
}
if let Some(cache) = self
.metadata
.context
.get_resource::<InMemoryAggregationCache>()
{
invalidate_aggregation_cache_namespace(cache, entity);
}
}
pub(crate) fn resolve_final_comment(
&self,
trace_chain: &[teaql_core::TraceNode],
comment: Option<String>,
) -> Option<String> {
let chain_str = (!trace_chain.is_empty()).then(|| {
let mut chain = String::with_capacity(trace_chain.len().saturating_mul(64));
for (index, node) in trace_chain.iter().enumerate() {
if index > 0 {
chain.push_str(" -> ");
}
match node.entity_id {
Some(id) => {
let _ = write!(chain, "{}({id}): {}", node.entity_type, node.comment);
}
None => {
let _ = write!(chain, "{}(pending): {}", node.entity_type, node.comment);
}
}
}
chain
});
let business_comment = chain_str.or(comment);
let user_id = self
.metadata
.context
.user_identifier()
.map(|s| s.to_owned());
match (user_id, business_comment) {
(Some(user), Some(bus)) if !user.is_empty() && !bus.is_empty() => {
Some(format!("[{user}] {bus}"))
}
(Some(user), _) if !user.is_empty() => Some(format!("[{user}]")),
(_, Some(bus)) if !bus.is_empty() => Some(bus),
_ => None,
}
}
}