Skip to main content

teaql_runtime/data_service/
resolved.rs

1use std::sync::Arc;
2
3use teaql_core::{
4    AggregationCacheOptions, DeleteCommand, Entity, InsertCommand, Record, RecoverCommand,
5    RelationAggregate, SelectQuery, SmartList, UpdateCommand, Value,
6};
7
8use crate::{
9    CheckObjectStatus, DataServiceError, EntityDataServiceBehavior, PurposedSelectQuery,
10    RawAuditEvent, RuntimeError, clear_record_status, mark_record_status,
11};
12
13use super::{
14    AggregationCacheBackend, ContextDataService, EntityDataService, InMemoryAggregationCache,
15    UserContextMetadata, helpers::*,
16};
17
18impl<'a, E> EntityDataService<'a, E>
19where
20    E: teaql_data_service::QueryExecutor
21        + teaql_data_service::MutationExecutor
22        + Send
23        + Sync
24        + 'static,
25{
26    pub(super) fn query_behavior(
27        &self,
28        entity: &str,
29    ) -> Option<Arc<dyn EntityDataServiceBehavior>> {
30        self.data_service
31            .metadata
32            .context
33            .entity_data_service_behavior(entity)
34    }
35
36    pub(super) fn behavior(&self) -> Option<Arc<dyn EntityDataServiceBehavior>> {
37        self.data_service
38            .metadata
39            .context
40            .entity_data_service_behavior(&self.entity)
41    }
42
43    pub fn entity(&self) -> &str {
44        &self.entity
45    }
46
47    pub fn select(&self) -> SelectQuery {
48        SelectQuery::new(self.entity.clone())
49    }
50
51    pub fn insert_command(&self) -> InsertCommand {
52        InsertCommand::new(self.entity.clone())
53    }
54
55    fn enforce_insert_policy(&self, command: &mut InsertCommand) -> Result<(), RuntimeError> {
56        if let Some(policy) = self.data_service.metadata.context.request_policy.as_ref() {
57            policy.enforce_insert(self.data_service.metadata.context, command)?;
58        }
59        Ok(())
60    }
61
62    fn enforce_update_policy(&self, command: &mut UpdateCommand) -> Result<(), RuntimeError> {
63        if let Some(policy) = self.data_service.metadata.context.request_policy.as_ref() {
64            policy.enforce_update(self.data_service.metadata.context, command)?;
65        }
66        Ok(())
67    }
68
69    fn enforce_delete_policy(&self, command: &mut DeleteCommand) -> Result<(), RuntimeError> {
70        if let Some(policy) = self.data_service.metadata.context.request_policy.as_ref() {
71            policy.enforce_delete(self.data_service.metadata.context, command)?;
72        }
73        Ok(())
74    }
75
76    fn enforce_recover_policy(&self, command: &mut RecoverCommand) -> Result<(), RuntimeError> {
77        if let Some(policy) = self.data_service.metadata.context.request_policy.as_ref() {
78            policy.enforce_recover(self.data_service.metadata.context, command)?;
79        }
80        Ok(())
81    }
82
83    fn prepare_select_query(&self, query: &SelectQuery) -> Result<SelectQuery, RuntimeError> {
84        let mut query = query.clone();
85
86        let mut full_trace = self.trace_context.clone();
87        full_trace.extend(query.trace_chain);
88        query.trace_chain = full_trace;
89
90        if let Some(behavior) = self.query_behavior(&query.entity) {
91            behavior.before_select(self.data_service.metadata.context, &mut query)?;
92        }
93        if let Some(policy) = self.data_service.metadata.context.request_policy.as_ref() {
94            policy.enforce_select(self.data_service.metadata.context, &mut query)?;
95        }
96        // Ensure local_key fields for relation loads are projected so that
97        // enhance_query_relations can match parent rows to child records.
98        if !query.relations.is_empty() {
99            if let Some(descriptor) = self.data_service.metadata.context.entity(&query.entity) {
100                for load in &query.relations {
101                    if let Some(relation) = descriptor.relation_by_name(&load.name) {
102                        if !query.projection.contains(&relation.local_key) {
103                            query.projection.push(relation.local_key.clone());
104                        }
105                    }
106                }
107            }
108        }
109        Ok(query)
110    }
111
112    pub fn prepare_insert_command(
113        &self,
114        command: &InsertCommand,
115    ) -> Result<InsertCommand, RuntimeError> {
116        let mut command = command.clone();
117        if let Some(behavior) = self.behavior() {
118            behavior.before_insert(self.data_service.metadata.context, &mut command)?;
119        }
120        self.enforce_insert_policy(&mut command)?;
121
122        let entity = self
123            .data_service
124            .metadata
125            .context
126            .require_entity(&command.entity)?;
127        if let Some(id_property) = entity.id_property() {
128            let needs_id = !command.values.contains_key(&id_property.name)
129                || is_unassigned_id(command.values.get(&id_property.name));
130            if needs_id {
131                let id = self
132                    .data_service
133                    .metadata
134                    .context
135                    .next_id(&command.entity)?;
136                command
137                    .values
138                    .insert(id_property.name.clone(), Value::U64(id));
139            }
140        }
141        ensure_initial_version(&mut command.values, entity);
142        mark_record_status(&mut command.values, CheckObjectStatus::Create);
143        let check_result = self
144            .data_service
145            .metadata
146            .context
147            .check_and_fix_record(&command.entity, &mut command.values);
148        clear_record_status(&mut command.values);
149        check_result?;
150
151        Ok(command)
152    }
153
154    pub fn update_command(&self, id: impl Into<Value>) -> UpdateCommand {
155        UpdateCommand::new(self.entity.clone(), id)
156    }
157
158    pub fn prepare_update_command(
159        &self,
160        command: &UpdateCommand,
161    ) -> Result<UpdateCommand, RuntimeError> {
162        let mut command = command.clone();
163        if let Some(behavior) = self.behavior() {
164            behavior.before_update(self.data_service.metadata.context, &mut command)?;
165        }
166        self.enforce_update_policy(&mut command)?;
167
168        Ok(command)
169    }
170
171    pub fn delete_command(&self, id: impl Into<Value>) -> DeleteCommand {
172        DeleteCommand::new(self.entity.clone(), id)
173    }
174
175    pub fn recover_command(&self, id: impl Into<Value>, expected_version: i64) -> RecoverCommand {
176        RecoverCommand::new(self.entity.clone(), id, expected_version)
177    }
178
179    pub(crate) async fn fetch_all_internal(
180        &self,
181        query: &SelectQuery,
182    ) -> Result<Vec<Record>, DataServiceError<E::Error>> {
183        let query = self
184            .prepare_select_query(query)
185            .map_err(DataServiceError::Runtime)?;
186        self.fetch_prepared_all(&query).await
187    }
188
189    /// Fetch records in streaming mode (chunked).
190    /// Returns a Vec of chunks, each chunk containing up to `chunk_size` rows.
191    /// Each chunk is enhanced (relations, children) before returning.
192    /// Requires E to implement StreamQueryExecutor.
193    pub(crate) async fn fetch_stream_internal(
194        &self,
195        query: &SelectQuery,
196    ) -> Result<Vec<teaql_data_service::StreamChunk>, DataServiceError<E::Error>>
197    where
198        E: teaql_data_service::StreamQueryExecutor,
199    {
200        let query = self
201            .prepare_select_query(query)
202            .map_err(DataServiceError::Runtime)?;
203
204        let chunk_size = query
205            .stream_config
206            .as_ref()
207            .map(|c| c.chunk_size)
208            .unwrap_or(1000);
209
210        let final_comment = self
211            .data_service
212            .resolve_final_comment(&query.trace_chain, query.comment.clone());
213        let mut query = query.clone();
214        query.comment = final_comment;
215
216        let request = teaql_data_service::QueryRequest {
217            query: query.clone(),
218            trace_chain: query.trace_chain.clone(),
219            comment: query.comment.clone(),
220        };
221
222        let chunks = self
223            .data_service
224            .executor
225            .query_stream(request, chunk_size)
226            .await
227            .map_err(DataServiceError::Executor)?;
228
229        // Enhance each chunk
230        let mut enhanced_chunks = Vec::with_capacity(chunks.len());
231        for mut chunk in chunks {
232            self.enhance_object_group_bys_internal(
233                &mut chunk.rows,
234                &query.object_group_bys,
235                &query.trace_chain,
236            )
237            .await?;
238            self.enhance_child_queries_internal(
239                &mut chunk.rows,
240                &query.child_enhancements,
241                &query.trace_chain,
242            )
243            .await?;
244            self.enhance_query_relations_internal(&mut chunk.rows, &query)
245                .await?;
246            enhanced_chunks.push(chunk);
247        }
248
249        Ok(enhanced_chunks)
250    }
251
252    async fn fetch_prepared_all(
253        &self,
254        query: &SelectQuery,
255    ) -> Result<Vec<Record>, DataServiceError<E::Error>> {
256        let mut rows = self.fetch_prepared_query(query).await?;
257        self.enhance_object_group_bys_internal(
258            &mut rows,
259            &query.object_group_bys,
260            &query.trace_chain,
261        )
262        .await?;
263        self.enhance_child_queries_internal(
264            &mut rows,
265            &query.child_enhancements,
266            &query.trace_chain,
267        )
268        .await?;
269        self.enhance_query_relations_internal(&mut rows, query)
270            .await?;
271        Ok(rows)
272    }
273
274    async fn fetch_prepared_query(
275        &self,
276        query: &SelectQuery,
277    ) -> Result<Vec<Record>, DataServiceError<E::Error>> {
278        let final_comment = self
279            .data_service
280            .resolve_final_comment(&query.trace_chain, query.comment.clone());
281        let mut query = query.clone();
282        query.comment = final_comment;
283        if let Some(options) = query.aggregation_cache.filter(|options| options.enabled) {
284            if let Some(cache) = self
285                .data_service
286                .metadata
287                .context
288                .get_resource::<Arc<dyn AggregationCacheBackend>>()
289            {
290                return self
291                    .fetch_prepared_query_with_cache(&query, options, cache.as_ref())
292                    .await;
293            }
294            if let Some(cache) = self
295                .data_service
296                .metadata
297                .context
298                .get_resource::<InMemoryAggregationCache>()
299            {
300                return self
301                    .fetch_prepared_query_with_cache(&query, options, cache)
302                    .await;
303            }
304        }
305        let request = teaql_data_service::QueryRequest {
306            query: query.clone(),
307            trace_chain: query.trace_chain.clone(),
308            comment: query.comment.clone(),
309        };
310        let res = self
311            .data_service
312            .executor
313            .query(request)
314            .await
315            .map_err(DataServiceError::Executor)?;
316        self.data_service
317            .metadata
318            .context
319            .record_metadata_log(&res.metadata);
320        Ok(res.rows)
321    }
322
323    async fn fetch_prepared_query_with_cache(
324        &self,
325        query: &SelectQuery,
326        options: AggregationCacheOptions,
327        cache: &dyn AggregationCacheBackend,
328    ) -> Result<Vec<Record>, DataServiceError<E::Error>> {
329        let key = aggregation_cache_key(
330            cache.namespace(),
331            &aggregation_cache_namespace(&query.entity),
332            query,
333        );
334        if let Some(rows) = cache.get(&key, options.cache_expired_millis) {
335            return Ok(rows);
336        }
337        let request = teaql_data_service::QueryRequest {
338            query: query.clone(),
339            trace_chain: query.trace_chain.clone(),
340            comment: query.comment.clone(),
341        };
342        let res = self
343            .data_service
344            .executor
345            .query(request)
346            .await
347            .map_err(DataServiceError::Executor)?;
348        self.data_service
349            .metadata
350            .context
351            .record_metadata_log(&res.metadata);
352        let rows = res.rows;
353        cache.put(key, rows.clone());
354        Ok(rows)
355    }
356
357    pub(crate) async fn fetch_all_with_relation_aggregates_internal(
358        &self,
359        query: &SelectQuery,
360        relation_aggregates: &[RelationAggregate],
361    ) -> Result<Vec<Record>, DataServiceError<E::Error>> {
362        let query = self
363            .prepare_select_query(query)
364            .map_err(DataServiceError::Runtime)?;
365
366        let mut rows = self.fetch_prepared_all(&query).await?;
367        self.enhance_relation_aggregates_internal(
368            &mut rows,
369            relation_aggregates,
370            query.aggregation_cache,
371            &query.trace_chain,
372        )
373        .await?;
374        Ok(rows)
375    }
376
377    pub(crate) async fn fetch_smart_list_internal(
378        &self,
379        query: &SelectQuery,
380    ) -> Result<SmartList<Record>, DataServiceError<E::Error>> {
381        let query = self
382            .prepare_select_query(query)
383            .map_err(DataServiceError::Runtime)?;
384
385        self.data_service.fetch_smart_list(&query).await
386    }
387
388    pub(crate) async fn fetch_smart_list_with_relation_aggregates_internal(
389        &self,
390        query: &SelectQuery,
391        relation_aggregates: &[RelationAggregate],
392    ) -> Result<SmartList<Record>, DataServiceError<E::Error>> {
393        self.fetch_all_with_relation_aggregates_internal(query, relation_aggregates)
394            .await
395            .map(SmartList::from)
396    }
397
398    pub(crate) async fn fetch_entities_internal<T>(
399        &self,
400        query: &SelectQuery,
401    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
402    where
403        T: Entity,
404    {
405        let query = self
406            .prepare_select_query(query)
407            .map_err(DataServiceError::Runtime)?;
408
409        self.data_service.fetch_entities(&query).await
410    }
411
412    pub(crate) async fn fetch_entities_with_relation_aggregates_internal<T>(
413        &self,
414        query: &SelectQuery,
415        relation_aggregates: &[RelationAggregate],
416    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
417    where
418        T: Entity,
419    {
420        self.fetch_all_with_relation_aggregates_internal(query, relation_aggregates)
421            .await?
422            .into_iter()
423            .map(|record| {
424                let mut entity = T::from_record(record)?;
425                let root = crate::EntityRoot::default();
426                entity.on_loaded(&root as &dyn std::any::Any);
427                Ok(entity)
428            })
429            .collect::<Result<Vec<_>, _>>()
430            .map(SmartList::from)
431            .map_err(DataServiceError::Entity)
432    }
433
434    pub(crate) async fn fetch_enhanced_entities_with_relation_aggregates_internal<T>(
435        &self,
436        query: &SelectQuery,
437        relation_aggregates: &[RelationAggregate],
438    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
439    where
440        T: Entity,
441    {
442        let query = self
443            .prepare_select_query(query)
444            .map_err(DataServiceError::Runtime)?;
445
446        let mut rows = self.fetch_prepared_all(&query).await?;
447        self.enhance_relation_aggregates_internal(
448            &mut rows,
449            relation_aggregates,
450            query.aggregation_cache,
451            &query.trace_chain,
452        )
453        .await?;
454        self.enhance_relations_internal(&mut rows).await?;
455        rows.into_iter()
456            .map(|record| {
457                let mut entity = T::from_record(record)?;
458                let root = crate::EntityRoot::default();
459                entity.on_loaded(&root as &dyn std::any::Any);
460                Ok(entity)
461            })
462            .collect::<Result<Vec<_>, _>>()
463            .map(SmartList::from)
464            .map_err(DataServiceError::Entity)
465    }
466
467    pub(crate) async fn fetch_enhanced_entities_internal<T>(
468        &self,
469        query: &SelectQuery,
470    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
471    where
472        T: Entity,
473    {
474        let query = self
475            .prepare_select_query(query)
476            .map_err(DataServiceError::Runtime)?;
477
478        let mut rows = self.fetch_prepared_all(&query).await?;
479        self.enhance_relations_internal(&mut rows).await?;
480        let root = self
481            .data_service
482            .metadata
483            .context
484            .get_resource::<crate::EntityRoot>()
485            .cloned();
486        rows.into_iter()
487            .map(|record| {
488                let mut entity = T::from_record(record)?;
489                if let Some(ref root) = root {
490                    entity.on_loaded(root as &dyn std::any::Any);
491                }
492                Ok(entity)
493            })
494            .collect::<Result<Vec<_>, _>>()
495            .map(SmartList::from)
496            .map_err(DataServiceError::Entity)
497    }
498
499    #[doc(hidden)]
500    pub async fn fetch_all(
501        &self,
502        query: &PurposedSelectQuery,
503    ) -> Result<Vec<Record>, DataServiceError<E::Error>> {
504        self.fetch_all_internal(query.as_query()).await
505    }
506
507    #[doc(hidden)]
508    pub async fn fetch_stream(
509        &self,
510        query: &PurposedSelectQuery,
511    ) -> Result<Vec<teaql_data_service::StreamChunk>, DataServiceError<E::Error>>
512    where
513        E: teaql_data_service::StreamQueryExecutor,
514    {
515        self.fetch_stream_internal(query.as_query()).await
516    }
517
518    #[doc(hidden)]
519    pub async fn fetch_smart_list(
520        &self,
521        query: &PurposedSelectQuery,
522    ) -> Result<SmartList<Record>, DataServiceError<E::Error>> {
523        self.fetch_smart_list_internal(query.as_query()).await
524    }
525
526    #[doc(hidden)]
527    pub async fn fetch_smart_list_with_relation_aggregates(
528        &self,
529        query: &PurposedSelectQuery,
530        relation_aggregates: &[RelationAggregate],
531    ) -> Result<SmartList<Record>, DataServiceError<E::Error>> {
532        self.fetch_smart_list_with_relation_aggregates_internal(
533            query.as_query(),
534            relation_aggregates,
535        )
536        .await
537    }
538
539    #[doc(hidden)]
540    pub async fn fetch_entities<T>(
541        &self,
542        query: &PurposedSelectQuery,
543    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
544    where
545        T: Entity,
546    {
547        self.fetch_entities_internal(query.as_query()).await
548    }
549
550    #[doc(hidden)]
551    pub async fn fetch_enhanced_entities<T>(
552        &self,
553        query: &PurposedSelectQuery,
554    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
555    where
556        T: Entity,
557    {
558        self.fetch_enhanced_entities_internal(query.as_query())
559            .await
560    }
561
562    #[doc(hidden)]
563    pub async fn fetch_enhanced_entities_with_relation_aggregates<T>(
564        &self,
565        query: &PurposedSelectQuery,
566        relation_aggregates: &[RelationAggregate],
567    ) -> Result<SmartList<T>, DataServiceError<E::Error>>
568    where
569        T: Entity,
570    {
571        self.fetch_enhanced_entities_with_relation_aggregates_internal(
572            query.as_query(),
573            relation_aggregates,
574        )
575        .await
576    }
577
578    pub(crate) async fn insert_internal(
579        &self,
580        command: &InsertCommand,
581    ) -> Result<u64, DataServiceError<E::Error>> {
582        let command = self
583            .prepare_insert_command(command)
584            .map_err(DataServiceError::Runtime)?;
585        self.execute_prepared_insert_with_comment(command, self.trace_context.clone())
586            .await
587    }
588
589    pub(crate) async fn update_internal(
590        &self,
591        command: &UpdateCommand,
592    ) -> Result<u64, DataServiceError<E::Error>> {
593        let command = self
594            .prepare_update_command(command)
595            .map_err(DataServiceError::Runtime)?;
596        self.execute_prepared_update_with_comment(command, self.trace_context.clone())
597            .await
598    }
599
600    pub(crate) async fn delete_internal(
601        &self,
602        command: &DeleteCommand,
603    ) -> Result<u64, DataServiceError<E::Error>> {
604        self.delete_scoped_internal(command, self.trace_context.clone())
605            .await
606    }
607
608    pub(crate) async fn delete_scoped_internal(
609        &self,
610        command: &DeleteCommand,
611        trace_chain: Vec<teaql_core::TraceNode>,
612    ) -> Result<u64, DataServiceError<E::Error>> {
613        let mut command = command.clone();
614        command.trace_chain = trace_chain.clone();
615        if let Some(behavior) = self.behavior() {
616            behavior
617                .before_delete(self.data_service.metadata.context, &mut command)
618                .map_err(DataServiceError::Runtime)?;
619        }
620        self.enforce_delete_policy(&mut command)
621            .map_err(DataServiceError::Runtime)?;
622
623        let old_values =
624            self.fetch_current_event_row(&command.entity, &command.id, trace_chain.clone())?;
625        let affected = self.data_service.delete(&command).await?;
626
627        let mut event = RawAuditEvent::deleted_with_old_values(
628            command.entity,
629            command.id,
630            command.expected_version,
631            old_values,
632        );
633        event.trace_chain = trace_chain;
634        self.emit_event(event).map_err(DataServiceError::Runtime)?;
635        Ok(affected)
636    }
637
638    pub(crate) async fn recover_internal(
639        &self,
640        command: &RecoverCommand,
641    ) -> Result<u64, DataServiceError<E::Error>> {
642        let mut command = command.clone();
643        command.trace_chain = self.trace_context.clone();
644        if let Some(behavior) = self.behavior() {
645            behavior
646                .before_recover(self.data_service.metadata.context, &mut command)
647                .map_err(DataServiceError::Runtime)?;
648        }
649        self.enforce_recover_policy(&mut command)
650            .map_err(DataServiceError::Runtime)?;
651        let old_values = self.fetch_current_event_row(
652            &command.entity,
653            &command.id,
654            command.trace_chain.clone(),
655        )?;
656        let affected = self.data_service.recover(&command).await?;
657        let event = RawAuditEvent::recovered_with_old_values(
658            command.entity,
659            command.id,
660            command.expected_version,
661            old_values,
662        );
663        self.emit_event(event).map_err(DataServiceError::Runtime)?;
664        Ok(affected)
665    }
666
667    fn emit_event(&self, event: RawAuditEvent) -> Result<(), RuntimeError> {
668        self.data_service.metadata.context.send_event(event)
669    }
670
671    #[allow(dead_code)]
672    pub(super) async fn execute_prepared_insert(
673        &self,
674        command: InsertCommand,
675    ) -> Result<u64, DataServiceError<E::Error>> {
676        self.execute_prepared_insert_with_comment(command, Vec::new())
677            .await
678    }
679
680    pub(super) async fn execute_prepared_insert_with_comment(
681        &self,
682        mut command: InsertCommand,
683        trace_chain: Vec<teaql_core::TraceNode>,
684    ) -> Result<u64, DataServiceError<E::Error>> {
685        command.trace_chain = trace_chain.clone();
686        let affected = self.data_service.insert(&command).await?;
687        let mut event = RawAuditEvent::created(command.entity, command.values);
688        event.trace_chain = trace_chain;
689        self.emit_event(event).map_err(DataServiceError::Runtime)?;
690        Ok(affected)
691    }
692
693    pub(super) async fn execute_prepared_batch_insert(
694        &self,
695        command: teaql_core::BatchInsertCommand,
696    ) -> Result<u64, DataServiceError<E::Error>> {
697        if command.batch_values.is_empty() {
698            return Ok(0);
699        }
700        let affected = self.data_service.batch_insert(&command).await?;
701
702        let entity = command.entity.clone();
703        for (i, values) in command.batch_values.into_iter().enumerate() {
704            let mut event = RawAuditEvent::created(entity.clone(), values);
705            if i < command.trace_chains.len() {
706                event.trace_chain = command.trace_chains[i].clone();
707            }
708            self.emit_event(event).map_err(DataServiceError::Runtime)?;
709        }
710        Ok(affected)
711    }
712
713    #[allow(dead_code)]
714    pub(super) async fn execute_prepared_update(
715        &self,
716        command: UpdateCommand,
717    ) -> Result<u64, DataServiceError<E::Error>> {
718        self.execute_prepared_update_with_comment(command, Vec::new())
719            .await
720    }
721
722    pub(super) async fn execute_prepared_update_with_comment(
723        &self,
724        mut command: UpdateCommand,
725        trace_chain: Vec<teaql_core::TraceNode>,
726    ) -> Result<u64, DataServiceError<E::Error>> {
727        command.trace_chain = trace_chain.clone();
728
729        let mut old_values = command.old_values.clone();
730        let needs_fetch = match &old_values {
731            Some(snapshot) => !command.values.keys().all(|k| snapshot.contains_key(k)),
732            None => true,
733        };
734        if needs_fetch {
735            old_values =
736                self.fetch_current_event_row(&command.entity, &command.id, trace_chain.clone())?;
737        }
738
739        let affected = self.data_service.update(&command).await?;
740        let updated_fields = command.values.keys().cloned().collect();
741        let mut values = command.values.clone();
742        values.insert("id".to_owned(), command.id.clone());
743        if let Some(version) = command.expected_version {
744            values.insert("version".to_owned(), Value::I64(version + 1));
745        }
746        let mut new_values = old_values.clone().unwrap_or_default();
747        for (field, value) in &values {
748            new_values.insert(field.clone(), value.clone());
749        }
750        let mut event = RawAuditEvent::updated_with_old_values(
751            command.entity,
752            values,
753            old_values,
754            new_values,
755            updated_fields,
756        );
757        event.trace_chain = trace_chain;
758        self.emit_event(event).map_err(DataServiceError::Runtime)?;
759        Ok(affected)
760    }
761
762    pub(super) async fn execute_prepared_batch_update(
763        &self,
764        command: teaql_core::BatchUpdateCommand,
765    ) -> Result<u64, DataServiceError<E::Error>> {
766        if command.batch_values.is_empty() {
767            return Ok(0);
768        }
769        let affected = self.data_service.batch_update(&command).await?;
770
771        let entity = command.entity.clone();
772        for (i, values) in command.batch_values.into_iter().enumerate() {
773            let mut full_values = values.clone();
774            full_values.insert("id".to_owned(), command.batch_ids[i].clone());
775            if let Some(Some(version)) = command.batch_expected_versions.get(i) {
776                full_values.insert("version".to_owned(), teaql_core::Value::I64(*version + 1));
777            }
778
779            let old_values = command.batch_old_values.get(i).cloned().unwrap_or(None);
780            let mut new_values = old_values.clone().unwrap_or_default();
781            for (field, value) in &full_values {
782                new_values.insert(field.clone(), value.clone());
783            }
784
785            let mut event = RawAuditEvent::updated_with_old_values(
786                entity.clone(),
787                full_values,
788                old_values,
789                new_values,
790                command.update_fields.clone(),
791            );
792            if i < command.trace_chains.len() {
793                event.trace_chain = command.trace_chains[i].clone();
794            }
795            self.emit_event(event).map_err(DataServiceError::Runtime)?;
796        }
797        Ok(affected)
798    }
799
800    fn fetch_current_event_row(
801        &self,
802        _entity: &str,
803        _id: &Value,
804        _trace_chain: Vec<teaql_core::TraceNode>,
805    ) -> Result<Option<Record>, DataServiceError<E::Error>> {
806        // PER THE USER: "我们不需要在审计的时候去抓旧的值"
807        // Avoid DB queries during event emission. We rely on in-memory `original_values`.
808        Ok(None)
809    }
810
811    pub(crate) fn scoped_data_service_internal(&self, entity: String) -> EntityDataService<'a, E> {
812        EntityDataService {
813            entity,
814            data_service: ContextDataService {
815                metadata: UserContextMetadata {
816                    context: self.data_service.metadata.context,
817                },
818                executor: self.data_service.executor,
819            },
820            trace_context: Vec::new(),
821        }
822    }
823}