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