// MessagePack framing protocol // 4-byte big-endian length prefix + payload // Implements D-020 IPC transport layer use std::io::{self, Read, Write}; /// Maximum message size: 16 MB. /// Sized for ObserverSnapshot with ~1000 entities (each ~40 bytes serialized), /// plus generous headroom for future field additions. A full-map dump of 10k /// entities would be ~400 KB, well within this limit. const MAX_MESSAGE_SIZE: u32 = 16 * 1024 * 1024; /// Write a length-prefixed message to a writer. /// Format: [4-byte BE length][payload] pub fn write_framed(writer: &mut impl Write, payload: &[u8]) -> io::Result<()> { let len = payload.len() as u32; if len > MAX_MESSAGE_SIZE { return Err(io::Error::new( io::ErrorKind::InvalidInput, format!( "message too large: {} bytes (max {})", len, MAX_MESSAGE_SIZE ), )); } writer.write_all(&len.to_be_bytes())?; writer.write_all(payload)?; writer.flush()?; Ok(()) } /// Read a length-prefixed message from a reader. /// Returns Ok(None) on clean EOF (connection closed). /// Returns error on incomplete/corrupted reads. pub fn read_framed(reader: &mut impl Read) -> io::Result>> { // Read 4-byte length prefix let mut len_bytes = [0u8; 4]; match reader.read_exact(&mut len_bytes) { Ok(()) => {} Err(e) if e.kind() == io::ErrorKind::UnexpectedEof => return Ok(None), Err(e) => return Err(e), } let len = u32::from_be_bytes(len_bytes); if len > MAX_MESSAGE_SIZE { return Err(io::Error::new( io::ErrorKind::InvalidData, format!( "message too large: {} bytes (max {})", len, MAX_MESSAGE_SIZE ), )); } // Read payload let mut payload = vec![0u8; len as usize]; reader.read_exact(&mut payload)?; Ok(Some(payload)) } #[cfg(test)] mod tests { use super::*; use std::io::Cursor; #[test] fn write_then_read_roundtrip() { let payload = b"hello, world!"; let mut buffer = Vec::new(); write_framed(&mut buffer, payload).expect("write failed"); let mut cursor = Cursor::new(buffer); let result = read_framed(&mut cursor).expect("read failed"); assert_eq!(result.unwrap(), payload); } #[test] fn empty_payload_roundtrip() { let payload = b""; let mut buffer = Vec::new(); write_framed(&mut buffer, payload).expect("write failed"); let mut cursor = Cursor::new(buffer); let result = read_framed(&mut cursor).expect("read failed"); assert_eq!(result.unwrap(), payload); } #[test] fn rejects_oversized_message() { let mut buffer = Vec::new(); let oversized_payload = vec![0u8; (MAX_MESSAGE_SIZE + 1) as usize]; let result = write_framed(&mut buffer, &oversized_payload); assert!(result.is_err()); assert!(result.unwrap_err().to_string().contains("too large")); } #[test] fn eof_returns_none() { let buffer = Vec::new(); let mut cursor = Cursor::new(buffer); let result = read_framed(&mut cursor).expect("read failed"); assert_eq!(result, None); } }