1use std::io::{self, Read};
42use std::ops::ControlFlow;
43
44use serde::{Deserialize, Serialize};
45
46use crate::run_meta::{ContextSnapshot, RegionEntrySnapshot, RegionSnapshot, RunMeta, RunStatus};
47
48#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
55pub struct RunIdentity {
56 pub run_id: String,
58 pub machine_id: String,
60 pub world_id: String,
62 pub created_at: i64,
64}
65
66#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
68pub struct MessageRecord {
69 pub role: String,
71 pub content: String,
73}
74
75#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
77pub struct ToolCallRecord {
78 pub id: String,
80 pub name: String,
82 pub arguments: String,
84 pub result: Option<crate::region::EntryContent>,
89 #[serde(default, skip_serializing_if = "Option::is_none")]
93 pub thought_signature: Option<String>,
94}
95
96#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
99pub struct InferenceRequestRecord {
100 pub model: String,
102 pub system: Vec<String>,
104 pub messages: Vec<MessageRecord>,
106 pub tool_names: Vec<String>,
108 pub temperature: f32,
110 pub max_tokens: usize,
112}
113
114#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
116pub struct InferenceResponseRecord {
117 pub content: String,
119 pub tool_calls: Vec<ToolCallRecord>,
121 pub prompt_tokens: usize,
123 pub completion_tokens: usize,
125 pub cached_tokens: usize,
127 pub cache_write_tokens: usize,
129}
130
131#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
133pub enum RegionDelta {
134 Set(RegionSnapshot),
137 Append {
140 name: String,
142 entries: Vec<RegionEntrySnapshot>,
144 current_tokens: usize,
146 },
147 Clear {
149 name: String,
151 },
152 Remove {
154 name: String,
156 },
157}
158
159#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
161pub struct ContextDelta {
162 pub stage_name: String,
164 pub total_tokens: usize,
166 pub max_tokens: usize,
168 pub regions: Vec<RegionDelta>,
170}
171
172#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
180#[serde(rename_all = "snake_case")]
181pub enum InferenceKind {
182 #[default]
186 Stage,
187 Compaction,
189 Title,
191 Routing,
193}
194
195impl InferenceKind {
196 pub fn label(&self) -> &'static str {
198 match self {
199 InferenceKind::Stage => "stage",
200 InferenceKind::Compaction => "compaction",
201 InferenceKind::Title => "title",
202 InferenceKind::Routing => "routing",
203 }
204 }
205
206 pub fn is_stage_work(&self) -> bool {
209 matches!(self, InferenceKind::Stage)
210 }
211}
212
213#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
215pub enum RunRecord {
216 Header {
218 identity: RunIdentity,
220 meta: Box<RunMeta>,
222 },
223 OwnershipChanged {
225 machine_id: String,
227 world_id: String,
229 at: i64,
231 },
232 Inference {
234 stage: String,
236 iteration: usize,
238 request: InferenceRequestRecord,
240 response: InferenceResponseRecord,
242 at: i64,
244 },
245 InferenceUsage {
260 #[serde(default)]
262 kind: InferenceKind,
263 stage: String,
266 iteration: usize,
268 provider: String,
270 model: String,
272 prompt_tokens: usize,
274 completion_tokens: usize,
276 cached_tokens: usize,
278 cache_write_tokens: usize,
280 #[serde(default, skip_serializing_if = "Option::is_none")]
291 cost_usd: Option<f64>,
292 #[serde(default, skip_serializing_if = "Option::is_none")]
295 cost_reported_by_provider: Option<bool>,
296 at: i64,
298 },
299 ToolBatch {
307 calls: Vec<ToolCallRecord>,
309 at: i64,
311 #[serde(default)]
313 stage_index: usize,
314 #[serde(default)]
317 iteration: usize,
318 #[serde(default)]
320 response: String,
321 },
322 ToolCallDone {
324 iteration: usize,
326 call_id: String,
328 result: crate::region::EntryContent,
330 at: i64,
332 },
333 ContextCheckpoint {
335 snapshot: ContextSnapshot,
337 at: i64,
339 },
340 ContextDiff {
342 delta: ContextDelta,
344 at: i64,
346 },
347 Message {
349 message: MessageRecord,
351 at: i64,
353 },
354 StatusChanged {
356 status: RunStatus,
358 at: i64,
360 },
361 Checkpoint {
364 meta: Box<RunMeta>,
366 context: ContextSnapshot,
368 at: i64,
370 },
371 Progress {
376 meta: Box<RunMeta>,
378 delta: ContextDelta,
380 at: i64,
382 },
383}
384
385fn is_prefix(prev: &[RegionEntrySnapshot], next: &[RegionEntrySnapshot]) -> bool {
389 prev.len() <= next.len() && next[..prev.len()] == *prev
390}
391
392trait PriorRegion {
400 fn name(&self) -> &str;
402 fn unchanged(&self, next: &RegionSnapshot) -> bool;
404 fn appended_to(&self, next: &RegionSnapshot) -> bool;
406 fn entry_count(&self) -> usize;
408}
409
410impl PriorRegion for RegionSnapshot {
411 fn name(&self) -> &str {
412 &self.name
413 }
414
415 fn unchanged(&self, next: &RegionSnapshot) -> bool {
416 self == next
417 }
418
419 fn appended_to(&self, next: &RegionSnapshot) -> bool {
420 self.kind == next.kind
421 && self.max_tokens == next.max_tokens
422 && is_prefix(&self.entries, &next.entries)
423 }
424
425 fn entry_count(&self) -> usize {
426 self.entries.len()
427 }
428}
429
430fn diff_regions<P: PriorRegion>(prev: &[P], next: &ContextSnapshot) -> ContextDelta {
435 let mut regions = Vec::new();
436 for nr in &next.regions {
437 match prev.iter().find(|r| r.name() == nr.name) {
438 None => regions.push(RegionDelta::Set(nr.clone())),
439 Some(pr) => {
440 if pr.unchanged(nr) {
441 } else if nr.entries.is_empty() && pr.entry_count() > 0 {
443 regions.push(RegionDelta::Clear {
444 name: nr.name.clone(),
445 });
446 } else if pr.appended_to(nr) {
447 regions.push(RegionDelta::Append {
448 name: nr.name.clone(),
449 entries: nr.entries[pr.entry_count()..].to_vec(),
450 current_tokens: nr.current_tokens,
451 });
452 } else {
453 regions.push(RegionDelta::Set(nr.clone()));
454 }
455 }
456 }
457 }
458 for pr in prev {
459 if !next.regions.iter().any(|r| r.name == pr.name()) {
460 regions.push(RegionDelta::Remove {
461 name: pr.name().to_string(),
462 });
463 }
464 }
465 ContextDelta {
466 stage_name: next.stage_name.clone(),
467 total_tokens: next.total_tokens,
468 max_tokens: next.max_tokens,
469 regions,
470 }
471}
472
473pub fn diff_context(prev: &ContextSnapshot, next: &ContextSnapshot) -> ContextDelta {
477 diff_regions(&prev.regions, next)
478}
479
480#[derive(Debug, Clone, PartialEq)]
491pub struct RegionDigest {
492 pub name: String,
494 pub kind: String,
496 pub current_tokens: usize,
498 pub max_tokens: usize,
500 pub entries: Vec<u64>,
502}
503
504#[derive(Debug, Clone, PartialEq, Default)]
506pub struct ContextDigest {
507 pub regions: Vec<RegionDigest>,
509}
510
511fn entry_digest(entry: &RegionEntrySnapshot) -> u64 {
515 use std::hash::{Hash, Hasher};
516 let mut hasher = std::collections::hash_map::DefaultHasher::new();
517 entry.content.hash(&mut hasher);
518 entry.tokens.hash(&mut hasher);
519 entry.key.hash(&mut hasher);
520 serde_json::to_string(&entry.kind)
523 .expect("EntryKind always serializes")
524 .hash(&mut hasher);
525 serde_json::to_string(&entry.metadata)
526 .expect("entry metadata always serializes")
527 .hash(&mut hasher);
528 serde_json::to_string(&entry.taint)
529 .expect("taint always serializes")
530 .hash(&mut hasher);
531 hasher.finish()
532}
533
534pub fn digest_context(snapshot: &ContextSnapshot) -> ContextDigest {
536 ContextDigest {
537 regions: snapshot
538 .regions
539 .iter()
540 .map(|r| RegionDigest {
541 name: r.name.clone(),
542 kind: r.kind.clone(),
543 current_tokens: r.current_tokens,
544 max_tokens: r.max_tokens,
545 entries: r.entries.iter().map(entry_digest).collect(),
546 })
547 .collect(),
548 }
549}
550
551fn is_prefix_digest(prev: &[u64], next: &[RegionEntrySnapshot]) -> bool {
553 prev.len() <= next.len()
554 && prev
555 .iter()
556 .zip(next)
557 .all(|(hash, entry)| *hash == entry_digest(entry))
558}
559
560impl PriorRegion for RegionDigest {
561 fn name(&self) -> &str {
562 &self.name
563 }
564
565 fn unchanged(&self, next: &RegionSnapshot) -> bool {
566 self.kind == next.kind
567 && self.max_tokens == next.max_tokens
568 && self.current_tokens == next.current_tokens
569 && self.entries.len() == next.entries.len()
570 && is_prefix_digest(&self.entries, &next.entries)
571 }
572
573 fn appended_to(&self, next: &RegionSnapshot) -> bool {
574 self.kind == next.kind
575 && self.max_tokens == next.max_tokens
576 && is_prefix_digest(&self.entries, &next.entries)
577 }
578
579 fn entry_count(&self) -> usize {
580 self.entries.len()
581 }
582}
583
584pub fn diff_context_digest(prev: &ContextDigest, next: &ContextSnapshot) -> ContextDelta {
589 diff_regions(&prev.regions, next)
590}
591
592pub fn apply_delta(base: &mut ContextSnapshot, delta: &ContextDelta) {
596 base.stage_name = delta.stage_name.clone();
597 base.total_tokens = delta.total_tokens;
598 base.max_tokens = delta.max_tokens;
599 for region_delta in &delta.regions {
600 match region_delta {
601 RegionDelta::Set(snapshot) => {
602 match base.regions.iter_mut().find(|r| r.name == snapshot.name) {
603 Some(existing) => *existing = snapshot.clone(),
604 None => base.regions.push(snapshot.clone()),
605 }
606 }
607 RegionDelta::Append {
608 name,
609 entries,
610 current_tokens,
611 } => {
612 if let Some(region) = base.regions.iter_mut().find(|r| &r.name == name) {
613 region.entries.extend(entries.iter().cloned());
614 region.current_tokens = *current_tokens;
615 }
616 }
617 RegionDelta::Clear { name } => {
618 if let Some(region) = base.regions.iter_mut().find(|r| &r.name == name) {
619 region.entries.clear();
620 region.current_tokens = 0;
621 }
622 }
623 RegionDelta::Remove { name } => {
624 base.regions.retain(|r| &r.name != name);
625 }
626 }
627 }
628}
629
630mod codec;
631
632pub use codec::{
633 Frame, RUN_ARCHIVE_MAGIC, RUN_ARCHIVE_VERSION, read_archive, read_archive_lenient,
634 read_archive_start, read_frame, read_record, write_archive_start, write_record,
635};
636
637#[derive(Debug, Clone, PartialEq)]
644pub struct PendingToolBatch {
645 pub stage_index: usize,
647 pub iteration: usize,
649 pub response: String,
651 pub calls: Vec<ToolCallRecord>,
653}
654
655#[derive(Debug, Clone, PartialEq)]
662pub struct InferenceUsageRecord {
663 pub kind: InferenceKind,
665 pub stage: String,
667 pub iteration: usize,
669 pub provider: String,
671 pub model: String,
673 pub prompt_tokens: usize,
675 pub completion_tokens: usize,
677 pub cached_tokens: usize,
679 pub cache_write_tokens: usize,
681 pub cost_usd: Option<f64>,
684 pub cost_reported_by_provider: Option<bool>,
686 pub at: i64,
688}
689
690#[derive(Debug, Clone, PartialEq)]
693pub struct FoldedRun {
694 pub identity: RunIdentity,
696 pub meta: RunMeta,
698 pub context: ContextSnapshot,
700 pub messages: Vec<MessageRecord>,
702 pub inference_count: usize,
704 pub inference_usage: Vec<InferenceUsageRecord>,
711 pub tool_call_count: usize,
713 pub pending_batch: Option<PendingToolBatch>,
717}
718
719pub fn context_contains_batch(context: &ContextSnapshot, batch: &PendingToolBatch) -> bool {
723 let Some(first_id) = batch.calls.first().map(|c| c.id.as_str()) else {
724 return false;
725 };
726 context.regions.iter().any(|region| {
727 region.entries.iter().any(|entry| {
728 matches!(
729 &entry.kind,
730 crate::region::EntryKind::AssistantTurn { tool_calls }
731 if tool_calls.iter().any(|tc| tc.id == first_id)
732 )
733 })
734 })
735}
736
737pub fn fold(records: &[RunRecord]) -> Option<FoldedRun> {
740 let mut iter = records.iter();
741 let (identity, meta) = match iter.next() {
742 Some(RunRecord::Header { identity, meta }) => (identity.clone(), (**meta).clone()),
743 _ => return None,
744 };
745 let mut folded = FoldedRun {
746 identity,
747 meta,
748 context: ContextSnapshot {
749 stage_name: String::new(),
750 total_tokens: 0,
751 max_tokens: 0,
752 regions: Vec::new(),
753 },
754 messages: Vec::new(),
755 inference_count: 0,
756 inference_usage: Vec::new(),
757 tool_call_count: 0,
758 pending_batch: None,
759 };
760 for record in iter {
761 match record {
762 RunRecord::Header { identity, meta } => {
763 folded.identity = identity.clone();
764 folded.meta = (**meta).clone();
765 }
766 RunRecord::OwnershipChanged {
767 machine_id,
768 world_id,
769 ..
770 } => {
771 folded.identity.machine_id = machine_id.clone();
772 folded.identity.world_id = world_id.clone();
773 }
774 RunRecord::Inference { .. } => folded.inference_count += 1,
775 RunRecord::InferenceUsage {
776 kind,
777 stage,
778 iteration,
779 provider,
780 model,
781 prompt_tokens,
782 completion_tokens,
783 cached_tokens,
784 cache_write_tokens,
785 cost_usd,
786 cost_reported_by_provider,
787 at,
788 } => {
789 folded.inference_count += 1;
793 folded.inference_usage.push(InferenceUsageRecord {
794 kind: *kind,
795 stage: stage.clone(),
796 iteration: *iteration,
797 provider: provider.clone(),
798 model: model.clone(),
799 prompt_tokens: *prompt_tokens,
800 completion_tokens: *completion_tokens,
801 cached_tokens: *cached_tokens,
802 cache_write_tokens: *cache_write_tokens,
803 cost_usd: *cost_usd,
804 cost_reported_by_provider: *cost_reported_by_provider,
805 at: *at,
806 });
807 }
808 RunRecord::ToolBatch {
809 calls,
810 stage_index,
811 iteration,
812 response,
813 ..
814 } => {
815 folded.tool_call_count += calls.len();
816 folded.pending_batch = Some(PendingToolBatch {
819 stage_index: *stage_index,
820 iteration: *iteration,
821 response: response.clone(),
822 calls: calls.clone(),
823 });
824 }
825 RunRecord::ToolCallDone {
826 iteration,
827 call_id,
828 result,
829 ..
830 } => {
831 if let Some(batch) = folded
834 .pending_batch
835 .as_mut()
836 .filter(|b| b.iteration == *iteration)
837 && let Some(call) = batch.calls.iter_mut().find(|c| c.id == *call_id)
838 {
839 call.result = Some(result.clone());
840 }
841 }
842 RunRecord::ContextCheckpoint { snapshot, .. } => folded.context = snapshot.clone(),
843 RunRecord::ContextDiff { delta, .. } => apply_delta(&mut folded.context, delta),
844 RunRecord::Message { message, .. } => folded.messages.push(message.clone()),
845 RunRecord::StatusChanged { status, .. } => folded.meta.status = status.clone(),
846 RunRecord::Checkpoint { meta, context, .. } => {
847 folded.meta = (**meta).clone();
848 folded.context = context.clone();
849 }
850 RunRecord::Progress { meta, delta, .. } => {
851 folded.meta = (**meta).clone();
852 apply_delta(&mut folded.context, delta);
853 }
854 }
855 }
856 if let Some(batch) = &folded.pending_batch
861 && (folded.meta.iteration != batch.iteration
862 || context_contains_batch(&folded.context, batch))
863 {
864 folded.pending_batch = None;
865 }
866 Some(folded)
867}
868
869#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
872pub struct RunPoint {
873 pub meta: RunMeta,
875 pub context: ContextSnapshot,
877 pub at: i64,
879}
880
881#[derive(Debug)]
884pub struct PointRef<'a> {
885 pub index: usize,
889 pub at: i64,
891 pub meta: &'a RunMeta,
893 pub context: &'a ContextSnapshot,
895}
896
897pub fn visit_points(records: &[RunRecord], visit: &mut dyn FnMut(PointRef<'_>) -> ControlFlow<()>) {
916 let mut iter = records.iter();
917 let Some(mut folder) = (match iter.next() {
918 Some(first) => PointFolder::start(first),
919 None => None,
920 }) else {
921 return;
922 };
923 for record in iter {
924 if folder.push(record, visit).is_break() {
925 return;
926 }
927 }
928}
929
930pub fn visit_archive_points(
940 r: &mut dyn Read,
941 visit: &mut dyn FnMut(PointRef<'_>) -> ControlFlow<()>,
942) -> io::Result<()> {
943 read_archive_start(r)?;
944 let mut folder = match read_frame(r) {
947 Ok(Some(Frame::Record(first))) => match PointFolder::start(&first) {
948 Some(folder) => folder,
949 None => return Ok(()),
950 },
951 _ => return Ok(()),
952 };
953 while let Ok(Some(frame)) = read_frame(r) {
957 let Frame::Record(record) = frame else {
958 continue;
959 };
960 if folder.push(&record, visit).is_break() {
961 return Ok(());
962 }
963 }
964 Ok(())
965}
966
967struct PointFolder {
972 meta: RunMeta,
973 context: ContextSnapshot,
974 index: usize,
975}
976
977impl PointFolder {
978 fn start(first: &RunRecord) -> Option<Self> {
982 match first {
983 RunRecord::Header { meta, .. } => Some(Self {
984 meta: (**meta).clone(),
985 context: ContextSnapshot {
986 stage_name: String::new(),
987 total_tokens: 0,
988 max_tokens: 0,
989 regions: Vec::new(),
990 },
991 index: 0,
992 }),
993 _ => None,
994 }
995 }
996
997 fn push(
999 &mut self,
1000 record: &RunRecord,
1001 visit: &mut dyn FnMut(PointRef<'_>) -> ControlFlow<()>,
1002 ) -> ControlFlow<()> {
1003 let at = match record {
1004 RunRecord::Header { meta: m, .. } => {
1005 self.meta = (**m).clone();
1006 return ControlFlow::Continue(());
1007 }
1008 RunRecord::StatusChanged { status, .. } => {
1009 self.meta.status = status.clone();
1010 return ControlFlow::Continue(());
1011 }
1012 RunRecord::ContextCheckpoint { snapshot, at } => {
1013 self.context = snapshot.clone();
1014 *at
1015 }
1016 RunRecord::ContextDiff { delta, at } => {
1017 apply_delta(&mut self.context, delta);
1018 *at
1019 }
1020 RunRecord::Checkpoint {
1021 meta: m,
1022 context: c,
1023 at,
1024 } => {
1025 self.meta = (**m).clone();
1026 self.context = c.clone();
1027 *at
1028 }
1029 RunRecord::Progress { meta: m, delta, at } => {
1030 self.meta = (**m).clone();
1031 apply_delta(&mut self.context, delta);
1032 *at
1033 }
1034 RunRecord::OwnershipChanged { .. }
1039 | RunRecord::Inference { .. }
1040 | RunRecord::InferenceUsage { .. }
1041 | RunRecord::ToolBatch { .. }
1042 | RunRecord::ToolCallDone { .. }
1043 | RunRecord::Message { .. } => return ControlFlow::Continue(()),
1044 };
1045 let flow = visit(PointRef {
1046 index: self.index,
1047 at,
1048 meta: &self.meta,
1049 context: &self.context,
1050 });
1051 self.index += 1;
1052 flow
1053 }
1054}
1055
1056pub fn replay_points(records: &[RunRecord]) -> Vec<RunPoint> {
1066 let mut points = Vec::new();
1067 visit_points(records, &mut |point| {
1068 points.push(RunPoint {
1069 meta: point.meta.clone(),
1070 context: point.context.clone(),
1071 at: point.at,
1072 });
1073 ControlFlow::Continue(())
1074 });
1075 points
1076}
1077
1078#[cfg(test)]
1079mod tests {
1080 use super::*;
1081 use crate::run_meta::RunStatus;
1083 use std::io::Write;
1084
1085 fn identity() -> RunIdentity {
1086 RunIdentity {
1087 run_id: "run-1".to_string(),
1088 machine_id: "machine-a".to_string(),
1089 world_id: "world-x".to_string(),
1090 created_at: 100,
1091 }
1092 }
1093
1094 const FIXTURE_NOW: i64 = 1_700_000_000;
1102
1103 fn meta() -> RunMeta {
1104 let mut meta = RunMeta::new(
1105 "run-1".to_string(),
1106 "coder".to_string(),
1107 "/agents/coder".to_string(),
1108 "do it".to_string(),
1109 Some("anthropic/claude".to_string()),
1110 "/work".to_string(),
1111 2,
1112 );
1113 meta.started_at = FIXTURE_NOW;
1114 meta.updated_at = FIXTURE_NOW;
1115 meta
1116 }
1117
1118 #[test]
1122 fn the_fixture_does_not_move_with_the_clock() {
1123 let first = meta();
1124 let mut later = meta();
1125 assert_eq!(first, later, "the fixture is rebuilt identically");
1129 later.started_at += 1;
1130 assert_ne!(
1131 first, later,
1132 "and the comparison is sensitive to the field that used to drift"
1133 );
1134 }
1135
1136 fn entry(content: &str, tokens: usize) -> RegionEntrySnapshot {
1137 RegionEntrySnapshot {
1138 content: content.into(),
1139 tokens,
1140 kind: crate::region::EntryKind::Text,
1141 metadata: None,
1142 key: None,
1143 taint: Default::default(),
1144 reasoning: None,
1145 }
1146 }
1147
1148 fn region(name: &str, entries: Vec<RegionEntrySnapshot>) -> RegionSnapshot {
1149 let current = entries.iter().map(|e| e.tokens).sum();
1150 RegionSnapshot {
1151 name: name.to_string(),
1152 kind: "clearable".to_string(),
1153 current_tokens: current,
1154 max_tokens: 1000,
1155 entries,
1156 description: None,
1157 }
1158 }
1159
1160 fn snapshot(stage: &str, regions: Vec<RegionSnapshot>) -> ContextSnapshot {
1161 let total = regions.iter().map(|r| r.current_tokens).sum();
1162 ContextSnapshot {
1163 stage_name: stage.to_string(),
1164 total_tokens: total,
1165 max_tokens: 10_000,
1166 regions,
1167 }
1168 }
1169
1170 fn header() -> RunRecord {
1171 RunRecord::Header {
1172 identity: identity(),
1173 meta: Box::new(meta()),
1174 }
1175 }
1176
1177 fn region_delta_kind(d: &RegionDelta) -> &'static str {
1181 match d {
1182 RegionDelta::Set(_) => "set",
1183 RegionDelta::Append { .. } => "append",
1184 RegionDelta::Clear { .. } => "clear",
1185 RegionDelta::Remove { .. } => "remove",
1186 }
1187 }
1188
1189 fn assert_diff_roundtrip(a: &ContextSnapshot, b: &ContextSnapshot) {
1194 let delta = diff_context(a, b);
1195 let mut base = a.clone();
1196 apply_delta(&mut base, &delta);
1197 assert_eq!(&base, b);
1198 }
1199
1200 #[test]
1201 fn diff_append_only_growth_is_compact() {
1202 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1203 let b = snapshot(
1204 "s1",
1205 vec![region("conv", vec![entry("hi", 1), entry("there", 2)])],
1206 );
1207 let delta = diff_context(&a, &b);
1208 assert_eq!(region_delta_kind(&delta.regions[0]), "append");
1209 assert_diff_roundtrip(&a, &b);
1210 }
1211
1212 #[test]
1213 fn diff_new_region_is_set() {
1214 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1215 let b = snapshot(
1216 "s1",
1217 vec![
1218 region("conv", vec![entry("hi", 1)]),
1219 region("plan", vec![entry("p", 3)]),
1220 ],
1221 );
1222 let delta = diff_context(&a, &b);
1223 assert!(delta.regions.iter().any(|d| region_delta_kind(d) == "set"));
1224 assert_diff_roundtrip(&a, &b);
1225 }
1226
1227 #[test]
1228 fn diff_cleared_region() {
1229 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1230 let b = snapshot("s1", vec![region("conv", vec![])]);
1231 let delta = diff_context(&a, &b);
1232 assert_eq!(region_delta_kind(&delta.regions[0]), "clear");
1233 assert_diff_roundtrip(&a, &b);
1234 }
1235
1236 #[test]
1237 fn diff_removed_region() {
1238 let a = snapshot(
1239 "s1",
1240 vec![
1241 region("conv", vec![entry("hi", 1)]),
1242 region("plan", vec![entry("p", 3)]),
1243 ],
1244 );
1245 let b = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1246 let delta = diff_context(&a, &b);
1247 assert!(
1248 delta
1249 .regions
1250 .iter()
1251 .any(|d| region_delta_kind(d) == "remove")
1252 );
1253 assert_diff_roundtrip(&a, &b);
1254 }
1255
1256 #[test]
1257 fn diff_non_prefix_rewrite_is_set() {
1258 let a = snapshot("s1", vec![region("conv", vec![entry("old", 1)])]);
1260 let b = snapshot("s1", vec![region("conv", vec![entry("new", 1)])]);
1261 let delta = diff_context(&a, &b);
1262 assert_eq!(region_delta_kind(&delta.regions[0]), "set");
1263 assert_diff_roundtrip(&a, &b);
1264 }
1265
1266 #[test]
1267 fn diff_kind_change_is_set_not_append() {
1268 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1270 let mut grown = region("conv", vec![entry("hi", 1), entry("more", 1)]);
1271 grown.kind = "sliding".to_string();
1272 let b = snapshot("s1", vec![grown]);
1273 let delta = diff_context(&a, &b);
1274 assert_eq!(region_delta_kind(&delta.regions[0]), "set");
1275 assert_diff_roundtrip(&a, &b);
1276 }
1277
1278 fn framed(records: &[RunRecord]) -> Vec<u8> {
1282 let mut buf = Vec::new();
1283 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1284 for record in records {
1285 write_record(&mut buf, record).unwrap();
1286 }
1287 buf
1288 }
1289
1290 fn try_collect_streamed(bytes: &[u8]) -> io::Result<Vec<(usize, i64, usize)>> {
1294 let mut seen = Vec::new();
1295 visit_archive_points(&mut &bytes[..], &mut |p| {
1296 seen.push((p.index, p.at, p.context.total_tokens));
1297 ControlFlow::Continue(())
1298 })?;
1299 Ok(seen)
1300 }
1301
1302 fn collect_streamed(bytes: &[u8]) -> Vec<(usize, i64, usize)> {
1304 try_collect_streamed(bytes).unwrap()
1305 }
1306
1307 #[test]
1308 fn visit_archive_points_matches_visit_points() {
1309 let records = vec![
1310 header(),
1311 RunRecord::ContextCheckpoint {
1312 snapshot: snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1313 at: 10,
1314 },
1315 RunRecord::StatusChanged {
1316 status: RunStatus::Running,
1317 at: 11,
1318 },
1319 RunRecord::Progress {
1320 meta: Box::new(meta()),
1321 delta: diff_context(
1322 &snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1323 &snapshot(
1324 "s1",
1325 vec![region("conv", vec![entry("hi", 1), entry("more", 2)])],
1326 ),
1327 ),
1328 at: 12,
1329 },
1330 ];
1331 let mut in_memory = Vec::new();
1332 visit_points(&records, &mut |p| {
1333 in_memory.push((p.index, p.at, p.context.total_tokens));
1334 ControlFlow::Continue(())
1335 });
1336 assert_eq!(collect_streamed(&framed(&records)), in_memory);
1337 assert_eq!(in_memory.len(), 2, "checkpoint + progress = two points");
1338 }
1339
1340 #[test]
1341 fn visit_archive_points_rejects_a_bad_preamble() {
1342 assert!(try_collect_streamed(b"not an archive at all").is_err());
1343 }
1344
1345 #[test]
1346 fn visit_archive_points_is_lenient_about_a_torn_tail() {
1347 let records = vec![
1348 header(),
1349 RunRecord::ContextCheckpoint {
1350 snapshot: snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1351 at: 10,
1352 },
1353 ];
1354 let mut bytes = framed(&records);
1355 bytes.extend_from_slice(&1000u64.to_be_bytes());
1357 bytes.extend_from_slice(b"partial");
1358 assert_eq!(collect_streamed(&bytes).len(), 1, "points before the tear");
1359 }
1360
1361 #[test]
1362 fn visit_archive_points_visits_nothing_without_a_header() {
1363 let records = vec![RunRecord::ContextCheckpoint {
1364 snapshot: snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1365 at: 10,
1366 }];
1367 assert!(collect_streamed(&framed(&records)).is_empty());
1368 assert!(collect_streamed(&framed(&[])).is_empty());
1370 }
1371
1372 #[test]
1373 fn visit_archive_points_stops_on_break() {
1374 let records = vec![
1375 header(),
1376 RunRecord::ContextCheckpoint {
1377 snapshot: snapshot("s1", vec![region("conv", vec![entry("a", 1)])]),
1378 at: 10,
1379 },
1380 RunRecord::ContextCheckpoint {
1381 snapshot: snapshot("s1", vec![region("conv", vec![entry("b", 2)])]),
1382 at: 11,
1383 },
1384 ];
1385 let bytes = framed(&records);
1386 let mut seen = 0;
1387 visit_archive_points(&mut &bytes[..], &mut |_| {
1388 seen += 1;
1389 ControlFlow::Break(())
1390 })
1391 .unwrap();
1392 assert_eq!(seen, 1);
1393 }
1394
1395 fn assert_digest_matches_full_diff(a: &ContextSnapshot, b: &ContextSnapshot) {
1402 let via_digest = diff_context_digest(&digest_context(a), b);
1403 assert_eq!(via_digest, diff_context(a, b));
1404 let mut base = a.clone();
1406 apply_delta(&mut base, &via_digest);
1407 assert_eq!(&base, b);
1408 }
1409
1410 #[test]
1411 fn digest_diff_append_only_growth_is_compact() {
1412 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1413 let b = snapshot(
1414 "s1",
1415 vec![region("conv", vec![entry("hi", 1), entry("there", 2)])],
1416 );
1417 let delta = diff_context_digest(&digest_context(&a), &b);
1418 assert_eq!(region_delta_kind(&delta.regions[0]), "append");
1419 assert_digest_matches_full_diff(&a, &b);
1420 }
1421
1422 #[test]
1423 fn digest_diff_new_cleared_removed_and_rewritten_regions() {
1424 let a = snapshot(
1425 "s1",
1426 vec![
1427 region("conv", vec![entry("hi", 1)]),
1428 region("gone", vec![entry("bye", 1)]),
1429 region("wiped", vec![entry("w", 1)]),
1430 region("rewritten", vec![entry("old", 1)]),
1431 ],
1432 );
1433 let b = snapshot(
1434 "s1",
1435 vec![
1436 region("conv", vec![entry("hi", 1)]),
1437 region("wiped", vec![]),
1438 region("rewritten", vec![entry("new", 1)]),
1439 region("fresh", vec![entry("f", 2)]),
1440 ],
1441 );
1442 let delta = diff_context_digest(&digest_context(&a), &b);
1443 let kinds: Vec<_> = delta.regions.iter().map(region_delta_kind).collect();
1444 assert_eq!(kinds, vec!["clear", "set", "set", "remove"]);
1445 assert_digest_matches_full_diff(&a, &b);
1446 }
1447
1448 #[test]
1449 fn digest_diff_unchanged_region_emits_nothing() {
1450 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1451 let delta = diff_context_digest(&digest_context(&a), &a.clone());
1452 assert!(delta.regions.is_empty());
1453 assert_digest_matches_full_diff(&a, &a.clone());
1454 }
1455
1456 #[test]
1457 fn digest_diff_kind_change_is_set_not_append() {
1458 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1459 let mut grown = region("conv", vec![entry("hi", 1), entry("more", 1)]);
1460 grown.kind = "sliding".to_string();
1461 let b = snapshot("s1", vec![grown]);
1462 let delta = diff_context_digest(&digest_context(&a), &b);
1463 assert_eq!(region_delta_kind(&delta.regions[0]), "set");
1464 assert_digest_matches_full_diff(&a, &b);
1465 }
1466
1467 #[test]
1470 fn digest_diff_token_recount_is_an_empty_append() {
1471 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1472 let mut recounted = region("conv", vec![entry("hi", 1)]);
1473 recounted.current_tokens = 42;
1474 let b = snapshot("s1", vec![recounted]);
1475 let delta = diff_context_digest(&digest_context(&a), &b);
1476 assert_eq!(region_delta_kind(&delta.regions[0]), "append");
1477 assert_digest_matches_full_diff(&a, &b);
1478 }
1479
1480 #[test]
1483 fn entry_digest_covers_every_field() {
1484 let base = entry("text", 1);
1485 let variants = [
1486 entry("other", 1),
1487 entry("text", 2),
1488 RegionEntrySnapshot {
1489 key: Some("k".to_string()),
1490 ..entry("text", 1)
1491 },
1492 RegionEntrySnapshot {
1493 metadata: Some(serde_json::json!({"a": 1})),
1494 ..entry("text", 1)
1495 },
1496 RegionEntrySnapshot {
1497 kind: crate::region::EntryKind::ToolResult {
1498 tool_call_id: "c1".to_string(),
1499 tool_name: "shell".to_string(),
1500 is_error: false,
1501 },
1502 ..entry("text", 1)
1503 },
1504 ];
1505 let base_hash = entry_digest(&base);
1506 for variant in &variants {
1507 assert_ne!(
1508 entry_digest(variant),
1509 base_hash,
1510 "field change must change the digest: {variant:?}"
1511 );
1512 }
1513 assert_eq!(entry_digest(&base), entry_digest(&entry("text", 1)));
1515 }
1516
1517 #[test]
1518 fn diff_unchanged_region_emits_nothing() {
1519 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1520 let b = a.clone();
1521 let delta = diff_context(&a, &b);
1522 assert!(delta.regions.is_empty());
1523 assert_diff_roundtrip(&a, &b);
1524 }
1525
1526 #[test]
1527 fn apply_delta_skips_unknown_regions_leniently() {
1528 let mut base = snapshot("s1", vec![]);
1530 let delta = ContextDelta {
1531 stage_name: "s1".to_string(),
1532 total_tokens: 0,
1533 max_tokens: 10_000,
1534 regions: vec![
1535 RegionDelta::Append {
1536 name: "ghost".to_string(),
1537 entries: vec![entry("x", 1)],
1538 current_tokens: 1,
1539 },
1540 RegionDelta::Clear {
1541 name: "ghost".to_string(),
1542 },
1543 RegionDelta::Remove {
1544 name: "ghost".to_string(),
1545 },
1546 ],
1547 };
1548 apply_delta(&mut base, &delta);
1549 assert!(base.regions.is_empty());
1550 }
1551
1552 fn all_record_kinds() -> Vec<RunRecord> {
1555 vec![
1556 header(),
1557 RunRecord::OwnershipChanged {
1558 machine_id: "machine-b".to_string(),
1559 world_id: "world-y".to_string(),
1560 at: 101,
1561 },
1562 RunRecord::Inference {
1563 stage: "plan".to_string(),
1564 iteration: 0,
1565 request: InferenceRequestRecord {
1566 model: "m".to_string(),
1567 system: vec!["sys".to_string()],
1568 messages: vec![MessageRecord {
1569 role: "user".to_string(),
1570 content: "hi".to_string(),
1571 }],
1572 tool_names: vec!["read_file".to_string()],
1573 temperature: 0.7,
1574 max_tokens: 1024,
1575 },
1576 response: InferenceResponseRecord {
1577 content: "ok".to_string(),
1578 tool_calls: vec![],
1579 prompt_tokens: 10,
1580 completion_tokens: 5,
1581 cached_tokens: 0,
1582 cache_write_tokens: 0,
1583 },
1584 at: 102,
1585 },
1586 RunRecord::InferenceUsage {
1587 kind: InferenceKind::Compaction,
1588 stage: "plan".to_string(),
1589 iteration: 2,
1590 provider: "anthropic".to_string(),
1591 model: "claude-sonnet-5".to_string(),
1592 prompt_tokens: 7000,
1593 completion_tokens: 70,
1594 cached_tokens: 12,
1595 cache_write_tokens: 34,
1596 cost_usd: None,
1597 cost_reported_by_provider: None,
1598 at: 102,
1599 },
1600 RunRecord::ToolBatch {
1601 calls: vec![ToolCallRecord {
1602 id: "c1".to_string(),
1603 name: "read_file".to_string(),
1604 arguments: "{}".to_string(),
1605 result: Some("body".to_string().into()),
1606 thought_signature: Some("sig".to_string()),
1607 }],
1608 at: 103,
1609 stage_index: 0,
1610 iteration: 0,
1611 response: "reading".to_string(),
1612 },
1613 RunRecord::ToolCallDone {
1614 iteration: 0,
1615 call_id: "c1".to_string(),
1616 result: "body".to_string().into(),
1617 at: 103,
1618 },
1619 RunRecord::ContextCheckpoint {
1620 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1621 at: 104,
1622 },
1623 RunRecord::ContextDiff {
1624 delta: ContextDelta {
1625 stage_name: "plan".to_string(),
1626 total_tokens: 3,
1627 max_tokens: 10_000,
1628 regions: vec![RegionDelta::Append {
1629 name: "conv".to_string(),
1630 entries: vec![entry("more", 2)],
1631 current_tokens: 3,
1632 }],
1633 },
1634 at: 105,
1635 },
1636 RunRecord::Message {
1637 message: MessageRecord {
1638 role: "user".to_string(),
1639 content: "another".to_string(),
1640 },
1641 at: 106,
1642 },
1643 RunRecord::StatusChanged {
1644 status: RunStatus::Complete,
1645 at: 107,
1646 },
1647 RunRecord::Checkpoint {
1648 meta: Box::new(meta()),
1649 context: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1650 at: 108,
1651 },
1652 RunRecord::Progress {
1653 meta: Box::new(meta()),
1654 delta: ContextDelta {
1655 stage_name: "plan".to_string(),
1656 total_tokens: 3,
1657 max_tokens: 10_000,
1658 regions: vec![RegionDelta::Append {
1659 name: "conv".to_string(),
1660 entries: vec![entry("step", 2)],
1661 current_tokens: 3,
1662 }],
1663 },
1664 at: 109,
1665 },
1666 ]
1667 }
1668
1669 #[test]
1670 fn archive_write_then_read_roundtrips_every_record_kind() {
1671 let records = all_record_kinds();
1672 let mut buf = Vec::new();
1673 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1674 for r in &records {
1675 write_record(&mut buf, r).unwrap();
1676 }
1677 let (version, read) = read_archive(&mut buf.as_slice()).unwrap();
1678 assert_eq!(version, RUN_ARCHIVE_VERSION);
1679 assert_eq!(read, records);
1680 }
1681
1682 #[test]
1683 fn read_archive_start_rejects_bad_magic() {
1684 let mut bytes: &[u8] = b"XXXX\x00\x01";
1685 let err = read_archive_start(&mut bytes).unwrap_err();
1686 assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1687 }
1688
1689 #[test]
1694 fn read_archive_start_reports_version() {
1695 let mut buf = Vec::new();
1696 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1697 assert_eq!(
1698 read_archive_start(&mut buf.as_slice()).unwrap(),
1699 RUN_ARCHIVE_VERSION
1700 );
1701 }
1702
1703 #[test]
1704 fn read_record_returns_none_at_clean_eof() {
1705 let empty: &[u8] = &[];
1706 assert!(read_record(&mut { empty }).unwrap().is_none());
1707 }
1708
1709 #[test]
1710 fn read_record_errors_on_truncated_length_prefix() {
1711 let mut bytes: &[u8] = &[0, 0];
1713 let err = read_record(&mut bytes).unwrap_err();
1714 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1715 }
1716
1717 #[test]
1718 fn read_record_errors_on_truncated_payload() {
1719 let mut bytes: &[u8] = &[0, 0, 0, 0, 0, 0, 0, 10, 1, 2];
1721 let err = read_record(&mut bytes).unwrap_err();
1722 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1723 }
1724
1725 #[test]
1726 fn read_record_errors_on_empty_payload_at_boundary() {
1727 let mut bytes: &[u8] = &[0, 0, 0, 0, 0, 0, 0, 10];
1730 let err = read_record(&mut bytes).unwrap_err();
1731 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1732 }
1733
1734 #[test]
1735 fn read_record_errors_on_invalid_json_payload() {
1736 let mut buf = Vec::new();
1738 let bad = b"not json";
1739 buf.extend_from_slice(&(bad.len() as u64).to_be_bytes());
1740 buf.extend_from_slice(bad);
1741 let err = read_record(&mut buf.as_slice()).unwrap_err();
1742 assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1743 }
1744
1745 struct FailingReader;
1748 impl Read for FailingReader {
1749 fn read(&mut self, _buf: &mut [u8]) -> io::Result<usize> {
1750 Err(io::Error::other("device error"))
1751 }
1752 }
1753
1754 #[test]
1755 fn read_record_propagates_reader_errors() {
1756 let err = read_record(&mut FailingReader).unwrap_err();
1757 assert_eq!(err.kind(), io::ErrorKind::Other);
1758 }
1759
1760 #[test]
1761 fn read_archive_propagates_a_bad_preamble() {
1762 let mut bytes: &[u8] = b"LV";
1764 assert!(read_archive(&mut bytes).is_err());
1765 }
1766
1767 #[test]
1768 fn read_archive_propagates_a_bad_frame() {
1769 let mut buf = Vec::new();
1771 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1772 buf.extend_from_slice(&[0, 0, 0, 0, 0, 0, 0, 5, 1, 2]); let err = read_archive(&mut buf.as_slice()).unwrap_err();
1774 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1775 }
1776
1777 struct FailAfter {
1779 remaining: usize,
1780 }
1781 impl Write for FailAfter {
1782 fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
1783 if self.remaining == 0 {
1784 return Err(io::Error::other("disk full"));
1785 }
1786 let n = buf.len().min(self.remaining);
1787 self.remaining -= n;
1788 Ok(n)
1789 }
1790 fn flush(&mut self) -> io::Result<()> {
1791 Ok(())
1792 }
1793 }
1794
1795 #[test]
1796 fn fail_after_writer_flush_is_a_noop() {
1797 assert!(FailAfter { remaining: 1 }.flush().is_ok());
1798 }
1799
1800 #[test]
1801 fn write_archive_start_propagates_write_errors() {
1802 assert!(write_archive_start(&mut FailAfter { remaining: 0 }, 1).is_err());
1804 assert!(write_archive_start(&mut FailAfter { remaining: 4 }, 1).is_err());
1805 }
1806
1807 #[test]
1808 fn write_record_propagates_write_errors() {
1809 let rec = header();
1810 assert!(write_record(&mut FailAfter { remaining: 0 }, &rec).is_err());
1812 assert!(write_record(&mut FailAfter { remaining: 8 }, &rec).is_err());
1813 }
1814
1815 #[test]
1821 fn an_absurd_frame_length_is_an_error_not_an_allocation() {
1822 let mut buf = Vec::new();
1823 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1824 write_record(&mut buf, &header()).unwrap();
1825 buf.extend_from_slice(&u64::MAX.to_be_bytes());
1827
1828 let err = read_archive(&mut buf.as_slice())
1829 .expect_err("the strict reader must refuse an impossible frame");
1830 assert_eq!(err.kind(), io::ErrorKind::InvalidData, "{err}");
1831
1832 let (_, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
1835 assert_eq!(records, vec![header()]);
1836 }
1837
1838 #[test]
1839 fn read_archive_lenient_matches_strict_on_a_clean_archive() {
1840 let records = all_record_kinds();
1843 let mut buf = Vec::new();
1844 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1845 for r in &records {
1846 write_record(&mut buf, r).unwrap();
1847 }
1848 let (version, read) = read_archive_lenient(&mut buf.as_slice()).unwrap();
1849 assert_eq!(version, RUN_ARCHIVE_VERSION);
1850 assert_eq!(read, records);
1851 }
1852
1853 #[test]
1854 fn read_archive_lenient_keeps_valid_prefix_before_a_torn_tail() {
1855 let mut buf = Vec::new();
1859 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1860 write_record(&mut buf, &header()).unwrap();
1861 write_record(
1862 &mut buf,
1863 &RunRecord::ContextCheckpoint {
1864 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1865 at: 1,
1866 },
1867 )
1868 .unwrap();
1869 buf.extend_from_slice(&[0, 0, 0, 0, 0, 0, 0, 10, 1, 2]);
1871
1872 assert!(read_archive(&mut buf.as_slice()).is_err());
1874 let (version, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
1876 assert_eq!(version, RUN_ARCHIVE_VERSION);
1877 assert_eq!(records.len(), 2);
1878 let folded = fold(&records).expect("prefix starts with a Header");
1879 assert_eq!(folded.context.regions[0].entries.len(), 1);
1880 }
1881
1882 #[test]
1883 fn read_archive_lenient_still_errors_on_a_bad_preamble() {
1884 let mut bad_magic: &[u8] = b"XXXX\x00\x01";
1887 assert!(read_archive_lenient(&mut bad_magic).is_err());
1888 let mut short: &[u8] = b"LVR1";
1890 assert!(read_archive_lenient(&mut short).is_err());
1891 }
1892
1893 #[test]
1894 fn read_archive_start_errors_on_truncated_version() {
1895 let mut bytes: &[u8] = b"LVR1";
1897 let err = read_archive_start(&mut bytes).unwrap_err();
1898 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1899 }
1900
1901 #[test]
1904 fn fold_requires_a_header_first() {
1905 assert!(fold(&[]).is_none());
1906 assert!(
1907 fold(&[RunRecord::StatusChanged {
1908 status: RunStatus::Complete,
1909 at: 1
1910 }])
1911 .is_none()
1912 );
1913 }
1914
1915 #[test]
1916 fn fold_reconstructs_state_from_the_journal() {
1917 let records = all_record_kinds();
1918 let folded = fold(&records).expect("has header");
1919 assert_eq!(folded.identity.machine_id, "machine-b");
1921 assert_eq!(folded.identity.world_id, "world-y");
1922 assert_eq!(folded.inference_count, 2);
1925 assert_eq!(folded.inference_usage.len(), 1);
1926 assert_eq!(folded.tool_call_count, 1);
1927 assert_eq!(folded.messages.len(), 1);
1929 assert_eq!(folded.messages[0].content, "another");
1930 assert_eq!(folded.context.regions[0].name, "conv");
1933 assert_eq!(folded.context.regions[0].entries.len(), 2);
1934 assert_eq!(folded.context.total_tokens, 3);
1935 assert_eq!(folded.meta.run_id, "run-1");
1936 let pending = folded.pending_batch.expect("batch never applied");
1939 assert_eq!(pending.calls[0].result.as_deref(), Some("body"));
1940 }
1941
1942 #[test]
1943 fn fold_applies_context_diffs_over_a_checkpoint() {
1944 let records = vec![
1946 header(),
1947 RunRecord::ContextCheckpoint {
1948 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1949 at: 1,
1950 },
1951 RunRecord::ContextDiff {
1952 delta: ContextDelta {
1953 stage_name: "plan".to_string(),
1954 total_tokens: 3,
1955 max_tokens: 10_000,
1956 regions: vec![RegionDelta::Append {
1957 name: "conv".to_string(),
1958 entries: vec![entry("there", 2)],
1959 current_tokens: 3,
1960 }],
1961 },
1962 at: 2,
1963 },
1964 ];
1965 let folded = fold(&records).unwrap();
1966 assert_eq!(folded.context.regions[0].entries.len(), 2);
1967 assert_eq!(folded.context.total_tokens, 3);
1968 }
1969
1970 #[test]
1971 fn fold_later_header_updates_identity_and_meta() {
1972 let mut second_meta = meta();
1974 second_meta.status = RunStatus::Running;
1975 let records = vec![
1976 header(),
1977 RunRecord::Header {
1978 identity: RunIdentity {
1979 run_id: "run-1".to_string(),
1980 machine_id: "machine-c".to_string(),
1981 world_id: "world-z".to_string(),
1982 created_at: 200,
1983 },
1984 meta: Box::new(second_meta),
1985 },
1986 ];
1987 let folded = fold(&records).unwrap();
1988 assert_eq!(folded.identity.machine_id, "machine-c");
1989 assert_eq!(folded.meta.status, RunStatus::Running);
1990 }
1991
1992 #[test]
1993 fn fold_progress_applies_meta_and_context_diff() {
1994 let mut advanced = meta();
1995 advanced.status = RunStatus::Running;
1996 advanced.iteration = 5;
1997 let records = vec![
1998 header(),
1999 RunRecord::ContextCheckpoint {
2000 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
2001 at: 1,
2002 },
2003 RunRecord::Progress {
2004 meta: Box::new(advanced),
2005 delta: ContextDelta {
2006 stage_name: "plan".to_string(),
2007 total_tokens: 3,
2008 max_tokens: 10_000,
2009 regions: vec![RegionDelta::Append {
2010 name: "conv".to_string(),
2011 entries: vec![entry("there", 2)],
2012 current_tokens: 3,
2013 }],
2014 },
2015 at: 2,
2016 },
2017 ];
2018 let folded = fold(&records).unwrap();
2019 assert_eq!(folded.meta.iteration, 5);
2020 assert_eq!(folded.meta.status, RunStatus::Running);
2021 assert_eq!(folded.context.regions[0].entries.len(), 2);
2022 }
2023
2024 #[test]
2028 fn fold_carries_a_submitted_final_output_through_progress() {
2029 let mut answered = meta();
2030 answered.final_output = Some(
2031 crate::output::FinalOutput::new(
2032 "renamed two helpers",
2033 Some("markdown".to_string()),
2034 "summary".to_string(),
2035 9,
2036 )
2037 .descriptor(),
2038 );
2039 answered.output_request = Some(crate::output::OutputSpec {
2040 format: Some("a2ui".to_string()),
2041 ..Default::default()
2042 });
2043 let records = vec![
2044 header(),
2045 RunRecord::Progress {
2046 meta: Box::new(answered),
2047 delta: ContextDelta {
2048 stage_name: "summary".to_string(),
2049 total_tokens: 0,
2050 max_tokens: 10_000,
2051 regions: vec![],
2052 },
2053 at: 2,
2054 },
2055 ];
2056 let folded = fold(&records).unwrap();
2057 let output = folded.meta.final_output.expect("the answer folded through");
2058 assert_eq!(output.bytes, "renamed two helpers".len());
2061 assert_eq!(output.stage, "summary");
2062 assert_eq!(
2063 folded.meta.output_request.and_then(|s| s.format).as_deref(),
2064 Some("a2ui")
2065 );
2066 }
2067
2068 fn call(id: &str, result: Option<&str>) -> ToolCallRecord {
2071 ToolCallRecord {
2072 id: id.to_string(),
2073 name: "shell".to_string(),
2074 arguments: "{}".to_string(),
2075 result: result.map(Into::into),
2076 thought_signature: None,
2077 }
2078 }
2079
2080 fn batch(iteration: usize, calls: Vec<ToolCallRecord>) -> RunRecord {
2081 RunRecord::ToolBatch {
2082 calls,
2083 at: 10,
2084 stage_index: 0,
2085 iteration,
2086 response: "running tools".to_string(),
2087 }
2088 }
2089
2090 fn turn_entry(call_ids: &[&str]) -> RegionEntrySnapshot {
2092 let mut e = entry("turn", 1);
2093 e.kind = crate::region::EntryKind::AssistantTurn {
2094 tool_calls: call_ids
2095 .iter()
2096 .map(|id| crate::region::SerializedToolCall {
2097 id: id.to_string(),
2098 name: "shell".to_string(),
2099 arguments: serde_json::Value::Null,
2100 thought_signature: None,
2101 })
2102 .collect(),
2103 };
2104 e
2105 }
2106
2107 #[test]
2108 fn fold_surfaces_a_pending_batch_with_merged_results() {
2109 let records = vec![
2114 header(),
2115 batch(
2116 0,
2117 vec![
2118 call("c1", None),
2119 call("c2", Some("inline")),
2120 call("c3", None),
2121 ],
2122 ),
2123 RunRecord::ToolCallDone {
2124 iteration: 0,
2125 call_id: "c1".to_string(),
2126 result: "ran".to_string().into(),
2127 at: 11,
2128 },
2129 ];
2130 let folded = fold(&records).unwrap();
2131 let pending = folded.pending_batch.expect("batch is pending");
2132 assert_eq!(pending.iteration, 0);
2133 assert_eq!(pending.response, "running tools");
2134 assert_eq!(pending.calls[0].result.as_deref(), Some("ran"));
2135 assert_eq!(pending.calls[1].result.as_deref(), Some("inline"));
2136 assert_eq!(pending.calls[2].result, None);
2137 assert_eq!(folded.tool_call_count, 3);
2138 }
2139
2140 #[test]
2141 fn fold_keeps_only_the_latest_batch_and_ignores_stale_done_records() {
2142 let mut advanced = meta();
2145 advanced.iteration = 1;
2146 let records = vec![
2147 header(),
2148 batch(0, vec![call("c1", None)]),
2149 RunRecord::Progress {
2150 meta: Box::new(advanced),
2151 delta: ContextDelta {
2152 stage_name: "plan".to_string(),
2153 total_tokens: 0,
2154 max_tokens: 10_000,
2155 regions: vec![],
2156 },
2157 at: 11,
2158 },
2159 batch(1, vec![call("c2", None)]),
2160 RunRecord::ToolCallDone {
2161 iteration: 0,
2162 call_id: "c1".to_string(),
2163 result: "stale".to_string().into(),
2164 at: 12,
2165 },
2166 RunRecord::ToolCallDone {
2167 iteration: 1,
2168 call_id: "unknown".to_string(),
2169 result: "nowhere to land".to_string().into(),
2170 at: 13,
2171 },
2172 ];
2173 let folded = fold(&records).unwrap();
2174 let pending = folded.pending_batch.expect("latest batch is pending");
2175 assert_eq!(pending.iteration, 1);
2176 assert_eq!(pending.calls.len(), 1);
2177 assert_eq!(pending.calls[0].id, "c2");
2178 assert_eq!(pending.calls[0].result, None, "stale/unknown dones ignored");
2179 }
2180
2181 #[test]
2182 fn fold_clears_a_batch_once_the_iteration_moves_on() {
2183 let mut advanced = meta();
2186 advanced.iteration = 1;
2187 let records = vec![
2188 header(),
2189 batch(0, vec![call("c1", Some("done"))]),
2190 RunRecord::Progress {
2191 meta: Box::new(advanced),
2192 delta: ContextDelta {
2193 stage_name: "plan".to_string(),
2194 total_tokens: 0,
2195 max_tokens: 10_000,
2196 regions: vec![],
2197 },
2198 at: 11,
2199 },
2200 ];
2201 assert_eq!(fold(&records).unwrap().pending_batch, None);
2202 }
2203
2204 #[test]
2205 fn fold_clears_a_batch_whose_turn_already_landed_in_the_window() {
2206 let records = vec![
2209 header(),
2210 batch(0, vec![call("c1", Some("done"))]),
2211 RunRecord::ContextCheckpoint {
2212 snapshot: snapshot("plan", vec![region("conv", vec![turn_entry(&["c1"])])]),
2213 at: 11,
2214 },
2215 ];
2216 assert_eq!(fold(&records).unwrap().pending_batch, None);
2217 }
2218
2219 #[test]
2220 fn context_contains_batch_matches_only_the_batch_turn() {
2221 let pending = PendingToolBatch {
2222 stage_index: 0,
2223 iteration: 0,
2224 response: String::new(),
2225 calls: vec![call("c1", None)],
2226 };
2227 let other = snapshot("plan", vec![region("conv", vec![turn_entry(&["zz"])])]);
2229 assert!(!context_contains_batch(&other, &pending));
2230 let own = snapshot(
2232 "plan",
2233 vec![region("conv", vec![turn_entry(&["c1", "c2"])])],
2234 );
2235 assert!(context_contains_batch(&own, &pending));
2236 let empty = PendingToolBatch {
2238 calls: vec![],
2239 ..pending
2240 };
2241 assert!(!context_contains_batch(&own, &empty));
2242 }
2243
2244 #[test]
2245 fn old_shape_tool_batch_json_still_parses() {
2246 let json = br#"{"ToolBatch":{"calls":[{"id":"c1","name":"shell","arguments":"{}","result":"ok"}],"at":9}}"#;
2250 let mut buf = Vec::new();
2251 buf.extend_from_slice(&(json.len() as u64).to_be_bytes());
2252 buf.extend_from_slice(json);
2253 let record = read_record(&mut buf.as_slice()).unwrap().unwrap();
2254 assert_eq!(
2255 record,
2256 RunRecord::ToolBatch {
2257 calls: vec![call("c1", Some("ok"))],
2258 at: 9,
2259 stage_index: 0,
2260 iteration: 0,
2261 response: String::new(),
2262 }
2263 );
2264 }
2265
2266 fn three_point_records() -> Vec<RunRecord> {
2270 let mut running = meta();
2271 running.status = RunStatus::Running;
2272 vec![
2273 header(),
2274 RunRecord::ContextCheckpoint {
2275 snapshot: snapshot("plan", vec![region("conv", vec![entry("first", 1)])]),
2276 at: 10,
2277 },
2278 RunRecord::ContextDiff {
2279 delta: ContextDelta {
2280 stage_name: "plan".to_string(),
2281 total_tokens: 2,
2282 max_tokens: 10_000,
2283 regions: vec![RegionDelta::Append {
2284 name: "conv".to_string(),
2285 entries: vec![entry("second", 1)],
2286 current_tokens: 2,
2287 }],
2288 },
2289 at: 20,
2290 },
2291 RunRecord::Progress {
2292 meta: Box::new(running),
2293 delta: ContextDelta {
2294 stage_name: "code".to_string(),
2295 total_tokens: 3,
2296 max_tokens: 10_000,
2297 regions: vec![RegionDelta::Append {
2298 name: "conv".to_string(),
2299 entries: vec![entry("third", 1)],
2300 current_tokens: 3,
2301 }],
2302 },
2303 at: 30,
2304 },
2305 ]
2306 }
2307
2308 #[test]
2309 fn visit_points_indexes_points_in_order_and_carries_the_running_window() {
2310 let records = three_point_records();
2311 let mut seen: Vec<(usize, i64, usize)> = Vec::new();
2312 visit_points(&records, &mut |point| {
2313 seen.push((
2314 point.index,
2315 point.at,
2316 point.context.regions[0].entries.len(),
2317 ));
2318 ControlFlow::Continue(())
2319 });
2320 assert_eq!(seen, vec![(0, 10, 1), (1, 20, 2), (2, 30, 3)]);
2322 }
2323
2324 #[test]
2327 fn visit_points_stops_at_the_first_break() {
2328 let records = three_point_records();
2329 let mut visits = 0;
2330 visit_points(&records, &mut |point| {
2331 visits += 1;
2332 if point.index == 1 {
2333 ControlFlow::Break(())
2334 } else {
2335 ControlFlow::Continue(())
2336 }
2337 });
2338 assert_eq!(
2339 visits, 2,
2340 "stopped at the breaking point, did not run the third"
2341 );
2342 }
2343
2344 #[test]
2345 fn visit_points_without_a_header_visits_nothing() {
2346 let mut visits = 0;
2347 {
2348 let mut count = |_: PointRef<'_>| {
2349 visits += 1;
2350 ControlFlow::Continue(())
2351 };
2352
2353 visit_points(&three_point_records(), &mut count);
2357 visit_points(&[], &mut count);
2360 visit_points(
2361 &[RunRecord::ContextCheckpoint {
2362 snapshot: snapshot("plan", vec![]),
2363 at: 1,
2364 }],
2365 &mut count,
2366 );
2367 }
2368 assert_eq!(visits, 3, "only the well-formed journal produced points");
2369 }
2370
2371 #[test]
2375 fn visit_points_and_replay_points_agree() {
2376 for records in [
2377 three_point_records(),
2378 vec![header()],
2379 vec![],
2380 vec![RunRecord::Message {
2381 message: MessageRecord {
2382 role: "user".to_string(),
2383 content: "x".to_string(),
2384 },
2385 at: 1,
2386 }],
2387 ] {
2388 let collected: Vec<RunPoint> = {
2389 let mut out = Vec::new();
2390 visit_points(&records, &mut |point| {
2391 out.push(RunPoint {
2392 meta: point.meta.clone(),
2393 context: point.context.clone(),
2394 at: point.at,
2395 });
2396 ControlFlow::Continue(())
2397 });
2398 out
2399 };
2400 assert_eq!(collected, replay_points(&records));
2401 }
2402 }
2403
2404 #[test]
2405 fn replay_points_requires_a_header() {
2406 assert!(replay_points(&[]).is_empty());
2407 assert!(
2408 replay_points(&[RunRecord::Message {
2409 message: MessageRecord {
2410 role: "user".to_string(),
2411 content: "x".to_string(),
2412 },
2413 at: 1,
2414 }])
2415 .is_empty()
2416 );
2417 }
2418
2419 #[test]
2420 fn replay_points_emits_a_snapshot_per_context_change() {
2421 let mut running = meta();
2424 running.status = RunStatus::Running;
2425 let records = vec![
2426 header(),
2427 RunRecord::Inference {
2428 stage: "plan".to_string(),
2429 iteration: 0,
2430 request: InferenceRequestRecord {
2431 model: "m".to_string(),
2432 system: vec![],
2433 messages: vec![],
2434 tool_names: vec![],
2435 temperature: 0.7,
2436 max_tokens: 10,
2437 },
2438 response: InferenceResponseRecord {
2439 content: "ok".to_string(),
2440 tool_calls: vec![],
2441 prompt_tokens: 1,
2442 completion_tokens: 1,
2443 cached_tokens: 0,
2444 cache_write_tokens: 0,
2445 },
2446 at: 1,
2447 },
2448 batch(0, vec![call("c1", None)]),
2449 RunRecord::ToolCallDone {
2450 iteration: 0,
2451 call_id: "c1".to_string(),
2452 result: "ran".to_string().into(),
2453 at: 1,
2454 },
2455 RunRecord::ContextCheckpoint {
2456 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
2457 at: 2,
2458 },
2459 RunRecord::StatusChanged {
2460 status: RunStatus::Running,
2461 at: 3,
2462 },
2463 RunRecord::Progress {
2464 meta: Box::new(running),
2465 delta: ContextDelta {
2466 stage_name: "implement".to_string(),
2467 total_tokens: 3,
2468 max_tokens: 10_000,
2469 regions: vec![RegionDelta::Append {
2470 name: "conv".to_string(),
2471 entries: vec![entry("more", 2)],
2472 current_tokens: 3,
2473 }],
2474 },
2475 at: 4,
2476 },
2477 ];
2478 let points = replay_points(&records);
2479 assert_eq!(points.len(), 2, "one point per context change");
2480 assert_eq!(points[0].at, 2);
2482 assert_eq!(points[0].context.regions[0].entries.len(), 1);
2483 assert_eq!(points[1].at, 4);
2486 assert_eq!(points[1].context.regions[0].entries.len(), 2);
2487 assert_eq!(points[1].context.stage_name, "implement");
2488 assert_eq!(points[1].meta.status, RunStatus::Running);
2489 }
2490
2491 #[test]
2492 fn replay_points_handles_context_diff_and_a_later_header() {
2493 let mut relabeled = meta();
2496 relabeled.agent_name = "renamed".to_string();
2497 let records = vec![
2498 header(),
2499 RunRecord::ContextCheckpoint {
2500 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
2501 at: 1,
2502 },
2503 RunRecord::Header {
2504 identity: identity(),
2505 meta: Box::new(relabeled),
2506 },
2507 RunRecord::ContextDiff {
2508 delta: ContextDelta {
2509 stage_name: "plan".to_string(),
2510 total_tokens: 3,
2511 max_tokens: 10_000,
2512 regions: vec![RegionDelta::Append {
2513 name: "conv".to_string(),
2514 entries: vec![entry("more", 2)],
2515 current_tokens: 3,
2516 }],
2517 },
2518 at: 2,
2519 },
2520 ];
2521 let points = replay_points(&records);
2522 assert_eq!(points.len(), 2); assert_eq!(points[1].context.regions[0].entries.len(), 2);
2524 assert_eq!(points[1].meta.agent_name, "renamed");
2526 }
2527
2528 #[test]
2529 fn replay_points_over_a_full_checkpoint() {
2530 let records = vec![
2532 header(),
2533 RunRecord::Checkpoint {
2534 meta: Box::new(meta()),
2535 context: snapshot("review", vec![region("conv", vec![entry("x", 4)])]),
2536 at: 9,
2537 },
2538 ];
2539 let points = replay_points(&records);
2540 assert_eq!(points.len(), 1);
2541 assert_eq!(points[0].context.stage_name, "review");
2542 assert_eq!(points[0].context.regions[0].entries[0].tokens, 4);
2543 }
2544
2545 #[test]
2549 fn every_inference_kind_has_a_distinct_label_and_serialized_name() {
2550 let all = [
2551 (InferenceKind::Stage, "stage"),
2552 (InferenceKind::Compaction, "compaction"),
2553 (InferenceKind::Title, "title"),
2554 (InferenceKind::Routing, "routing"),
2555 ];
2556 for (kind, label) in all {
2557 assert_eq!(kind.label(), label);
2558 assert_eq!(serde_json::to_value(kind).unwrap(), label);
2559 }
2560 let labels: std::collections::HashSet<_> = all.iter().map(|(k, _)| k.label()).collect();
2561 assert_eq!(labels.len(), all.len(), "labels must not collide");
2562 }
2563
2564 #[test]
2568 fn only_a_stage_turn_counts_as_stage_work() {
2569 assert!(InferenceKind::Stage.is_stage_work());
2570 for kind in [
2571 InferenceKind::Compaction,
2572 InferenceKind::Title,
2573 InferenceKind::Routing,
2574 ] {
2575 assert!(
2576 !kind.is_stage_work(),
2577 "{kind:?} is machinery, not stage work"
2578 );
2579 }
2580 }
2581
2582 #[test]
2586 fn a_usage_record_without_a_kind_reads_back_as_stage_work() {
2587 let json = serde_json::json!({
2588 "InferenceUsage": {
2589 "stage": "plan",
2590 "iteration": 1,
2591 "provider": "anthropic",
2592 "model": "claude-sonnet-5",
2593 "prompt_tokens": 10,
2594 "completion_tokens": 2,
2595 "cached_tokens": 0,
2596 "cache_write_tokens": 0,
2597 "at": 5,
2598 }
2599 });
2600 let record: RunRecord = serde_json::from_value(json).unwrap();
2605 assert_eq!(
2606 record,
2607 RunRecord::InferenceUsage {
2608 kind: InferenceKind::Stage,
2609 stage: "plan".to_string(),
2610 iteration: 1,
2611 provider: "anthropic".to_string(),
2612 model: "claude-sonnet-5".to_string(),
2613 prompt_tokens: 10,
2614 completion_tokens: 2,
2615 cached_tokens: 0,
2616 cache_write_tokens: 0,
2617 cost_usd: None,
2618 cost_reported_by_provider: None,
2619 at: 5,
2620 }
2621 );
2622 }
2623
2624 #[test]
2629 fn folding_keeps_each_call_separate_instead_of_summing_them() {
2630 let usage = |kind, prompt, at| RunRecord::InferenceUsage {
2631 kind,
2632 stage: "plan".to_string(),
2633 iteration: 1,
2634 provider: "anthropic".to_string(),
2635 model: "claude-sonnet-5".to_string(),
2636 prompt_tokens: prompt,
2637 completion_tokens: 1,
2638 cached_tokens: 0,
2639 cache_write_tokens: 0,
2640 cost_usd: None,
2641 cost_reported_by_provider: None,
2642 at,
2643 };
2644 let folded = fold(&[
2647 header(),
2648 usage(InferenceKind::Compaction, 7000, 1),
2649 usage(InferenceKind::Stage, 21_000, 2),
2650 ])
2651 .unwrap();
2652
2653 assert_eq!(folded.inference_count, 2);
2654 let seen: Vec<_> = folded
2655 .inference_usage
2656 .iter()
2657 .map(|u| (u.kind, u.prompt_tokens))
2658 .collect();
2659 assert_eq!(
2660 seen,
2661 vec![
2662 (InferenceKind::Compaction, 7000),
2663 (InferenceKind::Stage, 21_000)
2664 ]
2665 );
2666 assert!(
2670 folded
2671 .inference_usage
2672 .iter()
2673 .all(|u| u.prompt_tokens < 32_000),
2674 "no single call exceeded the window, and the journal can now prove it"
2675 );
2676 }
2677
2678 #[test]
2681 fn both_inference_record_kinds_count_as_one_call_each() {
2682 let records = all_record_kinds();
2683 let folded = fold(&records).unwrap();
2684 let written = records
2685 .iter()
2686 .filter(|r| {
2687 matches!(
2688 r,
2689 RunRecord::Inference { .. } | RunRecord::InferenceUsage { .. }
2690 )
2691 })
2692 .count();
2693 assert_eq!(folded.inference_count, written);
2694 assert_eq!(folded.inference_usage.len(), 1);
2695 }
2696
2697 #[test]
2706 fn probe_replay_matches_every_step() {
2707 fn region(name: &str, entries: &[(&str, usize)]) -> RegionSnapshot {
2708 RegionSnapshot {
2709 name: name.to_string(),
2710 kind: "temporary".to_string(),
2711 current_tokens: entries.iter().map(|(_, t)| *t).sum(),
2712 max_tokens: 1000,
2713 entries: entries
2714 .iter()
2715 .map(|(c, t)| RegionEntrySnapshot {
2716 content: (*c).into(),
2717 tokens: *t,
2718 key: None,
2719 kind: crate::region::EntryKind::Text,
2720 metadata: None,
2721 taint: crate::taint::TaintLevel::Public,
2722 reasoning: None,
2723 })
2724 .collect(),
2725 description: None,
2726 }
2727 }
2728
2729 fn snap(entries: &[(&str, usize)]) -> ContextSnapshot {
2730 let r = region("logs", entries);
2731 ContextSnapshot {
2732 stage_name: "s".to_string(),
2733 total_tokens: r.current_tokens,
2734 max_tokens: 1000,
2735 regions: vec![r],
2736 }
2737 }
2738
2739 let steps: Vec<ContextSnapshot> = vec![
2740 snap(&[]),
2741 snap(&[("a", 10)]),
2742 snap(&[("a", 10), ("b", 20)]),
2743 snap(&[("a", 10), ("b", 20), ("c", 30)]),
2744 snap(&[("b", 20), ("c", 30)]),
2746 snap(&[("c", 30), ("d", 40)]),
2749 snap(&[("d", 40)]),
2751 snap(&[("d", 40), ("x", 5)]),
2754 snap(&[("x", 5), ("x", 5)]),
2755 snap(&[("x", 5)]),
2756 snap(&[]),
2757 ];
2758
2759 let mut base = steps[0].clone();
2762 let mut digest = digest_context(&steps[0]);
2763 for (i, next) in steps.iter().enumerate().skip(1) {
2764 let delta = diff_context_digest(&digest, next);
2765 apply_delta(&mut base, &delta);
2766 assert_eq!(
2767 base, *next,
2768 "step {i}: replay drifted from the live state\n delta was {:?}",
2769 delta.regions
2770 );
2771 digest = digest_context(next);
2772 }
2773 }
2774
2775 fn trailing_message() -> RunRecord {
2778 RunRecord::Message {
2779 message: MessageRecord {
2780 role: "user".to_string(),
2781 content: "after the unknown".to_string(),
2782 },
2783 at: 9,
2784 }
2785 }
2786
2787 fn archive_with_an_unknown_record() -> (Vec<u8>, usize) {
2790 let mut buf = Vec::new();
2791 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
2792 write_record(&mut buf, &header()).unwrap();
2793 let payload =
2794 serde_json::to_vec(&serde_json::json!({ "SomethingNew": { "whatever": 1 } })).unwrap();
2795 buf.extend_from_slice(&(payload.len() as u64).to_be_bytes());
2796 buf.extend_from_slice(&payload);
2797 write_record(&mut buf, &trailing_message()).unwrap();
2798 (buf, payload.len())
2799 }
2800
2801 fn record_kind(record: &RunRecord) -> String {
2804 serde_json::to_value(record)
2805 .expect("a RunRecord always serializes")
2806 .as_object()
2807 .expect("externally tagged, so an object")
2808 .keys()
2809 .next()
2810 .expect("with exactly one key")
2811 .clone()
2812 }
2813
2814 #[test]
2820 fn an_unknown_record_kind_is_stepped_over_not_treated_as_the_end() {
2821 let (buf, _) = archive_with_an_unknown_record();
2822 let (version, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
2823 assert_eq!(version, RUN_ARCHIVE_VERSION);
2824 let kinds: Vec<String> = records.iter().map(record_kind).collect();
2825 assert_eq!(
2826 kinds,
2827 vec!["Header".to_string(), "Message".to_string()],
2828 "the header, and the readable record after the gap"
2829 );
2830 assert_eq!(records[1], trailing_message(), "intact, not just present");
2831 }
2832
2833 #[test]
2837 fn the_streaming_reader_also_steps_over_an_unknown_record() {
2838 let (mut buf, _) = archive_with_an_unknown_record();
2841 write_record(
2842 &mut buf,
2843 &RunRecord::ContextCheckpoint {
2844 snapshot: ContextSnapshot {
2845 stage_name: "s".to_string(),
2846 total_tokens: 1,
2847 max_tokens: 10,
2848 regions: vec![],
2849 },
2850 at: 11,
2851 },
2852 )
2853 .unwrap();
2854 let mut points = 0usize;
2855 visit_archive_points(&mut buf.as_slice(), &mut |_point| {
2856 points += 1;
2857 ControlFlow::Continue(())
2858 })
2859 .expect("a valid preamble");
2860 assert_eq!(points, 1, "the walk got past the unknown frame");
2861 }
2862
2863 #[test]
2867 fn a_frame_reports_whether_its_payload_was_readable() {
2868 let (buf, unknown_bytes) = archive_with_an_unknown_record();
2869 let mut r = buf.as_slice();
2870 read_archive_start(&mut r).unwrap();
2871
2872 let mut frames = Vec::new();
2873 while let Some(frame) = read_frame(&mut r).expect("no torn frames here") {
2874 frames.push(frame);
2875 }
2876 assert_eq!(
2877 frames,
2878 vec![
2879 Frame::Record(Box::new(header())),
2880 Frame::Unreadable {
2882 bytes: unknown_bytes
2883 },
2884 Frame::Record(Box::new(trailing_message())),
2885 ],
2886 "one frame per record, with the unreadable one accounted for rather than ending the read"
2887 );
2888 }
2889
2890 #[test]
2893 fn a_torn_frame_still_ends_the_read() {
2894 let mut buf = Vec::new();
2895 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
2896 write_record(&mut buf, &header()).unwrap();
2897 buf.extend_from_slice(&999u64.to_be_bytes());
2899 buf.extend_from_slice(b"not enough");
2900
2901 let (_, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
2902 assert_eq!(records.len(), 1, "everything intact before the tear");
2903 }
2904
2905 #[test]
2909 fn an_archive_from_a_newer_format_is_refused_with_both_versions_named() {
2910 let mut buf = Vec::new();
2911 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION + 1).unwrap();
2912 write_record(&mut buf, &header()).unwrap();
2913
2914 let err = read_archive_lenient(&mut buf.as_slice()).unwrap_err();
2915 let message = err.to_string();
2916 assert!(
2917 message.contains(&(RUN_ARCHIVE_VERSION + 1).to_string()),
2918 "{message}"
2919 );
2920 assert!(message.contains("upgrade leviath"), "{message}");
2921 }
2922
2923 #[test]
2926 fn an_archive_from_an_older_format_still_reads() {
2927 let mut buf = Vec::new();
2928 write_archive_start(&mut buf, 0).unwrap();
2929 write_record(&mut buf, &header()).unwrap();
2930 let (version, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
2931 assert_eq!(version, 0);
2932 assert_eq!(records.len(), 1);
2933 }
2934}