Skip to main content

teaql_runtime/data_service/
resolved.rs

1use std::collections::{BTreeMap, BTreeSet};
2use std::hash::{Hash, Hasher};
3use std::sync::Arc;
4use std::time::{Duration, SystemTime};
5
6use teaql_core::{
7    AggregationCacheOptions, DeleteCommand, Entity, Expr, InsertCommand, Record, RecoverCommand,
8    RelationAggregate, SelectQuery, SmartList, SortDirection, UpdateCommand, Value,
9};
10
11use crate::{
12    clear_entity_status, mark_entity_status, CheckObjectStatus, ContinuousPageCursor,
13    DataServiceError, EntityDataServiceBehavior, MetadataStore, PurposedSelectQuery, RawAuditEvent,
14    RuntimeError,
15};
16
17use super::{
18    helpers::*, AggregationCacheBackend, ContextDataService, EntityDataService,
19    InMemoryAggregationCache, UserContextMetadata,
20};
21
22#[derive(Debug, Clone)]
23struct ContinuousPageExecution {
24    query_key: String,
25    direction: SortDirection,
26    page_size: u64,
27    original_offset: u64,
28    ttl_seconds: u64,
29    optimized: bool,
30    seek_cursor_id: Option<String>,
31}
32
33impl<'a, E> EntityDataService<'a, E>
34where
35    E: teaql_data_service::QueryExecutor
36        + teaql_data_service::MutationExecutor
37        + Send
38        + Sync
39        + 'static,
40{
41    fn flatten_relation_graph(
42        &self,
43        entity_name: &str,
44        record: &mut BTreeMap<String, Value>,
45        root: &crate::EntityRoot,
46        graph: &mut crate::EntityGraphBuilder,
47        installed: &mut BTreeSet<(String, u64)>,
48    ) -> Result<(), teaql_core::EntityError> {
49        let context = self.data_service.metadata.context;
50        let relations = context
51            .entity(entity_name)
52            .map(|descriptor| descriptor.relations.clone())
53            .unwrap_or_default();
54
55        for relation in relations {
56            if !context.has_entity_graph_decoder(&relation.target_entity) {
57                continue;
58            }
59            let Some(value) = record.remove(&relation.name) else {
60                continue;
61            };
62            if !relation.many && matches!(value, Value::Null | Value::TypedNull(_)) {
63                record.insert(relation.name, value);
64                continue;
65            }
66            let mut child_records = match value {
67                Value::Object(child) => vec![child],
68                Value::List(values) => values
69                    .into_iter()
70                    .filter_map(|value| match value {
71                        Value::Object(child) => Some(child),
72                        _ => None,
73                    })
74                    .collect(),
75                Value::Null | Value::TypedNull(_) => Vec::new(),
76                other => {
77                    record.insert(relation.name, other);
78                    continue;
79                }
80            };
81
82            for child in &mut child_records {
83                self.flatten_relation_graph(
84                    &relation.target_entity,
85                    child,
86                    root,
87                    graph,
88                    installed,
89                )?;
90            }
91
92            if relation.many || relation.local_key == "id" {
93                let owner_id = record.get("id").and_then(Value::try_u64).ok_or_else(|| {
94                    teaql_core::EntityError::new(
95                        entity_name,
96                        "loaded reverse relation owner is missing its u64 id",
97                    )
98                })?;
99                if relation.many {
100                    context.decode_compact_entity_list_into_graph(
101                        &relation.target_entity,
102                        child_records
103                            .into_iter()
104                            .map(teaql_core::CompactRow::from_map)
105                            .collect(),
106                        root,
107                        graph,
108                        entity_name,
109                        owner_id,
110                        &relation.name,
111                    )?;
112                } else {
113                    context.decode_compact_entity_option_into_graph(
114                        &relation.target_entity,
115                        child_records
116                            .into_iter()
117                            .map(teaql_core::CompactRow::from_map)
118                            .collect(),
119                        root,
120                        graph,
121                        entity_name,
122                        owner_id,
123                        &relation.name,
124                    )?;
125                }
126                continue;
127            }
128
129            for child in child_records {
130                let id = child.get("id").and_then(Value::try_u64).ok_or_else(|| {
131                    teaql_core::EntityError::new(
132                        &relation.target_entity,
133                        "loaded relation is missing its u64 id",
134                    )
135                })?;
136                if installed.insert((relation.target_entity.clone(), id)) {
137                    context.decode_compact_entity_into_graph(
138                        &relation.target_entity,
139                        teaql_core::CompactRow::from_map(child),
140                        root,
141                        graph,
142                    )?;
143                }
144            }
145        }
146        Ok(())
147    }
148
149    fn attach_flat_relation_graph(
150        &self,
151        entity_name: &str,
152        rows: &mut [teaql_core::CompactRow],
153    ) -> Result<crate::EntityRoot, teaql_core::EntityError> {
154        let root = crate::EntityRoot::default();
155        let mut graph = crate::EntityGraphBuilder::default();
156        let mut installed = BTreeSet::new();
157        for row in rows {
158            let mut record = row.clone().into_map();
159            self.flatten_relation_graph(
160                entity_name,
161                &mut record,
162                &root,
163                &mut graph,
164                &mut installed,
165            )?;
166            *row = teaql_core::CompactRow::from_map(record);
167        }
168        root.freeze_graph(graph).map_err(|_| {
169            teaql_core::EntityError::new(entity_name, "identity graph was already frozen")
170        })?;
171        Ok(root)
172    }
173
174    pub(super) fn query_behavior(
175        &self,
176        entity: &str,
177    ) -> Option<Arc<dyn EntityDataServiceBehavior>> {
178        self.data_service
179            .metadata
180            .context
181            .entity_data_service_behavior(entity)
182    }
183
184    pub(super) fn behavior(&self) -> Option<Arc<dyn EntityDataServiceBehavior>> {
185        self.data_service
186            .metadata
187            .context
188            .entity_data_service_behavior(&self.entity)
189    }
190
191    pub fn entity(&self) -> &str {
192        &self.entity
193    }
194
195    pub fn select(&self) -> SelectQuery {
196        SelectQuery::new(self.entity.clone())
197    }
198
199    pub fn insert_command(&self) -> InsertCommand {
200        InsertCommand::new(self.entity.clone())
201    }
202
203    fn enforce_insert_policy(&self, command: &mut InsertCommand) -> Result<(), RuntimeError> {
204        if let Some(policy) = self.data_service.metadata.context.request_policy.as_ref() {
205            policy.enforce_insert(self.data_service.metadata.context, command)?;
206        }
207        Ok(())
208    }
209
210    fn enforce_update_policy(&self, command: &mut UpdateCommand) -> Result<(), RuntimeError> {
211        if let Some(policy) = self.data_service.metadata.context.request_policy.as_ref() {
212            policy.enforce_update(self.data_service.metadata.context, command)?;
213        }
214        Ok(())
215    }
216
217    fn enforce_delete_policy(&self, command: &mut DeleteCommand) -> Result<(), RuntimeError> {
218        if let Some(policy) = self.data_service.metadata.context.request_policy.as_ref() {
219            policy.enforce_delete(self.data_service.metadata.context, command)?;
220        }
221        Ok(())
222    }
223
224    fn enforce_recover_policy(&self, command: &mut RecoverCommand) -> Result<(), RuntimeError> {
225        if let Some(policy) = self.data_service.metadata.context.request_policy.as_ref() {
226            policy.enforce_recover(self.data_service.metadata.context, command)?;
227        }
228        Ok(())
229    }
230
231    fn prepare_select_query(&self, query: &SelectQuery) -> Result<SelectQuery, RuntimeError> {
232        self.prepare_select_query_owned(query.clone())
233    }
234
235    fn prepare_select_query_owned(
236        &self,
237        mut query: SelectQuery,
238    ) -> Result<SelectQuery, RuntimeError> {
239        let mut full_trace = self.trace_context.clone();
240        full_trace.extend(query.trace_chain);
241        query.trace_chain = full_trace;
242
243        if let Some(behavior) = self.query_behavior(&query.entity) {
244            behavior.before_select(self.data_service.metadata.context, &mut query)?;
245        }
246        if let Some(policy) = self.data_service.metadata.context.request_policy.as_ref() {
247            policy.enforce_select(self.data_service.metadata.context, &mut query)?;
248        }
249        // Ensure local_key fields for relation loads are projected so that
250        // enhance_query_relations can match parent rows to child records.
251        if !query.relations.is_empty() {
252            if let Some(descriptor) = self.data_service.metadata.context.entity(&query.entity) {
253                for load in &query.relations {
254                    if let Some(relation) = descriptor.relation_by_name(&load.name) {
255                        if !query.projection.contains(&relation.local_key) {
256                            query.projection.push(relation.local_key.clone());
257                        }
258                    }
259                }
260            }
261        }
262        Ok(query)
263    }
264
265    pub fn prepare_insert_command(
266        &self,
267        command: &InsertCommand,
268    ) -> Result<InsertCommand, RuntimeError> {
269        let mut command = command.clone();
270        if let Some(behavior) = self.behavior() {
271            behavior.before_insert(self.data_service.metadata.context, &mut command)?;
272        }
273        self.enforce_insert_policy(&mut command)?;
274
275        let entity = self
276            .data_service
277            .metadata
278            .context
279            .require_entity(&command.entity)?;
280        if let Some(id_property) = entity.id_property() {
281            let needs_id = !command.values.contains_key(&id_property.name)
282                || is_unassigned_id(command.values.get(&id_property.name));
283            if needs_id {
284                let id = self
285                    .data_service
286                    .metadata
287                    .context
288                    .next_id(&command.entity)?;
289                command
290                    .values
291                    .insert(id_property.name.clone(), Value::U64(id));
292            }
293        }
294        ensure_initial_version(&mut command.values, entity);
295        let mut checked_values: crate::EntityValues = command.values.into();
296        mark_entity_status(&mut checked_values, CheckObjectStatus::Create);
297        let check_result = self
298            .data_service
299            .metadata
300            .context
301            .check_and_fix_values(&command.entity, &mut checked_values);
302        clear_entity_status(&mut checked_values);
303        check_result?;
304        command.values = checked_values.into();
305
306        Ok(command)
307    }
308
309    pub fn update_command(&self, id: impl Into<Value>) -> UpdateCommand {
310        UpdateCommand::new(self.entity.clone(), id)
311    }
312
313    pub fn prepare_update_command(
314        &self,
315        command: &UpdateCommand,
316    ) -> Result<UpdateCommand, RuntimeError> {
317        let mut command = command.clone();
318        if let Some(behavior) = self.behavior() {
319            behavior.before_update(self.data_service.metadata.context, &mut command)?;
320        }
321        self.enforce_update_policy(&mut command)?;
322
323        Ok(command)
324    }
325
326    pub fn delete_command(&self, id: impl Into<Value>) -> DeleteCommand {
327        DeleteCommand::new(self.entity.clone(), id)
328    }
329
330    pub fn recover_command(&self, id: impl Into<Value>, expected_version: i64) -> RecoverCommand {
331        RecoverCommand::new(self.entity.clone(), id, expected_version)
332    }
333
334    pub(crate) async fn fetch_all_internal(
335        &self,
336        query: &SelectQuery,
337    ) -> Result<Vec<teaql_core::CompactRow>, DataServiceError<E::Error>> {
338        let query = self
339            .prepare_select_query(query)
340            .map_err(DataServiceError::Runtime)?;
341        let query = query
342            .prepare_for_list()
343            .map_err(|message| DataServiceError::Runtime(RuntimeError::Graph(message)))?;
344        if query.continuous_page_fetch.is_none()
345            && query.object_group_bys.is_empty()
346            && query.child_enhancements.is_empty()
347            && query.relations.is_empty()
348        {
349            return self.fetch_prepared_query_owned(query).await;
350        }
351        self.fetch_prepared_all(&query).await
352    }
353
354    async fn fetch_all_owned_internal(
355        &self,
356        query: SelectQuery,
357    ) -> Result<Vec<teaql_core::CompactRow>, DataServiceError<E::Error>> {
358        let query = self
359            .prepare_select_query_owned(query)
360            .map_err(DataServiceError::Runtime)?
361            .prepare_for_list()
362            .map_err(|message| DataServiceError::Runtime(RuntimeError::Graph(message)))?;
363        if query.continuous_page_fetch.is_none()
364            && query.object_group_bys.is_empty()
365            && query.child_enhancements.is_empty()
366            && query.relations.is_empty()
367        {
368            return self.fetch_prepared_query_owned(query).await;
369        }
370        self.fetch_prepared_all(&query).await
371    }
372
373    pub(crate) async fn fetch_compact_all_internal(
374        &self,
375        query: SelectQuery,
376    ) -> Result<Vec<teaql_core::CompactRow>, DataServiceError<E::Error>> {
377        let query = self
378            .prepare_select_query_owned(query)
379            .map_err(DataServiceError::Runtime)?
380            .prepare_for_list()
381            .map_err(|message| DataServiceError::Runtime(RuntimeError::Graph(message)))?;
382        self.fetch_prepared_all(&query).await
383    }
384
385    async fn prepare_continuous_page(
386        &self,
387        query: SelectQuery,
388    ) -> (SelectQuery, Option<ContinuousPageExecution>) {
389        let Some(options) = query.continuous_page_fetch.as_ref() else {
390            self.data_service
391                .metadata
392                .context
393                .observe_continuous_page("DISABLED", None);
394            return (query, None);
395        };
396        let Some(slice) = query.slice.as_ref() else {
397            self.data_service
398                .metadata
399                .context
400                .observe_continuous_page("OFFSET_FALLBACK:INVALID_SLICE", None);
401            return (query, None);
402        };
403        let Some(page_size) = slice.limit else {
404            self.data_service
405                .metadata
406                .context
407                .observe_continuous_page("OFFSET_FALLBACK:INVALID_SLICE", None);
408            return (query, None);
409        };
410        if query.partition_by.is_some()
411            || !query.aggregates.is_empty()
412            || !query.group_by.is_empty()
413        {
414            self.data_service
415                .metadata
416                .context
417                .observe_continuous_page("OFFSET_FALLBACK:UNSUPPORTED_QUERY_SHAPE", None);
418            return (query, None);
419        }
420        if query.order_by.len() != 1
421            || query.order_by[0].field != "id"
422            || query.order_by[0].expr.is_some()
423        {
424            self.data_service
425                .metadata
426                .context
427                .observe_continuous_page("OFFSET_FALLBACK:ORDER_NOT_SEEKABLE_ID", None);
428            return (query, None);
429        }
430        let direction = query.order_by[0].direction;
431        let query_key = self.continuous_page_query_key(&query, &options.namespace);
432        let execution = ContinuousPageExecution {
433            query_key: query_key.clone(),
434            direction,
435            page_size,
436            original_offset: slice.offset,
437            ttl_seconds: options.ttl_seconds,
438            optimized: false,
439            seek_cursor_id: None,
440        };
441        if slice.offset == 0 {
442            self.data_service
443                .metadata
444                .context
445                .observe_continuous_page("OFFSET_FALLBACK:FIRST_PAGE", None);
446            return (query, Some(execution));
447        }
448        let cursor = match self
449            .data_service
450            .metadata
451            .context
452            .continuous_page_cursor_store()
453            .get(&query_key, slice.offset)
454            .await
455        {
456            Ok(Some(cursor)) => cursor,
457            Ok(None) => {
458                self.data_service
459                    .metadata
460                    .context
461                    .observe_continuous_page("OFFSET_FALLBACK:CACHE_MISS", None);
462                return (query, Some(execution));
463            }
464            Err(_) => {
465                self.data_service
466                    .metadata
467                    .context
468                    .observe_continuous_page("OFFSET_FALLBACK:STORE_UNAVAILABLE", None);
469                return (query, Some(execution));
470            }
471        };
472        if cursor.entity != query.entity
473            || cursor.direction != direction
474            || cursor.page_size != page_size
475            || cursor.next_offset != slice.offset
476            || cursor.expires_at <= SystemTime::now()
477        {
478            self.data_service
479                .metadata
480                .context
481                .observe_continuous_page("OFFSET_FALLBACK:CURSOR_INVALID", None);
482            return (query, Some(execution));
483        }
484        let mut optimized = query;
485        optimized.slice.as_mut().expect("validated slice").offset = 0;
486        optimized = optimized.and_filter(match direction {
487            SortDirection::Asc => Expr::gt("id", cursor.boundary.clone()),
488            SortDirection::Desc => Expr::lt("id", cursor.boundary.clone()),
489        });
490        let seek_cursor_id = cursor.cursor_id;
491        self.data_service
492            .metadata
493            .context
494            .observe_continuous_page("CURSOR_SEEK", Some(seek_cursor_id.clone()));
495        (
496            optimized,
497            Some(ContinuousPageExecution {
498                optimized: true,
499                seek_cursor_id: Some(seek_cursor_id),
500                ..execution
501            }),
502        )
503    }
504
505    async fn register_continuous_page(
506        &self,
507        execution: &Option<ContinuousPageExecution>,
508        rows: &[teaql_core::CompactRow],
509    ) {
510        let Some(execution) = execution else { return };
511        if rows.len() as u64 != execution.page_size {
512            return;
513        }
514        let Some(boundary) = rows.last().and_then(|row| row.get("id")).cloned() else {
515            return;
516        };
517        let cursor_id = format!(
518            "cpg_{:x}",
519            SystemTime::now()
520                .duration_since(SystemTime::UNIX_EPOCH)
521                .unwrap_or_default()
522                .as_nanos()
523        );
524        let cursor = ContinuousPageCursor {
525            cursor_id,
526            query_key: execution.query_key.clone(),
527            entity: self.entity.clone(),
528            direction: execution.direction,
529            boundary,
530            page_size: execution.page_size,
531            next_offset: execution.original_offset + rows.len() as u64,
532            expires_at: SystemTime::now() + Duration::from_secs(execution.ttl_seconds),
533        };
534        if self
535            .data_service
536            .metadata
537            .context
538            .continuous_page_cursor_store()
539            .put(cursor)
540            .await
541            .is_err()
542        {
543            self.data_service
544                .metadata
545                .context
546                .observe_continuous_page("OFFSET_FALLBACK:STORE_UNAVAILABLE", None);
547        } else if execution.optimized {
548            self.data_service
549                .metadata
550                .context
551                .observe_continuous_page("CURSOR_SEEK", execution.seek_cursor_id.clone());
552        } else {
553            self.data_service
554                .metadata
555                .context
556                .observe_continuous_page("OFFSET_FALLBACK:FIRST_PAGE", None);
557        }
558    }
559
560    fn continuous_page_query_key(&self, query: &SelectQuery, namespace: &str) -> String {
561        let mut normalized = query.clone();
562        if let Some(slice) = normalized.slice.as_mut() {
563            slice.offset = 0;
564        }
565        normalized.comment = None;
566        normalized.trace_chain.clear();
567        let mut hasher = std::collections::hash_map::DefaultHasher::new();
568        namespace.hash(&mut hasher);
569        format!("{normalized:?}").hash(&mut hasher);
570        self.data_service
571            .metadata
572            .context
573            .user_identifier()
574            .hash(&mut hasher);
575        format!("teaql:continuous-page:v1:{:016x}", hasher.finish())
576    }
577
578    /// Fetch root records from the provider cursor without materializing them.
579    /// Relation and aggregate enhancement needs a separate batched protocol and
580    /// is rejected here instead of silently returning incomplete entities.
581    pub(crate) async fn fetch_stream_internal(
582        &self,
583        query: &SelectQuery,
584    ) -> Result<
585        std::pin::Pin<
586            Box<
587                dyn futures_core::Stream<
588                        Item = Result<teaql_data_service::StreamChunk, DataServiceError<E::Error>>,
589                    > + '_,
590            >,
591        >,
592        DataServiceError<E::Error>,
593    >
594    where
595        E: teaql_data_service::StreamQueryExecutor,
596    {
597        let query = self
598            .prepare_select_query(query)
599            .map_err(DataServiceError::Runtime)?;
600        let query = query
601            .prepare_for_list()
602            .map_err(|message| DataServiceError::Runtime(RuntimeError::Graph(message)))?;
603
604        if !query.relations.is_empty()
605            || !query.child_enhancements.is_empty()
606            || !query.object_group_bys.is_empty()
607        {
608            return Err(DataServiceError::Runtime(RuntimeError::Graph(
609                "streaming relation or aggregate enhancement is not supported; stream a root query or use execute_for_list"
610                    .to_owned(),
611            )));
612        }
613
614        let chunk_size = query
615            .stream_config
616            .as_ref()
617            .map(|c| c.chunk_size)
618            .unwrap_or(1000);
619
620        let final_comment = self
621            .data_service
622            .resolve_final_comment(&query.trace_chain, query.comment.clone());
623        let mut query = query.clone();
624        query.comment = final_comment;
625
626        let request = teaql_data_service::QueryRequest {
627            query: query.clone(),
628            trace_chain: query.trace_chain.clone(),
629            comment: query.comment.clone(),
630            capture_debug_query: self.data_service.metadata.capture_query_debug(),
631        };
632
633        let chunks = self.data_service.executor.query_stream(request, chunk_size);
634        use futures_util::StreamExt;
635        Ok(Box::pin(
636            chunks.map(|item| item.map_err(DataServiceError::Executor)),
637        ))
638    }
639
640    async fn fetch_prepared_all(
641        &self,
642        query: &SelectQuery,
643    ) -> Result<Vec<teaql_core::CompactRow>, DataServiceError<E::Error>> {
644        let query = query
645            .clone()
646            .prepare_for_list()
647            .map_err(|message| DataServiceError::Runtime(RuntimeError::Graph(message)))?;
648        if query.continuous_page_fetch.is_none()
649            && query.object_group_bys.is_empty()
650            && query.child_enhancements.is_empty()
651            && query.relations.is_empty()
652        {
653            return self.fetch_prepared_query(&query).await;
654        }
655        let (execution_query, continuous) = self.prepare_continuous_page(query).await;
656        let mut rows = self.fetch_prepared_query(&execution_query).await?;
657        self.enhance_object_group_bys_internal(
658            &mut rows,
659            &execution_query.object_group_bys,
660            &execution_query.trace_chain,
661        )
662        .await?;
663        self.enhance_child_queries_internal(
664            &mut rows,
665            &execution_query.child_enhancements,
666            &execution_query.trace_chain,
667        )
668        .await?;
669        self.enhance_query_relations_internal(&mut rows, &execution_query)
670            .await?;
671        self.register_continuous_page(&continuous, &rows).await;
672        Ok(rows)
673    }
674
675    async fn fetch_prepared_query(
676        &self,
677        query: &SelectQuery,
678    ) -> Result<Vec<teaql_core::CompactRow>, DataServiceError<E::Error>> {
679        let final_comment = self
680            .data_service
681            .resolve_final_comment(&query.trace_chain, query.comment.clone());
682        let mut query = query.clone();
683        query.comment = final_comment;
684        if let Some(options) = query.aggregation_cache.filter(|options| options.enabled) {
685            if let Some(cache) = self
686                .data_service
687                .metadata
688                .context
689                .get_resource::<Arc<dyn AggregationCacheBackend>>()
690            {
691                return self
692                    .fetch_prepared_query_with_cache(&query, options, cache.as_ref())
693                    .await;
694            }
695            if let Some(cache) = self
696                .data_service
697                .metadata
698                .context
699                .get_resource::<InMemoryAggregationCache>()
700            {
701                return self
702                    .fetch_prepared_query_with_cache(&query, options, cache)
703                    .await;
704            }
705        }
706        let request = teaql_data_service::QueryRequest {
707            query: query.clone(),
708            trace_chain: query.trace_chain.clone(),
709            comment: query.comment.clone(),
710            capture_debug_query: self.data_service.metadata.capture_query_debug(),
711        };
712        let res = self
713            .data_service
714            .executor
715            .query(request)
716            .await
717            .map_err(DataServiceError::Executor)?;
718        self.data_service
719            .metadata
720            .context
721            .record_metadata_log(&res.metadata);
722        Ok(res.rows)
723    }
724
725    async fn fetch_prepared_query_owned(
726        &self,
727        mut query: SelectQuery,
728    ) -> Result<Vec<teaql_core::CompactRow>, DataServiceError<E::Error>> {
729        if query
730            .aggregation_cache
731            .is_some_and(|options| options.enabled)
732        {
733            return self.fetch_prepared_query(&query).await;
734        }
735        query.comment = self
736            .data_service
737            .resolve_final_comment(&query.trace_chain, query.comment.take());
738        let trace_chain = std::mem::take(&mut query.trace_chain);
739        let request = teaql_data_service::QueryRequest {
740            trace_chain,
741            comment: query.comment.clone(),
742            capture_debug_query: self.data_service.metadata.capture_query_debug(),
743            query,
744        };
745        let res = self
746            .data_service
747            .executor
748            .query(request)
749            .await
750            .map_err(DataServiceError::Executor)?;
751        self.data_service
752            .metadata
753            .context
754            .record_metadata_log(&res.metadata);
755        Ok(res.rows)
756    }
757
758    pub(crate) async fn fetch_prepared_compact_owned(
759        &self,
760        mut query: SelectQuery,
761    ) -> Result<Vec<teaql_core::CompactRow>, DataServiceError<E::Error>> {
762        query.comment = self
763            .data_service
764            .resolve_final_comment(&query.trace_chain, query.comment.take());
765        let trace_chain = std::mem::take(&mut query.trace_chain);
766        let request = teaql_data_service::QueryRequest {
767            trace_chain,
768            comment: query.comment.clone(),
769            capture_debug_query: self.data_service.metadata.capture_query_debug(),
770            query,
771        };
772        let result = self
773            .data_service
774            .executor
775            .query(request)
776            .await
777            .map_err(DataServiceError::Executor)?;
778        self.data_service
779            .metadata
780            .context
781            .record_metadata_log(&result.metadata);
782        Ok(result.rows)
783    }
784
785    async fn fetch_prepared_query_with_cache(
786        &self,
787        query: &SelectQuery,
788        options: AggregationCacheOptions,
789        cache: &dyn AggregationCacheBackend,
790    ) -> Result<Vec<teaql_core::CompactRow>, DataServiceError<E::Error>> {
791        let key = aggregation_cache_key(
792            cache.namespace(),
793            &aggregation_cache_namespace(&query.entity),
794            query,
795        );
796        let scope = self.data_service.metadata.context.start_runtime_operation(
797            crate::RuntimeOperation::new("cache", format!("{}.aggregation.get", query.entity))
798                .attribute("teaql.cache.operation", "get"),
799        );
800        let result = scope
801            .run(async {
802                if let Some(rows) = cache.get(&key, options.cache_expired_millis) {
803                    return Ok((rows, "hit"));
804                }
805                let request = teaql_data_service::QueryRequest {
806                    query: query.clone(),
807                    trace_chain: query.trace_chain.clone(),
808                    comment: query.comment.clone(),
809                    capture_debug_query: self.data_service.metadata.capture_query_debug(),
810                };
811                let provider_kind = std::any::type_name::<E>().to_owned();
812                let provider_scope = self.data_service.metadata.context.start_runtime_operation(
813                    crate::RuntimeOperation::new("provider", format!("{provider_kind}.query"))
814                        .attribute("teaql.provider.kind", provider_kind)
815                        .attribute("teaql.provider.operation", "query"),
816                );
817                let provider_result = provider_scope
818                    .run(self.data_service.executor.query(request))
819                    .await;
820                let res = match provider_result {
821                    Ok(value) => {
822                        provider_scope.success(std::collections::BTreeMap::new());
823                        value
824                    }
825                    Err(error) => {
826                        provider_scope.failure("data_service_error");
827                        return Err(DataServiceError::Executor(error));
828                    }
829                };
830                self.data_service
831                    .metadata
832                    .context
833                    .record_metadata_log(&res.metadata);
834                let rows = res.rows;
835                cache.put(key, rows.clone());
836                Ok((rows, "miss"))
837            })
838            .await;
839        match result {
840            Ok((rows, cache_result)) => {
841                scope.success(std::collections::BTreeMap::from([(
842                    "teaql.cache.result".to_owned(),
843                    crate::RuntimeAttributeValue::from(cache_result),
844                )]));
845                Ok(rows)
846            }
847            Err(error) => {
848                scope.failure("cache_load_error");
849                Err(error)
850            }
851        }
852    }
853
854    pub(crate) async fn fetch_all_with_relation_aggregates_internal(
855        &self,
856        query: &SelectQuery,
857        relation_aggregates: &[RelationAggregate],
858    ) -> Result<Vec<teaql_core::CompactRow>, DataServiceError<E::Error>> {
859        let query = self
860            .prepare_select_query(query)
861            .map_err(DataServiceError::Runtime)?;
862
863        let mut rows = self.fetch_prepared_all(&query).await?;
864        self.enhance_relation_aggregates_internal(
865            &mut rows,
866            relation_aggregates,
867            query.aggregation_cache,
868            &query.trace_chain,
869        )
870        .await?;
871        Ok(rows)
872    }
873
874    pub(crate) async fn fetch_smart_list_internal(
875        &self,
876        query: &SelectQuery,
877    ) -> Result<SmartList<teaql_core::CompactRow>, DataServiceError<E::Error>> {
878        let query = self
879            .prepare_select_query(query)
880            .map_err(DataServiceError::Runtime)?;
881
882        self.data_service.fetch_smart_list(&query).await
883    }
884
885    pub(crate) async fn fetch_smart_list_with_relation_aggregates_internal(
886        &self,
887        query: &SelectQuery,
888        relation_aggregates: &[RelationAggregate],
889    ) -> Result<SmartList<teaql_core::CompactRow>, DataServiceError<E::Error>> {
890        self.fetch_all_with_relation_aggregates_internal(query, relation_aggregates)
891            .await
892            .map(SmartList::from)
893    }
894
895    pub(crate) async fn fetch_entities_internal<T>(
896        &self,
897        query: &SelectQuery,
898    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
899    where
900        T: Entity,
901    {
902        let query = self
903            .prepare_select_query(query)
904            .map_err(DataServiceError::Runtime)?;
905
906        self.data_service.fetch_entities(&query).await
907    }
908
909    pub(crate) async fn fetch_entities_with_relation_aggregates_internal<T>(
910        &self,
911        query: &SelectQuery,
912        relation_aggregates: &[RelationAggregate],
913    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
914    where
915        T: Entity,
916    {
917        let root = crate::EntityRoot::default();
918        self.fetch_all_with_relation_aggregates_internal(query, relation_aggregates)
919            .await?
920            .into_iter()
921            .map(|record| {
922                let mut entity = T::from_compact_row(record)?;
923                entity.on_loaded(&root as &dyn std::any::Any);
924                Ok(entity)
925            })
926            .collect::<Result<Vec<_>, _>>()
927            .map(SmartList::from)
928            .map_err(DataServiceError::Entity)
929    }
930
931    pub(crate) async fn fetch_enhanced_entities_with_relation_aggregates_internal<T>(
932        &self,
933        query: &SelectQuery,
934        relation_aggregates: &[RelationAggregate],
935    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
936    where
937        T: Entity,
938    {
939        let query = self
940            .prepare_select_query(query)
941            .map_err(DataServiceError::Runtime)?;
942        self.fetch_enhanced_entities_with_relation_aggregates_prepared(query, relation_aggregates)
943            .await
944    }
945
946    async fn fetch_enhanced_entities_with_relation_aggregates_owned_internal<T>(
947        &self,
948        query: SelectQuery,
949        relation_aggregates: &[RelationAggregate],
950    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
951    where
952        T: Entity,
953    {
954        let query = self
955            .prepare_select_query_owned(query)
956            .map_err(DataServiceError::Runtime)?;
957        self.fetch_enhanced_entities_with_relation_aggregates_prepared(query, relation_aggregates)
958            .await
959    }
960
961    async fn fetch_enhanced_entities_with_relation_aggregates_prepared<T>(
962        &self,
963        query: SelectQuery,
964        relation_aggregates: &[RelationAggregate],
965    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
966    where
967        T: Entity,
968    {
969        if relation_aggregates.is_empty()
970            && query.continuous_page_fetch.is_none()
971            && query.object_group_bys.is_empty()
972            && query.child_enhancements.is_empty()
973            && query.relations.is_empty()
974        {
975            let query = query
976                .prepare_for_list()
977                .map_err(|message| DataServiceError::Runtime(RuntimeError::Graph(message)))?;
978            let root = crate::EntityRoot::default();
979            return self
980                .fetch_prepared_query_owned(query)
981                .await?
982                .into_iter()
983                .map(|record| {
984                let mut entity = T::from_compact_row(record)?;
985                    entity.on_loaded(&root as &dyn std::any::Any);
986                    Ok(entity)
987                })
988                .collect::<Result<Vec<_>, _>>()
989                .map(SmartList::from)
990                .map_err(DataServiceError::Entity);
991        }
992
993        let flat_plans = self
994            .flat_relation_plans(&query)
995            .map_err(DataServiceError::Runtime)?;
996        let use_flat_hydration = flat_plans.is_some();
997        let mut root_query = query.clone();
998        if use_flat_hydration {
999            root_query.relations.clear();
1000        }
1001        if relation_aggregates.is_empty()
1002            && root_query.continuous_page_fetch.is_none()
1003            && root_query.object_group_bys.is_empty()
1004            && root_query.child_enhancements.is_empty()
1005        {
1006            if let Some((query_plans, behavior_plans)) = flat_plans.as_ref() {
1007                let root_query = root_query
1008                    .prepare_for_list()
1009                    .map_err(|message| DataServiceError::Runtime(RuntimeError::Graph(message)))?;
1010                let rows = self.fetch_prepared_compact_owned(root_query).await?;
1011                let root = crate::EntityRoot::default();
1012                let mut graph = crate::EntityGraphBuilder::default();
1013                self.hydrate_compact_flat_plans_internal(&rows, query_plans, &root, &mut graph)
1014                    .await?;
1015                self.hydrate_compact_flat_plans_internal(&rows, behavior_plans, &root, &mut graph)
1016                    .await?;
1017                root.freeze_graph(graph).map_err(|_| {
1018                    DataServiceError::Entity(teaql_core::EntityError::new(
1019                        &query.entity,
1020                        "identity graph was already frozen",
1021                    ))
1022                })?;
1023                return rows
1024                    .into_iter()
1025                    .map(|row| {
1026                        let mut entity = T::from_compact_row(row)?;
1027                        entity.on_loaded(&root as &dyn std::any::Any);
1028                        Ok(entity)
1029                    })
1030                    .collect::<Result<Vec<_>, _>>()
1031                    .map(SmartList::from)
1032                    .map_err(DataServiceError::Entity);
1033            }
1034        }
1035        let mut rows = self.fetch_prepared_all(&root_query).await?;
1036        self.enhance_relation_aggregates_internal(
1037            &mut rows,
1038            relation_aggregates,
1039            query.aggregation_cache,
1040            &query.trace_chain,
1041        )
1042        .await?;
1043        let root = if let Some((query_plans, behavior_plans)) = flat_plans {
1044            let root = crate::EntityRoot::default();
1045            let mut graph = crate::EntityGraphBuilder::default();
1046            self.hydrate_flat_plans_internal(&mut rows, &query_plans, &root, &mut graph)
1047                .await?;
1048            self.hydrate_flat_plans_internal(&mut rows, &behavior_plans, &root, &mut graph)
1049                .await?;
1050            root.freeze_graph(graph).map_err(|_| {
1051                DataServiceError::Entity(teaql_core::EntityError::new(
1052                    &query.entity,
1053                    "identity graph was already frozen",
1054                ))
1055            })?;
1056            root
1057        } else {
1058            self.enhance_relations_internal(&mut rows).await?;
1059            self.attach_flat_relation_graph(&query.entity, &mut rows)
1060                .map_err(DataServiceError::Entity)?
1061        };
1062        rows.into_iter()
1063            .map(|record| {
1064                let mut entity = T::from_compact_row(record)?;
1065                entity.on_loaded(&root as &dyn std::any::Any);
1066                Ok(entity)
1067            })
1068            .collect::<Result<Vec<_>, _>>()
1069            .map(SmartList::from)
1070            .map_err(DataServiceError::Entity)
1071    }
1072
1073    pub(crate) async fn fetch_enhanced_entities_internal<T>(
1074        &self,
1075        query: &SelectQuery,
1076    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
1077    where
1078        T: Entity,
1079    {
1080        let query = self
1081            .prepare_select_query(query)
1082            .map_err(DataServiceError::Runtime)?;
1083
1084        let flat_plans = self
1085            .flat_relation_plans(&query)
1086            .map_err(DataServiceError::Runtime)?;
1087        let use_flat_hydration = flat_plans.is_some();
1088        let mut root_query = query.clone();
1089        if use_flat_hydration {
1090            root_query.relations.clear();
1091        }
1092        if root_query.continuous_page_fetch.is_none()
1093            && root_query.object_group_bys.is_empty()
1094            && root_query.child_enhancements.is_empty()
1095        {
1096            if let Some((query_plans, behavior_plans)) = flat_plans.as_ref() {
1097                let root_query = root_query
1098                    .prepare_for_list()
1099                    .map_err(|message| DataServiceError::Runtime(RuntimeError::Graph(message)))?;
1100                let rows = self.fetch_prepared_compact_owned(root_query).await?;
1101                let root = crate::EntityRoot::default();
1102                let mut graph = crate::EntityGraphBuilder::default();
1103                self.hydrate_compact_flat_plans_internal(&rows, query_plans, &root, &mut graph)
1104                    .await?;
1105                self.hydrate_compact_flat_plans_internal(&rows, behavior_plans, &root, &mut graph)
1106                    .await?;
1107                root.freeze_graph(graph).map_err(|_| {
1108                    DataServiceError::Entity(teaql_core::EntityError::new(
1109                        &query.entity,
1110                        "identity graph was already frozen",
1111                    ))
1112                })?;
1113                return rows
1114                    .into_iter()
1115                    .map(|row| {
1116                        let mut entity = T::from_compact_row(row)?;
1117                        entity.on_loaded(&root as &dyn std::any::Any);
1118                        Ok(entity)
1119                    })
1120                    .collect::<Result<Vec<_>, _>>()
1121                    .map(SmartList::from)
1122                    .map_err(DataServiceError::Entity);
1123            }
1124        }
1125        let mut rows = self.fetch_prepared_all(&root_query).await?;
1126        let root = if let Some((query_plans, behavior_plans)) = flat_plans {
1127            let root = crate::EntityRoot::default();
1128            let mut graph = crate::EntityGraphBuilder::default();
1129            self.hydrate_flat_plans_internal(&mut rows, &query_plans, &root, &mut graph)
1130                .await?;
1131            self.hydrate_flat_plans_internal(&mut rows, &behavior_plans, &root, &mut graph)
1132                .await?;
1133            root.freeze_graph(graph).map_err(|_| {
1134                DataServiceError::Entity(teaql_core::EntityError::new(
1135                    &query.entity,
1136                    "identity graph was already frozen",
1137                ))
1138            })?;
1139            root
1140        } else {
1141            self.enhance_relations_internal(&mut rows).await?;
1142            self.attach_flat_relation_graph(&query.entity, &mut rows)
1143                .map_err(DataServiceError::Entity)?
1144        };
1145        rows.into_iter()
1146            .map(|record| {
1147                let mut entity = T::from_compact_row(record)?;
1148                entity.on_loaded(&root as &dyn std::any::Any);
1149                Ok(entity)
1150            })
1151            .collect::<Result<Vec<_>, _>>()
1152            .map(SmartList::from)
1153            .map_err(DataServiceError::Entity)
1154    }
1155
1156    #[doc(hidden)]
1157    pub async fn fetch_all(
1158        &self,
1159        query: &PurposedSelectQuery,
1160    ) -> Result<Vec<teaql_core::CompactRow>, DataServiceError<E::Error>> {
1161        self.fetch_all_internal(query.as_query()).await
1162    }
1163
1164    #[doc(hidden)]
1165    pub async fn fetch_all_owned(
1166        &self,
1167        query: PurposedSelectQuery,
1168    ) -> Result<Vec<teaql_core::CompactRow>, DataServiceError<E::Error>> {
1169        self.fetch_all_owned_internal(query.into_query()).await
1170    }
1171
1172    #[doc(hidden)]
1173    pub async fn fetch_stream(
1174        &self,
1175        query: &PurposedSelectQuery,
1176    ) -> Result<
1177        std::pin::Pin<
1178            Box<
1179                dyn futures_core::Stream<
1180                        Item = Result<teaql_data_service::StreamChunk, DataServiceError<E::Error>>,
1181                    > + '_,
1182            >,
1183        >,
1184        DataServiceError<E::Error>,
1185    >
1186    where
1187        E: teaql_data_service::StreamQueryExecutor,
1188    {
1189        self.fetch_stream_internal(query.as_query()).await
1190    }
1191
1192    #[doc(hidden)]
1193    pub async fn fetch_smart_list(
1194        &self,
1195        query: &PurposedSelectQuery,
1196    ) -> Result<SmartList<teaql_core::CompactRow>, DataServiceError<E::Error>> {
1197        self.fetch_smart_list_internal(query.as_query()).await
1198    }
1199
1200    #[doc(hidden)]
1201    pub async fn fetch_smart_list_with_relation_aggregates(
1202        &self,
1203        query: &PurposedSelectQuery,
1204        relation_aggregates: &[RelationAggregate],
1205    ) -> Result<SmartList<teaql_core::CompactRow>, DataServiceError<E::Error>> {
1206        self.fetch_smart_list_with_relation_aggregates_internal(
1207            query.as_query(),
1208            relation_aggregates,
1209        )
1210        .await
1211    }
1212
1213    #[doc(hidden)]
1214    pub async fn fetch_entities<T>(
1215        &self,
1216        query: &PurposedSelectQuery,
1217    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
1218    where
1219        T: Entity,
1220    {
1221        self.fetch_entities_internal(query.as_query()).await
1222    }
1223
1224    #[doc(hidden)]
1225    pub async fn fetch_enhanced_entities<T>(
1226        &self,
1227        query: &PurposedSelectQuery,
1228    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
1229    where
1230        T: Entity,
1231    {
1232        self.fetch_enhanced_entities_internal(query.as_query())
1233            .await
1234    }
1235
1236    #[doc(hidden)]
1237    pub async fn fetch_enhanced_entities_with_relation_aggregates<T>(
1238        &self,
1239        query: &PurposedSelectQuery,
1240        relation_aggregates: &[RelationAggregate],
1241    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
1242    where
1243        T: Entity,
1244    {
1245        self.fetch_enhanced_entities_with_relation_aggregates_internal(
1246            query.as_query(),
1247            relation_aggregates,
1248        )
1249        .await
1250    }
1251
1252    #[doc(hidden)]
1253    pub async fn fetch_enhanced_entities_with_relation_aggregates_owned<T>(
1254        &self,
1255        query: PurposedSelectQuery,
1256        relation_aggregates: &[RelationAggregate],
1257    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
1258    where
1259        T: Entity,
1260    {
1261        self.fetch_enhanced_entities_with_relation_aggregates_owned_internal(
1262            query.into_query(),
1263            relation_aggregates,
1264        )
1265        .await
1266    }
1267
1268    pub(crate) async fn insert_internal(
1269        &self,
1270        command: &InsertCommand,
1271    ) -> Result<u64, DataServiceError<E::Error>> {
1272        let command = self
1273            .prepare_insert_command(command)
1274            .map_err(DataServiceError::Runtime)?;
1275        self.execute_prepared_insert_with_comment(command, self.trace_context.clone())
1276            .await
1277    }
1278
1279    pub(crate) async fn update_internal(
1280        &self,
1281        command: &UpdateCommand,
1282    ) -> Result<u64, DataServiceError<E::Error>> {
1283        let command = self
1284            .prepare_update_command(command)
1285            .map_err(DataServiceError::Runtime)?;
1286        self.execute_prepared_update_with_comment(command, self.trace_context.clone())
1287            .await
1288    }
1289
1290    pub(crate) async fn delete_internal(
1291        &self,
1292        command: &DeleteCommand,
1293    ) -> Result<u64, DataServiceError<E::Error>> {
1294        self.delete_scoped_internal(command, self.trace_context.clone())
1295            .await
1296    }
1297
1298    pub(crate) async fn delete_scoped_internal(
1299        &self,
1300        command: &DeleteCommand,
1301        trace_chain: Vec<teaql_core::TraceNode>,
1302    ) -> Result<u64, DataServiceError<E::Error>> {
1303        let mut command = command.clone();
1304        command.trace_chain = trace_chain.clone();
1305        if let Some(behavior) = self.behavior() {
1306            behavior
1307                .before_delete(self.data_service.metadata.context, &mut command)
1308                .map_err(DataServiceError::Runtime)?;
1309        }
1310        self.enforce_delete_policy(&mut command)
1311            .map_err(DataServiceError::Runtime)?;
1312
1313        let old_values =
1314            self.fetch_current_event_row(&command.entity, &command.id, trace_chain.clone())?;
1315        let affected = self.data_service.delete(&command).await?;
1316
1317        let mut event = RawAuditEvent::deleted_with_old_values(
1318            command.entity,
1319            command.id,
1320            command.expected_version,
1321            old_values,
1322        );
1323        event.trace_chain = trace_chain;
1324        self.emit_event(event).map_err(DataServiceError::Runtime)?;
1325        Ok(affected)
1326    }
1327
1328    pub(crate) async fn recover_internal(
1329        &self,
1330        command: &RecoverCommand,
1331    ) -> Result<u64, DataServiceError<E::Error>> {
1332        let mut command = command.clone();
1333        command.trace_chain = self.trace_context.clone();
1334        if let Some(behavior) = self.behavior() {
1335            behavior
1336                .before_recover(self.data_service.metadata.context, &mut command)
1337                .map_err(DataServiceError::Runtime)?;
1338        }
1339        self.enforce_recover_policy(&mut command)
1340            .map_err(DataServiceError::Runtime)?;
1341        let old_values = self.fetch_current_event_row(
1342            &command.entity,
1343            &command.id,
1344            command.trace_chain.clone(),
1345        )?;
1346        let affected = self.data_service.recover(&command).await?;
1347        let event = RawAuditEvent::recovered_with_old_values(
1348            command.entity,
1349            command.id,
1350            command.expected_version,
1351            old_values,
1352        );
1353        self.emit_event(event).map_err(DataServiceError::Runtime)?;
1354        Ok(affected)
1355    }
1356
1357    fn emit_event(&self, event: RawAuditEvent) -> Result<(), RuntimeError> {
1358        self.data_service.metadata.context.send_event(event)
1359    }
1360
1361    #[allow(dead_code)]
1362    pub(super) async fn execute_prepared_insert(
1363        &self,
1364        command: InsertCommand,
1365    ) -> Result<u64, DataServiceError<E::Error>> {
1366        self.execute_prepared_insert_with_comment(command, Vec::new())
1367            .await
1368    }
1369
1370    pub(super) async fn execute_prepared_insert_with_comment(
1371        &self,
1372        mut command: InsertCommand,
1373        trace_chain: Vec<teaql_core::TraceNode>,
1374    ) -> Result<u64, DataServiceError<E::Error>> {
1375        command.trace_chain = trace_chain.clone();
1376        let affected = self.data_service.insert(&command).await?;
1377        let mut event = RawAuditEvent::created(command.entity, command.values.into());
1378        event.trace_chain = trace_chain;
1379        self.emit_event(event).map_err(DataServiceError::Runtime)?;
1380        Ok(affected)
1381    }
1382
1383    pub(super) async fn execute_prepared_batch_insert(
1384        &self,
1385        command: teaql_core::BatchInsertCommand,
1386    ) -> Result<u64, DataServiceError<E::Error>> {
1387        if command.batch_values.is_empty() {
1388            return Ok(0);
1389        }
1390        let affected = self.data_service.batch_insert(&command).await?;
1391
1392        let entity = command.entity.clone();
1393        for (i, values) in command.batch_values.into_iter().enumerate() {
1394            let mut event = RawAuditEvent::created(entity.clone(), values.into());
1395            if i < command.trace_chains.len() {
1396                event.trace_chain = command.trace_chains[i].clone();
1397            }
1398            self.emit_event(event).map_err(DataServiceError::Runtime)?;
1399        }
1400        Ok(affected)
1401    }
1402
1403    #[allow(dead_code)]
1404    pub(super) async fn execute_prepared_update(
1405        &self,
1406        command: UpdateCommand,
1407    ) -> Result<u64, DataServiceError<E::Error>> {
1408        self.execute_prepared_update_with_comment(command, Vec::new())
1409            .await
1410    }
1411
1412    pub(super) async fn execute_prepared_update_with_comment(
1413        &self,
1414        mut command: UpdateCommand,
1415        trace_chain: Vec<teaql_core::TraceNode>,
1416    ) -> Result<u64, DataServiceError<E::Error>> {
1417        command.trace_chain = trace_chain.clone();
1418
1419        let mut old_values = command.old_values.clone();
1420        let needs_fetch = match &old_values {
1421            Some(snapshot) => !command.values.keys().all(|k| snapshot.contains_key(k)),
1422            None => true,
1423        };
1424        if needs_fetch {
1425            old_values = self
1426                .fetch_current_event_row(&command.entity, &command.id, trace_chain.clone())?
1427                .map(Into::into);
1428        }
1429
1430        let affected = self.data_service.update(&command).await?;
1431        let updated_fields = command.values.keys().cloned().collect();
1432        let mut values = command.values.clone();
1433        values.insert("id".to_owned(), command.id.clone());
1434        if let Some(version) = command.expected_version {
1435            values.insert("version".to_owned(), Value::I64(version + 1));
1436        }
1437        let mut new_values = old_values.clone().unwrap_or_default();
1438        for (field, value) in &values {
1439            new_values.insert(field.clone(), value.clone());
1440        }
1441        let mut event = RawAuditEvent::updated_with_old_values(
1442            command.entity,
1443            values.into(),
1444            old_values.map(Into::into),
1445            new_values.into(),
1446            updated_fields,
1447        );
1448        event.trace_chain = trace_chain;
1449        self.emit_event(event).map_err(DataServiceError::Runtime)?;
1450        Ok(affected)
1451    }
1452
1453    pub(super) async fn execute_prepared_batch_update(
1454        &self,
1455        command: teaql_core::BatchUpdateCommand,
1456    ) -> Result<u64, DataServiceError<E::Error>> {
1457        if command.batch_values.is_empty() {
1458            return Ok(0);
1459        }
1460        let affected = self.data_service.batch_update(&command).await?;
1461
1462        let entity = command.entity.clone();
1463        for (i, values) in command.batch_values.into_iter().enumerate() {
1464            let mut full_values = values.clone();
1465            full_values.insert("id".to_owned(), command.batch_ids[i].clone());
1466            if let Some(Some(version)) = command.batch_expected_versions.get(i) {
1467                full_values.insert("version".to_owned(), teaql_core::Value::I64(*version + 1));
1468            }
1469
1470            let old_values = command.batch_old_values.get(i).cloned().unwrap_or(None);
1471            let mut new_values = old_values.clone().unwrap_or_default();
1472            for (field, value) in &full_values {
1473                new_values.insert(field.clone(), value.clone());
1474            }
1475
1476            let mut event = RawAuditEvent::updated_with_old_values(
1477                entity.clone(),
1478                full_values.into(),
1479                old_values.map(Into::into),
1480                new_values.into(),
1481                command.update_fields.clone(),
1482            );
1483            if i < command.trace_chains.len() {
1484                event.trace_chain = command.trace_chains[i].clone();
1485            }
1486            self.emit_event(event).map_err(DataServiceError::Runtime)?;
1487        }
1488        Ok(affected)
1489    }
1490
1491    fn fetch_current_event_row(
1492        &self,
1493        _entity: &str,
1494        _id: &Value,
1495        _trace_chain: Vec<teaql_core::TraceNode>,
1496    ) -> Result<Option<Record>, DataServiceError<E::Error>> {
1497        // PER THE USER: "我们不需要在审计的时候去抓旧的值"
1498        // Avoid DB queries during event emission. We rely on in-memory `original_values`.
1499        Ok(None)
1500    }
1501
1502    pub(crate) fn scoped_data_service_internal(&self, entity: String) -> EntityDataService<'a, E> {
1503        EntityDataService {
1504            entity,
1505            data_service: ContextDataService {
1506                metadata: UserContextMetadata {
1507                    context: self.data_service.metadata.context,
1508                },
1509                executor: self.data_service.executor,
1510            },
1511            trace_context: Vec::new(),
1512        }
1513    }
1514}