appcore_sync/sync/
message_json.rs1use crate::sync::SyncMessage;
14use std::io::{self, Write};
15
16const EVENT_BUFFER_BYTES: usize = 16 * 1024;
17
18pub 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}