Skip to main content

teaql_runtime/data_service/
context.rs

1// Compatibility adapters retained while generated callers migrate to EntityDataService.
2#![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    /// Register a data-service executor and automatically set up the
63    /// type-erased graph saver so that
64    /// [`Audited::save`](crate::AuditedSaveExt::save) works.
65    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}