Skip to main content

appcore_sync_sqlite/
outbox.rs

1// =============================================================================
2//        #######
3//     ###       ###     F: outbox.rs
4//    ##   ## ##   ##    P: AppCore-Runtime
5//         ## ##
6//                       C: 2026/08/26 00:00:00 by dnettoRaw
7//    ##   ## ##   ##    U: 2026/08/26 00:00:00 by dnettoRaw
8//      ###########      S: 2.0.0
9// =============================================================================
10
11//! Defines bounded outbox contracts and behavior for this crate.
12
13use 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/// SQLite-backed bounded synchronization outbox.
26#[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}