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