1use std::fmt::Write;
2use std::sync::Arc;
3use std::{collections::BTreeMap, future::Future};
4
5use teaql_core::{
6 CompactRow, DeleteCommand, Entity, InsertCommand, RecoverCommand, SelectQuery, SmartList,
7 UpdateCommand,
8};
9
10use crate::{
11 ContextError, DataServiceError, GraphMutationPlan, GraphNode, RuntimeError, UserContext,
12};
13
14use super::{
15 AggregationCacheBackend, ContextDataService, EntityDataService, InMemoryAggregationCache,
16 RuntimeDataService, UserContextMetadata, helpers::invalidate_aggregation_cache_namespace,
17};
18
19impl UserContext {
20 pub(crate) fn data_service_internal<E>(&self) -> Result<ContextDataService<'_, E>, ContextError>
21 where
22 E: teaql_data_service::QueryExecutor
23 + teaql_data_service::MutationExecutor
24 + Send
25 + Sync
26 + 'static,
27 {
28 if self.metadata.is_none() {
29 return Err(ContextError::MissingResource("metadata".to_owned()));
30 }
31
32 let executor = self.require_resource::<E>()?;
33 Ok(ContextDataService {
34 metadata: UserContextMetadata { context: self },
35 executor,
36 })
37 }
38
39 pub fn entity_data_service<E>(
40 &self,
41 entity: impl Into<String>,
42 ) -> Result<EntityDataService<'_, E>, ContextError>
43 where
44 E: teaql_data_service::QueryExecutor
45 + teaql_data_service::MutationExecutor
46 + Send
47 + Sync
48 + 'static,
49 {
50 let entity = entity.into();
51 if !self.has_entity_data_service(&entity) {
52 return Err(ContextError::MissingEntityDataService(entity));
53 }
54 Ok(EntityDataService {
55 entity,
56 data_service: self.data_service_internal::<E>()?,
57 trace_context: Vec::new(),
58 })
59 }
60
61 pub fn register_executor<E>(&mut self, executor: E)
65 where
66 E: teaql_data_service::QueryExecutor
67 + teaql_data_service::MutationExecutor
68 + teaql_data_service::TransactionExecutor
69 + Send
70 + Sync
71 + 'static,
72 for<'tx> <E as teaql_data_service::TransactionExecutor>::Tx<'tx>: Send + Sync,
73 {
74 use std::sync::Arc;
75 self.insert_resource::<Arc<dyn crate::entity_save::DynGraphSaver>>(Arc::new(
76 crate::entity_save::GraphSaverFor::<E>::new(),
77 ));
78 self.insert_resource(executor);
79 }
80}
81
82impl<'a, E> ContextDataService<'a, E>
83where
84 E: teaql_data_service::QueryExecutor + teaql_data_service::MutationExecutor + Send + Sync,
85{
86 async fn observe<T, F, N>(
87 &self,
88 family: &str,
89 name: N,
90 entity: &str,
91 work: F,
92 ) -> Result<T, DataServiceError<E::Error>>
93 where
94 F: Future<Output = Result<T, DataServiceError<E::Error>>>,
95 N: FnOnce() -> String,
96 {
97 if self.metadata.context.runtime_telemetry_is_noop() {
98 return work.await;
99 }
100 let operation = crate::RuntimeOperation::new(family, name())
101 .attribute("teaql.entity.type", entity.to_owned());
102 let scope = self.metadata.context.start_runtime_operation(operation);
103 let provider_kind = std::any::type_name::<E>().to_owned();
104 let provider_operation = family.to_owned();
105 let result = scope
106 .run(async {
107 let provider_scope = self.metadata.context.start_runtime_operation(
108 crate::RuntimeOperation::new(
109 "provider",
110 format!("{provider_kind}.{provider_operation}"),
111 )
112 .attribute("teaql.provider.kind", provider_kind)
113 .attribute("teaql.provider.operation", provider_operation),
114 );
115 let result = provider_scope.run(work).await;
116 match &result {
117 Ok(_) => provider_scope.success(BTreeMap::new()),
118 Err(_) => provider_scope.failure("data_service_error"),
119 }
120 result
121 })
122 .await;
123 match result {
124 Ok(value) => {
125 scope.success(BTreeMap::new());
126 Ok(value)
127 }
128 Err(error) => {
129 scope.failure("data_service_error");
130 Err(error)
131 }
132 }
133 }
134
135 fn data_service(&self) -> RuntimeDataService<'_, UserContextMetadata<'_>, E> {
136 RuntimeDataService::new(&self.metadata, self.executor)
137 }
138
139 pub(crate) async fn fetch_all(
140 &self,
141 mut query: SelectQuery,
142 ) -> Result<Vec<CompactRow>, DataServiceError<E::Error>> {
143 let final_comment = self.resolve_final_comment(&query.trace_chain, query.comment.clone());
144 query.comment = final_comment;
145 self.observe(
146 "query",
147 || format!("{}.list", query.entity),
148 &query.entity,
149 self.data_service().fetch_all(&query),
150 )
151 .await
152 }
153
154 pub(crate) async fn fetch_smart_list(
155 &self,
156 query: &SelectQuery,
157 ) -> Result<SmartList<CompactRow>, DataServiceError<E::Error>> {
158 self.observe(
159 "query",
160 || format!("{}.list", query.entity),
161 &query.entity,
162 self.data_service().fetch_smart_list(query),
163 )
164 .await
165 }
166
167 pub(crate) async fn fetch_entities<T>(
168 &self,
169 query: &SelectQuery,
170 ) -> Result<SmartList<T>, DataServiceError<E::Error>>
171 where
172 T: Entity,
173 {
174 self.observe(
175 "query",
176 || format!("{}.list", query.entity),
177 &query.entity,
178 self.data_service().fetch_entities(query),
179 )
180 .await
181 }
182
183 pub(crate) async fn fetch_enhanced_entities<T>(
184 &self,
185 query: &SelectQuery,
186 ) -> Result<SmartList<T>, DataServiceError<E::Error>>
187 where
188 T: Entity,
189 {
190 self.observe(
191 "query",
192 || format!("{}.list", query.entity),
193 &query.entity,
194 self.data_service().fetch_enhanced_entities(query),
195 )
196 .await
197 }
198
199 pub(crate) async fn insert(
200 &self,
201 command: &InsertCommand,
202 ) -> Result<u64, DataServiceError<E::Error>> {
203 let affected = self
204 .observe(
205 "mutation",
206 || format!("{}.insert", command.entity),
207 &command.entity,
208 self.data_service().insert(command),
209 )
210 .await?;
211 self.invalidate_aggregation_cache_for(&command.entity);
212 Ok(affected)
213 }
214
215 pub(crate) async fn update(
216 &self,
217 command: &UpdateCommand,
218 ) -> Result<u64, DataServiceError<E::Error>> {
219 let affected = self
220 .observe(
221 "mutation",
222 || format!("{}.update", command.entity),
223 &command.entity,
224 self.data_service().update(command),
225 )
226 .await?;
227 self.invalidate_aggregation_cache_for(&command.entity);
228 Ok(affected)
229 }
230
231 pub(crate) async fn batch_insert(
232 &self,
233 command: &teaql_core::BatchInsertCommand,
234 ) -> Result<u64, DataServiceError<E::Error>> {
235 let affected = self
236 .observe(
237 "mutation",
238 || format!("{}.batch_insert", command.entity),
239 &command.entity,
240 self.data_service().batch_insert(command),
241 )
242 .await?;
243 self.invalidate_aggregation_cache_for(&command.entity);
244 Ok(affected)
245 }
246
247 pub(crate) async fn batch_update(
248 &self,
249 command: &teaql_core::BatchUpdateCommand,
250 ) -> Result<u64, DataServiceError<E::Error>> {
251 let affected = self
252 .observe(
253 "mutation",
254 || format!("{}.batch_update", command.entity),
255 &command.entity,
256 self.data_service().batch_update(command),
257 )
258 .await?;
259 self.invalidate_aggregation_cache_for(&command.entity);
260 Ok(affected)
261 }
262
263 pub(crate) async fn delete(
264 &self,
265 command: &DeleteCommand,
266 ) -> Result<u64, DataServiceError<E::Error>> {
267 let affected = self
268 .observe(
269 "mutation",
270 || format!("{}.delete", command.entity),
271 &command.entity,
272 self.data_service().delete(command),
273 )
274 .await?;
275 self.invalidate_aggregation_cache_for(&command.entity);
276 Ok(affected)
277 }
278
279 pub(crate) async fn recover(
280 &self,
281 command: &RecoverCommand,
282 ) -> Result<u64, DataServiceError<E::Error>> {
283 let affected = self
284 .observe(
285 "mutation",
286 || format!("{}.recover", command.entity),
287 &command.entity,
288 self.data_service().recover(command),
289 )
290 .await?;
291 self.invalidate_aggregation_cache_for(&command.entity);
292 Ok(affected)
293 }
294
295 pub(super) fn invalidate_aggregation_cache_for(&self, entity: &str) {
296 if let Some(cache) = self
297 .metadata
298 .context
299 .get_resource::<Arc<dyn AggregationCacheBackend>>()
300 {
301 invalidate_aggregation_cache_namespace(cache.as_ref(), entity);
302 }
303 if let Some(cache) = self
304 .metadata
305 .context
306 .get_resource::<InMemoryAggregationCache>()
307 {
308 invalidate_aggregation_cache_namespace(cache, entity);
309 }
310 }
311
312 pub(crate) fn resolve_final_comment(
313 &self,
314 trace_chain: &[teaql_core::TraceNode],
315 comment: Option<String>,
316 ) -> Option<String> {
317 let chain_str = (!trace_chain.is_empty()).then(|| {
318 let mut chain = String::with_capacity(trace_chain.len().saturating_mul(64));
319 for (index, node) in trace_chain.iter().enumerate() {
320 if index > 0 {
321 chain.push_str(" -> ");
322 }
323 match node.entity_id {
324 Some(id) => {
325 let _ = write!(chain, "{}({id}): {}", node.entity_type, node.comment);
326 }
327 None => {
328 let _ = write!(chain, "{}(pending): {}", node.entity_type, node.comment);
329 }
330 }
331 }
332 chain
333 });
334
335 let business_comment = chain_str.or(comment);
336 let user_id = self
337 .metadata
338 .context
339 .user_identifier()
340 .map(|s| s.to_owned());
341
342 match (user_id, business_comment) {
343 (Some(user), Some(bus)) if !user.is_empty() && !bus.is_empty() => {
344 Some(format!("[{user}] {bus}"))
345 }
346 (Some(user), _) if !user.is_empty() => Some(format!("[{user}]")),
347 (_, Some(bus)) if !bus.is_empty() => Some(bus),
348 _ => None,
349 }
350 }
351}