use std::io::{self, Read, Write};
use serde::{Deserialize, Serialize};
use crate::run_meta::{ContextSnapshot, RegionEntrySnapshot, RegionSnapshot, RunMeta, RunStatus};
pub const RUN_ARCHIVE_MAGIC: &[u8; 4] = b"LVR1";
pub const RUN_ARCHIVE_VERSION: u16 = 1;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RunIdentity {
pub run_id: String,
pub machine_id: String,
pub world_id: String,
pub created_at: i64,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct MessageRecord {
pub role: String,
pub content: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ToolCallRecord {
pub id: String,
pub name: String,
pub arguments: String,
pub result: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct InferenceRequestRecord {
pub model: String,
pub system: Vec<String>,
pub messages: Vec<MessageRecord>,
pub tool_names: Vec<String>,
pub temperature: f32,
pub max_tokens: usize,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct InferenceResponseRecord {
pub content: String,
pub tool_calls: Vec<ToolCallRecord>,
pub prompt_tokens: usize,
pub completion_tokens: usize,
pub cached_tokens: usize,
pub cache_write_tokens: usize,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub enum RegionDelta {
Set(RegionSnapshot),
Append {
name: String,
entries: Vec<RegionEntrySnapshot>,
current_tokens: usize,
},
Clear {
name: String,
},
Remove {
name: String,
},
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ContextDelta {
pub stage_name: String,
pub total_tokens: usize,
pub max_tokens: usize,
pub regions: Vec<RegionDelta>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub enum RunRecord {
Header {
identity: RunIdentity,
meta: Box<RunMeta>,
},
OwnershipChanged {
machine_id: String,
world_id: String,
at: i64,
},
Inference {
stage: String,
iteration: usize,
request: InferenceRequestRecord,
response: InferenceResponseRecord,
at: i64,
},
ToolBatch {
calls: Vec<ToolCallRecord>,
at: i64,
},
ContextCheckpoint {
snapshot: ContextSnapshot,
at: i64,
},
ContextDiff {
delta: ContextDelta,
at: i64,
},
Message {
message: MessageRecord,
at: i64,
},
StatusChanged {
status: RunStatus,
at: i64,
},
Checkpoint {
meta: Box<RunMeta>,
context: ContextSnapshot,
at: i64,
},
Progress {
meta: Box<RunMeta>,
delta: ContextDelta,
at: i64,
},
}
fn is_prefix(prev: &[RegionEntrySnapshot], next: &[RegionEntrySnapshot]) -> bool {
prev.len() <= next.len() && next[..prev.len()] == *prev
}
pub fn diff_context(prev: &ContextSnapshot, next: &ContextSnapshot) -> ContextDelta {
let mut regions = Vec::new();
for nr in &next.regions {
match prev.regions.iter().find(|r| r.name == nr.name) {
None => regions.push(RegionDelta::Set(nr.clone())),
Some(pr) => {
if pr == nr {
} else if nr.entries.is_empty() && !pr.entries.is_empty() {
regions.push(RegionDelta::Clear {
name: nr.name.clone(),
});
} else if pr.kind == nr.kind
&& pr.max_tokens == nr.max_tokens
&& is_prefix(&pr.entries, &nr.entries)
{
regions.push(RegionDelta::Append {
name: nr.name.clone(),
entries: nr.entries[pr.entries.len()..].to_vec(),
current_tokens: nr.current_tokens,
});
} else {
regions.push(RegionDelta::Set(nr.clone()));
}
}
}
}
for pr in &prev.regions {
if !next.regions.iter().any(|r| r.name == pr.name) {
regions.push(RegionDelta::Remove {
name: pr.name.clone(),
});
}
}
ContextDelta {
stage_name: next.stage_name.clone(),
total_tokens: next.total_tokens,
max_tokens: next.max_tokens,
regions,
}
}
pub fn apply_delta(base: &mut ContextSnapshot, delta: &ContextDelta) {
base.stage_name = delta.stage_name.clone();
base.total_tokens = delta.total_tokens;
base.max_tokens = delta.max_tokens;
for region_delta in &delta.regions {
match region_delta {
RegionDelta::Set(snapshot) => {
match base.regions.iter_mut().find(|r| r.name == snapshot.name) {
Some(existing) => *existing = snapshot.clone(),
None => base.regions.push(snapshot.clone()),
}
}
RegionDelta::Append {
name,
entries,
current_tokens,
} => {
if let Some(region) = base.regions.iter_mut().find(|r| &r.name == name) {
region.entries.extend(entries.iter().cloned());
region.current_tokens = *current_tokens;
}
}
RegionDelta::Clear { name } => {
if let Some(region) = base.regions.iter_mut().find(|r| &r.name == name) {
region.entries.clear();
region.current_tokens = 0;
}
}
RegionDelta::Remove { name } => {
base.regions.retain(|r| &r.name != name);
}
}
}
}
pub fn write_archive_start(w: &mut dyn Write, version: u16) -> io::Result<()> {
w.write_all(RUN_ARCHIVE_MAGIC)?;
w.write_all(&version.to_be_bytes())?;
Ok(())
}
pub fn read_archive_start(r: &mut dyn Read) -> io::Result<u16> {
let mut magic = [0u8; 4];
r.read_exact(&mut magic)?;
if &magic != RUN_ARCHIVE_MAGIC {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"not a leviath run archive (bad magic)",
));
}
let mut version = [0u8; 2];
r.read_exact(&mut version)?;
Ok(u16::from_be_bytes(version))
}
pub fn write_record(w: &mut dyn Write, record: &RunRecord) -> io::Result<()> {
let payload = serde_json::to_vec(record).expect("a RunRecord always serializes to JSON");
let len = payload.len() as u64;
w.write_all(&len.to_be_bytes())?;
w.write_all(&payload)?;
Ok(())
}
fn read_exact_or_eof(r: &mut dyn Read, buf: &mut [u8]) -> io::Result<bool> {
let mut filled = 0;
while filled < buf.len() {
match r.read(&mut buf[filled..])? {
0 => {
if filled == 0 {
return Ok(false); }
return Err(io::Error::new(
io::ErrorKind::UnexpectedEof,
"truncated run-archive frame",
));
}
n => filled += n,
}
}
Ok(true)
}
pub fn read_record(r: &mut dyn Read) -> io::Result<Option<RunRecord>> {
let mut len_bytes = [0u8; 8];
if !read_exact_or_eof(r, &mut len_bytes)? {
return Ok(None);
}
let len = u64::from_be_bytes(len_bytes) as usize;
let mut payload = vec![0u8; len];
if !read_exact_or_eof(r, &mut payload)? {
return Err(io::Error::new(
io::ErrorKind::UnexpectedEof,
"truncated run-archive frame",
));
}
let record = serde_json::from_slice(&payload)
.map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
Ok(Some(record))
}
pub fn read_archive(r: &mut dyn Read) -> io::Result<(u16, Vec<RunRecord>)> {
let version = read_archive_start(r)?;
let mut records = Vec::new();
while let Some(record) = read_record(r)? {
records.push(record);
}
Ok((version, records))
}
pub fn read_archive_lenient(r: &mut dyn Read) -> io::Result<(u16, Vec<RunRecord>)> {
let version = read_archive_start(r)?;
let mut records = Vec::new();
while let Ok(Some(record)) = read_record(r) {
records.push(record);
}
Ok((version, records))
}
#[derive(Debug, Clone, PartialEq)]
pub struct FoldedRun {
pub identity: RunIdentity,
pub meta: RunMeta,
pub context: ContextSnapshot,
pub messages: Vec<MessageRecord>,
pub inference_count: usize,
pub tool_call_count: usize,
}
pub fn fold(records: &[RunRecord]) -> Option<FoldedRun> {
let mut iter = records.iter();
let (identity, meta) = match iter.next() {
Some(RunRecord::Header { identity, meta }) => (identity.clone(), (**meta).clone()),
_ => return None,
};
let mut folded = FoldedRun {
identity,
meta,
context: ContextSnapshot {
stage_name: String::new(),
total_tokens: 0,
max_tokens: 0,
regions: Vec::new(),
},
messages: Vec::new(),
inference_count: 0,
tool_call_count: 0,
};
for record in iter {
match record {
RunRecord::Header { identity, meta } => {
folded.identity = identity.clone();
folded.meta = (**meta).clone();
}
RunRecord::OwnershipChanged {
machine_id,
world_id,
..
} => {
folded.identity.machine_id = machine_id.clone();
folded.identity.world_id = world_id.clone();
}
RunRecord::Inference { .. } => folded.inference_count += 1,
RunRecord::ToolBatch { calls, .. } => folded.tool_call_count += calls.len(),
RunRecord::ContextCheckpoint { snapshot, .. } => folded.context = snapshot.clone(),
RunRecord::ContextDiff { delta, .. } => apply_delta(&mut folded.context, delta),
RunRecord::Message { message, .. } => folded.messages.push(message.clone()),
RunRecord::StatusChanged { status, .. } => folded.meta.status = status.clone(),
RunRecord::Checkpoint { meta, context, .. } => {
folded.meta = (**meta).clone();
folded.context = context.clone();
}
RunRecord::Progress { meta, delta, .. } => {
folded.meta = (**meta).clone();
apply_delta(&mut folded.context, delta);
}
}
}
Some(folded)
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RunPoint {
pub meta: RunMeta,
pub context: ContextSnapshot,
pub at: i64,
}
pub fn replay_points(records: &[RunRecord]) -> Vec<RunPoint> {
let mut iter = records.iter();
let mut meta = match iter.next() {
Some(RunRecord::Header { meta, .. }) => (**meta).clone(),
_ => return Vec::new(),
};
let mut context = ContextSnapshot {
stage_name: String::new(),
total_tokens: 0,
max_tokens: 0,
regions: Vec::new(),
};
let mut points = Vec::new();
for record in iter {
match record {
RunRecord::Header { meta: m, .. } => meta = (**m).clone(),
RunRecord::StatusChanged { status, .. } => meta.status = status.clone(),
RunRecord::ContextCheckpoint { snapshot, at } => {
context = snapshot.clone();
points.push(RunPoint {
meta: meta.clone(),
context: context.clone(),
at: *at,
});
}
RunRecord::ContextDiff { delta, at } => {
apply_delta(&mut context, delta);
points.push(RunPoint {
meta: meta.clone(),
context: context.clone(),
at: *at,
});
}
RunRecord::Checkpoint {
meta: m,
context: c,
at,
} => {
meta = (**m).clone();
context = c.clone();
points.push(RunPoint {
meta: meta.clone(),
context: context.clone(),
at: *at,
});
}
RunRecord::Progress { meta: m, delta, at } => {
meta = (**m).clone();
apply_delta(&mut context, delta);
points.push(RunPoint {
meta: meta.clone(),
context: context.clone(),
at: *at,
});
}
RunRecord::OwnershipChanged { .. }
| RunRecord::Inference { .. }
| RunRecord::ToolBatch { .. }
| RunRecord::Message { .. } => {}
}
}
points
}
#[cfg(test)]
mod tests {
use super::*;
use crate::run_meta::RunStatus;
fn identity() -> RunIdentity {
RunIdentity {
run_id: "run-1".to_string(),
machine_id: "machine-a".to_string(),
world_id: "world-x".to_string(),
created_at: 100,
}
}
fn meta() -> RunMeta {
RunMeta::new(
"run-1".to_string(),
"coder".to_string(),
"/agents/coder".to_string(),
"do it".to_string(),
Some("anthropic/claude".to_string()),
"/work".to_string(),
2,
)
}
fn entry(content: &str, tokens: usize) -> RegionEntrySnapshot {
RegionEntrySnapshot {
content: content.to_string(),
tokens,
kind: crate::region::EntryKind::Text,
metadata: None,
key: None,
taint: Default::default(),
}
}
fn region(name: &str, entries: Vec<RegionEntrySnapshot>) -> RegionSnapshot {
let current = entries.iter().map(|e| e.tokens).sum();
RegionSnapshot {
name: name.to_string(),
kind: "clearable".to_string(),
current_tokens: current,
max_tokens: 1000,
entries,
}
}
fn snapshot(stage: &str, regions: Vec<RegionSnapshot>) -> ContextSnapshot {
let total = regions.iter().map(|r| r.current_tokens).sum();
ContextSnapshot {
stage_name: stage.to_string(),
total_tokens: total,
max_tokens: 10_000,
regions,
}
}
fn header() -> RunRecord {
RunRecord::Header {
identity: identity(),
meta: Box::new(meta()),
}
}
fn region_delta_kind(d: &RegionDelta) -> &'static str {
match d {
RegionDelta::Set(_) => "set",
RegionDelta::Append { .. } => "append",
RegionDelta::Clear { .. } => "clear",
RegionDelta::Remove { .. } => "remove",
}
}
fn assert_diff_roundtrip(a: &ContextSnapshot, b: &ContextSnapshot) {
let delta = diff_context(a, b);
let mut base = a.clone();
apply_delta(&mut base, &delta);
assert_eq!(&base, b);
}
#[test]
fn diff_append_only_growth_is_compact() {
let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
let b = snapshot(
"s1",
vec![region("conv", vec![entry("hi", 1), entry("there", 2)])],
);
let delta = diff_context(&a, &b);
assert_eq!(region_delta_kind(&delta.regions[0]), "append");
assert_diff_roundtrip(&a, &b);
}
#[test]
fn diff_new_region_is_set() {
let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
let b = snapshot(
"s1",
vec![
region("conv", vec![entry("hi", 1)]),
region("plan", vec![entry("p", 3)]),
],
);
let delta = diff_context(&a, &b);
assert!(delta.regions.iter().any(|d| region_delta_kind(d) == "set"));
assert_diff_roundtrip(&a, &b);
}
#[test]
fn diff_cleared_region() {
let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
let b = snapshot("s1", vec![region("conv", vec![])]);
let delta = diff_context(&a, &b);
assert_eq!(region_delta_kind(&delta.regions[0]), "clear");
assert_diff_roundtrip(&a, &b);
}
#[test]
fn diff_removed_region() {
let a = snapshot(
"s1",
vec![
region("conv", vec![entry("hi", 1)]),
region("plan", vec![entry("p", 3)]),
],
);
let b = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
let delta = diff_context(&a, &b);
assert!(
delta
.regions
.iter()
.any(|d| region_delta_kind(d) == "remove")
);
assert_diff_roundtrip(&a, &b);
}
#[test]
fn diff_non_prefix_rewrite_is_set() {
let a = snapshot("s1", vec![region("conv", vec![entry("old", 1)])]);
let b = snapshot("s1", vec![region("conv", vec![entry("new", 1)])]);
let delta = diff_context(&a, &b);
assert_eq!(region_delta_kind(&delta.regions[0]), "set");
assert_diff_roundtrip(&a, &b);
}
#[test]
fn diff_kind_change_is_set_not_append() {
let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
let mut grown = region("conv", vec![entry("hi", 1), entry("more", 1)]);
grown.kind = "sliding".to_string();
let b = snapshot("s1", vec![grown]);
let delta = diff_context(&a, &b);
assert_eq!(region_delta_kind(&delta.regions[0]), "set");
assert_diff_roundtrip(&a, &b);
}
#[test]
fn diff_unchanged_region_emits_nothing() {
let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
let b = a.clone();
let delta = diff_context(&a, &b);
assert!(delta.regions.is_empty());
assert_diff_roundtrip(&a, &b);
}
#[test]
fn apply_delta_skips_unknown_regions_leniently() {
let mut base = snapshot("s1", vec![]);
let delta = ContextDelta {
stage_name: "s1".to_string(),
total_tokens: 0,
max_tokens: 10_000,
regions: vec![
RegionDelta::Append {
name: "ghost".to_string(),
entries: vec![entry("x", 1)],
current_tokens: 1,
},
RegionDelta::Clear {
name: "ghost".to_string(),
},
RegionDelta::Remove {
name: "ghost".to_string(),
},
],
};
apply_delta(&mut base, &delta);
assert!(base.regions.is_empty());
}
fn all_record_kinds() -> Vec<RunRecord> {
vec![
header(),
RunRecord::OwnershipChanged {
machine_id: "machine-b".to_string(),
world_id: "world-y".to_string(),
at: 101,
},
RunRecord::Inference {
stage: "plan".to_string(),
iteration: 0,
request: InferenceRequestRecord {
model: "m".to_string(),
system: vec!["sys".to_string()],
messages: vec![MessageRecord {
role: "user".to_string(),
content: "hi".to_string(),
}],
tool_names: vec!["read_file".to_string()],
temperature: 0.7,
max_tokens: 1024,
},
response: InferenceResponseRecord {
content: "ok".to_string(),
tool_calls: vec![],
prompt_tokens: 10,
completion_tokens: 5,
cached_tokens: 0,
cache_write_tokens: 0,
},
at: 102,
},
RunRecord::ToolBatch {
calls: vec![ToolCallRecord {
id: "c1".to_string(),
name: "read_file".to_string(),
arguments: "{}".to_string(),
result: Some("body".to_string()),
}],
at: 103,
},
RunRecord::ContextCheckpoint {
snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
at: 104,
},
RunRecord::ContextDiff {
delta: ContextDelta {
stage_name: "plan".to_string(),
total_tokens: 3,
max_tokens: 10_000,
regions: vec![RegionDelta::Append {
name: "conv".to_string(),
entries: vec![entry("more", 2)],
current_tokens: 3,
}],
},
at: 105,
},
RunRecord::Message {
message: MessageRecord {
role: "user".to_string(),
content: "another".to_string(),
},
at: 106,
},
RunRecord::StatusChanged {
status: RunStatus::Complete,
at: 107,
},
RunRecord::Checkpoint {
meta: Box::new(meta()),
context: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
at: 108,
},
RunRecord::Progress {
meta: Box::new(meta()),
delta: ContextDelta {
stage_name: "plan".to_string(),
total_tokens: 3,
max_tokens: 10_000,
regions: vec![RegionDelta::Append {
name: "conv".to_string(),
entries: vec![entry("step", 2)],
current_tokens: 3,
}],
},
at: 109,
},
]
}
#[test]
fn archive_write_then_read_roundtrips_every_record_kind() {
let records = all_record_kinds();
let mut buf = Vec::new();
write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
for r in &records {
write_record(&mut buf, r).unwrap();
}
let (version, read) = read_archive(&mut buf.as_slice()).unwrap();
assert_eq!(version, RUN_ARCHIVE_VERSION);
assert_eq!(read, records);
}
#[test]
fn read_archive_start_rejects_bad_magic() {
let mut bytes: &[u8] = b"XXXX\x00\x01";
let err = read_archive_start(&mut bytes).unwrap_err();
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
}
#[test]
fn read_archive_start_reports_version() {
let mut buf = Vec::new();
write_archive_start(&mut buf, 7).unwrap();
assert_eq!(read_archive_start(&mut buf.as_slice()).unwrap(), 7);
}
#[test]
fn read_record_returns_none_at_clean_eof() {
let empty: &[u8] = &[];
assert!(read_record(&mut { empty }).unwrap().is_none());
}
#[test]
fn read_record_errors_on_truncated_length_prefix() {
let mut bytes: &[u8] = &[0, 0];
let err = read_record(&mut bytes).unwrap_err();
assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
}
#[test]
fn read_record_errors_on_truncated_payload() {
let mut bytes: &[u8] = &[0, 0, 0, 0, 0, 0, 0, 10, 1, 2];
let err = read_record(&mut bytes).unwrap_err();
assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
}
#[test]
fn read_record_errors_on_empty_payload_at_boundary() {
let mut bytes: &[u8] = &[0, 0, 0, 0, 0, 0, 0, 10];
let err = read_record(&mut bytes).unwrap_err();
assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
}
#[test]
fn read_record_errors_on_invalid_json_payload() {
let mut buf = Vec::new();
let bad = b"not json";
buf.extend_from_slice(&(bad.len() as u64).to_be_bytes());
buf.extend_from_slice(bad);
let err = read_record(&mut buf.as_slice()).unwrap_err();
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
}
struct FailingReader;
impl Read for FailingReader {
fn read(&mut self, _buf: &mut [u8]) -> io::Result<usize> {
Err(io::Error::other("device error"))
}
}
#[test]
fn read_record_propagates_reader_errors() {
let err = read_record(&mut FailingReader).unwrap_err();
assert_eq!(err.kind(), io::ErrorKind::Other);
}
#[test]
fn read_archive_propagates_a_bad_preamble() {
let mut bytes: &[u8] = b"LV";
assert!(read_archive(&mut bytes).is_err());
}
#[test]
fn read_archive_propagates_a_bad_frame() {
let mut buf = Vec::new();
write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
buf.extend_from_slice(&[0, 0, 0, 0, 0, 0, 0, 5, 1, 2]); let err = read_archive(&mut buf.as_slice()).unwrap_err();
assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
}
struct FailAfter {
remaining: usize,
}
impl Write for FailAfter {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
if self.remaining == 0 {
return Err(io::Error::other("disk full"));
}
let n = buf.len().min(self.remaining);
self.remaining -= n;
Ok(n)
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
#[test]
fn fail_after_writer_flush_is_a_noop() {
assert!(FailAfter { remaining: 1 }.flush().is_ok());
}
#[test]
fn write_archive_start_propagates_write_errors() {
assert!(write_archive_start(&mut FailAfter { remaining: 0 }, 1).is_err());
assert!(write_archive_start(&mut FailAfter { remaining: 4 }, 1).is_err());
}
#[test]
fn write_record_propagates_write_errors() {
let rec = header();
assert!(write_record(&mut FailAfter { remaining: 0 }, &rec).is_err());
assert!(write_record(&mut FailAfter { remaining: 8 }, &rec).is_err());
}
#[test]
fn read_archive_lenient_matches_strict_on_a_clean_archive() {
let records = all_record_kinds();
let mut buf = Vec::new();
write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
for r in &records {
write_record(&mut buf, r).unwrap();
}
let (version, read) = read_archive_lenient(&mut buf.as_slice()).unwrap();
assert_eq!(version, RUN_ARCHIVE_VERSION);
assert_eq!(read, records);
}
#[test]
fn read_archive_lenient_keeps_valid_prefix_before_a_torn_tail() {
let mut buf = Vec::new();
write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
write_record(&mut buf, &header()).unwrap();
write_record(
&mut buf,
&RunRecord::ContextCheckpoint {
snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
at: 1,
},
)
.unwrap();
buf.extend_from_slice(&[0, 0, 0, 0, 0, 0, 0, 10, 1, 2]);
assert!(read_archive(&mut buf.as_slice()).is_err());
let (version, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
assert_eq!(version, RUN_ARCHIVE_VERSION);
assert_eq!(records.len(), 2);
let folded = fold(&records).expect("prefix starts with a Header");
assert_eq!(folded.context.regions[0].entries.len(), 1);
}
#[test]
fn read_archive_lenient_still_errors_on_a_bad_preamble() {
let mut bad_magic: &[u8] = b"XXXX\x00\x01";
assert!(read_archive_lenient(&mut bad_magic).is_err());
let mut short: &[u8] = b"LVR1";
assert!(read_archive_lenient(&mut short).is_err());
}
#[test]
fn read_archive_start_errors_on_truncated_version() {
let mut bytes: &[u8] = b"LVR1";
let err = read_archive_start(&mut bytes).unwrap_err();
assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
}
#[test]
fn fold_requires_a_header_first() {
assert!(fold(&[]).is_none());
assert!(
fold(&[RunRecord::StatusChanged {
status: RunStatus::Complete,
at: 1
}])
.is_none()
);
}
#[test]
fn fold_reconstructs_state_from_the_journal() {
let records = all_record_kinds();
let folded = fold(&records).expect("has header");
assert_eq!(folded.identity.machine_id, "machine-b");
assert_eq!(folded.identity.world_id, "world-y");
assert_eq!(folded.inference_count, 1);
assert_eq!(folded.tool_call_count, 1);
assert_eq!(folded.messages.len(), 1);
assert_eq!(folded.messages[0].content, "another");
assert_eq!(folded.context.regions[0].name, "conv");
assert_eq!(folded.context.regions[0].entries.len(), 2);
assert_eq!(folded.context.total_tokens, 3);
assert_eq!(folded.meta.run_id, "run-1");
}
#[test]
fn fold_applies_context_diffs_over_a_checkpoint() {
let records = vec![
header(),
RunRecord::ContextCheckpoint {
snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
at: 1,
},
RunRecord::ContextDiff {
delta: ContextDelta {
stage_name: "plan".to_string(),
total_tokens: 3,
max_tokens: 10_000,
regions: vec![RegionDelta::Append {
name: "conv".to_string(),
entries: vec![entry("there", 2)],
current_tokens: 3,
}],
},
at: 2,
},
];
let folded = fold(&records).unwrap();
assert_eq!(folded.context.regions[0].entries.len(), 2);
assert_eq!(folded.context.total_tokens, 3);
}
#[test]
fn fold_later_header_updates_identity_and_meta() {
let mut second_meta = meta();
second_meta.status = RunStatus::Running;
let records = vec![
header(),
RunRecord::Header {
identity: RunIdentity {
run_id: "run-1".to_string(),
machine_id: "machine-c".to_string(),
world_id: "world-z".to_string(),
created_at: 200,
},
meta: Box::new(second_meta),
},
];
let folded = fold(&records).unwrap();
assert_eq!(folded.identity.machine_id, "machine-c");
assert_eq!(folded.meta.status, RunStatus::Running);
}
#[test]
fn fold_progress_applies_meta_and_context_diff() {
let mut advanced = meta();
advanced.status = RunStatus::Running;
advanced.iteration = 5;
let records = vec![
header(),
RunRecord::ContextCheckpoint {
snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
at: 1,
},
RunRecord::Progress {
meta: Box::new(advanced),
delta: ContextDelta {
stage_name: "plan".to_string(),
total_tokens: 3,
max_tokens: 10_000,
regions: vec![RegionDelta::Append {
name: "conv".to_string(),
entries: vec![entry("there", 2)],
current_tokens: 3,
}],
},
at: 2,
},
];
let folded = fold(&records).unwrap();
assert_eq!(folded.meta.iteration, 5);
assert_eq!(folded.meta.status, RunStatus::Running);
assert_eq!(folded.context.regions[0].entries.len(), 2);
}
#[test]
fn replay_points_requires_a_header() {
assert!(replay_points(&[]).is_empty());
assert!(
replay_points(&[RunRecord::Message {
message: MessageRecord {
role: "user".to_string(),
content: "x".to_string(),
},
at: 1,
}])
.is_empty()
);
}
#[test]
fn replay_points_emits_a_snapshot_per_context_change() {
let mut running = meta();
running.status = RunStatus::Running;
let records = vec![
header(),
RunRecord::Inference {
stage: "plan".to_string(),
iteration: 0,
request: InferenceRequestRecord {
model: "m".to_string(),
system: vec![],
messages: vec![],
tool_names: vec![],
temperature: 0.7,
max_tokens: 10,
},
response: InferenceResponseRecord {
content: "ok".to_string(),
tool_calls: vec![],
prompt_tokens: 1,
completion_tokens: 1,
cached_tokens: 0,
cache_write_tokens: 0,
},
at: 1,
},
RunRecord::ContextCheckpoint {
snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
at: 2,
},
RunRecord::StatusChanged {
status: RunStatus::Running,
at: 3,
},
RunRecord::Progress {
meta: Box::new(running),
delta: ContextDelta {
stage_name: "implement".to_string(),
total_tokens: 3,
max_tokens: 10_000,
regions: vec![RegionDelta::Append {
name: "conv".to_string(),
entries: vec![entry("more", 2)],
current_tokens: 3,
}],
},
at: 4,
},
];
let points = replay_points(&records);
assert_eq!(points.len(), 2, "one point per context change");
assert_eq!(points[0].at, 2);
assert_eq!(points[0].context.regions[0].entries.len(), 1);
assert_eq!(points[1].at, 4);
assert_eq!(points[1].context.regions[0].entries.len(), 2);
assert_eq!(points[1].context.stage_name, "implement");
assert_eq!(points[1].meta.status, RunStatus::Running);
}
#[test]
fn replay_points_handles_context_diff_and_a_later_header() {
let mut relabeled = meta();
relabeled.agent_name = "renamed".to_string();
let records = vec![
header(),
RunRecord::ContextCheckpoint {
snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
at: 1,
},
RunRecord::Header {
identity: identity(),
meta: Box::new(relabeled),
},
RunRecord::ContextDiff {
delta: ContextDelta {
stage_name: "plan".to_string(),
total_tokens: 3,
max_tokens: 10_000,
regions: vec![RegionDelta::Append {
name: "conv".to_string(),
entries: vec![entry("more", 2)],
current_tokens: 3,
}],
},
at: 2,
},
];
let points = replay_points(&records);
assert_eq!(points.len(), 2); assert_eq!(points[1].context.regions[0].entries.len(), 2);
assert_eq!(points[1].meta.agent_name, "renamed");
}
#[test]
fn replay_points_over_a_full_checkpoint() {
let records = vec![
header(),
RunRecord::Checkpoint {
meta: Box::new(meta()),
context: snapshot("review", vec![region("conv", vec![entry("x", 4)])]),
at: 9,
},
];
let points = replay_points(&records);
assert_eq!(points.len(), 1);
assert_eq!(points[0].context.stage_name, "review");
assert_eq!(points[0].context.regions[0].entries[0].tokens, 4);
}
}