Skip to main content

teaql_runtime/data_service/
resolved.rs

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