Skip to main content

dovecote_sqlx_sqlite/
page.rs

1//! `SQLite` live and finite snapshot paging.
2
3use crate::{
4    begin_read, commit_transaction,
5    error::PageError,
6    hydrate::{DurableRow, hydrate_page},
7    install_foreign_keys,
8};
9use dovecote::{Limit, PagedEvent, RowId, TenantId};
10use sqlx::{Sqlite, SqlitePool, Transaction, query_as, query_scalar};
11use std::marker::PhantomData;
12
13pub(crate) async fn page_for_scope(
14    pool: &SqlitePool,
15    tenant_id: Option<&TenantId>,
16    after_row_id: Option<RowId>,
17    limit: Limit,
18) -> Result<Vec<PagedEvent>, PageError> {
19    let mut connection = pool
20        .acquire()
21        .await
22        .map_err(|source| PageError::sql("acquire live page connection", source))?;
23    install_foreign_keys(&mut connection)
24        .await
25        .map_err(|source| PageError::sql("enable live-page foreign keys", source))?;
26    read_page_scoped(
27        &mut *connection,
28        tenant_id,
29        after_row_id.map_or(0, RowId::get),
30        None,
31        limit,
32    )
33    .await
34}
35
36pub(crate) async fn begin_snapshot_for_scope(
37    pool: &SqlitePool,
38    tenant_id: Option<&TenantId>,
39) -> Result<SnapshotPager, PageError> {
40    let mut transaction = begin_read(pool)
41        .await
42        .map_err(|source| PageError::sql("begin snapshot transaction", source))?;
43    let upper_bound = match query_scalar::<_, Option<i64>>(
44        "SELECT MAX(row_id) FROM dovecote_events WHERE (? IS NULL OR tenant_id = ?)",
45    )
46    .bind(tenant_id.map(TenantId::as_str))
47    .bind(tenant_id.map(TenantId::as_str))
48    .fetch_one(&mut *transaction)
49    .await
50    {
51        Ok(value) => value,
52        Err(source) => {
53            let _ = transaction.rollback().await;
54            return Err(PageError::sql("read snapshot upper row ID", source));
55        }
56    };
57    let upper_bound = match upper_bound
58        .map(|value| RowId::new(value).map_err(|error| PageError::serialization(error.to_string())))
59        .transpose()
60    {
61        Ok(value) => value,
62        Err(error) => {
63            let _ = transaction.rollback().await;
64            return Err(error);
65        }
66    };
67    Ok(SnapshotPager {
68        transaction: Some(transaction),
69        upper_bound,
70        cursor: None,
71        exhausted: upper_bound.is_none(),
72        tenant_id: tenant_id.cloned(),
73        _not_send: PhantomData,
74    })
75}
76
77/// A finite pager retaining one `SQLite` read transaction. The explicit marker
78/// makes accidental movement across unrelated executors a compile-time error.
79///
80/// ```compile_fail
81/// use dovecote_sqlx_sqlite::SnapshotPager;
82///
83/// fn requires_send<T: Send>() {}
84///
85/// fn main() {
86///     requires_send::<SnapshotPager>();
87/// }
88/// ```
89pub struct SnapshotPager {
90    transaction: Option<Transaction<'static, Sqlite>>,
91    upper_bound: Option<RowId>,
92    cursor: Option<RowId>,
93    exhausted: bool,
94    _not_send: PhantomData<*mut ()>,
95    tenant_id: Option<TenantId>,
96}
97
98impl SnapshotPager {
99    /// Returns the last row ID returned by a non-empty page.
100    #[must_use]
101    pub const fn cursor(&self) -> Option<RowId> {
102        self.cursor
103    }
104    /// Returns the maximum row ID visible to this pager.
105    #[must_use]
106    pub const fn upper_bound(&self) -> Option<RowId> {
107        self.upper_bound
108    }
109    /// Returns whether the pager has returned its final page.
110    #[must_use]
111    pub const fn is_exhausted(&self) -> bool {
112        self.exhausted
113    }
114
115    /// Reads the next bounded page from the retained snapshot.
116    ///
117    /// # Errors
118    /// Returns an error if a stored event or delivery cannot be validated or the
119    /// database read fails. Roll back or drop the pager after a failed read.
120    pub async fn next_page(&mut self, limit: Limit) -> Result<Vec<PagedEvent>, PageError> {
121        if self.exhausted {
122            return Ok(Vec::new());
123        }
124
125        let transaction = self.transaction.as_mut().ok_or(PageError::Closed)?;
126        let upper = self
127            .upper_bound
128            .expect("non-exhausted pager has an upper bound");
129        let result = read_page_scoped(
130            &mut **transaction,
131            self.tenant_id.as_ref(),
132            self.cursor.map_or(0, RowId::get),
133            Some(upper.get()),
134            limit,
135        )
136        .await;
137        let rows = match result {
138            Ok(rows) => rows,
139            Err(error) => {
140                if let Some(transaction) = self.transaction.take() {
141                    let _ = transaction.rollback().await;
142                }
143
144                return Err(error);
145            }
146        };
147        if let Some(last) = rows.last() {
148            self.cursor = Some(last.row_id());
149
150            if rows.len() < limit.get() as usize || self.cursor == self.upper_bound {
151                self.exhausted = true;
152            }
153        } else {
154            self.exhausted = true;
155        }
156
157        Ok(rows)
158    }
159
160    /// Commits the read-only snapshot transaction and releases its connection.
161    ///
162    /// # Errors
163    /// Returns a database error if committing the read transaction fails. A lost
164    /// commit response does not establish whether the server committed.
165    pub async fn finish(mut self) -> Result<(), PageError> {
166        let Some(transaction) = self.transaction.take() else {
167            return Ok(());
168        };
169        commit_transaction(transaction)
170            .await
171            .map_err(|source| PageError::sql("finish snapshot transaction", source))
172    }
173    /// Rolls back the snapshot transaction and releases its connection.
174    ///
175    /// # Errors
176    /// Returns a database error if rolling back the read transaction fails.
177    pub async fn rollback(mut self) -> Result<(), PageError> {
178        let Some(transaction) = self.transaction.take() else {
179            return Ok(());
180        };
181        transaction
182            .rollback()
183            .await
184            .map_err(|source| PageError::sql("rollback snapshot transaction", source))
185    }
186    /// Closes the pager by rolling back its transaction.
187    ///
188    /// # Errors
189    /// Returns a database error if rolling back the read transaction fails.
190    pub async fn close(self) -> Result<(), PageError> {
191        self.rollback().await
192    }
193}
194
195async fn read_page_scoped<'c, E>(
196    executor: E,
197    tenant_id: Option<&TenantId>,
198    after_row_id: i64,
199    upper_bound: Option<i64>,
200    limit: Limit,
201) -> Result<Vec<PagedEvent>, PageError>
202where
203    E: sqlx::Executor<'c, Database = Sqlite>,
204{
205    let rows = match (tenant_id, upper_bound) {
206        (Some(tenant_id), Some(upper)) => {
207            query_as::<_, DurableRow>(SCOPED_PAGE_SNAPSHOT_SQL)
208                .bind(after_row_id)
209                .bind(upper)
210                .bind(tenant_id.as_str())
211                .bind(i64::from(limit.get()))
212                .fetch_all(executor)
213                .await
214        }
215        (Some(tenant_id), None) => {
216            query_as::<_, DurableRow>(SCOPED_PAGE_SQL)
217                .bind(after_row_id)
218                .bind(tenant_id.as_str())
219                .bind(i64::from(limit.get()))
220                .fetch_all(executor)
221                .await
222        }
223        (None, Some(upper)) => {
224            query_as::<_, DurableRow>(PAGE_SNAPSHOT_SQL)
225                .bind(after_row_id)
226                .bind(upper)
227                .bind(i64::from(limit.get()))
228                .fetch_all(executor)
229                .await
230        }
231        (None, None) => {
232            query_as::<_, DurableRow>(PAGE_SQL)
233                .bind(after_row_id)
234                .bind(i64::from(limit.get()))
235                .fetch_all(executor)
236                .await
237        }
238    }
239    .map_err(|source| PageError::sql("read event page", source))?;
240    rows.into_iter()
241        .map(hydrate_page)
242        .collect::<Result<Vec<_>, _>>()
243        .map_err(PageError::serialization)
244}
245
246const PAGE_SQL: &str = "SELECT e.row_id, e.tenant_id, e.stream, e.specversion, e.event_id, e.source, e.event_type, e.subject, e.occurred_at, e.enqueued_at, e.datacontenttype, e.dataschema, e.partitionkey, e.extensions, e.data_kind, e.data, d.state, d.available_at, d.attempts, d.claim_token, d.claimed_by, d.claim_expires_at, d.last_failure_code, d.last_failure_detail, d.delivered_at, d.quarantined_at, d.quarantine_reason FROM dovecote_events AS e LEFT JOIN dovecote_deliveries AS d ON d.tenant_id = e.tenant_id AND d.event_row_id = e.row_id WHERE e.row_id > ? ORDER BY e.row_id ASC LIMIT ?";
247const PAGE_SNAPSHOT_SQL: &str = "SELECT e.row_id, e.tenant_id, e.stream, e.specversion, e.event_id, e.source, e.event_type, e.subject, e.occurred_at, e.enqueued_at, e.datacontenttype, e.dataschema, e.partitionkey, e.extensions, e.data_kind, e.data, d.state, d.available_at, d.attempts, d.claim_token, d.claimed_by, d.claim_expires_at, d.last_failure_code, d.last_failure_detail, d.delivered_at, d.quarantined_at, d.quarantine_reason FROM dovecote_events AS e LEFT JOIN dovecote_deliveries AS d ON d.tenant_id = e.tenant_id AND d.event_row_id = e.row_id WHERE e.row_id > ? AND e.row_id <= ? ORDER BY e.row_id ASC LIMIT ?";
248const SCOPED_PAGE_SQL: &str = "SELECT e.row_id, e.tenant_id, e.stream, e.specversion, e.event_id, e.source, e.event_type, e.subject, e.occurred_at, e.enqueued_at, e.datacontenttype, e.dataschema, e.partitionkey, e.extensions, e.data_kind, e.data, d.state, d.available_at, d.attempts, d.claim_token, d.claimed_by, d.claim_expires_at, d.last_failure_code, d.last_failure_detail, d.delivered_at, d.quarantined_at, d.quarantine_reason FROM dovecote_events AS e LEFT JOIN dovecote_deliveries AS d ON d.tenant_id = e.tenant_id AND d.event_row_id = e.row_id WHERE e.row_id > ? AND e.tenant_id = ? ORDER BY e.row_id ASC LIMIT ?";
249const SCOPED_PAGE_SNAPSHOT_SQL: &str = "SELECT e.row_id, e.tenant_id, e.stream, e.specversion, e.event_id, e.source, e.event_type, e.subject, e.occurred_at, e.enqueued_at, e.datacontenttype, e.dataschema, e.partitionkey, e.extensions, e.data_kind, e.data, d.state, d.available_at, d.attempts, d.claim_token, d.claimed_by, d.claim_expires_at, d.last_failure_code, d.last_failure_detail, d.delivered_at, d.quarantined_at, d.quarantine_reason FROM dovecote_events AS e LEFT JOIN dovecote_deliveries AS d ON d.tenant_id = e.tenant_id AND d.event_row_id = e.row_id WHERE e.row_id > ? AND e.row_id <= ? AND e.tenant_id = ? ORDER BY e.row_id ASC LIMIT ?";