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