Skip to main content

teaql_runtime/data_service/
resolved.rs

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        // Ensure local_key fields for relation loads are projected so that
283        // enhance_query_relations can match parent rows to child records.
284        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        // An empty projection means all mapped properties. Once callers use a
300        // reduced projection, entity hydration must still retain identity and
301        // the optimistic-lock token. Aggregate/DTO paths do not call this
302        // helper, so their SQL shape is left untouched.
303        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        // Relation probes build a fresh, runtime-owned query for every foreign key.
435        // Keep that ownership all the way into the provider when no runtime
436        // enhancement needs to retain the query. Going through
437        // `fetch_prepared_all(&query)` cloned the complete SelectQuery (including
438        // projections and expression trees) twice for every probe.
439        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        // A unique final key makes the retained order deterministic even when
710        // the caller's preceding sort fields contain ties.
711        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    /// Fetch root records from the provider cursor without materializing them.
867    /// Relation and aggregate enhancement needs a separate batched protocol and
868    /// is rejected here instead of silently returning incomplete entities.
869    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        // PER THE USER: "我们不需要在审计的时候去抓旧的值"
1754        // Avoid DB queries during event emission. We rely on in-memory `original_values`.
1755        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}