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