Skip to main content

teaql_runtime/data_service/
resolved.rs

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