#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u8)]
#[non_exhaustive]
pub enum SyncMessageType {
Handshake = 0x01,
HandshakeAck = 0x02,
DeltaPush = 0x10,
DeltaAck = 0x11,
DeltaReject = 0x12,
CollectionSchema = 0x13,
CollectionPurged = 0x14,
ShapeSubscribe = 0x20,
ShapeSnapshot = 0x21,
ShapeDelta = 0x22,
ShapeUnsubscribe = 0x23,
VectorClockSync = 0x30,
TimeseriesPush = 0x40,
TimeseriesAck = 0x41,
ResyncRequest = 0x50,
Throttle = 0x52,
TokenRefresh = 0x60,
TokenRefreshAck = 0x61,
DefinitionSync = 0x70,
PresenceUpdate = 0x80,
PresenceBroadcast = 0x81,
PresenceLeave = 0x82,
ArrayDelta = 0x90,
ArrayDeltaBatch = 0x91,
ArraySnapshot = 0x92,
ArraySnapshotChunk = 0x93,
ArraySchema = 0x94,
ArrayAck = 0x95,
ArrayReject = 0x96,
ArrayCatchupRequest = 0x97,
ColumnarInsert = 0xA0,
ColumnarInsertAck = 0xA1,
VectorInsert = 0xA2,
VectorInsertAck = 0xA3,
VectorDelete = 0xA4,
VectorDeleteAck = 0xA5,
FtsIndex = 0xA6,
FtsIndexAck = 0xA7,
FtsDelete = 0xA8,
FtsDeleteAck = 0xA9,
SpatialInsert = 0xAA,
SpatialInsertAck = 0xAB,
SpatialDelete = 0xAC,
SpatialDeleteAck = 0xAD,
PingPong = 0xFF,
}
impl SyncMessageType {
pub fn from_u8(v: u8) -> Option<Self> {
match v {
0x01 => Some(Self::Handshake),
0x02 => Some(Self::HandshakeAck),
0x10 => Some(Self::DeltaPush),
0x11 => Some(Self::DeltaAck),
0x12 => Some(Self::DeltaReject),
0x13 => Some(Self::CollectionSchema),
0x14 => Some(Self::CollectionPurged),
0x20 => Some(Self::ShapeSubscribe),
0x21 => Some(Self::ShapeSnapshot),
0x22 => Some(Self::ShapeDelta),
0x23 => Some(Self::ShapeUnsubscribe),
0x30 => Some(Self::VectorClockSync),
0x40 => Some(Self::TimeseriesPush),
0x41 => Some(Self::TimeseriesAck),
0x50 => Some(Self::ResyncRequest),
0x52 => Some(Self::Throttle),
0x60 => Some(Self::TokenRefresh),
0x61 => Some(Self::TokenRefreshAck),
0x70 => Some(Self::DefinitionSync),
0x80 => Some(Self::PresenceUpdate),
0x81 => Some(Self::PresenceBroadcast),
0x82 => Some(Self::PresenceLeave),
0x90 => Some(Self::ArrayDelta),
0x91 => Some(Self::ArrayDeltaBatch),
0x92 => Some(Self::ArraySnapshot),
0x93 => Some(Self::ArraySnapshotChunk),
0x94 => Some(Self::ArraySchema),
0x95 => Some(Self::ArrayAck),
0x96 => Some(Self::ArrayReject),
0x97 => Some(Self::ArrayCatchupRequest),
0xA0 => Some(Self::ColumnarInsert),
0xA1 => Some(Self::ColumnarInsertAck),
0xA2 => Some(Self::VectorInsert),
0xA3 => Some(Self::VectorInsertAck),
0xA4 => Some(Self::VectorDelete),
0xA5 => Some(Self::VectorDeleteAck),
0xA6 => Some(Self::FtsIndex),
0xA7 => Some(Self::FtsIndexAck),
0xA8 => Some(Self::FtsDelete),
0xA9 => Some(Self::FtsDeleteAck),
0xAA => Some(Self::SpatialInsert),
0xAB => Some(Self::SpatialInsertAck),
0xAC => Some(Self::SpatialDelete),
0xAD => Some(Self::SpatialDeleteAck),
0xFF => Some(Self::PingPong),
_ => None,
}
}
}
#[non_exhaustive]
#[derive(Clone)]
pub struct SyncFrame {
pub msg_type: SyncMessageType,
pub body: Vec<u8>,
}
impl SyncFrame {
pub const FORMAT_VERSION: u8 = 1;
pub const HEADER_SIZE: usize = 10;
pub fn to_bytes(&self) -> Vec<u8> {
let len = self.body.len() as u32;
let crc = crc32c::crc32c(&self.body);
let mut buf = Vec::with_capacity(Self::HEADER_SIZE + self.body.len());
buf.push(Self::FORMAT_VERSION);
buf.push(self.msg_type as u8);
buf.extend_from_slice(&len.to_le_bytes());
buf.extend_from_slice(&crc.to_le_bytes());
buf.extend_from_slice(&self.body);
buf
}
pub fn from_bytes(data: &[u8]) -> Option<Self> {
if data.len() < Self::HEADER_SIZE {
return None;
}
let version = data[0];
if version != Self::FORMAT_VERSION {
return None;
}
let msg_type = SyncMessageType::from_u8(data[1])?;
let len = u32::from_le_bytes(data[2..6].try_into().ok()?) as usize;
let expected_crc = u32::from_le_bytes(data[6..10].try_into().ok()?);
if data.len() < Self::HEADER_SIZE + len {
return None;
}
let body = data[Self::HEADER_SIZE..Self::HEADER_SIZE + len].to_vec();
let actual_crc = crc32c::crc32c(&body);
if actual_crc != expected_crc {
tracing::warn!(
msg_type = data[1],
expected_crc,
actual_crc,
"sync frame CRC32C mismatch; dropping corrupt frame"
);
return None;
}
Some(Self { msg_type, body })
}
pub fn new_msgpack<T: zerompk::ToMessagePack>(
msg_type: SyncMessageType,
value: &T,
) -> Option<Self> {
let body = zerompk::to_msgpack_vec(value).ok()?;
Some(Self { msg_type, body })
}
pub fn try_encode<T: zerompk::ToMessagePack>(
msg_type: SyncMessageType,
value: &T,
) -> Option<Self> {
match zerompk::to_msgpack_vec(value) {
Ok(body) => Some(Self { msg_type, body }),
Err(e) => {
tracing::error!(
msg_type = msg_type as u8,
error = %e,
"failed to encode sync frame body; dropping response"
);
None
}
}
}
pub fn decode_body<T: zerompk::FromMessagePackOwned>(&self) -> Option<T> {
zerompk::from_msgpack(&self.body).ok()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn make_frame(msg_type: SyncMessageType, body: Vec<u8>) -> SyncFrame {
SyncFrame { msg_type, body }
}
#[test]
fn roundtrip_preserves_msg_type_and_body() {
let body = b"hello sync world".to_vec();
let frame = make_frame(SyncMessageType::PingPong, body.clone());
let bytes = frame.to_bytes();
let decoded = SyncFrame::from_bytes(&bytes).unwrap();
assert_eq!(decoded.msg_type, SyncMessageType::PingPong);
assert_eq!(decoded.body, body);
}
#[test]
fn flipped_body_byte_returns_none() {
let body = b"integrity check".to_vec();
let frame = make_frame(SyncMessageType::DeltaPush, body);
let mut bytes = frame.to_bytes();
bytes[SyncFrame::HEADER_SIZE] ^= 0xFF;
assert!(SyncFrame::from_bytes(&bytes).is_none());
}
#[test]
fn truncated_buffer_returns_none() {
assert!(SyncFrame::from_bytes(&[]).is_none());
assert!(SyncFrame::from_bytes(&[1u8; SyncFrame::HEADER_SIZE - 1]).is_none());
let frame = make_frame(SyncMessageType::PingPong, b"abcdef".to_vec());
let bytes = frame.to_bytes();
let truncated = &bytes[..bytes.len() - 1];
assert!(SyncFrame::from_bytes(truncated).is_none());
}
#[test]
fn wrong_version_returns_none() {
let frame = make_frame(SyncMessageType::PingPong, b"version test".to_vec());
let mut bytes = frame.to_bytes();
bytes[0] = SyncFrame::FORMAT_VERSION.wrapping_add(1);
assert!(SyncFrame::from_bytes(&bytes).is_none());
}
#[test]
fn header_size_is_ten_and_total_length_is_correct() {
assert_eq!(SyncFrame::HEADER_SIZE, 10);
let body = b"nodedb".to_vec();
let frame = make_frame(SyncMessageType::PingPong, body.clone());
let bytes = frame.to_bytes();
assert_eq!(bytes.len(), SyncFrame::HEADER_SIZE + body.len());
}
#[test]
fn crc32c_field_matches_crate_output() {
let body = b"crc check".to_vec();
let frame = make_frame(SyncMessageType::Handshake, body.clone());
let bytes = frame.to_bytes();
let stored = u32::from_le_bytes(bytes[6..10].try_into().unwrap());
let expected = crc32c::crc32c(&body);
assert_eq!(stored, expected);
}
}