1use 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
77pub 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 #[must_use]
101 pub const fn cursor(&self) -> Option<RowId> {
102 self.cursor
103 }
104 #[must_use]
106 pub const fn upper_bound(&self) -> Option<RowId> {
107 self.upper_bound
108 }
109 #[must_use]
111 pub const fn is_exhausted(&self) -> bool {
112 self.exhausted
113 }
114
115 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 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 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 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 ?";