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