use serde::{Deserialize, Serialize};
use crate::run_meta::{ContextSnapshot, RegionEntrySnapshot, RegionSnapshot, RunMeta, RunStatus};
#[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,
#[serde(default, skip_serializing_if = "String::is_empty")]
pub execution_id: String,
pub name: String,
pub arguments: String,
pub result: Option<crate::region::EntryContent>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub thought_signature: 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, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum InferenceKind {
#[default]
Stage,
Compaction,
Title,
Routing,
}
impl InferenceKind {
pub fn label(&self) -> &'static str {
match self {
InferenceKind::Stage => "stage",
InferenceKind::Compaction => "compaction",
InferenceKind::Title => "title",
InferenceKind::Routing => "routing",
}
}
pub fn is_stage_work(&self) -> bool {
matches!(self, InferenceKind::Stage)
}
}
#[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,
},
InferenceUsage {
#[serde(default)]
kind: InferenceKind,
stage: String,
iteration: usize,
provider: String,
model: String,
prompt_tokens: usize,
completion_tokens: usize,
cached_tokens: usize,
cache_write_tokens: usize,
#[serde(default, skip_serializing_if = "Option::is_none")]
cost_usd: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
cost_reported_by_provider: Option<bool>,
at: i64,
},
ToolBatch {
calls: Vec<ToolCallRecord>,
at: i64,
#[serde(default)]
stage_index: usize,
#[serde(default)]
iteration: usize,
#[serde(default, skip_serializing_if = "String::is_empty")]
visit_id: String,
#[serde(default, skip_serializing_if = "String::is_empty")]
requested_by: String,
#[serde(default)]
response: String,
},
ToolCallDone {
iteration: usize,
call_id: String,
#[serde(default, skip_serializing_if = "String::is_empty")]
execution_id: String,
result: crate::region::EntryContent,
#[serde(default, skip_serializing_if = "Option::is_none")]
outcome: Option<crate::execution::ToolOutcome>,
at: i64,
},
ArtifactsProduced {
execution_id: String,
artifacts: Vec<crate::output::Artifact>,
at: i64,
},
Interaction {
request_id: String,
kind: crate::interaction::InteractionKind,
#[serde(default, skip_serializing_if = "Option::is_none")]
tool: Option<String>,
prompt: String,
stage: String,
settlement: crate::interaction::Settlement,
asked_at: i64,
at: i64,
},
InferenceAttempt(AttemptRecord),
InferenceFailover(FailoverRecord),
ContextCheckpoint {
snapshot: ContextSnapshot,
at: i64,
},
ContextChange {
region: String,
cause: crate::ContextCause,
entries_added: usize,
entries_removed: usize,
token_delta: i64,
at: i64,
},
ContextTransaction {
revision_before: String,
revision_after: String,
cause: crate::ContextCause,
regions: Vec<RegionCommit>,
#[serde(default, skip_serializing_if = "String::is_empty")]
execution_id: String,
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
}
trait PriorRegion {
fn name(&self) -> &str;
fn unchanged(&self, next: &RegionSnapshot) -> bool;
fn appended_to(&self, next: &RegionSnapshot) -> bool;
fn entry_count(&self) -> usize;
}
impl PriorRegion for RegionSnapshot {
fn name(&self) -> &str {
&self.name
}
fn unchanged(&self, next: &RegionSnapshot) -> bool {
self == next
}
fn appended_to(&self, next: &RegionSnapshot) -> bool {
self.kind == next.kind
&& self.max_tokens == next.max_tokens
&& is_prefix(&self.entries, &next.entries)
}
fn entry_count(&self) -> usize {
self.entries.len()
}
}
fn diff_regions<P: PriorRegion>(prev: &[P], next: &ContextSnapshot) -> ContextDelta {
let mut regions = Vec::new();
for nr in &next.regions {
match prev.iter().find(|r| r.name() == nr.name) {
None => regions.push(RegionDelta::Set(nr.clone())),
Some(pr) => {
if pr.unchanged(nr) {
} else if nr.entries.is_empty() && pr.entry_count() > 0 {
regions.push(RegionDelta::Clear {
name: nr.name.clone(),
});
} else if pr.appended_to(nr) {
regions.push(RegionDelta::Append {
name: nr.name.clone(),
entries: nr.entries[pr.entry_count()..].to_vec(),
current_tokens: nr.current_tokens,
});
} else {
regions.push(RegionDelta::Set(nr.clone()));
}
}
}
}
for pr in prev {
if !next.regions.iter().any(|r| r.name == pr.name()) {
regions.push(RegionDelta::Remove {
name: pr.name().to_string(),
});
}
}
ContextDelta {
stage_name: next.stage_name.clone(),
total_tokens: next.total_tokens,
max_tokens: next.max_tokens,
regions,
}
}
pub fn diff_context(prev: &ContextSnapshot, next: &ContextSnapshot) -> ContextDelta {
diff_regions(&prev.regions, next)
}
#[derive(Debug, Clone, PartialEq)]
pub struct RegionDigest {
pub name: String,
pub kind: String,
pub current_tokens: usize,
pub max_tokens: usize,
pub entries: Vec<u64>,
}
#[derive(Debug, Clone, PartialEq, Default)]
pub struct ContextDigest {
pub regions: Vec<RegionDigest>,
}
impl ContextDigest {
pub fn fingerprint(&self) -> String {
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
for region in &self.regions {
region.name.hash(&mut hasher);
region.kind.hash(&mut hasher);
region.current_tokens.hash(&mut hasher);
region.max_tokens.hash(&mut hasher);
region.entries.hash(&mut hasher);
}
format!("{:016x}", hasher.finish())
}
}
fn entry_digest(entry: &RegionEntrySnapshot) -> u64 {
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
entry.content.hash(&mut hasher);
entry.tokens.hash(&mut hasher);
entry.key.hash(&mut hasher);
serde_json::to_string(&entry.kind)
.expect("EntryKind always serializes")
.hash(&mut hasher);
serde_json::to_string(&entry.metadata)
.expect("entry metadata always serializes")
.hash(&mut hasher);
serde_json::to_string(&entry.taint)
.expect("taint always serializes")
.hash(&mut hasher);
hasher.finish()
}
pub fn digest_context(snapshot: &ContextSnapshot) -> ContextDigest {
ContextDigest {
regions: snapshot
.regions
.iter()
.map(|r| RegionDigest {
name: r.name.clone(),
kind: r.kind.clone(),
current_tokens: r.current_tokens,
max_tokens: r.max_tokens,
entries: r.entries.iter().map(entry_digest).collect(),
})
.collect(),
}
}
fn is_prefix_digest(prev: &[u64], next: &[RegionEntrySnapshot]) -> bool {
prev.len() <= next.len()
&& prev
.iter()
.zip(next)
.all(|(hash, entry)| *hash == entry_digest(entry))
}
impl PriorRegion for RegionDigest {
fn name(&self) -> &str {
&self.name
}
fn unchanged(&self, next: &RegionSnapshot) -> bool {
self.kind == next.kind
&& self.max_tokens == next.max_tokens
&& self.current_tokens == next.current_tokens
&& self.entries.len() == next.entries.len()
&& is_prefix_digest(&self.entries, &next.entries)
}
fn appended_to(&self, next: &RegionSnapshot) -> bool {
self.kind == next.kind
&& self.max_tokens == next.max_tokens
&& is_prefix_digest(&self.entries, &next.entries)
}
fn entry_count(&self) -> usize {
self.entries.len()
}
}
pub fn diff_context_digest(prev: &ContextDigest, next: &ContextSnapshot) -> ContextDelta {
diff_regions(&prev.regions, next)
}
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);
}
}
}
}
mod attempt;
mod codec;
mod executions;
mod points;
mod transaction;
pub use attempt::{
AttemptOutcome, AttemptRecord, CaptureStatus, FailoverRecord, ModelInput, RequestDigest, Retry,
};
pub use codec::{
Frame, Frames, RUN_ARCHIVE_MAGIC, RUN_ARCHIVE_VERSION, read_archive, read_archive_lenient,
read_archive_start, read_frame, read_record, write_archive_start, write_record,
};
pub use executions::{Execution, SeekRead, read_archive_executions, read_result_at};
pub use points::{PointRef, RunPoint, replay_points, visit_archive_points, visit_points};
pub use transaction::{
ContextChangeRecord, IndexedChange, RegionCommit, RegionTransition, read_archive_changes,
};
#[derive(Debug, Clone, PartialEq)]
pub struct PendingToolBatch {
pub stage_index: usize,
pub iteration: usize,
pub response: String,
pub calls: Vec<ToolCallRecord>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct InteractionRecord {
pub request_id: String,
pub kind: crate::interaction::InteractionKind,
pub tool: Option<String>,
pub prompt: String,
pub stage: String,
pub settlement: crate::interaction::Settlement,
pub asked_at: i64,
pub at: i64,
}
#[derive(Debug, Clone, PartialEq)]
pub struct InferenceUsageRecord {
pub kind: InferenceKind,
pub stage: String,
pub iteration: usize,
pub provider: String,
pub model: String,
pub prompt_tokens: usize,
pub completion_tokens: usize,
pub cached_tokens: usize,
pub cache_write_tokens: usize,
pub cost_usd: Option<f64>,
pub cost_reported_by_provider: Option<bool>,
pub at: i64,
}
#[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 inference_usage: Vec<InferenceUsageRecord>,
pub tool_call_count: usize,
pub interactions: Vec<InteractionRecord>,
pub attempts: Vec<AttemptRecord>,
pub failovers: Vec<FailoverRecord>,
pub context_changes: Vec<ContextChangeRecord>,
pub pending_batch: Option<PendingToolBatch>,
}
pub fn context_contains_batch(context: &ContextSnapshot, batch: &PendingToolBatch) -> bool {
let Some(first_id) = batch.calls.first().map(|c| c.id.as_str()) else {
return false;
};
context.regions.iter().any(|region| {
region.entries.iter().any(|entry| {
matches!(
&entry.kind,
crate::region::EntryKind::AssistantTurn { tool_calls }
if tool_calls.iter().any(|tc| tc.id == first_id)
)
})
})
}
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,
inference_usage: Vec::new(),
tool_call_count: 0,
interactions: Vec::new(),
attempts: Vec::new(),
failovers: Vec::new(),
context_changes: Vec::new(),
pending_batch: None,
};
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::ArtifactsProduced { .. } => {}
RunRecord::InferenceAttempt(attempt) => folded.attempts.push(attempt.clone()),
RunRecord::InferenceFailover(failover) => folded.failovers.push(failover.clone()),
RunRecord::Interaction {
request_id,
kind,
tool,
prompt,
stage,
settlement,
asked_at,
at,
} => folded.interactions.push(InteractionRecord {
request_id: request_id.clone(),
kind: kind.clone(),
tool: tool.clone(),
prompt: prompt.clone(),
stage: stage.clone(),
settlement: settlement.clone(),
asked_at: *asked_at,
at: *at,
}),
RunRecord::InferenceUsage {
kind,
stage,
iteration,
provider,
model,
prompt_tokens,
completion_tokens,
cached_tokens,
cache_write_tokens,
cost_usd,
cost_reported_by_provider,
at,
} => {
folded.inference_count += 1;
folded.inference_usage.push(InferenceUsageRecord {
kind: *kind,
stage: stage.clone(),
iteration: *iteration,
provider: provider.clone(),
model: model.clone(),
prompt_tokens: *prompt_tokens,
completion_tokens: *completion_tokens,
cached_tokens: *cached_tokens,
cache_write_tokens: *cache_write_tokens,
cost_usd: *cost_usd,
cost_reported_by_provider: *cost_reported_by_provider,
at: *at,
});
}
RunRecord::ToolBatch {
calls,
stage_index,
iteration,
response,
..
} => {
folded.tool_call_count += calls.len();
folded.pending_batch =
calls
.iter()
.any(|call| call.result.is_none())
.then(|| PendingToolBatch {
stage_index: *stage_index,
iteration: *iteration,
response: response.clone(),
calls: calls.clone(),
});
}
RunRecord::ToolCallDone {
iteration,
call_id,
result,
..
} => {
if let Some(batch) = folded
.pending_batch
.as_mut()
.filter(|b| b.iteration == *iteration)
&& let Some(call) = batch.calls.iter_mut().find(|c| c.id == *call_id)
{
call.result = Some(result.clone());
}
}
RunRecord::ContextChange { .. } | RunRecord::ContextTransaction { .. } => {
folded
.context_changes
.extend(transaction::change_of(record));
}
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);
}
}
}
if let Some(batch) = &folded.pending_batch
&& (folded.meta.iteration != batch.iteration
|| context_contains_batch(&folded.context, batch))
{
folded.pending_batch = None;
}
Some(folded)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ContextCause;
use crate::run_meta::RunStatus;
use std::io::{self, Read, Write};
use std::ops::ControlFlow;
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,
}
}
const FIXTURE_NOW: i64 = 1_700_000_000;
fn meta() -> RunMeta {
let mut meta = 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,
);
meta.started_at = FIXTURE_NOW;
meta.updated_at = FIXTURE_NOW;
meta
}
#[test]
fn the_fixture_does_not_move_with_the_clock() {
let first = meta();
let mut later = meta();
assert_eq!(first, later, "the fixture is rebuilt identically");
later.started_at += 1;
assert_ne!(
first, later,
"and the comparison is sensitive to the field that used to drift"
);
}
fn entry(content: &str, tokens: usize) -> RegionEntrySnapshot {
RegionEntrySnapshot {
content: content.into(),
tokens,
kind: crate::region::EntryKind::Text,
metadata: None,
key: None,
taint: Default::default(),
reasoning: None,
}
}
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,
description: None,
}
}
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);
}
fn framed(records: &[RunRecord]) -> Vec<u8> {
let mut buf = Vec::new();
write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
for record in records {
write_record(&mut buf, record).unwrap();
}
buf
}
fn try_collect_streamed(bytes: &[u8]) -> io::Result<Vec<(usize, i64, usize)>> {
let mut seen = Vec::new();
visit_archive_points(&mut &bytes[..], &mut |p| {
seen.push((p.index, p.at, p.context.total_tokens));
ControlFlow::Continue(())
})?;
Ok(seen)
}
fn collect_streamed(bytes: &[u8]) -> Vec<(usize, i64, usize)> {
try_collect_streamed(bytes).unwrap()
}
#[test]
fn visit_archive_points_matches_visit_points() {
let records = vec![
header(),
RunRecord::ContextCheckpoint {
snapshot: snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
at: 10,
},
RunRecord::StatusChanged {
status: RunStatus::Running,
at: 11,
},
RunRecord::Progress {
meta: Box::new(meta()),
delta: diff_context(
&snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
&snapshot(
"s1",
vec![region("conv", vec![entry("hi", 1), entry("more", 2)])],
),
),
at: 12,
},
];
let mut in_memory = Vec::new();
visit_points(&records, &mut |p| {
in_memory.push((p.index, p.at, p.context.total_tokens));
ControlFlow::Continue(())
});
assert_eq!(collect_streamed(&framed(&records)), in_memory);
assert_eq!(in_memory.len(), 2, "checkpoint + progress = two points");
}
#[test]
fn visit_archive_points_rejects_a_bad_preamble() {
assert!(try_collect_streamed(b"not an archive at all").is_err());
}
#[test]
fn visit_archive_points_is_lenient_about_a_torn_tail() {
let records = vec![
header(),
RunRecord::ContextCheckpoint {
snapshot: snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
at: 10,
},
];
let mut bytes = framed(&records);
bytes.extend_from_slice(&1000u64.to_be_bytes());
bytes.extend_from_slice(b"partial");
assert_eq!(collect_streamed(&bytes).len(), 1, "points before the tear");
}
#[test]
fn visit_archive_points_visits_nothing_without_a_header() {
let records = vec![RunRecord::ContextCheckpoint {
snapshot: snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
at: 10,
}];
assert!(collect_streamed(&framed(&records)).is_empty());
assert!(collect_streamed(&framed(&[])).is_empty());
}
#[test]
fn visit_archive_points_stops_on_break() {
let records = vec![
header(),
RunRecord::ContextCheckpoint {
snapshot: snapshot("s1", vec![region("conv", vec![entry("a", 1)])]),
at: 10,
},
RunRecord::ContextCheckpoint {
snapshot: snapshot("s1", vec![region("conv", vec![entry("b", 2)])]),
at: 11,
},
];
let bytes = framed(&records);
let mut seen = 0;
visit_archive_points(&mut &bytes[..], &mut |_| {
seen += 1;
ControlFlow::Break(())
})
.unwrap();
assert_eq!(seen, 1);
}
fn assert_digest_matches_full_diff(a: &ContextSnapshot, b: &ContextSnapshot) {
let via_digest = diff_context_digest(&digest_context(a), b);
assert_eq!(via_digest, diff_context(a, b));
let mut base = a.clone();
apply_delta(&mut base, &via_digest);
assert_eq!(&base, b);
}
#[test]
fn the_folded_fingerprint_follows_the_digest_it_is_folded_from() {
let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
let same = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
let grown = snapshot(
"s1",
vec![region("conv", vec![entry("hi", 1), entry("there", 2)])],
);
let renamed = snapshot("s1", vec![region("plan", vec![entry("hi", 1)])]);
let print = |snap: &ContextSnapshot| digest_context(snap).fingerprint();
assert_eq!(print(&a).len(), 16);
assert_eq!(print(&a), print(&same));
assert_ne!(print(&a), print(&grown));
assert_ne!(print(&a), print(&renamed));
assert_eq!(ContextDigest::default().fingerprint().len(), 16);
}
#[test]
fn digest_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_digest(&digest_context(&a), &b);
assert_eq!(region_delta_kind(&delta.regions[0]), "append");
assert_digest_matches_full_diff(&a, &b);
}
#[test]
fn digest_diff_new_cleared_removed_and_rewritten_regions() {
let a = snapshot(
"s1",
vec![
region("conv", vec![entry("hi", 1)]),
region("gone", vec![entry("bye", 1)]),
region("wiped", vec![entry("w", 1)]),
region("rewritten", vec![entry("old", 1)]),
],
);
let b = snapshot(
"s1",
vec![
region("conv", vec![entry("hi", 1)]),
region("wiped", vec![]),
region("rewritten", vec![entry("new", 1)]),
region("fresh", vec![entry("f", 2)]),
],
);
let delta = diff_context_digest(&digest_context(&a), &b);
let kinds: Vec<_> = delta.regions.iter().map(region_delta_kind).collect();
assert_eq!(kinds, vec!["clear", "set", "set", "remove"]);
assert_digest_matches_full_diff(&a, &b);
}
#[test]
fn digest_diff_unchanged_region_emits_nothing() {
let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
let delta = diff_context_digest(&digest_context(&a), &a.clone());
assert!(delta.regions.is_empty());
assert_digest_matches_full_diff(&a, &a.clone());
}
#[test]
fn digest_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_digest(&digest_context(&a), &b);
assert_eq!(region_delta_kind(&delta.regions[0]), "set");
assert_digest_matches_full_diff(&a, &b);
}
#[test]
fn digest_diff_token_recount_is_an_empty_append() {
let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
let mut recounted = region("conv", vec![entry("hi", 1)]);
recounted.current_tokens = 42;
let b = snapshot("s1", vec![recounted]);
let delta = diff_context_digest(&digest_context(&a), &b);
assert_eq!(region_delta_kind(&delta.regions[0]), "append");
assert_digest_matches_full_diff(&a, &b);
}
#[test]
fn entry_digest_covers_every_field() {
let base = entry("text", 1);
let variants = [
entry("other", 1),
entry("text", 2),
RegionEntrySnapshot {
key: Some("k".to_string()),
..entry("text", 1)
},
RegionEntrySnapshot {
metadata: Some(serde_json::json!({"a": 1})),
..entry("text", 1)
},
RegionEntrySnapshot {
kind: crate::region::EntryKind::ToolResult {
tool_call_id: "c1".to_string(),
tool_name: "shell".to_string(),
is_error: false,
},
..entry("text", 1)
},
];
let base_hash = entry_digest(&base);
for variant in &variants {
assert_ne!(
entry_digest(variant),
base_hash,
"field change must change the digest: {variant:?}"
);
}
assert_eq!(entry_digest(&base), entry_digest(&entry("text", 1)));
}
#[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::InferenceUsage {
kind: InferenceKind::Compaction,
stage: "plan".to_string(),
iteration: 2,
provider: "anthropic".to_string(),
model: "claude-sonnet-5".to_string(),
prompt_tokens: 7000,
completion_tokens: 70,
cached_tokens: 12,
cache_write_tokens: 34,
cost_usd: None,
cost_reported_by_provider: None,
at: 102,
},
RunRecord::InferenceAttempt(AttemptRecord {
id: "a0001".to_string(),
stage: "plan".to_string(),
attempt: 1,
provider: "anthropic".to_string(),
model: "claude-sonnet-5".to_string(),
outcome: AttemptOutcome::Failed {
kind: "server-error".to_string(),
transient: true,
capacity: false,
next: Retry::SameModel,
},
finish_reason: String::new(),
stopped_for: None,
duration_ms: 1_200,
backoff_ms: 0,
digest: RequestDigest {
system_hash: 99,
messages: 4,
tools: 1,
max_tokens: 1024,
temperature: 0.7,
},
model_input: Some(ModelInput {
capture_status: CaptureStatus::Retained,
request: Some(serde_json::json!({ "model": "claude-sonnet-5" })),
bytes: 31,
source_context_digest: "0123456789abcdef".to_string(),
parameters: [("temperature".to_string(), serde_json::json!(0.7))]
.into_iter()
.collect(),
tool_catalog_version: "fedcba9876543210".to_string(),
assembly_version: "1".to_string(),
}),
at: 102,
}),
RunRecord::InferenceFailover(FailoverRecord {
stage: "plan".to_string(),
iteration: 2,
from_provider: "anthropic".to_string(),
from_model: "claude-sonnet-5".to_string(),
to_provider: "openai".to_string(),
to_model: "gpt-5.5".to_string(),
reason: "unreachable".to_string(),
kind: "timeout".to_string(),
at: 102,
}),
RunRecord::ToolBatch {
calls: vec![ToolCallRecord {
execution_id: String::new(),
id: "c1".to_string(),
name: "read_file".to_string(),
arguments: "{}".to_string(),
result: None,
thought_signature: Some("sig".to_string()),
}],
at: 103,
stage_index: 0,
iteration: 0,
visit_id: String::new(),
requested_by: String::new(),
response: "reading".to_string(),
},
RunRecord::ToolCallDone {
execution_id: String::new(),
outcome: None,
iteration: 0,
call_id: "c1".to_string(),
result: "body".to_string().into(),
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, RUN_ARCHIVE_VERSION).unwrap();
assert_eq!(
read_archive_start(&mut buf.as_slice()).unwrap(),
RUN_ARCHIVE_VERSION
);
}
#[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 an_absurd_frame_length_is_an_error_not_an_allocation() {
let mut buf = Vec::new();
write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
write_record(&mut buf, &header()).unwrap();
buf.extend_from_slice(&u64::MAX.to_be_bytes());
let err = read_archive(&mut buf.as_slice())
.expect_err("the strict reader must refuse an impossible frame");
assert_eq!(err.kind(), io::ErrorKind::InvalidData, "{err}");
let (_, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
assert_eq!(records, vec![header()]);
}
#[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, 2);
assert_eq!(folded.inference_usage.len(), 1);
assert_eq!(folded.tool_call_count, 1);
assert_eq!(folded.attempts.len(), 1);
assert_eq!(folded.attempts[0].attempt, 1);
assert_eq!(folded.attempts[0].digest.messages, 4);
assert_eq!(folded.failovers.len(), 1);
assert_eq!(folded.failovers[0].from_provider, "anthropic");
assert_eq!(folded.failovers[0].to_provider, "openai");
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");
let pending = folded.pending_batch.expect("batch never applied");
assert_eq!(pending.calls[0].result.as_deref(), Some("body"));
}
#[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 fold_carries_a_submitted_final_output_through_progress() {
let mut answered = meta();
answered.final_output = Some(
crate::output::FinalOutput::new(
"renamed two helpers",
Some("markdown".to_string()),
"summary".to_string(),
9,
)
.descriptor(),
);
answered.output_request = Some(crate::output::OutputSpec {
format: Some("a2ui".to_string()),
..Default::default()
});
let records = vec![
header(),
RunRecord::Progress {
meta: Box::new(answered),
delta: ContextDelta {
stage_name: "summary".to_string(),
total_tokens: 0,
max_tokens: 10_000,
regions: vec![],
},
at: 2,
},
];
let folded = fold(&records).unwrap();
let output = folded.meta.final_output.expect("the answer folded through");
assert_eq!(output.bytes, "renamed two helpers".len());
assert_eq!(output.stage, "summary");
assert_eq!(
folded.meta.output_request.and_then(|s| s.format).as_deref(),
Some("a2ui")
);
}
fn call(id: &str, result: Option<&str>) -> ToolCallRecord {
ToolCallRecord {
execution_id: String::new(),
id: id.to_string(),
name: "shell".to_string(),
arguments: "{}".to_string(),
result: result.map(Into::into),
thought_signature: None,
}
}
fn batch(iteration: usize, calls: Vec<ToolCallRecord>) -> RunRecord {
RunRecord::ToolBatch {
calls,
at: 10,
stage_index: 0,
iteration,
visit_id: String::new(),
requested_by: String::new(),
response: "running tools".to_string(),
}
}
fn turn_entry(call_ids: &[&str]) -> RegionEntrySnapshot {
let mut e = entry("turn", 1);
e.kind = crate::region::EntryKind::AssistantTurn {
tool_calls: call_ids
.iter()
.map(|id| crate::region::SerializedToolCall {
id: id.to_string(),
name: "shell".to_string(),
arguments: serde_json::Value::Null,
thought_signature: None,
})
.collect(),
};
e
}
#[test]
fn fold_surfaces_a_pending_batch_with_merged_results() {
let records = vec![
header(),
batch(
0,
vec![
call("c1", None),
call("c2", Some("inline")),
call("c3", None),
],
),
RunRecord::ToolCallDone {
execution_id: String::new(),
outcome: None,
iteration: 0,
call_id: "c1".to_string(),
result: "ran".to_string().into(),
at: 11,
},
];
let folded = fold(&records).unwrap();
let pending = folded.pending_batch.expect("batch is pending");
assert_eq!(pending.iteration, 0);
assert_eq!(pending.response, "running tools");
assert_eq!(pending.calls[0].result.as_deref(), Some("ran"));
assert_eq!(pending.calls[1].result.as_deref(), Some("inline"));
assert_eq!(pending.calls[2].result, None);
assert_eq!(folded.tool_call_count, 3);
}
#[test]
fn folding_gathers_why_each_region_moved() {
use crate::ContextCause;
let moved = |region: &str, before: usize, after: usize, added| RegionCommit {
region: region.to_string(),
digest_before: format!("rg1-{before:032x}"),
digest_after: format!("rg1-{after:032x}"),
tokens_before: before * 10,
tokens_after: after * 10,
entries_before: before,
entries_after: after,
entries_added: added,
};
let committed = |cause, regions, at| RunRecord::ContextTransaction {
revision_before: format!("cw1-{at:032x}"),
revision_after: format!("cw1-{:032x}", at + 1),
cause,
regions,
execution_id: String::new(),
at,
};
let records = vec![
header(),
committed(ContextCause::Seed, vec![moved("plan", 0, 1, 1)], 20),
committed(
ContextCause::ToolResult,
vec![moved("conversation", 0, 2, 2)],
21,
),
committed(
ContextCause::Compaction,
vec![moved("plan", 6, 0, 0), moved("plan_history", 0, 1, 1)],
22,
),
];
let folded = fold(&records).expect("a journal with a header folds");
assert_eq!(folded.context_changes.len(), 3);
assert_eq!(folded.context_changes[0].regions[0].region, "plan");
assert_eq!(folded.context_changes[0].cause, ContextCause::Seed);
assert_eq!(folded.context_changes[0].regions[0].entries_added, 1);
assert_eq!(folded.context_changes[0].regions[0].token_delta, 10);
assert_eq!(folded.context_changes[0].at, 20);
assert_eq!(
folded.context_changes[0].revision_after.as_deref(),
Some(format!("cw1-{:032x}", 21).as_str()),
"a transaction names the window it produced"
);
assert_eq!(folded.context_changes[1].cause, ContextCause::ToolResult);
let compacted = &folded.context_changes[2];
assert_eq!(compacted.cause, ContextCause::Compaction);
assert_eq!(compacted.regions.len(), 2, "both halves, in one record");
assert_eq!(compacted.regions[0].entries_removed, 6);
assert_eq!(
compacted.regions[0].token_delta, -60,
"a region that shrank reads as a loss, not as an absence"
);
assert_eq!(compacted.regions[1].region, "plan_history");
assert_eq!(folded.inference_count, 0);
assert_eq!(folded.tool_call_count, 0);
}
#[test]
fn folding_gathers_every_question_the_run_asked() {
use crate::interaction::{ApprovalScope, InteractionKind, Settlement};
let asked = |id: &str, settlement: Settlement, at: i64| RunRecord::Interaction {
request_id: id.to_string(),
kind: InteractionKind::ToolApproval,
tool: Some("shell".to_string()),
prompt: format!("Run {id}?"),
stage: "plan".to_string(),
settlement,
asked_at: at,
at: at + 1,
};
let records = vec![
header(),
asked(
"approve-1",
Settlement::Answered {
approved: Some(true),
scope: Some(ApprovalScope::Stage),
choice: Some(0),
text: None,
feedback: None,
},
20,
),
asked("approve-2", Settlement::TimedOut, 30),
asked("approve-3", Settlement::Cancelled, 40),
];
let folded = fold(&records).expect("a journal with a header folds");
assert_eq!(folded.interactions.len(), 3);
assert_eq!(folded.interactions[0].request_id, "approve-1");
assert_eq!(folded.interactions[0].tool.as_deref(), Some("shell"));
assert_eq!(folded.interactions[0].stage, "plan");
assert_eq!(folded.interactions[0].prompt, "Run approve-1?");
assert_eq!(folded.interactions[0].asked_at, 20);
assert_eq!(folded.interactions[0].at, 21);
assert!(matches!(
folded.interactions[0].settlement,
Settlement::Answered {
approved: Some(true),
scope: Some(ApprovalScope::Stage),
..
}
));
assert_eq!(folded.interactions[1].settlement, Settlement::TimedOut);
assert_eq!(folded.interactions[2].settlement, Settlement::Cancelled);
assert_eq!(folded.tool_call_count, 0);
}
#[test]
fn folding_reads_a_journal_of_single_region_changes() {
let changed = |region: &str, cause: ContextCause, added, removed, delta, at| {
RunRecord::ContextChange {
region: region.to_string(),
cause,
entries_added: added,
entries_removed: removed,
token_delta: delta,
at,
}
};
let records = vec![
header(),
changed("plan", ContextCause::Seed, 1, 0, 40, 10),
changed("plan", ContextCause::ContextTool, 1, 1, -5, 20),
changed("plan", ContextCause::Compaction, 0, 3, -120, 30),
];
let folded = fold(&records).expect("a journal with a header folds");
let causes: Vec<ContextCause> = folded.context_changes.iter().map(|c| c.cause).collect();
assert_eq!(
causes,
vec![
ContextCause::Seed,
ContextCause::ContextTool,
ContextCause::Compaction
]
);
let compaction = &folded.context_changes[2];
assert_eq!(compaction.regions.len(), 1, "one region is all it recorded");
assert_eq!(compaction.regions[0].region, "plan");
assert_eq!(compaction.regions[0].entries_added, 0);
assert_eq!(compaction.regions[0].entries_removed, 3);
assert_eq!(compaction.regions[0].token_delta, -120);
assert_eq!(compaction.at, 30);
assert_eq!(compaction.revision_before, None);
assert_eq!(compaction.revision_after, None);
assert_eq!(compaction.regions[0].digest_after, None);
assert_eq!(compaction.regions[0].tokens_after, None);
assert!(replay_points(&records).is_empty());
}
#[test]
fn fold_keeps_only_the_latest_batch_and_ignores_stale_done_records() {
let mut advanced = meta();
advanced.iteration = 1;
let records = vec![
header(),
batch(0, vec![call("c1", None)]),
RunRecord::Progress {
meta: Box::new(advanced),
delta: ContextDelta {
stage_name: "plan".to_string(),
total_tokens: 0,
max_tokens: 10_000,
regions: vec![],
},
at: 11,
},
batch(1, vec![call("c2", None)]),
RunRecord::ToolCallDone {
execution_id: String::new(),
outcome: None,
iteration: 0,
call_id: "c1".to_string(),
result: "stale".to_string().into(),
at: 12,
},
RunRecord::ToolCallDone {
execution_id: String::new(),
outcome: None,
iteration: 1,
call_id: "unknown".to_string(),
result: "nowhere to land".to_string().into(),
at: 13,
},
];
let folded = fold(&records).unwrap();
let pending = folded.pending_batch.expect("latest batch is pending");
assert_eq!(pending.iteration, 1);
assert_eq!(pending.calls.len(), 1);
assert_eq!(pending.calls[0].id, "c2");
assert_eq!(pending.calls[0].result, None, "stale/unknown dones ignored");
}
#[test]
fn folding_passes_over_the_files_an_execution_produced() {
let records = vec![
header(),
RunRecord::ArtifactsProduced {
execution_id: "x1".to_string(),
artifacts: vec![crate::output::Artifact {
name: "report".to_string(),
path: "out/report.md".to_string(),
mime_type: crate::mime::MimeType::parse("text/markdown").expect("a type"),
size: 12,
sha256: "beef".to_string(),
}],
at: 30,
},
];
let folded = fold(&records).expect("a journal with a header folds");
assert_eq!(folded.tool_call_count, 0, "a file is not a call");
assert!(folded.context_changes.is_empty());
assert!(folded.pending_batch.is_none());
}
#[test]
fn fold_does_not_make_a_batch_it_resolved_itself_pending() {
let records = vec![header(), batch(0, vec![call("c1", Some("wrote the plan"))])];
assert_eq!(fold(&records).unwrap().pending_batch, None);
assert_eq!(
fold(&records).unwrap().tool_call_count,
1,
"it is still a call the run made"
);
let after = vec![
header(),
batch(0, vec![call("c1", None)]),
batch(0, vec![call("c2", Some("wrote the plan"))]),
];
assert_eq!(fold(&after).unwrap().pending_batch, None);
}
#[test]
fn fold_clears_a_batch_once_the_iteration_moves_on() {
let mut advanced = meta();
advanced.iteration = 1;
let records = vec![
header(),
batch(0, vec![call("c1", None)]),
RunRecord::ToolCallDone {
execution_id: String::new(),
outcome: None,
iteration: 0,
call_id: "c1".to_string(),
result: "done".to_string().into(),
at: 10,
},
RunRecord::Progress {
meta: Box::new(advanced),
delta: ContextDelta {
stage_name: "plan".to_string(),
total_tokens: 0,
max_tokens: 10_000,
regions: vec![],
},
at: 11,
},
];
assert_eq!(fold(&records).unwrap().pending_batch, None);
}
#[test]
fn fold_clears_a_batch_whose_turn_already_landed_in_the_window() {
let records = vec![
header(),
batch(0, vec![call("c1", None)]),
RunRecord::ToolCallDone {
execution_id: String::new(),
outcome: None,
iteration: 0,
call_id: "c1".to_string(),
result: "done".to_string().into(),
at: 10,
},
RunRecord::ContextCheckpoint {
snapshot: snapshot("plan", vec![region("conv", vec![turn_entry(&["c1"])])]),
at: 11,
},
];
assert_eq!(fold(&records).unwrap().pending_batch, None);
}
#[test]
fn context_contains_batch_matches_only_the_batch_turn() {
let pending = PendingToolBatch {
stage_index: 0,
iteration: 0,
response: String::new(),
calls: vec![call("c1", None)],
};
let other = snapshot("plan", vec![region("conv", vec![turn_entry(&["zz"])])]);
assert!(!context_contains_batch(&other, &pending));
let own = snapshot(
"plan",
vec![region("conv", vec![turn_entry(&["c1", "c2"])])],
);
assert!(context_contains_batch(&own, &pending));
let empty = PendingToolBatch {
calls: vec![],
..pending
};
assert!(!context_contains_batch(&own, &empty));
}
#[test]
fn old_shape_tool_batch_json_still_parses() {
let json = br#"{"ToolBatch":{"calls":[{"id":"c1","name":"shell","arguments":"{}","result":"ok"}],"at":9}}"#;
let mut buf = Vec::new();
buf.extend_from_slice(&(json.len() as u64).to_be_bytes());
buf.extend_from_slice(json);
let record = read_record(&mut buf.as_slice()).unwrap().unwrap();
assert_eq!(
record,
RunRecord::ToolBatch {
calls: vec![call("c1", Some("ok"))],
at: 9,
stage_index: 0,
iteration: 0,
visit_id: String::new(),
requested_by: String::new(),
response: String::new(),
}
);
}
fn three_point_records() -> Vec<RunRecord> {
let mut running = meta();
running.status = RunStatus::Running;
vec![
header(),
RunRecord::ContextCheckpoint {
snapshot: snapshot("plan", vec![region("conv", vec![entry("first", 1)])]),
at: 10,
},
RunRecord::ContextDiff {
delta: ContextDelta {
stage_name: "plan".to_string(),
total_tokens: 2,
max_tokens: 10_000,
regions: vec![RegionDelta::Append {
name: "conv".to_string(),
entries: vec![entry("second", 1)],
current_tokens: 2,
}],
},
at: 20,
},
RunRecord::Progress {
meta: Box::new(running),
delta: ContextDelta {
stage_name: "code".to_string(),
total_tokens: 3,
max_tokens: 10_000,
regions: vec![RegionDelta::Append {
name: "conv".to_string(),
entries: vec![entry("third", 1)],
current_tokens: 3,
}],
},
at: 30,
},
]
}
#[test]
fn visit_points_indexes_points_in_order_and_carries_the_running_window() {
let records = three_point_records();
let mut seen: Vec<(usize, i64, usize)> = Vec::new();
visit_points(&records, &mut |point| {
seen.push((
point.index,
point.at,
point.context.regions[0].entries.len(),
));
ControlFlow::Continue(())
});
assert_eq!(seen, vec![(0, 10, 1), (1, 20, 2), (2, 30, 3)]);
}
#[test]
fn visit_points_stops_at_the_first_break() {
let records = three_point_records();
let mut visits = 0;
visit_points(&records, &mut |point| {
visits += 1;
if point.index == 1 {
ControlFlow::Break(())
} else {
ControlFlow::Continue(())
}
});
assert_eq!(
visits, 2,
"stopped at the breaking point, did not run the third"
);
}
#[test]
fn visit_points_without_a_header_visits_nothing() {
let mut visits = 0;
{
let mut count = |_: PointRef<'_>| {
visits += 1;
ControlFlow::Continue(())
};
visit_points(&three_point_records(), &mut count);
visit_points(&[], &mut count);
visit_points(
&[RunRecord::ContextCheckpoint {
snapshot: snapshot("plan", vec![]),
at: 1,
}],
&mut count,
);
}
assert_eq!(visits, 3, "only the well-formed journal produced points");
}
#[test]
fn visit_points_and_replay_points_agree() {
for records in [
three_point_records(),
vec![header()],
vec![],
vec![RunRecord::Message {
message: MessageRecord {
role: "user".to_string(),
content: "x".to_string(),
},
at: 1,
}],
] {
let collected: Vec<RunPoint> = {
let mut out = Vec::new();
visit_points(&records, &mut |point| {
out.push(RunPoint {
meta: point.meta.clone(),
context: point.context.clone(),
at: point.at,
});
ControlFlow::Continue(())
});
out
};
assert_eq!(collected, replay_points(&records));
}
}
#[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,
},
batch(0, vec![call("c1", None)]),
RunRecord::ToolCallDone {
execution_id: String::new(),
outcome: None,
iteration: 0,
call_id: "c1".to_string(),
result: "ran".to_string().into(),
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);
}
#[test]
fn every_inference_kind_has_a_distinct_label_and_serialized_name() {
let all = [
(InferenceKind::Stage, "stage"),
(InferenceKind::Compaction, "compaction"),
(InferenceKind::Title, "title"),
(InferenceKind::Routing, "routing"),
];
for (kind, label) in all {
assert_eq!(kind.label(), label);
assert_eq!(serde_json::to_value(kind).unwrap(), label);
}
let labels: std::collections::HashSet<_> = all.iter().map(|(k, _)| k.label()).collect();
assert_eq!(labels.len(), all.len(), "labels must not collide");
}
#[test]
fn only_a_stage_turn_counts_as_stage_work() {
assert!(InferenceKind::Stage.is_stage_work());
for kind in [
InferenceKind::Compaction,
InferenceKind::Title,
InferenceKind::Routing,
] {
assert!(
!kind.is_stage_work(),
"{kind:?} is machinery, not stage work"
);
}
}
#[test]
fn a_usage_record_without_a_kind_reads_back_as_stage_work() {
let json = serde_json::json!({
"InferenceUsage": {
"stage": "plan",
"iteration": 1,
"provider": "anthropic",
"model": "claude-sonnet-5",
"prompt_tokens": 10,
"completion_tokens": 2,
"cached_tokens": 0,
"cache_write_tokens": 0,
"at": 5,
}
});
let record: RunRecord = serde_json::from_value(json).unwrap();
assert_eq!(
record,
RunRecord::InferenceUsage {
kind: InferenceKind::Stage,
stage: "plan".to_string(),
iteration: 1,
provider: "anthropic".to_string(),
model: "claude-sonnet-5".to_string(),
prompt_tokens: 10,
completion_tokens: 2,
cached_tokens: 0,
cache_write_tokens: 0,
cost_usd: None,
cost_reported_by_provider: None,
at: 5,
}
);
}
#[test]
fn folding_keeps_each_call_separate_instead_of_summing_them() {
let usage = |kind, prompt, at| RunRecord::InferenceUsage {
kind,
stage: "plan".to_string(),
iteration: 1,
provider: "anthropic".to_string(),
model: "claude-sonnet-5".to_string(),
prompt_tokens: prompt,
completion_tokens: 1,
cached_tokens: 0,
cache_write_tokens: 0,
cost_usd: None,
cost_reported_by_provider: None,
at,
};
let folded = fold(&[
header(),
usage(InferenceKind::Compaction, 7000, 1),
usage(InferenceKind::Stage, 21_000, 2),
])
.unwrap();
assert_eq!(folded.inference_count, 2);
let seen: Vec<_> = folded
.inference_usage
.iter()
.map(|u| (u.kind, u.prompt_tokens))
.collect();
assert_eq!(
seen,
vec![
(InferenceKind::Compaction, 7000),
(InferenceKind::Stage, 21_000)
]
);
assert!(
folded
.inference_usage
.iter()
.all(|u| u.prompt_tokens < 32_000),
"no single call exceeded the window, and the journal can now prove it"
);
}
#[test]
fn an_attempts_finish_reason_is_kept_and_an_old_record_reads_without_one() {
let record = |finish_reason: &str, stopped_for: Option<&str>| AttemptRecord {
id: "a1".to_string(),
stage: "plan".to_string(),
attempt: 1,
provider: "anthropic".to_string(),
model: "claude".to_string(),
outcome: AttemptOutcome::Succeeded,
finish_reason: finish_reason.to_string(),
stopped_for: stopped_for.map(str::to_string),
duration_ms: 10,
backoff_ms: 0,
digest: RequestDigest {
system_hash: 1,
messages: 1,
tools: 0,
max_tokens: 10,
temperature: 0.0,
},
model_input: None,
at: 1,
};
let unknown = record("unknown", Some("content_filter"));
let json = serde_json::to_string(&unknown).unwrap();
assert!(json.contains("\"finish_reason\":\"unknown\""), "{json}");
assert!(
json.contains("\"stopped_for\":\"content_filter\""),
"{json}"
);
let back: AttemptRecord = serde_json::from_str(&json).unwrap();
assert_eq!(back, unknown);
let plain = record("", None);
let json = serde_json::to_string(&plain).unwrap();
assert!(!json.contains("finish_reason"), "{json}");
assert!(!json.contains("stopped_for"), "{json}");
let old: AttemptRecord = serde_json::from_str(
r#"{"stage":"plan","attempt":1,"provider":"anthropic","model":"claude",
"outcome":"succeeded","duration_ms":10,"backoff_ms":0,
"digest":{"system_hash":1,"messages":1,"tools":0,"max_tokens":10,"temperature":0.0},
"at":1}"#,
)
.unwrap();
assert_eq!(old.finish_reason, "");
assert_eq!(old.stopped_for, None);
}
#[test]
fn an_attempts_outcomes_and_follow_ups_keep_their_wire_names() {
for (outcome, json) in [
(AttemptOutcome::Succeeded, "\"succeeded\"".to_string()),
(
AttemptOutcome::Failed {
kind: "timeout".to_string(),
transient: true,
capacity: false,
next: Retry::Reported,
},
"{\"failed\":{\"kind\":\"timeout\",\"transient\":true,\"capacity\":false,\
\"next\":\"reported\"}}"
.to_string(),
),
] {
let wire = serde_json::to_string(&outcome).expect("an outcome serializes");
assert_eq!(wire, json);
assert_eq!(
serde_json::from_str::<AttemptOutcome>(&wire).expect("and reads back"),
outcome
);
}
for (next, json) in [
(Retry::Reported, "\"reported\""),
(Retry::SameModel, "\"same_model\""),
(Retry::RenewedFiles, "\"renewed_files\""),
] {
assert_eq!(serde_json::to_string(&next).expect("serializes"), json);
assert_eq!(
serde_json::from_str::<Retry>(json).expect("and reads back"),
next
);
}
}
#[test]
fn both_inference_record_kinds_count_as_one_call_each() {
let records = all_record_kinds();
let folded = fold(&records).unwrap();
let written = records
.iter()
.filter(|r| {
matches!(
r,
RunRecord::Inference { .. } | RunRecord::InferenceUsage { .. }
)
})
.count();
assert_eq!(folded.inference_count, written);
assert_eq!(folded.inference_usage.len(), 1);
}
#[test]
fn probe_replay_matches_every_step() {
fn region(name: &str, entries: &[(&str, usize)]) -> RegionSnapshot {
RegionSnapshot {
name: name.to_string(),
kind: "temporary".to_string(),
current_tokens: entries.iter().map(|(_, t)| *t).sum(),
max_tokens: 1000,
entries: entries
.iter()
.map(|(c, t)| RegionEntrySnapshot {
content: (*c).into(),
tokens: *t,
key: None,
kind: crate::region::EntryKind::Text,
metadata: None,
taint: crate::taint::TaintLevel::Public,
reasoning: None,
})
.collect(),
description: None,
}
}
fn snap(entries: &[(&str, usize)]) -> ContextSnapshot {
let r = region("logs", entries);
ContextSnapshot {
stage_name: "s".to_string(),
total_tokens: r.current_tokens,
max_tokens: 1000,
regions: vec![r],
}
}
let steps: Vec<ContextSnapshot> = vec![
snap(&[]),
snap(&[("a", 10)]),
snap(&[("a", 10), ("b", 20)]),
snap(&[("a", 10), ("b", 20), ("c", 30)]),
snap(&[("b", 20), ("c", 30)]),
snap(&[("c", 30), ("d", 40)]),
snap(&[("d", 40)]),
snap(&[("d", 40), ("x", 5)]),
snap(&[("x", 5), ("x", 5)]),
snap(&[("x", 5)]),
snap(&[]),
];
let mut base = steps[0].clone();
let mut digest = digest_context(&steps[0]);
for (i, next) in steps.iter().enumerate().skip(1) {
let delta = diff_context_digest(&digest, next);
apply_delta(&mut base, &delta);
assert_eq!(
base, *next,
"step {i}: replay drifted from the live state\n delta was {:?}",
delta.regions
);
digest = digest_context(next);
}
}
fn trailing_message() -> RunRecord {
RunRecord::Message {
message: MessageRecord {
role: "user".to_string(),
content: "after the unknown".to_string(),
},
at: 9,
}
}
fn archive_with_an_unknown_record() -> (Vec<u8>, usize) {
let mut buf = Vec::new();
write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
write_record(&mut buf, &header()).unwrap();
let payload =
serde_json::to_vec(&serde_json::json!({ "SomethingNew": { "whatever": 1 } })).unwrap();
buf.extend_from_slice(&(payload.len() as u64).to_be_bytes());
buf.extend_from_slice(&payload);
write_record(&mut buf, &trailing_message()).unwrap();
(buf, payload.len())
}
fn record_kind(record: &RunRecord) -> String {
serde_json::to_value(record)
.expect("a RunRecord always serializes")
.as_object()
.expect("externally tagged, so an object")
.keys()
.next()
.expect("with exactly one key")
.clone()
}
#[test]
fn an_unknown_record_kind_is_stepped_over_not_treated_as_the_end() {
let (buf, _) = archive_with_an_unknown_record();
let (version, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
assert_eq!(version, RUN_ARCHIVE_VERSION);
let kinds: Vec<String> = records.iter().map(record_kind).collect();
assert_eq!(
kinds,
vec!["Header".to_string(), "Message".to_string()],
"the header, and the readable record after the gap"
);
assert_eq!(records[1], trailing_message(), "intact, not just present");
}
#[test]
fn the_streaming_reader_also_steps_over_an_unknown_record() {
let (mut buf, _) = archive_with_an_unknown_record();
write_record(
&mut buf,
&RunRecord::ContextCheckpoint {
snapshot: ContextSnapshot {
stage_name: "s".to_string(),
total_tokens: 1,
max_tokens: 10,
regions: vec![],
},
at: 11,
},
)
.unwrap();
let mut points = 0usize;
visit_archive_points(&mut buf.as_slice(), &mut |_point| {
points += 1;
ControlFlow::Continue(())
})
.expect("a valid preamble");
assert_eq!(points, 1, "the walk got past the unknown frame");
}
#[test]
fn a_frame_reports_whether_its_payload_was_readable() {
let (buf, unknown_bytes) = archive_with_an_unknown_record();
let mut r = buf.as_slice();
read_archive_start(&mut r).unwrap();
let mut frames = Vec::new();
while let Some(frame) = read_frame(&mut r).expect("no torn frames here") {
frames.push(frame);
}
assert_eq!(
frames,
vec![
Frame::Record(Box::new(header())),
Frame::Unreadable {
bytes: unknown_bytes
},
Frame::Record(Box::new(trailing_message())),
],
"one frame per record, with the unreadable one accounted for rather than ending the read"
);
}
#[test]
fn a_torn_frame_still_ends_the_read() {
let mut buf = Vec::new();
write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
write_record(&mut buf, &header()).unwrap();
buf.extend_from_slice(&999u64.to_be_bytes());
buf.extend_from_slice(b"not enough");
let (_, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
assert_eq!(records.len(), 1, "everything intact before the tear");
}
#[test]
fn an_archive_from_a_newer_format_is_refused_with_both_versions_named() {
let mut buf = Vec::new();
write_archive_start(&mut buf, RUN_ARCHIVE_VERSION + 1).unwrap();
write_record(&mut buf, &header()).unwrap();
let err = read_archive_lenient(&mut buf.as_slice()).unwrap_err();
let message = err.to_string();
assert!(
message.contains(&(RUN_ARCHIVE_VERSION + 1).to_string()),
"{message}"
);
assert!(message.contains("upgrade leviath"), "{message}");
}
#[test]
fn an_archive_from_an_older_format_still_reads() {
let mut buf = Vec::new();
write_archive_start(&mut buf, 0).unwrap();
write_record(&mut buf, &header()).unwrap();
let (version, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
assert_eq!(version, 0);
assert_eq!(records.len(), 1);
}
}