//! Length-prefixed, checksummed CBOR session framing (`SMH-SPEC-SPEC0001`, //! Sessions). //! //! Wire format, little endian: //! //! ```text //! [len u32][kind u16][version u16][crc32 u32][CBOR payload] //! ``` //! //! `len` counts every byte after itself, so a damaged frame is skipped by its //! declared length instead of scanning payload bytes for a plausible header. //! `kind` and `version` say what the payload is; `crc32` covers kind, version, //! and payload, which is what keeps corruption distinguishable from a frame //! kind that a future build introduced. use crate::session::{Entry, EntryFrame, EntryWire, SessionHeader}; use smith::error::{FrameFault, Result, SmithError}; use std::io::{Read, Write}; /// Default frame size ceiling when no configuration is supplied. pub const DEFAULT_MAX_FRAME_BYTES: usize = 1024 * 1024; /// Length of the frame length prefix. pub const HEADER_LEN: usize = 4; /// Length of the body prefix: kind, version, and checksum. pub const BODY_PREFIX_LEN: usize = 8; /// Frame kind of a session entry payload. pub const KIND_ENTRY: u16 = 1; /// Supported version of the [`KIND_ENTRY`] payload schema. /// Version 2 carries the content as an opaque canonical CBOR byte string. pub const VERSION_ENTRY: u16 = 2; /// Frame kind of the session header payload. pub const KIND_SESSION_HEADER: u16 = 2; /// Supported version of the [`KIND_SESSION_HEADER`] payload schema. pub const VERSION_SESSION_HEADER: u16 = 1; /// A decoded session frame. #[derive(Clone, PartialEq, Eq, Debug)] pub enum Frame { /// Durable session identity; the first frame of a session file. Header { /// Decoded header. header: SessionHeader, }, /// A known session entry. Known { /// Decoded entry frame; its content stays canonical CBOR bytes. entry: EntryFrame, }, /// A checksum-valid frame of a kind this build does not know. /// Kind, version, and payload are preserved so it round-trips unchanged. Unknown { /// Frame kind tag. kind: u16, /// Frame version tag. version: u16, /// Exact payload bytes as read. payload: Vec, }, } impl Frame { /// Serialize a payload value, rejecting sizes above `max_frame_bytes` /// before unbounded allocation occurs. /// /// # Errors /// /// Returns a [`FrameFault::TooLarge`] or [`FrameFault::Encode`] error. pub fn encode_payload( value: &T, max_frame_bytes: usize, ) -> Result> { serialize_limited(value, Vec::new(), max_frame_bytes) } /// Encode the full frame (length prefix, body prefix, payload) enforcing /// `max_frame_bytes` on the payload. /// /// # Errors /// /// Returns a [`FrameFault::TooLarge`] or [`FrameFault::Encode`] error. pub fn encode(&self, max_frame_bytes: usize) -> Result> { match self { Self::Header { header } => encode_frame( KIND_SESSION_HEADER, VERSION_SESSION_HEADER, header, 0, max_frame_bytes, ), Self::Known { entry } => encode_frame( KIND_ENTRY, VERSION_ENTRY, entry, entry.content.len(), max_frame_bytes, ), Self::Unknown { kind, version, payload, } => { if payload.len() > max_frame_bytes { return Err(too_large(max_frame_bytes, payload.len())); } let mut out = Vec::with_capacity(FRAME_PREFIX_LEN + payload.len()); out.resize(FRAME_PREFIX_LEN, 0); out.extend_from_slice(payload); seal(*kind, *version, out, max_frame_bytes) } } } } /// Encode a session record and its arena content bytes as a full entry /// frame, without copying the content out of the arena first. /// /// # Errors /// /// Returns a [`FrameFault::TooLarge`] or [`FrameFault::Encode`] error. pub fn encode_entry(entry: &Entry, content: &[u8], max_frame_bytes: usize) -> Result> { encode_frame( KIND_ENTRY, VERSION_ENTRY, &EntryWire::of(entry, content), content.len(), max_frame_bytes, ) } /// Bytes before a frame payload: length prefix, kind, version, checksum. const FRAME_PREFIX_LEN: usize = HEADER_LEN + BODY_PREFIX_LEN; /// Capacity hint for the payload bytes of a frame beyond its entry content; /// an underestimate costs one reallocation, never correctness. const PAYLOAD_OVERHEAD_HINT: usize = 192; /// Serialize `value` as a frame payload directly after a reserved prefix, so /// one buffer holds the whole frame. fn encode_frame( kind: u16, version: u16, value: &T, content_len: usize, max_frame_bytes: usize, ) -> Result> { let capacity = (content_len + PAYLOAD_OVERHEAD_HINT).min(max_frame_bytes); let mut out = Vec::with_capacity(FRAME_PREFIX_LEN + capacity); out.resize(FRAME_PREFIX_LEN, 0); let out = serialize_limited(value, out, max_frame_bytes)?; seal(kind, version, out, max_frame_bytes) } /// Serialize `value` after the bytes already in `buf`, rejecting payloads /// above `max_frame_bytes` before unbounded allocation occurs. fn serialize_limited( value: &T, buf: Vec, max_frame_bytes: usize, ) -> Result> { let mut limited = LimitedWriter::new(buf, max_frame_bytes); let result = ciborium::ser::into_writer(value, &mut limited); if limited.overflowed() { return Err(too_large(max_frame_bytes, max_frame_bytes + 1)); } match result { Ok(()) => Ok(limited.into_inner()), Err(e) => Err(fault(FrameFault::Encode { message: e.to_string(), })), } } /// Fill the reserved prefix of `out` (length, kind, version, checksum) for /// the payload that follows it. fn seal(kind: u16, version: u16, mut out: Vec, max_frame_bytes: usize) -> Result> { let payload_len = out.len().saturating_sub(FRAME_PREFIX_LEN); let Ok(declared) = u32::try_from(BODY_PREFIX_LEN + payload_len) else { return Err(too_large(max_frame_bytes, payload_len)); }; let checksum = crc32(kind, version, &out[FRAME_PREFIX_LEN..]); out[..HEADER_LEN].copy_from_slice(&declared.to_le_bytes()); out[HEADER_LEN..HEADER_LEN + 2].copy_from_slice(&kind.to_le_bytes()); out[HEADER_LEN + 2..HEADER_LEN + 4].copy_from_slice(&version.to_le_bytes()); out[HEADER_LEN + 4..FRAME_PREFIX_LEN].copy_from_slice(&checksum.to_le_bytes()); Ok(out) } /// Classify one frame body: verify the checksum first, then the kind and /// version, then the payload schema. fn decode_body(body: &[u8]) -> std::result::Result { if body.len() < BODY_PREFIX_LEN { return Err(FrameFault::Malformed { declared_length: body.len(), }); } let kind = u16::from_le_bytes([body[0], body[1]]); let version = u16::from_le_bytes([body[2], body[3]]); let checksum = u32::from_le_bytes([body[4], body[5], body[6], body[7]]); let payload = &body[BODY_PREFIX_LEN..]; if checksum != crc32(kind, version, payload) { return Err(FrameFault::Checksum { declared_length: body.len(), }); } match (kind, version) { (KIND_ENTRY, VERSION_ENTRY) => { decode_payload(payload).map(|entry: EntryFrame| Frame::Known { entry }) } (KIND_SESSION_HEADER, VERSION_SESSION_HEADER) => { decode_payload(payload).map(|header| Frame::Header { header }) } (KIND_ENTRY | KIND_SESSION_HEADER, _) => { Err(FrameFault::UnsupportedVersion { kind, version }) } _ => Ok(Frame::Unknown { kind, version, payload: payload.to_vec(), }), } } fn decode_payload( payload: &[u8], ) -> std::result::Result { ciborium::de::from_reader::(payload).map_err(|e| FrameFault::SchemaDecode { message: e.to_string(), }) } /// Frames recovered from a frame stream, plus explicit damage. /// /// Damage never discards the surrounding frames: a corrupt frame is reported /// and skipped by its declared length, and the frames after it stay readable. #[derive(Clone, PartialEq, Eq, Debug, Default)] pub struct Recovery { /// Frames recovered in stream order. pub frames: Vec, /// Damage reports, each locating its frame header by byte offset. pub damage: Vec, /// Offset just past the last complete frame. pub good_end: u64, /// Offset where an incomplete trailing frame starts, if any. pub torn_tail_at: Option, } /// Read every complete frame from `reader` with per-frame bounded allocation. /// /// A torn tail ends the walk and is reported through /// [`Recovery::torn_tail_at`]; complete frames before it are always returned. pub fn read_frames(reader: &mut R, max_frame_bytes: usize) -> Recovery { let mut recovery = Recovery::default(); let declared_limit = max_frame_bytes.saturating_add(BODY_PREFIX_LEN); loop { let offset = recovery.good_end; let mut header = [0u8; HEADER_LEN]; match read_exact_or_eof(reader, &mut header) { Ok(true) => {} Ok(false) => break, Err(_) => { recovery.torn_tail_at = Some(offset); break; } } let declared = u32::from_le_bytes(header) as usize; let total = HEADER_LEN as u64 + declared as u64; if declared > declared_limit { if !skip_exact(reader, declared) { recovery.torn_tail_at = Some(offset); break; } recovery.damage.push(located( FrameFault::TooLarge { max_bytes: max_frame_bytes, actual: declared.saturating_sub(BODY_PREFIX_LEN), }, offset, )); recovery.good_end = offset + total; continue; } let mut body = vec![0u8; declared]; if reader.read_exact(&mut body).is_err() { recovery.torn_tail_at = Some(offset); break; } match decode_body(&body) { Ok(frame) => recovery.frames.push(frame), Err(f) => recovery.damage.push(located(f, offset)), } recovery.good_end = offset + total; } recovery } /// Write a full frame to a writer with an enforced payload size limit. /// /// # Errors /// /// Returns a framing error when encoding or writing fails. pub fn write_frame(writer: &mut impl Write, frame: &Frame, max_frame_bytes: usize) -> Result<()> { let bytes = frame.encode(max_frame_bytes)?; writer.write_all(&bytes).map_err(|e| { fault(FrameFault::Encode { message: e.to_string(), }) })?; Ok(()) } const fn fault(f: FrameFault) -> SmithError { SmithError::Frame { fault: f, offset: None, } } const fn located(f: FrameFault, offset: u64) -> SmithError { SmithError::Frame { fault: f, offset: Some(offset), } } const fn too_large(max_bytes: usize, actual: usize) -> SmithError { fault(FrameFault::TooLarge { max_bytes, actual }) } /// Read exactly `buf.len()` bytes. `Ok(false)` means clean end of stream, /// `Err` means a partial read, which is a torn tail. fn read_exact_or_eof(reader: &mut R, buf: &mut [u8]) -> std::io::Result { let mut filled = 0; while filled < buf.len() { match reader.read(&mut buf[filled..]) { Ok(0) => { if filled == 0 { return Ok(false); } return Err(std::io::Error::new( std::io::ErrorKind::UnexpectedEof, "torn frame", )); } Ok(n) => filled += n, Err(e) if e.kind() == std::io::ErrorKind::Interrupted => {} Err(e) => return Err(e), } } Ok(true) } /// Discard `len` bytes through a bounded copy buffer. `false` means the /// stream ended early. fn skip_exact(reader: &mut R, len: usize) -> bool { let len = len as u64; std::io::copy(&mut reader.take(len), &mut std::io::sink()).is_ok_and(|skipped| skipped == len) } /// CRC-32 (IEEE 802.3, reflected) over kind, version, and payload. fn crc32(kind: u16, version: u16, payload: &[u8]) -> u32 { let mut crc = 0xFFFF_FFFF_u32; let mut step = |byte: u8| { let index = ((crc ^ u32::from(byte)) & 0xFF) as usize; crc = (crc >> 8) ^ CRC32_TABLE[index]; }; for byte in kind.to_le_bytes() { step(byte); } for byte in version.to_le_bytes() { step(byte); } for byte in payload { step(*byte); } !crc } /// Recursion, not loops: a mutated step in a `const` loop makes rustc spin /// forever under capped lints; a mutated recursion fails to compile. const CRC32_TABLE: [u32; 256] = crc32_table([0; 256], 0); #[expect( clippy::large_types_passed_by_value, reason = "const evaluation threads the table through the recursion by value" )] const fn crc32_table(table: [u32; 256], row: u32) -> [u32; 256] { if row == 16 { return table; } crc32_table(crc32_row(table, row * 16, 0), row + 1) } /// Sixteen entries per row keep const-eval recursion well under its frame cap. const fn crc32_row(mut table: [u32; 256], start: u32, offset: u32) -> [u32; 256] { if offset == 16 { return table; } let index = start + offset; table[index as usize] = crc32_entry(index, 0); crc32_row(table, start, offset + 1) } const fn crc32_entry(value: u32, bit: u32) -> u32 { if bit == 8 { return value; } let next = if value & 1 == 1 { 0xEDB8_8320 ^ (value >> 1) } else { value >> 1 }; crc32_entry(next, bit + 1) } /// A `Vec` writer that stops accepting bytes beyond a limit, /// bounding allocation during encoding. struct LimitedWriter { buf: Vec, remaining: usize, overflow: bool, } impl LimitedWriter { /// Accept up to `limit` bytes after the ones already in `buf`. const fn new(buf: Vec, limit: usize) -> Self { Self { buf, remaining: limit, overflow: false, } } const fn overflowed(&self) -> bool { self.overflow } fn into_inner(self) -> Vec { self.buf } } impl Write for LimitedWriter { fn write(&mut self, buf: &[u8]) -> std::io::Result { if self.overflow { return Err(std::io::Error::new( std::io::ErrorKind::WriteZero, "frame size limit exceeded", )); } let keep = buf.len().min(self.remaining); self.buf.extend_from_slice(&buf[..keep]); self.remaining -= keep; if keep < buf.len() { self.overflow = true; return Err(std::io::Error::new( std::io::ErrorKind::WriteZero, "frame size limit exceeded", )); } Ok(keep) } fn flush(&mut self) -> std::io::Result<()> { Ok(()) } } #[cfg(test)] #[expect( clippy::unwrap_used, clippy::panic, reason = "tests may panic on invariant violations" )] mod tests { use super::*; use crate::session::{EntryContent, EntryKind, Session}; use smith::id::ToolCallId; use smith::message::{ContentBlock, Message, Role}; fn message_frame(message: Message) -> EntryFrame { EntryFrame::new(None, &EntryContent::Message(message)).unwrap() } fn sample_entry(thinking: bool) -> EntryFrame { let mut msg = Message::with_text(Role::User, "hello"); if thinking { msg.add_block(ContentBlock::thinking("internal")); } message_frame(msg) } fn known_bytes(text: &str) -> Vec { Frame::Known { entry: message_frame(Message::with_text(Role::User, text)), } .encode(DEFAULT_MAX_FRAME_BYTES) .unwrap() } fn text_of(frame: &Frame) -> String { match frame { Frame::Known { entry } => match EntryContent::decode(&entry.content).unwrap() { EntryContent::Message(m) => match &m.blocks[0] { ContentBlock::Text(t) => t.clone(), other => panic!("unexpected block {other:?}"), }, other => panic!("unexpected content {other:?}"), }, other => panic!("unexpected frame {other:?}"), } } fn read_all(bytes: &[u8], max_frame_bytes: usize) -> Recovery { read_frames(&mut &bytes[..], max_frame_bytes) } #[test] fn corrupt_frame_is_reported_and_neighbours_survive() { let mut bytes = known_bytes("before"); let corrupt_at = bytes.len() as u64; let mut corrupt = known_bytes("corrupt"); // Flip one payload byte: the length prefix stays intact, the checksum fails. let last = corrupt.len() - 1; corrupt[last] ^= 0xFF; bytes.extend_from_slice(&corrupt); bytes.extend_from_slice(&known_bytes("after")); let recovery = read_all(&bytes, DEFAULT_MAX_FRAME_BYTES); let texts: Vec = recovery.frames.iter().map(text_of).collect(); assert_eq!(texts, vec!["before".to_string(), "after".to_string()]); assert_eq!(recovery.damage.len(), 1); assert_eq!(recovery.damage[0].code(), "FRAME_CHECKSUM"); assert_eq!( recovery.damage[0], SmithError::Frame { fault: FrameFault::Checksum { declared_length: corrupt.len() - HEADER_LEN }, offset: Some(corrupt_at), } ); assert_eq!(recovery.good_end, bytes.len() as u64); assert_eq!(recovery.torn_tail_at, None); } #[test] fn future_frame_round_trips_byte_for_byte() { let future = Frame::Unknown { kind: 0xBEEF, version: 7, payload: vec![0xDE, 0xAD, 0x01, 0x02], }; let bytes = future.encode(DEFAULT_MAX_FRAME_BYTES).unwrap(); let recovery = read_all(&bytes, DEFAULT_MAX_FRAME_BYTES); assert!(recovery.damage.is_empty()); assert_eq!(recovery.frames, vec![future]); assert_eq!( recovery.frames[0].encode(DEFAULT_MAX_FRAME_BYTES).unwrap(), bytes ); } #[test] fn known_frame_round_trips_through_cbor() { let entry = sample_entry(true); let bytes = Frame::Known { entry: entry.clone(), } .encode(DEFAULT_MAX_FRAME_BYTES) .unwrap(); let recovery = read_all(&bytes, DEFAULT_MAX_FRAME_BYTES); assert_eq!(recovery.frames, vec![Frame::Known { entry }]); assert_eq!(recovery.good_end, bytes.len() as u64); } #[test] fn entry_frame_carries_kind_and_content_as_one_byte_string() { let call_id = ToolCallId::new(); let content = EntryContent::ToolResult { ok: true, output: "done".to_string(), call_id, }; let entry = EntryFrame::new(None, &content).unwrap(); assert_eq!(entry.kind, EntryKind::ToolResult { call_id, ok: true }); assert_eq!(EntryContent::decode(&entry.content).unwrap(), content); // The content is embedded verbatim after a CBOR byte-string header. let payload = Frame::encode_payload(&entry, DEFAULT_MAX_FRAME_BYTES).unwrap(); let at = payload .windows(entry.content.len()) .position(|window| window == entry.content.as_slice()) .unwrap(); let len = entry.content.len(); let header: &[u8] = if len < 24 { &[0x40 | len as u8] } else { &[0x58, len as u8] }; assert_eq!(&payload[at - header.len()..at], header); // A session record frames straight from arena bytes, identically. let mut session = Session::new(); let record = *session.append(&content).unwrap(); let from_arena = encode_entry(&record, session.bytes(&record), DEFAULT_MAX_FRAME_BYTES).unwrap(); let owned = Frame::Known { entry: EntryFrame { id: record.id, parent: record.parent, timestamp_ms: record.timestamp_ms, kind: record.kind, content: session.bytes(&record).to_vec(), }, }; assert_eq!(from_arena, owned.encode(DEFAULT_MAX_FRAME_BYTES).unwrap()); } #[test] fn torn_tail_is_located_and_earlier_frames_kept() { let mut bytes = known_bytes("kept"); let tail_at = bytes.len() as u64; bytes.extend_from_slice(&known_bytes("lost")[..6]); let recovery = read_all(&bytes, DEFAULT_MAX_FRAME_BYTES); assert_eq!( recovery.frames.iter().map(text_of).collect::>(), vec!["kept".to_string()] ); assert_eq!(recovery.torn_tail_at, Some(tail_at)); assert_eq!(recovery.good_end, tail_at); // A partial length prefix is torn too, never a fabricated header. let mut short = known_bytes("kept"); let short_tail = short.len() as u64; short.extend_from_slice(&[0xFF, 0xFF]); let recovery = read_all(&short, DEFAULT_MAX_FRAME_BYTES); assert_eq!(recovery.frames.len(), 1); assert_eq!(recovery.torn_tail_at, Some(short_tail)); } #[test] fn oversized_complete_frame_is_skipped_by_length() { let mut bytes = (5000_u32).to_le_bytes().to_vec(); bytes.extend_from_slice(&vec![0xAB_u8; 5000]); let skipped = bytes.len() as u64; bytes.extend_from_slice(&known_bytes("tail")); let recovery = read_all(&bytes, 1024); assert_eq!( recovery.frames.iter().map(text_of).collect::>(), vec!["tail".to_string()] ); assert_eq!(recovery.damage.len(), 1); assert_eq!(recovery.damage[0].code(), "FRAME_TOO_LARGE"); assert_eq!( recovery.damage[0], SmithError::Frame { fault: FrameFault::TooLarge { max_bytes: 1024, actual: 5000 - BODY_PREFIX_LEN }, offset: Some(0), } ); assert_eq!(recovery.good_end, bytes.len() as u64); assert_eq!(skipped, 5004); } #[test] fn oversized_truncated_frame_is_torn_not_scanned() { let mut bytes = 200_u32.to_le_bytes().to_vec(); bytes.extend_from_slice(&[0u8; 10]); let recovery = read_all(&bytes, 50); assert!(recovery.frames.is_empty()); assert!(recovery.damage.is_empty()); assert_eq!(recovery.torn_tail_at, Some(0)); } #[test] fn declared_length_below_body_prefix_is_malformed() { let mut bytes = 3_u32.to_le_bytes().to_vec(); bytes.extend_from_slice(&[1, 2, 3]); bytes.extend_from_slice(&known_bytes("after")); let recovery = read_all(&bytes, DEFAULT_MAX_FRAME_BYTES); assert_eq!(recovery.damage.len(), 1); assert_eq!(recovery.damage[0].code(), "FRAME_MALFORMED"); assert_eq!(recovery.frames.len(), 1); } #[test] fn unsupported_entry_version_is_reported_not_preserved() { let payload = vec![0xA0]; // empty CBOR map let frame = Frame::Unknown { kind: KIND_ENTRY, version: VERSION_ENTRY + 1, payload, }; let bytes = frame.encode(DEFAULT_MAX_FRAME_BYTES).unwrap(); let recovery = read_all(&bytes, DEFAULT_MAX_FRAME_BYTES); assert!(recovery.frames.is_empty()); assert_eq!(recovery.damage[0].code(), "FRAME_VERSION"); } #[test] fn checksum_valid_entry_with_bad_schema_is_schema_damage() { let frame = Frame::Unknown { kind: KIND_ENTRY, version: VERSION_ENTRY, payload: vec![0xA0], // valid CBOR, not an Entry }; let bytes = frame.encode(DEFAULT_MAX_FRAME_BYTES).unwrap(); let recovery = read_all(&bytes, DEFAULT_MAX_FRAME_BYTES); assert!(recovery.frames.is_empty()); assert_eq!(recovery.damage[0].code(), "FRAME_SCHEMA"); } #[test] fn version_one_frames_without_body_prefix_are_rejected() { // Old layout: [len][CBOR payload] with no kind, version, or checksum. let entry = sample_entry(false); let payload = Frame::encode_payload(&entry, DEFAULT_MAX_FRAME_BYTES).unwrap(); let mut bytes = u32::try_from(payload.len()).unwrap().to_le_bytes().to_vec(); bytes.extend_from_slice(&payload); let recovery = read_all(&bytes, DEFAULT_MAX_FRAME_BYTES); assert!(recovery.frames.is_empty()); assert_eq!(recovery.damage.len(), 1); } #[test] fn sequential_frames_consume_the_whole_stream() { let mut bytes = known_bytes("one"); bytes.extend_from_slice(&known_bytes("two")); let recovery = read_all(&bytes, DEFAULT_MAX_FRAME_BYTES); assert_eq!( recovery.frames.iter().map(text_of).collect::>(), vec!["one".to_string(), "two".to_string()] ); assert_eq!(recovery.good_end, bytes.len() as u64); } #[test] fn encoding_rejects_frames_above_limit() { let big = message_frame({ let mut m = Message::with_text(Role::User, "x"); m.add_block(ContentBlock::text("y".repeat(4096))); m }); let err = Frame::Known { entry: big }.encode(64).unwrap_err(); assert_eq!(err, too_large(64, 65)); } #[test] fn default_ceiling_admits_one_mebibyte_payloads_and_no_more() { let frame = |len| Frame::Unknown { kind: 9, version: 1, payload: vec![0; len], }; let at_ceiling = frame(1 << 20).encode(DEFAULT_MAX_FRAME_BYTES).unwrap(); assert_eq!(at_ceiling.len(), HEADER_LEN + BODY_PREFIX_LEN + (1 << 20)); assert_eq!( frame((1 << 20) + 1).encode(DEFAULT_MAX_FRAME_BYTES), Err(too_large(1 << 20, (1 << 20) + 1)) ); } #[test] fn equivalent_inputs_encode_identically_despite_key_order() { let a = serde_json::json!({"path": "p", "content": {"b": 1, "a": 2}}); let b = serde_json::json!({"content": {"a": 2, "b": 1}, "path": "p"}); let call_id = ToolCallId::new(); let e1 = EntryFrame::new( None, &EntryContent::ToolCall { name: "write".into(), input: a, call_id, }, ) .unwrap(); let mut e2 = EntryFrame::new( None, &EntryContent::ToolCall { name: "write".into(), input: b, call_id, }, ) .unwrap(); // Force identical ids/timestamps so only key order could differ. e2.id = e1.id; e2.timestamp_ms = e1.timestamp_ms; let b1 = Frame::Known { entry: e1 }.encode(1 << 20).unwrap(); let b2 = Frame::Known { entry: e2 }.encode(1 << 20).unwrap(); assert_eq!(b1, b2); } #[test] fn header_frame_round_trips() { let header = SessionHeader { session_id: smith::id::SessionId::new(), branch_id: smith::id::BranchId::new(), }; let bytes = Frame::Header { header } .encode(DEFAULT_MAX_FRAME_BYTES) .unwrap(); let recovery = read_all(&bytes, DEFAULT_MAX_FRAME_BYTES); assert_eq!(recovery.frames, vec![Frame::Header { header }]); assert!(recovery.damage.is_empty()); } #[test] fn checksum_matches_known_reference_value() { // CRC-32 of "123456789" is a published check value. assert_eq!(crc32(0x3231, 0x3433, b"56789"), 0xCBF4_3926); } } /// Frame stream round-trip law (`SMH-SPEC-SPEC0001`, Sessions). #[cfg(test)] mod properties { use super::*; use crate::session::EntryKind; use crate::session::properties::role; use proptest::collection::vec; use proptest::prelude::*; use smith::id::{BranchId, EntryId, MessageId, SessionId, ToolCallId}; fn kind() -> impl Strategy { prop_oneof![ (any::(), role()).prop_map(|(id, role)| EntryKind::Message { id: MessageId::from_u128(id), role, }), any::().prop_map(|id| EntryKind::ToolCall { call_id: ToolCallId::from_u128(id), }), (any::(), any::()).prop_map(|(id, ok)| EntryKind::ToolResult { call_id: ToolCallId::from_u128(id), ok, }), any::().prop_map(|compaction| EntryKind::Meta { compaction }), ] } /// Any frame this build writes: header, entry with opaque content bytes, /// or a frame of a kind it does not know. fn frame() -> impl Strategy { prop_oneof![ (any::(), any::()).prop_map(|(session, branch)| Frame::Header { header: SessionHeader { session_id: SessionId::from_u128(session), branch_id: BranchId::from_u128(branch), }, }), ( any::(), proptest::option::of(any::()), any::(), kind(), vec(any::(), 0..64), ) .prop_map(|(id, parent, timestamp_ms, kind, content)| Frame::Known { entry: EntryFrame { id: EntryId::from_u128(id), parent: parent.map(EntryId::from_u128), timestamp_ms, kind, content, }, }), ( any::().prop_filter("known kinds decode as known", |kind| { !matches!(*kind, KIND_ENTRY | KIND_SESSION_HEADER) }), any::(), vec(any::(), 0..64), ) .prop_map(|(kind, version, payload)| Frame::Unknown { kind, version, payload, }), ] } proptest! { #[test] fn encoded_frame_streams_decode_to_the_same_frames(frames in vec(frame(), 0..8)) { let mut bytes = Vec::new(); for frame in &frames { bytes.extend( frame .encode(DEFAULT_MAX_FRAME_BYTES) .map_err(|e| TestCaseError::fail(e.to_string()))?, ); } let recovery = read_frames(&mut bytes.as_slice(), DEFAULT_MAX_FRAME_BYTES); prop_assert!(recovery.damage.is_empty(), "damage: {:?}", recovery.damage); prop_assert_eq!(recovery.torn_tail_at, None); prop_assert_eq!(recovery.good_end, bytes.len() as u64); let mut again = Vec::new(); for frame in &recovery.frames { again.extend( frame .encode(DEFAULT_MAX_FRAME_BYTES) .map_err(|e| TestCaseError::fail(e.to_string()))?, ); } prop_assert_eq!(recovery.frames, frames); prop_assert_eq!(again, bytes); } #[test] fn a_frame_exactly_at_the_limit_round_trips_and_one_byte_less_rejects_it(frame in frame()) { let unlimited = frame .encode(DEFAULT_MAX_FRAME_BYTES) .map_err(|e| TestCaseError::fail(e.to_string()))?; let payload_len = unlimited.len() - HEADER_LEN - BODY_PREFIX_LEN; let declared = u32::from_le_bytes(unlimited[..HEADER_LEN].try_into().unwrap()); prop_assert_eq!(declared as usize, BODY_PREFIX_LEN + payload_len); let at_limit = frame .encode(payload_len) .map_err(|e| TestCaseError::fail(e.to_string()))?; prop_assert_eq!(&at_limit, &unlimited); let recovery = read_frames(&mut at_limit.as_slice(), payload_len); prop_assert!(recovery.damage.is_empty(), "damage: {:?}", recovery.damage); prop_assert_eq!(recovery.frames, vec![frame.clone()]); if let Some(below) = payload_len.checked_sub(1) { prop_assert_eq!(frame.encode(below), Err(too_large(below, payload_len))); let recovery = read_frames(&mut at_limit.as_slice(), below); prop_assert!(recovery.frames.is_empty()); prop_assert_eq!(recovery.good_end, at_limit.len() as u64); } } } }