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