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
11use 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/// SQLite-backed bounded synchronization outbox.
21#[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}