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 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 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 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 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}