1use crate::outbox_blob::{
14 encoded_message_bytes, insert_message_blob, message_blob_matches, read_message_blob,
15};
16use crate::{SqliteSyncError, SqliteSyncStore};
17use appcore_sync::{
18 SyncError, SyncMessage, SyncOutbox, SyncOutboxReceipt, SyncOutboxStats, SyncResult,
19 MAX_OUTBOX_PAGE_BYTES, MAX_OUTBOX_PAGE_MESSAGES,
20};
21use rusqlite::{params, OptionalExtension, TransactionBehavior};
22
23const MAX_BATCH_ID_BYTES: usize = 1_024;
24
25#[derive(Debug, Clone)]
27pub struct SqliteSyncOutbox {
28 store: SqliteSyncStore,
29}
30
31impl SqliteSyncOutbox {
32 pub(crate) fn new(store: SqliteSyncStore) -> Self {
33 Self { store }
34 }
35
36 fn read_messages(&self, limit: usize) -> SyncResult<Vec<SyncMessage>> {
37 self.store
38 .with_connection(|connection| {
39 let limit = i64::try_from(limit)
40 .map_err(|_| SqliteSyncError::CapacityExceeded("outbox"))?;
41 let mut statement = connection
42 .prepare(
43 "SELECT position, length(encoded) FROM appcore_sync_outbox
44 ORDER BY position LIMIT ?1",
45 )
46 .map_err(SqliteSyncError::database)?;
47 let rows = statement
48 .query_map([limit], |row| {
49 Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?))
50 })
51 .map_err(SqliteSyncError::database)?;
52 let mut messages = Vec::new();
53 let mut bytes = 0usize;
54 for row in rows {
55 let (position, encoded_bytes) = row.map_err(SqliteSyncError::database)?;
56 let encoded_bytes = checked_usize(encoded_bytes, "outbox byte")?;
57 if encoded_bytes > self.store.config().max_outbox_record_bytes {
58 return Err(SqliteSyncError::CapacityExceeded("outbox record"));
59 }
60 bytes = bytes
61 .checked_add(encoded_bytes)
62 .ok_or(SqliteSyncError::CapacityExceeded("outbox byte"))?;
63 if bytes as u64 > self.store.config().max_database_bytes {
64 return Err(SqliteSyncError::CapacityExceeded("outbox byte"));
65 }
66 messages.push(read_message_blob(connection, position, encoded_bytes)?);
67 }
68 Ok(messages)
69 })
70 .map_err(SqliteSyncError::sync)
71 }
72
73 fn read_page(
74 &self,
75 limit: usize,
76 max_bytes: usize,
77 ready_at_ms: Option<u64>,
78 ) -> SyncResult<Vec<SyncMessage>> {
79 validate_page_limits(limit, max_bytes)?;
80 self.store
81 .with_connection(|connection| {
82 let limit = limit.min(self.store.config().max_outbox_entries);
83 let sql_limit = i64::try_from(limit)
84 .map_err(|_| SqliteSyncError::CapacityExceeded("outbox page"))?;
85 let mut statement = connection
86 .prepare(
87 "SELECT position, length(encoded), next_ready_at_ms
88 FROM appcore_sync_outbox ORDER BY position LIMIT ?1",
89 )
90 .map_err(SqliteSyncError::database)?;
91 let rows = statement
92 .query_map([sql_limit], |row| {
93 Ok((
94 row.get::<_, i64>(0)?,
95 row.get::<_, i64>(1)?,
96 row.get::<_, i64>(2)?,
97 ))
98 })
99 .map_err(SqliteSyncError::database)?;
100 let mut selected = 0usize;
101 let mut selected_bytes = 0usize;
102 let mut last_position = None;
103 for row in rows {
104 let (position, encoded_bytes, next_ready) =
105 row.map_err(SqliteSyncError::database)?;
106 let encoded_bytes = usize::try_from(encoded_bytes)
107 .map_err(|_| SqliteSyncError::CorruptRecord("outbox byte"))?;
108 if ready_at_ms.is_some_and(|now| {
109 u64::try_from(next_ready).map_or(true, |ready| ready > now)
110 }) || selected_bytes
111 .checked_add(encoded_bytes)
112 .is_none_or(|bytes| bytes > max_bytes)
113 {
114 break;
115 }
116 selected_bytes += encoded_bytes;
117 selected += 1;
118 last_position = Some(position);
119 }
120 drop(statement);
121 let Some(last_position) = last_position else {
122 return Ok(Vec::new());
123 };
124 read_selected_messages(
125 connection,
126 last_position,
127 selected,
128 selected_bytes,
129 self.store.config().max_outbox_record_bytes,
130 )
131 })
132 .map_err(SqliteSyncError::sync)
133 }
134}
135
136impl SyncOutbox for SqliteSyncOutbox {
137 fn try_enqueue(&self, message: SyncMessage, max_len: usize) -> SyncResult<bool> {
138 validate_batch_id(&message.batch_id)?;
139 let encoded_bytes =
140 encoded_message_bytes(&message, self.store.config().max_outbox_record_bytes)?;
141 let limit = max_len.min(self.store.config().max_outbox_entries);
142 self.store
143 .with_connection(|connection| {
144 let transaction = connection
145 .transaction_with_behavior(TransactionBehavior::Immediate)
146 .map_err(SqliteSyncError::database)?;
147 let existing: Option<(i64, i64)> = transaction
148 .query_row(
149 "SELECT position, length(encoded) FROM appcore_sync_outbox
150 WHERE batch_id = ?1",
151 [&message.batch_id],
152 |row| Ok((row.get(0)?, row.get(1)?)),
153 )
154 .optional()
155 .map_err(SqliteSyncError::database)?;
156 if let Some((position, existing_bytes)) = existing {
157 let existing_bytes = usize::try_from(existing_bytes)
158 .map_err(|_| SqliteSyncError::CorruptRecord("outbox byte"))?;
159 return if existing_bytes == encoded_bytes
160 && message_blob_matches(&transaction, position, &message, encoded_bytes)?
161 {
162 Ok(false)
163 } else {
164 Err(SqliteSyncError::CorruptRecord("outbox conflict"))
165 };
166 }
167 let count = count_entries(&transaction)?;
168 if count >= limit {
169 return Ok(false);
170 }
171 insert_message_blob(&transaction, &message, encoded_bytes)?;
172 transaction.commit().map_err(SqliteSyncError::database)?;
173 Ok(true)
174 })
175 .map_err(|error| match error {
176 SqliteSyncError::CorruptRecord("outbox conflict") => {
177 SyncError::InvalidSyncMessage("outbox batch conflict")
178 }
179 other => other.sync(),
180 })
181 }
182
183 fn front(&self) -> SyncResult<Option<SyncMessage>> {
184 Ok(self.read_messages(1)?.into_iter().next())
185 }
186
187 fn acknowledge_front(&self, batch_id: &str) -> SyncResult<()> {
188 let receipt = SyncOutboxReceipt::new(vec![batch_id.to_string()])?;
189 self.acknowledge_receipt(&receipt).map(|_| ())
190 }
191
192 fn messages(&self) -> SyncResult<Vec<SyncMessage>> {
193 let length = self.len()?;
194 if length > self.store.config().max_outbox_entries {
195 return Err(SqliteSyncError::CapacityExceeded("outbox").sync());
196 }
197 self.read_messages(length)
198 }
199
200 fn len(&self) -> SyncResult<usize> {
201 self.store
202 .with_connection(|connection| count_entries(connection))
203 .map_err(SqliteSyncError::sync)
204 }
205
206 fn peek(&self, limit: usize, max_bytes: usize) -> SyncResult<Vec<SyncMessage>> {
207 self.read_page(limit, max_bytes, None)
208 }
209
210 fn stats(&self) -> SyncResult<SyncOutboxStats> {
211 self.store
212 .with_connection(|connection| {
213 let values: (i64, i64, i64, i64, Option<i64>) = connection
214 .query_row(
215 "SELECT COUNT(*), COALESCE(SUM(length(encoded)), 0),
216 COALESCE(SUM(CASE WHEN attempts > 0 THEN 1 ELSE 0 END), 0),
217 COALESCE(SUM(attempts), 0),
218 (SELECT next_ready_at_ms FROM appcore_sync_outbox
219 ORDER BY position LIMIT 1)
220 FROM appcore_sync_outbox",
221 [],
222 |row| {
223 Ok((
224 row.get(0)?,
225 row.get(1)?,
226 row.get(2)?,
227 row.get(3)?,
228 row.get(4)?,
229 ))
230 },
231 )
232 .map_err(SqliteSyncError::database)?;
233 Ok(SyncOutboxStats {
234 pending_messages: checked_usize(values.0, "outbox count")?,
235 pending_bytes: Some(checked_usize(values.1, "outbox byte")?),
236 attempted_messages: Some(checked_usize(values.2, "outbox attempt")?),
237 total_attempts: Some(
238 u64::try_from(values.3)
239 .map_err(|_| SqliteSyncError::CorruptRecord("outbox attempt"))?,
240 ),
241 next_ready_at_ms: values
242 .4
243 .map(|value| {
244 u64::try_from(value)
245 .map_err(|_| SqliteSyncError::CorruptRecord("outbox readiness"))
246 })
247 .transpose()?,
248 })
249 })
250 .map_err(SqliteSyncError::sync)
251 }
252
253 fn mark_attempt(&self, batch_id: &str, next_ready_at_ms: u64) -> SyncResult<u32> {
254 validate_batch_id(batch_id)?;
255 let next_ready_at_ms = i64::try_from(next_ready_at_ms)
256 .map_err(|_| SyncError::InvalidSyncMessage("outbox readiness overflow"))?;
257 self.store
258 .with_connection(|connection| {
259 let transaction = connection
260 .transaction_with_behavior(TransactionBehavior::Immediate)
261 .map_err(SqliteSyncError::database)?;
262 let front: Option<(i64, String, i64)> = transaction
263 .query_row(
264 "SELECT position, batch_id, attempts FROM appcore_sync_outbox
265 ORDER BY position LIMIT 1",
266 [],
267 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
268 )
269 .optional()
270 .map_err(SqliteSyncError::database)?;
271 let Some((position, current_id, attempts)) = front else {
272 return Err(SqliteSyncError::CorruptRecord("outbox attempt"));
273 };
274 if current_id != batch_id {
275 return Err(SqliteSyncError::CorruptRecord("outbox attempt"));
276 }
277 let attempts = attempts
278 .checked_add(1)
279 .and_then(|value| u32::try_from(value).ok())
280 .ok_or(SqliteSyncError::CapacityExceeded("outbox attempt"))?;
281 transaction
282 .execute(
283 "UPDATE appcore_sync_outbox
284 SET attempts = ?1, next_ready_at_ms = ?2 WHERE position = ?3",
285 params![attempts, next_ready_at_ms, position],
286 )
287 .map_err(SqliteSyncError::database)?;
288 transaction.commit().map_err(SqliteSyncError::database)?;
289 Ok(attempts)
290 })
291 .map_err(|error| match error {
292 SqliteSyncError::CorruptRecord("outbox attempt") => {
293 SyncError::InvalidSyncMessage("outbox attempt mismatch")
294 }
295 other => other.sync(),
296 })
297 }
298
299 fn next_ready(
300 &self,
301 now_ms: u64,
302 limit: usize,
303 max_bytes: usize,
304 ) -> SyncResult<Vec<SyncMessage>> {
305 self.read_page(limit, max_bytes, Some(now_ms))
306 }
307
308 fn acknowledge_receipt(&self, receipt: &SyncOutboxReceipt) -> SyncResult<usize> {
309 self.store
310 .with_connection(|connection| {
311 let transaction = connection
312 .transaction_with_behavior(TransactionBehavior::Immediate)
313 .map_err(SqliteSyncError::database)?;
314 let limit = i64::try_from(receipt.batch_ids().len())
315 .map_err(|_| SqliteSyncError::CapacityExceeded("outbox receipt"))?;
316 let mut statement = transaction
317 .prepare(
318 "SELECT position, batch_id FROM appcore_sync_outbox
319 ORDER BY position LIMIT ?1",
320 )
321 .map_err(SqliteSyncError::database)?;
322 let rows = statement
323 .query_map([limit], |row| {
324 Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?))
325 })
326 .map_err(SqliteSyncError::database)?;
327 let selected = rows
328 .collect::<Result<Vec<_>, _>>()
329 .map_err(SqliteSyncError::database)?;
330 if selected.len() != receipt.batch_ids().len()
331 || selected
332 .iter()
333 .zip(receipt.batch_ids())
334 .any(|((_, current), expected)| current != expected)
335 {
336 return Err(SqliteSyncError::CorruptRecord("outbox acknowledgement"));
337 }
338 let last_position = selected
339 .last()
340 .map(|(position, _)| *position)
341 .ok_or(SqliteSyncError::CorruptRecord("outbox acknowledgement"))?;
342 drop(statement);
343 let removed = transaction
344 .execute(
345 "DELETE FROM appcore_sync_outbox WHERE position <= ?1",
346 [last_position],
347 )
348 .map_err(SqliteSyncError::database)?;
349 if removed != receipt.batch_ids().len() {
350 return Err(SqliteSyncError::CorruptRecord("outbox acknowledgement"));
351 }
352 transaction.commit().map_err(SqliteSyncError::database)?;
353 Ok(removed)
354 })
355 .map_err(|error| match error {
356 SqliteSyncError::CorruptRecord("outbox acknowledgement") => {
357 SyncError::InvalidSyncMessage("outbox acknowledgement mismatch")
358 }
359 other => other.sync(),
360 })
361 }
362}
363
364fn read_selected_messages(
365 connection: &rusqlite::Connection,
366 last_position: i64,
367 selected: usize,
368 selected_bytes: usize,
369 max_record_bytes: usize,
370) -> Result<Vec<SyncMessage>, SqliteSyncError> {
371 let limit =
372 i64::try_from(selected).map_err(|_| SqliteSyncError::CapacityExceeded("outbox page"))?;
373 let mut statement = connection
374 .prepare(
375 "SELECT position, length(encoded) FROM appcore_sync_outbox
376 WHERE position <= ?1 ORDER BY position LIMIT ?2",
377 )
378 .map_err(SqliteSyncError::database)?;
379 let rows = statement
380 .query_map(params![last_position, limit], |row| {
381 Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?))
382 })
383 .map_err(SqliteSyncError::database)?;
384 let mut messages = Vec::with_capacity(selected);
385 let mut actual_bytes = 0usize;
386 for row in rows {
387 let (position, encoded_bytes) = row.map_err(SqliteSyncError::database)?;
388 let encoded_bytes = checked_usize(encoded_bytes, "outbox byte")?;
389 if encoded_bytes > max_record_bytes {
390 return Err(SqliteSyncError::CapacityExceeded("outbox record"));
391 }
392 actual_bytes = actual_bytes
393 .checked_add(encoded_bytes)
394 .ok_or(SqliteSyncError::CapacityExceeded("outbox byte"))?;
395 messages.push(read_message_blob(connection, position, encoded_bytes)?);
396 }
397 if messages.len() != selected || actual_bytes != selected_bytes {
398 return Err(SqliteSyncError::CorruptRecord("outbox page"));
399 }
400 Ok(messages)
401}
402
403fn validate_page_limits(limit: usize, max_bytes: usize) -> SyncResult<()> {
404 if limit > MAX_OUTBOX_PAGE_MESSAGES || max_bytes > MAX_OUTBOX_PAGE_BYTES {
405 return Err(SyncError::InvalidSyncMessage("invalid outbox page limits"));
406 }
407 Ok(())
408}
409
410fn checked_usize(value: i64, reason: &'static str) -> Result<usize, SqliteSyncError> {
411 usize::try_from(value).map_err(|_| SqliteSyncError::CorruptRecord(reason))
412}
413
414fn count_entries(connection: &rusqlite::Connection) -> Result<usize, SqliteSyncError> {
415 let count: i64 = connection
416 .query_row("SELECT COUNT(*) FROM appcore_sync_outbox", [], |row| {
417 row.get(0)
418 })
419 .map_err(SqliteSyncError::database)?;
420 usize::try_from(count).map_err(|_| SqliteSyncError::CorruptRecord("outbox count"))
421}
422
423pub(crate) fn validate_batch_id(batch_id: &str) -> SyncResult<()> {
424 if batch_id.is_empty()
425 || batch_id.len() > MAX_BATCH_ID_BYTES
426 || batch_id.chars().any(char::is_control)
427 {
428 return Err(SyncError::InvalidSyncMessage("invalid outbox batch id"));
429 }
430 Ok(())
431}