Skip to main content

es_entity/
query.rs

1//! Query execution infrastructure for event-sourced entities.
2//!
3//! This module provides the underlying query types used by the `es_query!` macro.
4//! **These types are not intended to be used directly** - instead, use the `es_query!`
5//! macro which provides a simpler interface for querying index tables and automatically
6//! hydrating entities from their events.
7//!
8//! # Example
9//!
10//! Instead of using these types directly, use the `es_query!` macro:
11//!
12//! ```rust,ignore
13//! es_query!(
14//!     "SELECT id FROM users WHERE name = $1",
15//!     name
16//! ).fetch_optional(&pool).await
17//! ```
18//!
19//! See the `es_query!` macro documentation for more details.
20
21use crate::{
22    db,
23    events::{EntityEvents, HydrationRow},
24    one_time_executor::IntoOneTimeExecutor,
25    snapshot::NO_SNAPSHOT_FINGERPRINT,
26    traits::*,
27    tree_query::{TreeQuerySource, build_tree_query, partition_by_tag, snapshot_fingerprints},
28};
29
30/// Query builder for event-sourced entities.
31///
32/// This type is generated by the `es_query!` macro and should not be constructed directly.
33/// It wraps a SQLx query and provides methods to fetch and hydrate entities from their events.
34///
35/// `R` is the row type the underlying `sqlx::query_as!` decodes into —
36/// `GenericEvent<Id>` for a plain repo, `SnapshotGenericEvent<Id>` for a
37/// `#[es_repo(snapshot)]` repo — normalised into a `HydrationRow<Id>` before
38/// an entity is built from it.
39pub struct EsQuery<
40    'q,
41    Repo,
42    Flavor,
43    F,
44    A,
45    R = crate::events::GenericEvent<
46        <<<Repo as EsRepo>::Entity as EsEntity>::Event as EsEvent>::EntityId,
47    >,
48> where
49    Repo: EsRepo,
50{
51    inner: sqlx::query::Map<'q, db::Db, F, A>,
52    source: TreeQuerySource<<<<Repo as EsRepo>::Entity as EsEntity>::Event as EsEvent>::EntityId>,
53    _repo: std::marker::PhantomData<Repo>,
54    _flavor: std::marker::PhantomData<Flavor>,
55    _row: std::marker::PhantomData<fn() -> R>,
56}
57
58/// Query flavor for flat entities without nested relationships.
59pub struct EsQueryFlavorFlat;
60
61/// Query flavor for entities with nested relationships that need to be loaded recursively.
62pub struct EsQueryFlavorNested;
63
64impl<'q, Repo, Flavor, F, A, R> EsQuery<'q, Repo, Flavor, F, A, R>
65where
66    Repo: EsRepo,
67    <<<Repo as EsRepo>::Entity as EsEntity>::Event as EsEvent>::EntityId: Unpin,
68    F: FnMut(db::Row) -> Result<R, sqlx::Error> + Send,
69    R: Into<HydrationRow<<<<Repo as EsRepo>::Entity as EsEntity>::Event as EsEvent>::EntityId>>
70        + Send
71        + Unpin,
72    A: 'q + Send + sqlx::IntoArguments<'q, db::Db>,
73{
74    pub fn new(
75        query: sqlx::query::Map<'q, db::Db, F, A>,
76        source: TreeQuerySource<
77            <<<Repo as EsRepo>::Entity as EsEntity>::Event as EsEvent>::EntityId,
78        >,
79    ) -> Self {
80        Self {
81            inner: query,
82            source,
83            _repo: std::marker::PhantomData,
84            _flavor: std::marker::PhantomData,
85            _row: std::marker::PhantomData,
86        }
87    }
88
89    async fn fetch_optional_inner(
90        self,
91        op: impl IntoOneTimeExecutor<'_>,
92    ) -> Result<Option<<Repo as EsRepo>::Entity>, crate::RepoFault> {
93        let executor = op.into_executor();
94        let rows = executor.fetch_all(self.inner).await?;
95        if rows.is_empty() {
96            return Ok(None);
97        }
98
99        EntityEvents::load_first(rows.into_iter()).map_err(Into::into)
100    }
101
102    async fn fetch_n_inner(
103        self,
104        op: impl IntoOneTimeExecutor<'_>,
105        first: usize,
106    ) -> Result<(Vec<<Repo as EsRepo>::Entity>, bool), crate::RepoFault> {
107        let executor = op.into_executor();
108        let rows = executor.fetch_all(self.inner).await?;
109        EntityEvents::load_n(rows.into_iter(), first).map_err(Into::into)
110    }
111
112    async fn fetch_tree_rows(
113        self,
114        op: impl IntoOneTimeExecutor<'_>,
115        include_deleted: bool,
116    ) -> Result<
117        (
118            Vec<HydrationRow<<<<Repo as EsRepo>::Entity as EsEntity>::Event as EsEvent>::EntityId>>,
119            std::collections::HashMap<i32, Vec<db::Row>>,
120        ),
121        crate::RepoFault,
122    > {
123        let executor = op.into_executor();
124        let spec = <Repo as EsRepo>::nested_tree_spec();
125        let sql = build_tree_query(
126            self.source.user_sql,
127            self.source.order_by_cols,
128            &spec,
129            include_deleted,
130            self.source.n_user_args + 1,
131        );
132
133        let mut inner = self.inner;
134        let mut args = sqlx::Execute::take_arguments(&mut inner)
135            .map_err(sqlx::Error::Encode)?
136            .unwrap_or_default();
137
138        // A full-history load binds `NO_SNAPSHOT_FINGERPRINT` as the root's
139        // own fingerprint, forcing every node in the tree to bind the
140        // sentinel too. The `snapshot_table_name` check guards against a
141        // false positive: `NoSnapshot::FINGERPRINT` *equals* that sentinel,
142        // so an ordinary non-snapshot root would otherwise always match.
143        let full_history = spec.snapshot_table_name.is_some()
144            && self.source.snapshot_fingerprint == NO_SNAPSHOT_FINGERPRINT;
145        // The root's own `es_query!` call site already bound its
146        // fingerprint, so only descendants need binding here — walked in
147        // the same DFS order `build_tree_query` assigned positions in.
148        let mut fingerprints = snapshot_fingerprints(&spec).into_iter();
149        if spec.snapshot_table_name.is_some() {
150            fingerprints.next();
151        }
152        for fp in fingerprints {
153            let bind = if full_history {
154                NO_SNAPSHOT_FINGERPRINT
155            } else {
156                fp
157            };
158            sqlx::Arguments::add(&mut args, bind).map_err(sqlx::Error::Encode)?;
159        }
160
161        let rows: Vec<db::Row> = sqlx::query_with::<db::Db, _>(&sql, args)
162            .fetch_all(executor)
163            .await?;
164        let mut by_tag = partition_by_tag(rows)?;
165        let root_rows = by_tag.remove(&0).unwrap_or_default();
166        let root = root_rows
167            .iter()
168            .map(self.source.decode)
169            .collect::<Result<Vec<_>, sqlx::Error>>()?;
170        Ok((root, by_tag))
171    }
172}
173
174impl<'q, Repo, F, A, R> EsQuery<'q, Repo, EsQueryFlavorFlat, F, A, R>
175where
176    Repo: EsRepo,
177    <<<Repo as EsRepo>::Entity as EsEntity>::Event as EsEvent>::EntityId: Unpin,
178    F: FnMut(db::Row) -> Result<R, sqlx::Error> + Send,
179    R: Into<HydrationRow<<<<Repo as EsRepo>::Entity as EsEntity>::Event as EsEvent>::EntityId>>
180        + Send
181        + Unpin,
182    A: 'q + Send + sqlx::IntoArguments<'q, db::Db>,
183{
184    /// Fetches at most one entity from the query results.
185    ///
186    /// Returns `Ok(None)` if no entities match the query, or `Ok(Some(entity))` if found.
187    pub async fn fetch_optional(
188        self,
189        op: impl IntoOneTimeExecutor<'_>,
190    ) -> Result<Option<<Repo as EsRepo>::Entity>, crate::RepoFault> {
191        self.fetch_optional_inner(op).await
192    }
193
194    /// Fetches up to `first` entities from the query results.
195    ///
196    /// Returns a tuple of (entities, has_more) where `has_more` indicates if there
197    /// were more entities available beyond the requested limit.
198    pub async fn fetch_n(
199        self,
200        op: impl IntoOneTimeExecutor<'_>,
201        first: usize,
202    ) -> Result<(Vec<<Repo as EsRepo>::Entity>, bool), crate::RepoFault> {
203        self.fetch_n_inner(op, first).await
204    }
205}
206
207impl<'q, Repo, F, A, R> EsQuery<'q, Repo, EsQueryFlavorNested, F, A, R>
208where
209    Repo: EsRepo,
210    <<<Repo as EsRepo>::Entity as EsEntity>::Event as EsEvent>::EntityId: Unpin,
211    F: FnMut(db::Row) -> Result<R, sqlx::Error> + Send,
212    R: Into<HydrationRow<<<<Repo as EsRepo>::Entity as EsEntity>::Event as EsEvent>::EntityId>>
213        + Send
214        + Unpin,
215    A: 'q + Send + sqlx::IntoArguments<'q, db::Db>,
216{
217    /// Fetches at most one entity and loads all nested relationships.
218    ///
219    /// Returns `Ok(None)` if no entities match, or `Ok(Some(entity))` with all
220    /// nested entities loaded if found.
221    pub async fn fetch_optional(
222        self,
223        op: impl IntoOneTimeExecutor<'_>,
224    ) -> Result<Option<<Repo as EsRepo>::Entity>, crate::RepoFault> {
225        self.fetch_optional_tree(op, false).await
226    }
227
228    /// Fetches up to `first` entities and loads all nested relationships.
229    ///
230    /// Returns a tuple of (entities, has_more) where all entities have their nested
231    /// relationships loaded, and `has_more` indicates if more entities were available.
232    pub async fn fetch_n(
233        self,
234        op: impl IntoOneTimeExecutor<'_>,
235        first: usize,
236    ) -> Result<(Vec<<Repo as EsRepo>::Entity>, bool), crate::RepoFault> {
237        self.fetch_n_tree(op, first, false).await
238    }
239
240    /// Like [`fetch_optional`](EsQuery::fetch_optional) but transitively includes
241    /// soft-deleted nested entities.
242    pub async fn fetch_optional_include_deleted(
243        self,
244        op: impl IntoOneTimeExecutor<'_>,
245    ) -> Result<Option<<Repo as EsRepo>::Entity>, crate::RepoFault> {
246        self.fetch_optional_tree(op, true).await
247    }
248
249    /// Like [`fetch_n`](EsQuery::fetch_n) but transitively includes soft-deleted
250    /// nested entities.
251    pub async fn fetch_n_include_deleted(
252        self,
253        op: impl IntoOneTimeExecutor<'_>,
254        first: usize,
255    ) -> Result<(Vec<<Repo as EsRepo>::Entity>, bool), crate::RepoFault> {
256        self.fetch_n_tree(op, first, true).await
257    }
258
259    async fn fetch_optional_tree(
260        self,
261        op: impl IntoOneTimeExecutor<'_>,
262        include_deleted: bool,
263    ) -> Result<Option<<Repo as EsRepo>::Entity>, crate::RepoFault> {
264        let (root, mut by_tag) = self.fetch_tree_rows(op, include_deleted).await?;
265        let Some(entity) = EntityEvents::load_first::<<Repo as EsRepo>::Entity>(root)? else {
266            return Ok(None);
267        };
268        let mut entities = [entity];
269        let mut cursor = 1i32;
270        <Repo as EsRepo>::hydrate_nested_from_rows(&mut by_tag, &mut cursor, &mut entities)?;
271        let [entity] = entities;
272        Ok(Some(entity))
273    }
274
275    async fn fetch_n_tree(
276        self,
277        op: impl IntoOneTimeExecutor<'_>,
278        first: usize,
279        include_deleted: bool,
280    ) -> Result<(Vec<<Repo as EsRepo>::Entity>, bool), crate::RepoFault> {
281        let (root, mut by_tag) = self.fetch_tree_rows(op, include_deleted).await?;
282        let (mut entities, more) = EntityEvents::load_n::<<Repo as EsRepo>::Entity>(root, first)?;
283        let mut cursor = 1i32;
284        <Repo as EsRepo>::hydrate_nested_from_rows(&mut by_tag, &mut cursor, &mut entities)?;
285        Ok((entities, more))
286    }
287}