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