Skip to main content

appcore_sync/sync/
message_json.rs

1// =============================================================================
2//        #######
3//     ###       ###     F: message_json.rs
4//    ##   ## ##   ##    P: AppCore-Runtime
5//         ## ##
6//                       C: 2026/09/03 00:00:00 by dnettoRaw
7//    ##   ## ##   ##    U: 2026/09/03 00:00:00 by dnettoRaw
8//      ###########      S: 1.0.1-rc.8
9// =============================================================================
10
11//! Writes canonical synchronization message JSON to bounded sinks.
12
13use crate::sync::SyncMessage;
14use std::io::{self, Write};
15
16const EVENT_BUFFER_BYTES: usize = 16 * 1024;
17
18/// Writes the compact JSON representation used by `SyncMessage` serialization.
19///
20/// Unlike `serde_json::to_vec`, this function does not require a second
21/// payload-sized allocation. It preserves the derived Serde field order and
22/// escaping so persistent providers can stream directly to bounded storage.
23pub fn write_sync_message_json(writer: &mut impl Write, message: &SyncMessage) -> io::Result<()> {
24    writer.write_all(b"{\"batch_id\":")?;
25    write_string(writer, &message.batch_id)?;
26    writer.write_all(b",\"source_node_id\":")?;
27    write_string(writer, message.source_node_id.as_str())?;
28    writer.write_all(b",\"sequence_start\":")?;
29    write_u64(writer, message.sequence_start)?;
30    writer.write_all(b",\"sequence_end\":")?;
31    write_u64(writer, message.sequence_end)?;
32    writer.write_all(b",\"event_count\":")?;
33    write_usize(writer, message.event_count)?;
34    writer.write_all(b",\"events_hash\":")?;
35    write_string(writer, &message.events_hash)?;
36    writer.write_all(b",\"created_at_ms\":")?;
37    write_u64(writer, message.created_at_ms)?;
38    writer.write_all(b",\"previous_batch_hash\":")?;
39    match message.previous_batch_hash.as_deref() {
40        Some(hash) => write_string(writer, hash)?,
41        None => writer.write_all(b"null")?,
42    }
43    writer.write_all(b",\"events\":[")?;
44    write_events(writer, &message.events)?;
45    writer.write_all(b"]}")
46}
47
48fn write_events(writer: &mut impl Write, events: &[Vec<u8>]) -> io::Result<()> {
49    let mut output = EventJsonWriter::new(writer);
50    for (event_index, event) in events.iter().enumerate() {
51        if event_index != 0 {
52            output.push(b',')?;
53        }
54        output.push(b'[')?;
55        for (byte_index, value) in event.iter().enumerate() {
56            if byte_index != 0 {
57                output.push(b',')?;
58            }
59            output.push_byte(*value)?;
60        }
61        output.push(b']')?;
62    }
63    output.finish()
64}
65
66fn write_string(writer: &mut impl Write, value: &str) -> io::Result<()> {
67    serde_json::to_writer(writer, value).map_err(io::Error::other)
68}
69
70fn write_u64(writer: &mut impl Write, value: u64) -> io::Result<()> {
71    write_decimal(writer, u128::from(value))
72}
73
74fn write_usize(writer: &mut impl Write, value: usize) -> io::Result<()> {
75    write_decimal(writer, value as u128)
76}
77
78fn write_decimal(writer: &mut impl Write, mut value: u128) -> io::Result<()> {
79    let mut buffer = [0u8; 39];
80    let mut cursor = buffer.len();
81    loop {
82        cursor -= 1;
83        buffer[cursor] = b'0' + (value % 10) as u8;
84        value /= 10;
85        if value == 0 {
86            return writer.write_all(&buffer[cursor..]);
87        }
88    }
89}
90
91struct EventJsonWriter<'a, W> {
92    writer: &'a mut W,
93    buffer: [u8; EVENT_BUFFER_BYTES],
94    used: usize,
95}
96
97impl<'a, W: Write> EventJsonWriter<'a, W> {
98    fn new(writer: &'a mut W) -> Self {
99        Self {
100            writer,
101            buffer: [0; EVENT_BUFFER_BYTES],
102            used: 0,
103        }
104    }
105
106    fn push(&mut self, byte: u8) -> io::Result<()> {
107        if self.used == self.buffer.len() {
108            self.flush()?;
109        }
110        self.buffer[self.used] = byte;
111        self.used += 1;
112        Ok(())
113    }
114
115    fn push_byte(&mut self, value: u8) -> io::Result<()> {
116        let hundreds = value / 100;
117        let tens = (value % 100) / 10;
118        let ones = value % 10;
119        if hundreds != 0 {
120            self.push(hundreds + b'0')?;
121        }
122        if hundreds != 0 || tens != 0 {
123            self.push(tens + b'0')?;
124        }
125        self.push(ones + b'0')
126    }
127
128    fn flush(&mut self) -> io::Result<()> {
129        self.writer.write_all(&self.buffer[..self.used])?;
130        self.used = 0;
131        Ok(())
132    }
133
134    fn finish(mut self) -> io::Result<()> {
135        self.flush()
136    }
137}