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