appcore_sync/sync/
outbox_size.rs1use crate::sync::{SyncError, SyncMessage, SyncResult};
14use std::io::{self, Write};
15
16const LENGTH_OVERFLOW: SyncError =
17 SyncError::InvalidSyncMessage("outbox serialization length overflow");
18
19pub fn encoded_sync_message_bytes(message: &SyncMessage) -> SyncResult<usize> {
24 let mut bytes = 0usize;
25 add(&mut bytes, b"{\"batch_id\":".len())?;
26 add(&mut bytes, encoded_string_bytes(&message.batch_id)?)?;
27 add(&mut bytes, b",\"source_node_id\":".len())?;
28 add(
29 &mut bytes,
30 encoded_string_bytes(message.source_node_id.as_str())?,
31 )?;
32 add(&mut bytes, b",\"sequence_start\":".len())?;
33 add(&mut bytes, decimal_digits_u64(message.sequence_start))?;
34 add(&mut bytes, b",\"sequence_end\":".len())?;
35 add(&mut bytes, decimal_digits_u64(message.sequence_end))?;
36 add(&mut bytes, b",\"event_count\":".len())?;
37 add(&mut bytes, decimal_digits_usize(message.event_count))?;
38 add(&mut bytes, b",\"events_hash\":".len())?;
39 add(&mut bytes, encoded_string_bytes(&message.events_hash)?)?;
40 add(&mut bytes, b",\"created_at_ms\":".len())?;
41 add(&mut bytes, decimal_digits_u64(message.created_at_ms))?;
42 add(&mut bytes, b",\"previous_batch_hash\":".len())?;
43 add_optional_string(&mut bytes, message.previous_batch_hash.as_deref())?;
44 add(&mut bytes, b",\"events\":[".len())?;
45 add_events(&mut bytes, &message.events)?;
46 add(&mut bytes, b"]}".len())?;
47 Ok(bytes)
48}
49
50fn add_events(total: &mut usize, events: &[Vec<u8>]) -> SyncResult<()> {
51 add(total, events.len().saturating_sub(1))?;
52 for event in events {
53 add(total, 2)?;
54 add(total, event.len().saturating_sub(1))?;
55 let digits = event.iter().try_fold(0usize, |bytes, value| {
56 bytes
57 .checked_add(decimal_digits_byte(*value))
58 .ok_or(LENGTH_OVERFLOW)
59 })?;
60 add(total, digits)?;
61 }
62 Ok(())
63}
64
65fn add_optional_string(total: &mut usize, value: Option<&str>) -> SyncResult<()> {
66 match value {
67 Some(value) => add(total, encoded_string_bytes(value)?),
68 None => add(total, b"null".len()),
69 }
70}
71
72fn encoded_string_bytes(value: &str) -> SyncResult<usize> {
73 let mut counter = JsonLengthCounter::default();
74 match serde_json::to_writer(&mut counter, value) {
75 Ok(()) => Ok(counter.bytes),
76 Err(_) if counter.overflowed => Err(LENGTH_OVERFLOW),
77 Err(_) => Err(SyncError::InvalidSyncMessage("outbox serialization failed")),
78 }
79}
80
81fn add(total: &mut usize, bytes: usize) -> SyncResult<()> {
82 *total = total.checked_add(bytes).ok_or(LENGTH_OVERFLOW)?;
83 Ok(())
84}
85
86const fn decimal_digits_byte(value: u8) -> usize {
87 if value < 10 {
88 1
89 } else if value < 100 {
90 2
91 } else {
92 3
93 }
94}
95
96const fn decimal_digits_u64(mut value: u64) -> usize {
97 let mut digits = 1;
98 while value >= 10 {
99 value /= 10;
100 digits += 1;
101 }
102 digits
103}
104
105const fn decimal_digits_usize(mut value: usize) -> usize {
106 let mut digits = 1;
107 while value >= 10 {
108 value /= 10;
109 digits += 1;
110 }
111 digits
112}
113
114#[derive(Default)]
115struct JsonLengthCounter {
116 bytes: usize,
117 overflowed: bool,
118}
119
120impl Write for JsonLengthCounter {
121 fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
122 let Some(total) = self.bytes.checked_add(bytes.len()) else {
123 self.overflowed = true;
124 return Err(io::Error::other("JSON length overflow"));
125 };
126 self.bytes = total;
127 Ok(bytes.len())
128 }
129
130 fn flush(&mut self) -> io::Result<()> {
131 Ok(())
132 }
133}