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