1#[derive(Debug, Clone, Copy, PartialEq, Eq)]
21#[repr(u8)]
22#[non_exhaustive]
23pub enum SyncMessageType {
24 Handshake = 0x01,
25 HandshakeAck = 0x02,
26 DeltaPush = 0x10,
27 DeltaAck = 0x11,
28 DeltaReject = 0x12,
29 CollectionSchema = 0x13,
33 CollectionPurged = 0x14,
39 ShapeSubscribe = 0x20,
40 ShapeSnapshot = 0x21,
41 ShapeDelta = 0x22,
42 ShapeUnsubscribe = 0x23,
43 VectorClockSync = 0x30,
44 TimeseriesPush = 0x40,
46 TimeseriesAck = 0x41,
48 ResyncRequest = 0x50,
51 Throttle = 0x52,
54 TokenRefresh = 0x60,
56 TokenRefreshAck = 0x61,
58 DefinitionSync = 0x70,
61 PresenceUpdate = 0x80,
63 PresenceBroadcast = 0x81,
65 PresenceLeave = 0x82,
67 ArrayDelta = 0x90,
69 ArrayDeltaBatch = 0x91,
71 ArraySnapshot = 0x92,
73 ArraySnapshotChunk = 0x93,
75 ArraySchema = 0x94,
77 ArrayAck = 0x95,
79 ArrayReject = 0x96,
81 ArrayCatchupRequest = 0x97,
83 ColumnarInsert = 0xA0,
85 ColumnarInsertAck = 0xA1,
87 VectorInsert = 0xA2,
89 VectorInsertAck = 0xA3,
91 VectorDelete = 0xA4,
93 VectorDeleteAck = 0xA5,
95 FtsIndex = 0xA6,
97 FtsIndexAck = 0xA7,
99 FtsDelete = 0xA8,
101 FtsDeleteAck = 0xA9,
103 SpatialInsert = 0xAA,
105 SpatialInsertAck = 0xAB,
107 SpatialDelete = 0xAC,
109 SpatialDeleteAck = 0xAD,
111 PingPong = 0xFF,
112}
113
114impl SyncMessageType {
115 pub fn from_u8(v: u8) -> Option<Self> {
116 match v {
117 0x01 => Some(Self::Handshake),
118 0x02 => Some(Self::HandshakeAck),
119 0x10 => Some(Self::DeltaPush),
120 0x11 => Some(Self::DeltaAck),
121 0x12 => Some(Self::DeltaReject),
122 0x13 => Some(Self::CollectionSchema),
123 0x14 => Some(Self::CollectionPurged),
124 0x20 => Some(Self::ShapeSubscribe),
125 0x21 => Some(Self::ShapeSnapshot),
126 0x22 => Some(Self::ShapeDelta),
127 0x23 => Some(Self::ShapeUnsubscribe),
128 0x30 => Some(Self::VectorClockSync),
129 0x40 => Some(Self::TimeseriesPush),
130 0x41 => Some(Self::TimeseriesAck),
131 0x50 => Some(Self::ResyncRequest),
132 0x52 => Some(Self::Throttle),
133 0x60 => Some(Self::TokenRefresh),
134 0x61 => Some(Self::TokenRefreshAck),
135 0x70 => Some(Self::DefinitionSync),
136 0x80 => Some(Self::PresenceUpdate),
137 0x81 => Some(Self::PresenceBroadcast),
138 0x82 => Some(Self::PresenceLeave),
139 0x90 => Some(Self::ArrayDelta),
140 0x91 => Some(Self::ArrayDeltaBatch),
141 0x92 => Some(Self::ArraySnapshot),
142 0x93 => Some(Self::ArraySnapshotChunk),
143 0x94 => Some(Self::ArraySchema),
144 0x95 => Some(Self::ArrayAck),
145 0x96 => Some(Self::ArrayReject),
146 0x97 => Some(Self::ArrayCatchupRequest),
147 0xA0 => Some(Self::ColumnarInsert),
148 0xA1 => Some(Self::ColumnarInsertAck),
149 0xA2 => Some(Self::VectorInsert),
150 0xA3 => Some(Self::VectorInsertAck),
151 0xA4 => Some(Self::VectorDelete),
152 0xA5 => Some(Self::VectorDeleteAck),
153 0xA6 => Some(Self::FtsIndex),
154 0xA7 => Some(Self::FtsIndexAck),
155 0xA8 => Some(Self::FtsDelete),
156 0xA9 => Some(Self::FtsDeleteAck),
157 0xAA => Some(Self::SpatialInsert),
158 0xAB => Some(Self::SpatialInsertAck),
159 0xAC => Some(Self::SpatialDelete),
160 0xAD => Some(Self::SpatialDeleteAck),
161 0xFF => Some(Self::PingPong),
162 _ => None,
163 }
164 }
165}
166
167#[non_exhaustive]
181#[derive(Clone)]
182pub struct SyncFrame {
183 pub msg_type: SyncMessageType,
184 pub body: Vec<u8>,
185}
186
187impl SyncFrame {
188 pub const FORMAT_VERSION: u8 = 1;
190
191 pub const HEADER_SIZE: usize = 10;
195
196 pub fn to_bytes(&self) -> Vec<u8> {
200 let len = self.body.len() as u32;
201 let crc = crc32c::crc32c(&self.body);
202 let mut buf = Vec::with_capacity(Self::HEADER_SIZE + self.body.len());
203 buf.push(Self::FORMAT_VERSION);
204 buf.push(self.msg_type as u8);
205 buf.extend_from_slice(&len.to_le_bytes());
206 buf.extend_from_slice(&crc.to_le_bytes());
207 buf.extend_from_slice(&self.body);
208 buf
209 }
210
211 pub fn from_bytes(data: &[u8]) -> Option<Self> {
220 if data.len() < Self::HEADER_SIZE {
221 return None;
222 }
223 let version = data[0];
224 if version != Self::FORMAT_VERSION {
225 return None;
226 }
227 let msg_type = SyncMessageType::from_u8(data[1])?;
228 let len = u32::from_le_bytes(data[2..6].try_into().ok()?) as usize;
229 let expected_crc = u32::from_le_bytes(data[6..10].try_into().ok()?);
230 if data.len() < Self::HEADER_SIZE + len {
231 return None;
232 }
233 let body = data[Self::HEADER_SIZE..Self::HEADER_SIZE + len].to_vec();
234 let actual_crc = crc32c::crc32c(&body);
235 if actual_crc != expected_crc {
236 tracing::warn!(
237 msg_type = data[1],
238 expected_crc,
239 actual_crc,
240 "sync frame CRC32C mismatch; dropping corrupt frame"
241 );
242 return None;
243 }
244 Some(Self { msg_type, body })
245 }
246
247 pub fn new_msgpack<T: zerompk::ToMessagePack>(
249 msg_type: SyncMessageType,
250 value: &T,
251 ) -> Option<Self> {
252 let body = zerompk::to_msgpack_vec(value).ok()?;
253 Some(Self { msg_type, body })
254 }
255
256 pub fn try_encode<T: zerompk::ToMessagePack>(
262 msg_type: SyncMessageType,
263 value: &T,
264 ) -> Option<Self> {
265 match zerompk::to_msgpack_vec(value) {
266 Ok(body) => Some(Self { msg_type, body }),
267 Err(e) => {
268 tracing::error!(
269 msg_type = msg_type as u8,
270 error = %e,
271 "failed to encode sync frame body; dropping response"
272 );
273 None
274 }
275 }
276 }
277
278 pub fn decode_body<T: zerompk::FromMessagePackOwned>(&self) -> Option<T> {
280 zerompk::from_msgpack(&self.body).ok()
281 }
282}
283
284#[cfg(test)]
285mod tests {
286 use super::*;
287
288 fn make_frame(msg_type: SyncMessageType, body: Vec<u8>) -> SyncFrame {
289 SyncFrame { msg_type, body }
290 }
291
292 #[test]
293 fn roundtrip_preserves_msg_type_and_body() {
294 let body = b"hello sync world".to_vec();
295 let frame = make_frame(SyncMessageType::PingPong, body.clone());
296 let bytes = frame.to_bytes();
297 let decoded = SyncFrame::from_bytes(&bytes).unwrap();
298 assert_eq!(decoded.msg_type, SyncMessageType::PingPong);
299 assert_eq!(decoded.body, body);
300 }
301
302 #[test]
303 fn flipped_body_byte_returns_none() {
304 let body = b"integrity check".to_vec();
305 let frame = make_frame(SyncMessageType::DeltaPush, body);
306 let mut bytes = frame.to_bytes();
307 bytes[SyncFrame::HEADER_SIZE] ^= 0xFF;
309 assert!(SyncFrame::from_bytes(&bytes).is_none());
310 }
311
312 #[test]
313 fn truncated_buffer_returns_none() {
314 assert!(SyncFrame::from_bytes(&[]).is_none());
316 assert!(SyncFrame::from_bytes(&[1u8; SyncFrame::HEADER_SIZE - 1]).is_none());
317
318 let frame = make_frame(SyncMessageType::PingPong, b"abcdef".to_vec());
320 let bytes = frame.to_bytes();
321 let truncated = &bytes[..bytes.len() - 1];
322 assert!(SyncFrame::from_bytes(truncated).is_none());
323 }
324
325 #[test]
326 fn wrong_version_returns_none() {
327 let frame = make_frame(SyncMessageType::PingPong, b"version test".to_vec());
328 let mut bytes = frame.to_bytes();
329 bytes[0] = SyncFrame::FORMAT_VERSION.wrapping_add(1);
331 assert!(SyncFrame::from_bytes(&bytes).is_none());
332 }
333
334 #[test]
335 fn header_size_is_ten_and_total_length_is_correct() {
336 assert_eq!(SyncFrame::HEADER_SIZE, 10);
337 let body = b"nodedb".to_vec();
338 let frame = make_frame(SyncMessageType::PingPong, body.clone());
339 let bytes = frame.to_bytes();
340 assert_eq!(bytes.len(), SyncFrame::HEADER_SIZE + body.len());
341 }
342
343 #[test]
344 fn crc32c_field_matches_crate_output() {
345 let body = b"crc check".to_vec();
346 let frame = make_frame(SyncMessageType::Handshake, body.clone());
347 let bytes = frame.to_bytes();
348 let stored = u32::from_le_bytes(bytes[6..10].try_into().unwrap());
350 let expected = crc32c::crc32c(&body);
351 assert_eq!(stored, expected);
352 }
353}