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