repositories / smith
smith
There are many coding harnesses - but this one is fast
owned by admin
smith-core/src/frame.rs
Raw//! 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<u8>,
},
}
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<T: serde::Serialize>(
value: &T,
max_frame_bytes: usize,
) -> Result<Vec<u8>> {
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<Vec<u8>> {
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<Vec<u8>> {
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<T: serde::Serialize>(
kind: u16,
version: u16,
value: &T,
content_len: usize,
max_frame_bytes: usize,
) -> Result<Vec<u8>> {
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<T: serde::Serialize>(
value: &T,
buf: Vec<u8>,
max_frame_bytes: usize,
) -> Result<Vec<u8>> {
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<u8>, max_frame_bytes: usize) -> Result<Vec<u8>> {
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<Frame, FrameFault> {
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<T: serde::de::DeserializeOwned>(
payload: &[u8],
) -> std::result::Result<T, FrameFault> {
ciborium::de::from_reader::<T, _>(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<Frame>,
/// Damage reports, each locating its frame header by byte offset.
pub damage: Vec<SmithError>,
/// 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<u64>,
}
/// 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<R: Read>(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<R: Read>(reader: &mut R, buf: &mut [u8]) -> std::io::Result<bool> {
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<R: Read>(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<u8>` writer that stops accepting bytes beyond a limit,
/// bounding allocation during encoding.
struct LimitedWriter {
buf: Vec<u8>,
remaining: usize,
overflow: bool,
}
impl LimitedWriter {
/// Accept up to `limit` bytes after the ones already in `buf`.
const fn new(buf: Vec<u8>, limit: usize) -> Self {
Self {
buf,
remaining: limit,
overflow: false,
}
}
const fn overflowed(&self) -> bool {
self.overflow
}
fn into_inner(self) -> Vec<u8> {
self.buf
}
}
impl Write for LimitedWriter {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
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<u8> {
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<String> = 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<_>>(),
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<_>>(),
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<_>>(),
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<Value = EntryKind> {
prop_oneof![
(any::<u128>(), role()).prop_map(|(id, role)| EntryKind::Message {
id: MessageId::from_u128(id),
role,
}),
any::<u128>().prop_map(|id| EntryKind::ToolCall {
call_id: ToolCallId::from_u128(id),
}),
(any::<u128>(), any::<bool>()).prop_map(|(id, ok)| EntryKind::ToolResult {
call_id: ToolCallId::from_u128(id),
ok,
}),
any::<bool>().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<Value = Frame> {
prop_oneof![
(any::<u128>(), any::<u128>()).prop_map(|(session, branch)| Frame::Header {
header: SessionHeader {
session_id: SessionId::from_u128(session),
branch_id: BranchId::from_u128(branch),
},
}),
(
any::<u128>(),
proptest::option::of(any::<u128>()),
any::<i64>(),
kind(),
vec(any::<u8>(), 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::<u16>().prop_filter("known kinds decode as known", |kind| {
!matches!(*kind, KIND_ENTRY | KIND_SESSION_HEADER)
}),
any::<u16>(),
vec(any::<u8>(), 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);
}
}
}
}