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<String>,
86 #[serde(default, skip_serializing_if = "Option::is_none")]
90 pub thought_signature: Option<String>,
91}
92
93#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
96pub struct InferenceRequestRecord {
97 pub model: String,
99 pub system: Vec<String>,
101 pub messages: Vec<MessageRecord>,
103 pub tool_names: Vec<String>,
105 pub temperature: f32,
107 pub max_tokens: usize,
109}
110
111#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
113pub struct InferenceResponseRecord {
114 pub content: String,
116 pub tool_calls: Vec<ToolCallRecord>,
118 pub prompt_tokens: usize,
120 pub completion_tokens: usize,
122 pub cached_tokens: usize,
124 pub cache_write_tokens: usize,
126}
127
128#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
130pub enum RegionDelta {
131 Set(RegionSnapshot),
134 Append {
137 name: String,
139 entries: Vec<RegionEntrySnapshot>,
141 current_tokens: usize,
143 },
144 Clear {
146 name: String,
148 },
149 Remove {
151 name: String,
153 },
154}
155
156#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
158pub struct ContextDelta {
159 pub stage_name: String,
161 pub total_tokens: usize,
163 pub max_tokens: usize,
165 pub regions: Vec<RegionDelta>,
167}
168
169#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
177#[serde(rename_all = "snake_case")]
178pub enum InferenceKind {
179 #[default]
183 Stage,
184 Compaction,
186 Title,
188 Routing,
190}
191
192impl InferenceKind {
193 pub fn label(&self) -> &'static str {
195 match self {
196 InferenceKind::Stage => "stage",
197 InferenceKind::Compaction => "compaction",
198 InferenceKind::Title => "title",
199 InferenceKind::Routing => "routing",
200 }
201 }
202
203 pub fn is_stage_work(&self) -> bool {
206 matches!(self, InferenceKind::Stage)
207 }
208}
209
210#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
212pub enum RunRecord {
213 Header {
215 identity: RunIdentity,
217 meta: Box<RunMeta>,
219 },
220 OwnershipChanged {
222 machine_id: String,
224 world_id: String,
226 at: i64,
228 },
229 Inference {
231 stage: String,
233 iteration: usize,
235 request: InferenceRequestRecord,
237 response: InferenceResponseRecord,
239 at: i64,
241 },
242 InferenceUsage {
257 #[serde(default)]
259 kind: InferenceKind,
260 stage: String,
263 iteration: usize,
265 provider: String,
267 model: String,
269 prompt_tokens: usize,
271 completion_tokens: usize,
273 cached_tokens: usize,
275 cache_write_tokens: usize,
277 at: i64,
279 },
280 ToolBatch {
288 calls: Vec<ToolCallRecord>,
290 at: i64,
292 #[serde(default)]
294 stage_index: usize,
295 #[serde(default)]
298 iteration: usize,
299 #[serde(default)]
301 response: String,
302 },
303 ToolCallDone {
305 iteration: usize,
307 call_id: String,
309 result: String,
311 at: i64,
313 },
314 ContextCheckpoint {
316 snapshot: ContextSnapshot,
318 at: i64,
320 },
321 ContextDiff {
323 delta: ContextDelta,
325 at: i64,
327 },
328 Message {
330 message: MessageRecord,
332 at: i64,
334 },
335 StatusChanged {
337 status: RunStatus,
339 at: i64,
341 },
342 Checkpoint {
345 meta: Box<RunMeta>,
347 context: ContextSnapshot,
349 at: i64,
351 },
352 Progress {
357 meta: Box<RunMeta>,
359 delta: ContextDelta,
361 at: i64,
363 },
364}
365
366fn is_prefix(prev: &[RegionEntrySnapshot], next: &[RegionEntrySnapshot]) -> bool {
370 prev.len() <= next.len() && next[..prev.len()] == *prev
371}
372
373pub fn diff_context(prev: &ContextSnapshot, next: &ContextSnapshot) -> ContextDelta {
377 let mut regions = Vec::new();
378 for nr in &next.regions {
379 match prev.regions.iter().find(|r| r.name == nr.name) {
380 None => regions.push(RegionDelta::Set(nr.clone())),
381 Some(pr) => {
382 if pr == nr {
383 } else if nr.entries.is_empty() && !pr.entries.is_empty() {
385 regions.push(RegionDelta::Clear {
386 name: nr.name.clone(),
387 });
388 } else if pr.kind == nr.kind
389 && pr.max_tokens == nr.max_tokens
390 && is_prefix(&pr.entries, &nr.entries)
391 {
392 regions.push(RegionDelta::Append {
393 name: nr.name.clone(),
394 entries: nr.entries[pr.entries.len()..].to_vec(),
395 current_tokens: nr.current_tokens,
396 });
397 } else {
398 regions.push(RegionDelta::Set(nr.clone()));
399 }
400 }
401 }
402 }
403 for pr in &prev.regions {
404 if !next.regions.iter().any(|r| r.name == pr.name) {
405 regions.push(RegionDelta::Remove {
406 name: pr.name.clone(),
407 });
408 }
409 }
410 ContextDelta {
411 stage_name: next.stage_name.clone(),
412 total_tokens: next.total_tokens,
413 max_tokens: next.max_tokens,
414 regions,
415 }
416}
417
418#[derive(Debug, Clone, PartialEq)]
429pub struct RegionDigest {
430 pub name: String,
432 pub kind: String,
434 pub current_tokens: usize,
436 pub max_tokens: usize,
438 pub entries: Vec<u64>,
440}
441
442#[derive(Debug, Clone, PartialEq, Default)]
444pub struct ContextDigest {
445 pub regions: Vec<RegionDigest>,
447}
448
449fn entry_digest(entry: &RegionEntrySnapshot) -> u64 {
453 use std::hash::{Hash, Hasher};
454 let mut hasher = std::collections::hash_map::DefaultHasher::new();
455 entry.content.hash(&mut hasher);
456 entry.tokens.hash(&mut hasher);
457 entry.key.hash(&mut hasher);
458 serde_json::to_string(&entry.kind)
461 .expect("EntryKind always serializes")
462 .hash(&mut hasher);
463 serde_json::to_string(&entry.metadata)
464 .expect("entry metadata always serializes")
465 .hash(&mut hasher);
466 serde_json::to_string(&entry.taint)
467 .expect("taint always serializes")
468 .hash(&mut hasher);
469 hasher.finish()
470}
471
472pub fn digest_context(snapshot: &ContextSnapshot) -> ContextDigest {
474 ContextDigest {
475 regions: snapshot
476 .regions
477 .iter()
478 .map(|r| RegionDigest {
479 name: r.name.clone(),
480 kind: r.kind.clone(),
481 current_tokens: r.current_tokens,
482 max_tokens: r.max_tokens,
483 entries: r.entries.iter().map(entry_digest).collect(),
484 })
485 .collect(),
486 }
487}
488
489fn is_prefix_digest(prev: &[u64], next: &[RegionEntrySnapshot]) -> bool {
491 prev.len() <= next.len()
492 && prev
493 .iter()
494 .zip(next)
495 .all(|(hash, entry)| *hash == entry_digest(entry))
496}
497
498pub fn diff_context_digest(prev: &ContextDigest, next: &ContextSnapshot) -> ContextDelta {
503 let mut regions = Vec::new();
504 for nr in &next.regions {
505 match prev.regions.iter().find(|r| r.name == nr.name) {
506 None => regions.push(RegionDelta::Set(nr.clone())),
507 Some(pr) => {
508 let unchanged = pr.kind == nr.kind
509 && pr.max_tokens == nr.max_tokens
510 && pr.current_tokens == nr.current_tokens
511 && pr.entries.len() == nr.entries.len()
512 && is_prefix_digest(&pr.entries, &nr.entries);
513 if unchanged {
514 } else if nr.entries.is_empty() && !pr.entries.is_empty() {
516 regions.push(RegionDelta::Clear {
517 name: nr.name.clone(),
518 });
519 } else if pr.kind == nr.kind
520 && pr.max_tokens == nr.max_tokens
521 && is_prefix_digest(&pr.entries, &nr.entries)
522 {
523 regions.push(RegionDelta::Append {
524 name: nr.name.clone(),
525 entries: nr.entries[pr.entries.len()..].to_vec(),
526 current_tokens: nr.current_tokens,
527 });
528 } else {
529 regions.push(RegionDelta::Set(nr.clone()));
530 }
531 }
532 }
533 }
534 for pr in &prev.regions {
535 if !next.regions.iter().any(|r| r.name == pr.name) {
536 regions.push(RegionDelta::Remove {
537 name: pr.name.clone(),
538 });
539 }
540 }
541 ContextDelta {
542 stage_name: next.stage_name.clone(),
543 total_tokens: next.total_tokens,
544 max_tokens: next.max_tokens,
545 regions,
546 }
547}
548
549pub fn apply_delta(base: &mut ContextSnapshot, delta: &ContextDelta) {
553 base.stage_name = delta.stage_name.clone();
554 base.total_tokens = delta.total_tokens;
555 base.max_tokens = delta.max_tokens;
556 for region_delta in &delta.regions {
557 match region_delta {
558 RegionDelta::Set(snapshot) => {
559 match base.regions.iter_mut().find(|r| r.name == snapshot.name) {
560 Some(existing) => *existing = snapshot.clone(),
561 None => base.regions.push(snapshot.clone()),
562 }
563 }
564 RegionDelta::Append {
565 name,
566 entries,
567 current_tokens,
568 } => {
569 if let Some(region) = base.regions.iter_mut().find(|r| &r.name == name) {
570 region.entries.extend(entries.iter().cloned());
571 region.current_tokens = *current_tokens;
572 }
573 }
574 RegionDelta::Clear { name } => {
575 if let Some(region) = base.regions.iter_mut().find(|r| &r.name == name) {
576 region.entries.clear();
577 region.current_tokens = 0;
578 }
579 }
580 RegionDelta::Remove { name } => {
581 base.regions.retain(|r| &r.name != name);
582 }
583 }
584 }
585}
586
587mod codec;
588
589pub use codec::{
590 Frame, RUN_ARCHIVE_MAGIC, RUN_ARCHIVE_VERSION, read_archive, read_archive_lenient,
591 read_archive_start, read_frame, read_record, write_archive_start, write_record,
592};
593
594#[derive(Debug, Clone, PartialEq)]
601pub struct PendingToolBatch {
602 pub stage_index: usize,
604 pub iteration: usize,
606 pub response: String,
608 pub calls: Vec<ToolCallRecord>,
610}
611
612#[derive(Debug, Clone, PartialEq, Eq)]
617pub struct InferenceUsageRecord {
618 pub kind: InferenceKind,
620 pub stage: String,
622 pub iteration: usize,
624 pub provider: String,
626 pub model: String,
628 pub prompt_tokens: usize,
630 pub completion_tokens: usize,
632 pub cached_tokens: usize,
634 pub cache_write_tokens: usize,
636 pub at: i64,
638}
639
640#[derive(Debug, Clone, PartialEq)]
643pub struct FoldedRun {
644 pub identity: RunIdentity,
646 pub meta: RunMeta,
648 pub context: ContextSnapshot,
650 pub messages: Vec<MessageRecord>,
652 pub inference_count: usize,
654 pub inference_usage: Vec<InferenceUsageRecord>,
661 pub tool_call_count: usize,
663 pub pending_batch: Option<PendingToolBatch>,
667}
668
669pub fn context_contains_batch(context: &ContextSnapshot, batch: &PendingToolBatch) -> bool {
673 let Some(first_id) = batch.calls.first().map(|c| c.id.as_str()) else {
674 return false;
675 };
676 context.regions.iter().any(|region| {
677 region.entries.iter().any(|entry| {
678 matches!(
679 &entry.kind,
680 crate::region::EntryKind::AssistantTurn { tool_calls }
681 if tool_calls.iter().any(|tc| tc.id == first_id)
682 )
683 })
684 })
685}
686
687pub fn fold(records: &[RunRecord]) -> Option<FoldedRun> {
690 let mut iter = records.iter();
691 let (identity, meta) = match iter.next() {
692 Some(RunRecord::Header { identity, meta }) => (identity.clone(), (**meta).clone()),
693 _ => return None,
694 };
695 let mut folded = FoldedRun {
696 identity,
697 meta,
698 context: ContextSnapshot {
699 stage_name: String::new(),
700 total_tokens: 0,
701 max_tokens: 0,
702 regions: Vec::new(),
703 },
704 messages: Vec::new(),
705 inference_count: 0,
706 inference_usage: Vec::new(),
707 tool_call_count: 0,
708 pending_batch: None,
709 };
710 for record in iter {
711 match record {
712 RunRecord::Header { identity, meta } => {
713 folded.identity = identity.clone();
714 folded.meta = (**meta).clone();
715 }
716 RunRecord::OwnershipChanged {
717 machine_id,
718 world_id,
719 ..
720 } => {
721 folded.identity.machine_id = machine_id.clone();
722 folded.identity.world_id = world_id.clone();
723 }
724 RunRecord::Inference { .. } => folded.inference_count += 1,
725 RunRecord::InferenceUsage {
726 kind,
727 stage,
728 iteration,
729 provider,
730 model,
731 prompt_tokens,
732 completion_tokens,
733 cached_tokens,
734 cache_write_tokens,
735 at,
736 } => {
737 folded.inference_count += 1;
741 folded.inference_usage.push(InferenceUsageRecord {
742 kind: *kind,
743 stage: stage.clone(),
744 iteration: *iteration,
745 provider: provider.clone(),
746 model: model.clone(),
747 prompt_tokens: *prompt_tokens,
748 completion_tokens: *completion_tokens,
749 cached_tokens: *cached_tokens,
750 cache_write_tokens: *cache_write_tokens,
751 at: *at,
752 });
753 }
754 RunRecord::ToolBatch {
755 calls,
756 stage_index,
757 iteration,
758 response,
759 ..
760 } => {
761 folded.tool_call_count += calls.len();
762 folded.pending_batch = Some(PendingToolBatch {
765 stage_index: *stage_index,
766 iteration: *iteration,
767 response: response.clone(),
768 calls: calls.clone(),
769 });
770 }
771 RunRecord::ToolCallDone {
772 iteration,
773 call_id,
774 result,
775 ..
776 } => {
777 if let Some(batch) = folded
780 .pending_batch
781 .as_mut()
782 .filter(|b| b.iteration == *iteration)
783 && let Some(call) = batch.calls.iter_mut().find(|c| c.id == *call_id)
784 {
785 call.result = Some(result.clone());
786 }
787 }
788 RunRecord::ContextCheckpoint { snapshot, .. } => folded.context = snapshot.clone(),
789 RunRecord::ContextDiff { delta, .. } => apply_delta(&mut folded.context, delta),
790 RunRecord::Message { message, .. } => folded.messages.push(message.clone()),
791 RunRecord::StatusChanged { status, .. } => folded.meta.status = status.clone(),
792 RunRecord::Checkpoint { meta, context, .. } => {
793 folded.meta = (**meta).clone();
794 folded.context = context.clone();
795 }
796 RunRecord::Progress { meta, delta, .. } => {
797 folded.meta = (**meta).clone();
798 apply_delta(&mut folded.context, delta);
799 }
800 }
801 }
802 if let Some(batch) = &folded.pending_batch
807 && (folded.meta.iteration != batch.iteration
808 || context_contains_batch(&folded.context, batch))
809 {
810 folded.pending_batch = None;
811 }
812 Some(folded)
813}
814
815#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
818pub struct RunPoint {
819 pub meta: RunMeta,
821 pub context: ContextSnapshot,
823 pub at: i64,
825}
826
827#[derive(Debug)]
830pub struct PointRef<'a> {
831 pub index: usize,
835 pub at: i64,
837 pub meta: &'a RunMeta,
839 pub context: &'a ContextSnapshot,
841}
842
843pub fn visit_points(records: &[RunRecord], visit: &mut dyn FnMut(PointRef<'_>) -> ControlFlow<()>) {
862 let mut iter = records.iter();
863 let Some(mut folder) = (match iter.next() {
864 Some(first) => PointFolder::start(first),
865 None => None,
866 }) else {
867 return;
868 };
869 for record in iter {
870 if folder.push(record, visit).is_break() {
871 return;
872 }
873 }
874}
875
876pub fn visit_archive_points(
886 r: &mut dyn Read,
887 visit: &mut dyn FnMut(PointRef<'_>) -> ControlFlow<()>,
888) -> io::Result<()> {
889 read_archive_start(r)?;
890 let mut folder = match read_frame(r) {
893 Ok(Some(Frame::Record(first))) => match PointFolder::start(&first) {
894 Some(folder) => folder,
895 None => return Ok(()),
896 },
897 _ => return Ok(()),
898 };
899 while let Ok(Some(frame)) = read_frame(r) {
903 let Frame::Record(record) = frame else {
904 continue;
905 };
906 if folder.push(&record, visit).is_break() {
907 return Ok(());
908 }
909 }
910 Ok(())
911}
912
913struct PointFolder {
918 meta: RunMeta,
919 context: ContextSnapshot,
920 index: usize,
921}
922
923impl PointFolder {
924 fn start(first: &RunRecord) -> Option<Self> {
928 match first {
929 RunRecord::Header { meta, .. } => Some(Self {
930 meta: (**meta).clone(),
931 context: ContextSnapshot {
932 stage_name: String::new(),
933 total_tokens: 0,
934 max_tokens: 0,
935 regions: Vec::new(),
936 },
937 index: 0,
938 }),
939 _ => None,
940 }
941 }
942
943 fn push(
945 &mut self,
946 record: &RunRecord,
947 visit: &mut dyn FnMut(PointRef<'_>) -> ControlFlow<()>,
948 ) -> ControlFlow<()> {
949 let at = match record {
950 RunRecord::Header { meta: m, .. } => {
951 self.meta = (**m).clone();
952 return ControlFlow::Continue(());
953 }
954 RunRecord::StatusChanged { status, .. } => {
955 self.meta.status = status.clone();
956 return ControlFlow::Continue(());
957 }
958 RunRecord::ContextCheckpoint { snapshot, at } => {
959 self.context = snapshot.clone();
960 *at
961 }
962 RunRecord::ContextDiff { delta, at } => {
963 apply_delta(&mut self.context, delta);
964 *at
965 }
966 RunRecord::Checkpoint {
967 meta: m,
968 context: c,
969 at,
970 } => {
971 self.meta = (**m).clone();
972 self.context = c.clone();
973 *at
974 }
975 RunRecord::Progress { meta: m, delta, at } => {
976 self.meta = (**m).clone();
977 apply_delta(&mut self.context, delta);
978 *at
979 }
980 RunRecord::OwnershipChanged { .. }
985 | RunRecord::Inference { .. }
986 | RunRecord::InferenceUsage { .. }
987 | RunRecord::ToolBatch { .. }
988 | RunRecord::ToolCallDone { .. }
989 | RunRecord::Message { .. } => return ControlFlow::Continue(()),
990 };
991 let flow = visit(PointRef {
992 index: self.index,
993 at,
994 meta: &self.meta,
995 context: &self.context,
996 });
997 self.index += 1;
998 flow
999 }
1000}
1001
1002pub fn replay_points(records: &[RunRecord]) -> Vec<RunPoint> {
1012 let mut points = Vec::new();
1013 visit_points(records, &mut |point| {
1014 points.push(RunPoint {
1015 meta: point.meta.clone(),
1016 context: point.context.clone(),
1017 at: point.at,
1018 });
1019 ControlFlow::Continue(())
1020 });
1021 points
1022}
1023
1024#[cfg(test)]
1025mod tests {
1026 use super::*;
1027 use crate::run_meta::RunStatus;
1029 use std::io::Write;
1030
1031 fn identity() -> RunIdentity {
1032 RunIdentity {
1033 run_id: "run-1".to_string(),
1034 machine_id: "machine-a".to_string(),
1035 world_id: "world-x".to_string(),
1036 created_at: 100,
1037 }
1038 }
1039
1040 const FIXTURE_NOW: i64 = 1_700_000_000;
1048
1049 fn meta() -> RunMeta {
1050 let mut meta = RunMeta::new(
1051 "run-1".to_string(),
1052 "coder".to_string(),
1053 "/agents/coder".to_string(),
1054 "do it".to_string(),
1055 Some("anthropic/claude".to_string()),
1056 "/work".to_string(),
1057 2,
1058 );
1059 meta.started_at = FIXTURE_NOW;
1060 meta.updated_at = FIXTURE_NOW;
1061 meta
1062 }
1063
1064 #[test]
1068 fn the_fixture_does_not_move_with_the_clock() {
1069 let first = meta();
1070 let mut later = meta();
1071 assert_eq!(first, later, "the fixture is rebuilt identically");
1075 later.started_at += 1;
1076 assert_ne!(
1077 first, later,
1078 "and the comparison is sensitive to the field that used to drift"
1079 );
1080 }
1081
1082 fn entry(content: &str, tokens: usize) -> RegionEntrySnapshot {
1083 RegionEntrySnapshot {
1084 content: content.to_string(),
1085 tokens,
1086 kind: crate::region::EntryKind::Text,
1087 metadata: None,
1088 key: None,
1089 taint: Default::default(),
1090 }
1091 }
1092
1093 fn region(name: &str, entries: Vec<RegionEntrySnapshot>) -> RegionSnapshot {
1094 let current = entries.iter().map(|e| e.tokens).sum();
1095 RegionSnapshot {
1096 name: name.to_string(),
1097 kind: "clearable".to_string(),
1098 current_tokens: current,
1099 max_tokens: 1000,
1100 entries,
1101 description: None,
1102 }
1103 }
1104
1105 fn snapshot(stage: &str, regions: Vec<RegionSnapshot>) -> ContextSnapshot {
1106 let total = regions.iter().map(|r| r.current_tokens).sum();
1107 ContextSnapshot {
1108 stage_name: stage.to_string(),
1109 total_tokens: total,
1110 max_tokens: 10_000,
1111 regions,
1112 }
1113 }
1114
1115 fn header() -> RunRecord {
1116 RunRecord::Header {
1117 identity: identity(),
1118 meta: Box::new(meta()),
1119 }
1120 }
1121
1122 fn region_delta_kind(d: &RegionDelta) -> &'static str {
1126 match d {
1127 RegionDelta::Set(_) => "set",
1128 RegionDelta::Append { .. } => "append",
1129 RegionDelta::Clear { .. } => "clear",
1130 RegionDelta::Remove { .. } => "remove",
1131 }
1132 }
1133
1134 fn assert_diff_roundtrip(a: &ContextSnapshot, b: &ContextSnapshot) {
1139 let delta = diff_context(a, b);
1140 let mut base = a.clone();
1141 apply_delta(&mut base, &delta);
1142 assert_eq!(&base, b);
1143 }
1144
1145 #[test]
1146 fn diff_append_only_growth_is_compact() {
1147 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1148 let b = snapshot(
1149 "s1",
1150 vec![region("conv", vec![entry("hi", 1), entry("there", 2)])],
1151 );
1152 let delta = diff_context(&a, &b);
1153 assert_eq!(region_delta_kind(&delta.regions[0]), "append");
1154 assert_diff_roundtrip(&a, &b);
1155 }
1156
1157 #[test]
1158 fn diff_new_region_is_set() {
1159 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1160 let b = snapshot(
1161 "s1",
1162 vec![
1163 region("conv", vec![entry("hi", 1)]),
1164 region("plan", vec![entry("p", 3)]),
1165 ],
1166 );
1167 let delta = diff_context(&a, &b);
1168 assert!(delta.regions.iter().any(|d| region_delta_kind(d) == "set"));
1169 assert_diff_roundtrip(&a, &b);
1170 }
1171
1172 #[test]
1173 fn diff_cleared_region() {
1174 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1175 let b = snapshot("s1", vec![region("conv", vec![])]);
1176 let delta = diff_context(&a, &b);
1177 assert_eq!(region_delta_kind(&delta.regions[0]), "clear");
1178 assert_diff_roundtrip(&a, &b);
1179 }
1180
1181 #[test]
1182 fn diff_removed_region() {
1183 let a = snapshot(
1184 "s1",
1185 vec![
1186 region("conv", vec![entry("hi", 1)]),
1187 region("plan", vec![entry("p", 3)]),
1188 ],
1189 );
1190 let b = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1191 let delta = diff_context(&a, &b);
1192 assert!(
1193 delta
1194 .regions
1195 .iter()
1196 .any(|d| region_delta_kind(d) == "remove")
1197 );
1198 assert_diff_roundtrip(&a, &b);
1199 }
1200
1201 #[test]
1202 fn diff_non_prefix_rewrite_is_set() {
1203 let a = snapshot("s1", vec![region("conv", vec![entry("old", 1)])]);
1205 let b = snapshot("s1", vec![region("conv", vec![entry("new", 1)])]);
1206 let delta = diff_context(&a, &b);
1207 assert_eq!(region_delta_kind(&delta.regions[0]), "set");
1208 assert_diff_roundtrip(&a, &b);
1209 }
1210
1211 #[test]
1212 fn diff_kind_change_is_set_not_append() {
1213 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1215 let mut grown = region("conv", vec![entry("hi", 1), entry("more", 1)]);
1216 grown.kind = "sliding".to_string();
1217 let b = snapshot("s1", vec![grown]);
1218 let delta = diff_context(&a, &b);
1219 assert_eq!(region_delta_kind(&delta.regions[0]), "set");
1220 assert_diff_roundtrip(&a, &b);
1221 }
1222
1223 fn framed(records: &[RunRecord]) -> Vec<u8> {
1227 let mut buf = Vec::new();
1228 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1229 for record in records {
1230 write_record(&mut buf, record).unwrap();
1231 }
1232 buf
1233 }
1234
1235 fn try_collect_streamed(bytes: &[u8]) -> io::Result<Vec<(usize, i64, usize)>> {
1239 let mut seen = Vec::new();
1240 visit_archive_points(&mut &bytes[..], &mut |p| {
1241 seen.push((p.index, p.at, p.context.total_tokens));
1242 ControlFlow::Continue(())
1243 })?;
1244 Ok(seen)
1245 }
1246
1247 fn collect_streamed(bytes: &[u8]) -> Vec<(usize, i64, usize)> {
1249 try_collect_streamed(bytes).unwrap()
1250 }
1251
1252 #[test]
1253 fn visit_archive_points_matches_visit_points() {
1254 let records = vec![
1255 header(),
1256 RunRecord::ContextCheckpoint {
1257 snapshot: snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1258 at: 10,
1259 },
1260 RunRecord::StatusChanged {
1261 status: RunStatus::Running,
1262 at: 11,
1263 },
1264 RunRecord::Progress {
1265 meta: Box::new(meta()),
1266 delta: diff_context(
1267 &snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1268 &snapshot(
1269 "s1",
1270 vec![region("conv", vec![entry("hi", 1), entry("more", 2)])],
1271 ),
1272 ),
1273 at: 12,
1274 },
1275 ];
1276 let mut in_memory = Vec::new();
1277 visit_points(&records, &mut |p| {
1278 in_memory.push((p.index, p.at, p.context.total_tokens));
1279 ControlFlow::Continue(())
1280 });
1281 assert_eq!(collect_streamed(&framed(&records)), in_memory);
1282 assert_eq!(in_memory.len(), 2, "checkpoint + progress = two points");
1283 }
1284
1285 #[test]
1286 fn visit_archive_points_rejects_a_bad_preamble() {
1287 assert!(try_collect_streamed(b"not an archive at all").is_err());
1288 }
1289
1290 #[test]
1291 fn visit_archive_points_is_lenient_about_a_torn_tail() {
1292 let records = vec![
1293 header(),
1294 RunRecord::ContextCheckpoint {
1295 snapshot: snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1296 at: 10,
1297 },
1298 ];
1299 let mut bytes = framed(&records);
1300 bytes.extend_from_slice(&1000u64.to_be_bytes());
1302 bytes.extend_from_slice(b"partial");
1303 assert_eq!(collect_streamed(&bytes).len(), 1, "points before the tear");
1304 }
1305
1306 #[test]
1307 fn visit_archive_points_visits_nothing_without_a_header() {
1308 let records = vec![RunRecord::ContextCheckpoint {
1309 snapshot: snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1310 at: 10,
1311 }];
1312 assert!(collect_streamed(&framed(&records)).is_empty());
1313 assert!(collect_streamed(&framed(&[])).is_empty());
1315 }
1316
1317 #[test]
1318 fn visit_archive_points_stops_on_break() {
1319 let records = vec![
1320 header(),
1321 RunRecord::ContextCheckpoint {
1322 snapshot: snapshot("s1", vec![region("conv", vec![entry("a", 1)])]),
1323 at: 10,
1324 },
1325 RunRecord::ContextCheckpoint {
1326 snapshot: snapshot("s1", vec![region("conv", vec![entry("b", 2)])]),
1327 at: 11,
1328 },
1329 ];
1330 let bytes = framed(&records);
1331 let mut seen = 0;
1332 visit_archive_points(&mut &bytes[..], &mut |_| {
1333 seen += 1;
1334 ControlFlow::Break(())
1335 })
1336 .unwrap();
1337 assert_eq!(seen, 1);
1338 }
1339
1340 fn assert_digest_matches_full_diff(a: &ContextSnapshot, b: &ContextSnapshot) {
1347 let via_digest = diff_context_digest(&digest_context(a), b);
1348 assert_eq!(via_digest, diff_context(a, b));
1349 let mut base = a.clone();
1351 apply_delta(&mut base, &via_digest);
1352 assert_eq!(&base, b);
1353 }
1354
1355 #[test]
1356 fn digest_diff_append_only_growth_is_compact() {
1357 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1358 let b = snapshot(
1359 "s1",
1360 vec![region("conv", vec![entry("hi", 1), entry("there", 2)])],
1361 );
1362 let delta = diff_context_digest(&digest_context(&a), &b);
1363 assert_eq!(region_delta_kind(&delta.regions[0]), "append");
1364 assert_digest_matches_full_diff(&a, &b);
1365 }
1366
1367 #[test]
1368 fn digest_diff_new_cleared_removed_and_rewritten_regions() {
1369 let a = snapshot(
1370 "s1",
1371 vec![
1372 region("conv", vec![entry("hi", 1)]),
1373 region("gone", vec![entry("bye", 1)]),
1374 region("wiped", vec![entry("w", 1)]),
1375 region("rewritten", vec![entry("old", 1)]),
1376 ],
1377 );
1378 let b = snapshot(
1379 "s1",
1380 vec![
1381 region("conv", vec![entry("hi", 1)]),
1382 region("wiped", vec![]),
1383 region("rewritten", vec![entry("new", 1)]),
1384 region("fresh", vec![entry("f", 2)]),
1385 ],
1386 );
1387 let delta = diff_context_digest(&digest_context(&a), &b);
1388 let kinds: Vec<_> = delta.regions.iter().map(region_delta_kind).collect();
1389 assert_eq!(kinds, vec!["clear", "set", "set", "remove"]);
1390 assert_digest_matches_full_diff(&a, &b);
1391 }
1392
1393 #[test]
1394 fn digest_diff_unchanged_region_emits_nothing() {
1395 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1396 let delta = diff_context_digest(&digest_context(&a), &a.clone());
1397 assert!(delta.regions.is_empty());
1398 assert_digest_matches_full_diff(&a, &a.clone());
1399 }
1400
1401 #[test]
1402 fn digest_diff_kind_change_is_set_not_append() {
1403 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1404 let mut grown = region("conv", vec![entry("hi", 1), entry("more", 1)]);
1405 grown.kind = "sliding".to_string();
1406 let b = snapshot("s1", vec![grown]);
1407 let delta = diff_context_digest(&digest_context(&a), &b);
1408 assert_eq!(region_delta_kind(&delta.regions[0]), "set");
1409 assert_digest_matches_full_diff(&a, &b);
1410 }
1411
1412 #[test]
1415 fn digest_diff_token_recount_is_an_empty_append() {
1416 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1417 let mut recounted = region("conv", vec![entry("hi", 1)]);
1418 recounted.current_tokens = 42;
1419 let b = snapshot("s1", vec![recounted]);
1420 let delta = diff_context_digest(&digest_context(&a), &b);
1421 assert_eq!(region_delta_kind(&delta.regions[0]), "append");
1422 assert_digest_matches_full_diff(&a, &b);
1423 }
1424
1425 #[test]
1428 fn entry_digest_covers_every_field() {
1429 let base = entry("text", 1);
1430 let variants = [
1431 entry("other", 1),
1432 entry("text", 2),
1433 RegionEntrySnapshot {
1434 key: Some("k".to_string()),
1435 ..entry("text", 1)
1436 },
1437 RegionEntrySnapshot {
1438 metadata: Some(serde_json::json!({"a": 1})),
1439 ..entry("text", 1)
1440 },
1441 RegionEntrySnapshot {
1442 kind: crate::region::EntryKind::ToolResult {
1443 tool_call_id: "c1".to_string(),
1444 tool_name: "shell".to_string(),
1445 is_error: false,
1446 },
1447 ..entry("text", 1)
1448 },
1449 ];
1450 let base_hash = entry_digest(&base);
1451 for variant in &variants {
1452 assert_ne!(
1453 entry_digest(variant),
1454 base_hash,
1455 "field change must change the digest: {variant:?}"
1456 );
1457 }
1458 assert_eq!(entry_digest(&base), entry_digest(&entry("text", 1)));
1460 }
1461
1462 #[test]
1463 fn diff_unchanged_region_emits_nothing() {
1464 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1465 let b = a.clone();
1466 let delta = diff_context(&a, &b);
1467 assert!(delta.regions.is_empty());
1468 assert_diff_roundtrip(&a, &b);
1469 }
1470
1471 #[test]
1472 fn apply_delta_skips_unknown_regions_leniently() {
1473 let mut base = snapshot("s1", vec![]);
1475 let delta = ContextDelta {
1476 stage_name: "s1".to_string(),
1477 total_tokens: 0,
1478 max_tokens: 10_000,
1479 regions: vec![
1480 RegionDelta::Append {
1481 name: "ghost".to_string(),
1482 entries: vec![entry("x", 1)],
1483 current_tokens: 1,
1484 },
1485 RegionDelta::Clear {
1486 name: "ghost".to_string(),
1487 },
1488 RegionDelta::Remove {
1489 name: "ghost".to_string(),
1490 },
1491 ],
1492 };
1493 apply_delta(&mut base, &delta);
1494 assert!(base.regions.is_empty());
1495 }
1496
1497 fn all_record_kinds() -> Vec<RunRecord> {
1500 vec![
1501 header(),
1502 RunRecord::OwnershipChanged {
1503 machine_id: "machine-b".to_string(),
1504 world_id: "world-y".to_string(),
1505 at: 101,
1506 },
1507 RunRecord::Inference {
1508 stage: "plan".to_string(),
1509 iteration: 0,
1510 request: InferenceRequestRecord {
1511 model: "m".to_string(),
1512 system: vec!["sys".to_string()],
1513 messages: vec![MessageRecord {
1514 role: "user".to_string(),
1515 content: "hi".to_string(),
1516 }],
1517 tool_names: vec!["read_file".to_string()],
1518 temperature: 0.7,
1519 max_tokens: 1024,
1520 },
1521 response: InferenceResponseRecord {
1522 content: "ok".to_string(),
1523 tool_calls: vec![],
1524 prompt_tokens: 10,
1525 completion_tokens: 5,
1526 cached_tokens: 0,
1527 cache_write_tokens: 0,
1528 },
1529 at: 102,
1530 },
1531 RunRecord::InferenceUsage {
1532 kind: InferenceKind::Compaction,
1533 stage: "plan".to_string(),
1534 iteration: 2,
1535 provider: "anthropic".to_string(),
1536 model: "claude-sonnet-5".to_string(),
1537 prompt_tokens: 7000,
1538 completion_tokens: 70,
1539 cached_tokens: 12,
1540 cache_write_tokens: 34,
1541 at: 102,
1542 },
1543 RunRecord::ToolBatch {
1544 calls: vec![ToolCallRecord {
1545 id: "c1".to_string(),
1546 name: "read_file".to_string(),
1547 arguments: "{}".to_string(),
1548 result: Some("body".to_string()),
1549 thought_signature: Some("sig".to_string()),
1550 }],
1551 at: 103,
1552 stage_index: 0,
1553 iteration: 0,
1554 response: "reading".to_string(),
1555 },
1556 RunRecord::ToolCallDone {
1557 iteration: 0,
1558 call_id: "c1".to_string(),
1559 result: "body".to_string(),
1560 at: 103,
1561 },
1562 RunRecord::ContextCheckpoint {
1563 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1564 at: 104,
1565 },
1566 RunRecord::ContextDiff {
1567 delta: ContextDelta {
1568 stage_name: "plan".to_string(),
1569 total_tokens: 3,
1570 max_tokens: 10_000,
1571 regions: vec![RegionDelta::Append {
1572 name: "conv".to_string(),
1573 entries: vec![entry("more", 2)],
1574 current_tokens: 3,
1575 }],
1576 },
1577 at: 105,
1578 },
1579 RunRecord::Message {
1580 message: MessageRecord {
1581 role: "user".to_string(),
1582 content: "another".to_string(),
1583 },
1584 at: 106,
1585 },
1586 RunRecord::StatusChanged {
1587 status: RunStatus::Complete,
1588 at: 107,
1589 },
1590 RunRecord::Checkpoint {
1591 meta: Box::new(meta()),
1592 context: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1593 at: 108,
1594 },
1595 RunRecord::Progress {
1596 meta: Box::new(meta()),
1597 delta: ContextDelta {
1598 stage_name: "plan".to_string(),
1599 total_tokens: 3,
1600 max_tokens: 10_000,
1601 regions: vec![RegionDelta::Append {
1602 name: "conv".to_string(),
1603 entries: vec![entry("step", 2)],
1604 current_tokens: 3,
1605 }],
1606 },
1607 at: 109,
1608 },
1609 ]
1610 }
1611
1612 #[test]
1613 fn archive_write_then_read_roundtrips_every_record_kind() {
1614 let records = all_record_kinds();
1615 let mut buf = Vec::new();
1616 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1617 for r in &records {
1618 write_record(&mut buf, r).unwrap();
1619 }
1620 let (version, read) = read_archive(&mut buf.as_slice()).unwrap();
1621 assert_eq!(version, RUN_ARCHIVE_VERSION);
1622 assert_eq!(read, records);
1623 }
1624
1625 #[test]
1626 fn read_archive_start_rejects_bad_magic() {
1627 let mut bytes: &[u8] = b"XXXX\x00\x01";
1628 let err = read_archive_start(&mut bytes).unwrap_err();
1629 assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1630 }
1631
1632 #[test]
1637 fn read_archive_start_reports_version() {
1638 let mut buf = Vec::new();
1639 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1640 assert_eq!(
1641 read_archive_start(&mut buf.as_slice()).unwrap(),
1642 RUN_ARCHIVE_VERSION
1643 );
1644 }
1645
1646 #[test]
1647 fn read_record_returns_none_at_clean_eof() {
1648 let empty: &[u8] = &[];
1649 assert!(read_record(&mut { empty }).unwrap().is_none());
1650 }
1651
1652 #[test]
1653 fn read_record_errors_on_truncated_length_prefix() {
1654 let mut bytes: &[u8] = &[0, 0];
1656 let err = read_record(&mut bytes).unwrap_err();
1657 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1658 }
1659
1660 #[test]
1661 fn read_record_errors_on_truncated_payload() {
1662 let mut bytes: &[u8] = &[0, 0, 0, 0, 0, 0, 0, 10, 1, 2];
1664 let err = read_record(&mut bytes).unwrap_err();
1665 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1666 }
1667
1668 #[test]
1669 fn read_record_errors_on_empty_payload_at_boundary() {
1670 let mut bytes: &[u8] = &[0, 0, 0, 0, 0, 0, 0, 10];
1673 let err = read_record(&mut bytes).unwrap_err();
1674 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1675 }
1676
1677 #[test]
1678 fn read_record_errors_on_invalid_json_payload() {
1679 let mut buf = Vec::new();
1681 let bad = b"not json";
1682 buf.extend_from_slice(&(bad.len() as u64).to_be_bytes());
1683 buf.extend_from_slice(bad);
1684 let err = read_record(&mut buf.as_slice()).unwrap_err();
1685 assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1686 }
1687
1688 struct FailingReader;
1691 impl Read for FailingReader {
1692 fn read(&mut self, _buf: &mut [u8]) -> io::Result<usize> {
1693 Err(io::Error::other("device error"))
1694 }
1695 }
1696
1697 #[test]
1698 fn read_record_propagates_reader_errors() {
1699 let err = read_record(&mut FailingReader).unwrap_err();
1700 assert_eq!(err.kind(), io::ErrorKind::Other);
1701 }
1702
1703 #[test]
1704 fn read_archive_propagates_a_bad_preamble() {
1705 let mut bytes: &[u8] = b"LV";
1707 assert!(read_archive(&mut bytes).is_err());
1708 }
1709
1710 #[test]
1711 fn read_archive_propagates_a_bad_frame() {
1712 let mut buf = Vec::new();
1714 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1715 buf.extend_from_slice(&[0, 0, 0, 0, 0, 0, 0, 5, 1, 2]); let err = read_archive(&mut buf.as_slice()).unwrap_err();
1717 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1718 }
1719
1720 struct FailAfter {
1722 remaining: usize,
1723 }
1724 impl Write for FailAfter {
1725 fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
1726 if self.remaining == 0 {
1727 return Err(io::Error::other("disk full"));
1728 }
1729 let n = buf.len().min(self.remaining);
1730 self.remaining -= n;
1731 Ok(n)
1732 }
1733 fn flush(&mut self) -> io::Result<()> {
1734 Ok(())
1735 }
1736 }
1737
1738 #[test]
1739 fn fail_after_writer_flush_is_a_noop() {
1740 assert!(FailAfter { remaining: 1 }.flush().is_ok());
1741 }
1742
1743 #[test]
1744 fn write_archive_start_propagates_write_errors() {
1745 assert!(write_archive_start(&mut FailAfter { remaining: 0 }, 1).is_err());
1747 assert!(write_archive_start(&mut FailAfter { remaining: 4 }, 1).is_err());
1748 }
1749
1750 #[test]
1751 fn write_record_propagates_write_errors() {
1752 let rec = header();
1753 assert!(write_record(&mut FailAfter { remaining: 0 }, &rec).is_err());
1755 assert!(write_record(&mut FailAfter { remaining: 8 }, &rec).is_err());
1756 }
1757
1758 #[test]
1764 fn an_absurd_frame_length_is_an_error_not_an_allocation() {
1765 let mut buf = Vec::new();
1766 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1767 write_record(&mut buf, &header()).unwrap();
1768 buf.extend_from_slice(&u64::MAX.to_be_bytes());
1770
1771 let err = read_archive(&mut buf.as_slice())
1772 .expect_err("the strict reader must refuse an impossible frame");
1773 assert_eq!(err.kind(), io::ErrorKind::InvalidData, "{err}");
1774
1775 let (_, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
1778 assert_eq!(records, vec![header()]);
1779 }
1780
1781 #[test]
1782 fn read_archive_lenient_matches_strict_on_a_clean_archive() {
1783 let records = all_record_kinds();
1786 let mut buf = Vec::new();
1787 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1788 for r in &records {
1789 write_record(&mut buf, r).unwrap();
1790 }
1791 let (version, read) = read_archive_lenient(&mut buf.as_slice()).unwrap();
1792 assert_eq!(version, RUN_ARCHIVE_VERSION);
1793 assert_eq!(read, records);
1794 }
1795
1796 #[test]
1797 fn read_archive_lenient_keeps_valid_prefix_before_a_torn_tail() {
1798 let mut buf = Vec::new();
1802 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1803 write_record(&mut buf, &header()).unwrap();
1804 write_record(
1805 &mut buf,
1806 &RunRecord::ContextCheckpoint {
1807 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1808 at: 1,
1809 },
1810 )
1811 .unwrap();
1812 buf.extend_from_slice(&[0, 0, 0, 0, 0, 0, 0, 10, 1, 2]);
1814
1815 assert!(read_archive(&mut buf.as_slice()).is_err());
1817 let (version, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
1819 assert_eq!(version, RUN_ARCHIVE_VERSION);
1820 assert_eq!(records.len(), 2);
1821 let folded = fold(&records).expect("prefix starts with a Header");
1822 assert_eq!(folded.context.regions[0].entries.len(), 1);
1823 }
1824
1825 #[test]
1826 fn read_archive_lenient_still_errors_on_a_bad_preamble() {
1827 let mut bad_magic: &[u8] = b"XXXX\x00\x01";
1830 assert!(read_archive_lenient(&mut bad_magic).is_err());
1831 let mut short: &[u8] = b"LVR1";
1833 assert!(read_archive_lenient(&mut short).is_err());
1834 }
1835
1836 #[test]
1837 fn read_archive_start_errors_on_truncated_version() {
1838 let mut bytes: &[u8] = b"LVR1";
1840 let err = read_archive_start(&mut bytes).unwrap_err();
1841 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1842 }
1843
1844 #[test]
1847 fn fold_requires_a_header_first() {
1848 assert!(fold(&[]).is_none());
1849 assert!(
1850 fold(&[RunRecord::StatusChanged {
1851 status: RunStatus::Complete,
1852 at: 1
1853 }])
1854 .is_none()
1855 );
1856 }
1857
1858 #[test]
1859 fn fold_reconstructs_state_from_the_journal() {
1860 let records = all_record_kinds();
1861 let folded = fold(&records).expect("has header");
1862 assert_eq!(folded.identity.machine_id, "machine-b");
1864 assert_eq!(folded.identity.world_id, "world-y");
1865 assert_eq!(folded.inference_count, 2);
1868 assert_eq!(folded.inference_usage.len(), 1);
1869 assert_eq!(folded.tool_call_count, 1);
1870 assert_eq!(folded.messages.len(), 1);
1872 assert_eq!(folded.messages[0].content, "another");
1873 assert_eq!(folded.context.regions[0].name, "conv");
1876 assert_eq!(folded.context.regions[0].entries.len(), 2);
1877 assert_eq!(folded.context.total_tokens, 3);
1878 assert_eq!(folded.meta.run_id, "run-1");
1879 let pending = folded.pending_batch.expect("batch never applied");
1882 assert_eq!(pending.calls[0].result.as_deref(), Some("body"));
1883 }
1884
1885 #[test]
1886 fn fold_applies_context_diffs_over_a_checkpoint() {
1887 let records = vec![
1889 header(),
1890 RunRecord::ContextCheckpoint {
1891 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1892 at: 1,
1893 },
1894 RunRecord::ContextDiff {
1895 delta: ContextDelta {
1896 stage_name: "plan".to_string(),
1897 total_tokens: 3,
1898 max_tokens: 10_000,
1899 regions: vec![RegionDelta::Append {
1900 name: "conv".to_string(),
1901 entries: vec![entry("there", 2)],
1902 current_tokens: 3,
1903 }],
1904 },
1905 at: 2,
1906 },
1907 ];
1908 let folded = fold(&records).unwrap();
1909 assert_eq!(folded.context.regions[0].entries.len(), 2);
1910 assert_eq!(folded.context.total_tokens, 3);
1911 }
1912
1913 #[test]
1914 fn fold_later_header_updates_identity_and_meta() {
1915 let mut second_meta = meta();
1917 second_meta.status = RunStatus::Running;
1918 let records = vec![
1919 header(),
1920 RunRecord::Header {
1921 identity: RunIdentity {
1922 run_id: "run-1".to_string(),
1923 machine_id: "machine-c".to_string(),
1924 world_id: "world-z".to_string(),
1925 created_at: 200,
1926 },
1927 meta: Box::new(second_meta),
1928 },
1929 ];
1930 let folded = fold(&records).unwrap();
1931 assert_eq!(folded.identity.machine_id, "machine-c");
1932 assert_eq!(folded.meta.status, RunStatus::Running);
1933 }
1934
1935 #[test]
1936 fn fold_progress_applies_meta_and_context_diff() {
1937 let mut advanced = meta();
1938 advanced.status = RunStatus::Running;
1939 advanced.iteration = 5;
1940 let records = vec![
1941 header(),
1942 RunRecord::ContextCheckpoint {
1943 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1944 at: 1,
1945 },
1946 RunRecord::Progress {
1947 meta: Box::new(advanced),
1948 delta: ContextDelta {
1949 stage_name: "plan".to_string(),
1950 total_tokens: 3,
1951 max_tokens: 10_000,
1952 regions: vec![RegionDelta::Append {
1953 name: "conv".to_string(),
1954 entries: vec![entry("there", 2)],
1955 current_tokens: 3,
1956 }],
1957 },
1958 at: 2,
1959 },
1960 ];
1961 let folded = fold(&records).unwrap();
1962 assert_eq!(folded.meta.iteration, 5);
1963 assert_eq!(folded.meta.status, RunStatus::Running);
1964 assert_eq!(folded.context.regions[0].entries.len(), 2);
1965 }
1966
1967 #[test]
1971 fn fold_carries_a_submitted_final_output_through_progress() {
1972 let mut answered = meta();
1973 answered.final_output = Some(
1974 crate::output::FinalOutput::new(
1975 "renamed two helpers",
1976 Some("markdown".to_string()),
1977 "summary".to_string(),
1978 9,
1979 )
1980 .descriptor(),
1981 );
1982 answered.output_request = Some(crate::output::OutputSpec {
1983 format: Some("a2ui".to_string()),
1984 ..Default::default()
1985 });
1986 let records = vec![
1987 header(),
1988 RunRecord::Progress {
1989 meta: Box::new(answered),
1990 delta: ContextDelta {
1991 stage_name: "summary".to_string(),
1992 total_tokens: 0,
1993 max_tokens: 10_000,
1994 regions: vec![],
1995 },
1996 at: 2,
1997 },
1998 ];
1999 let folded = fold(&records).unwrap();
2000 let output = folded.meta.final_output.expect("the answer folded through");
2001 assert_eq!(output.bytes, "renamed two helpers".len());
2004 assert_eq!(output.stage, "summary");
2005 assert_eq!(
2006 folded.meta.output_request.and_then(|s| s.format).as_deref(),
2007 Some("a2ui")
2008 );
2009 }
2010
2011 fn call(id: &str, result: Option<&str>) -> ToolCallRecord {
2014 ToolCallRecord {
2015 id: id.to_string(),
2016 name: "shell".to_string(),
2017 arguments: "{}".to_string(),
2018 result: result.map(str::to_string),
2019 thought_signature: None,
2020 }
2021 }
2022
2023 fn batch(iteration: usize, calls: Vec<ToolCallRecord>) -> RunRecord {
2024 RunRecord::ToolBatch {
2025 calls,
2026 at: 10,
2027 stage_index: 0,
2028 iteration,
2029 response: "running tools".to_string(),
2030 }
2031 }
2032
2033 fn turn_entry(call_ids: &[&str]) -> RegionEntrySnapshot {
2035 let mut e = entry("turn", 1);
2036 e.kind = crate::region::EntryKind::AssistantTurn {
2037 tool_calls: call_ids
2038 .iter()
2039 .map(|id| crate::region::SerializedToolCall {
2040 id: id.to_string(),
2041 name: "shell".to_string(),
2042 arguments: serde_json::Value::Null,
2043 thought_signature: None,
2044 })
2045 .collect(),
2046 };
2047 e
2048 }
2049
2050 #[test]
2051 fn fold_surfaces_a_pending_batch_with_merged_results() {
2052 let records = vec![
2057 header(),
2058 batch(
2059 0,
2060 vec![
2061 call("c1", None),
2062 call("c2", Some("inline")),
2063 call("c3", None),
2064 ],
2065 ),
2066 RunRecord::ToolCallDone {
2067 iteration: 0,
2068 call_id: "c1".to_string(),
2069 result: "ran".to_string(),
2070 at: 11,
2071 },
2072 ];
2073 let folded = fold(&records).unwrap();
2074 let pending = folded.pending_batch.expect("batch is pending");
2075 assert_eq!(pending.iteration, 0);
2076 assert_eq!(pending.response, "running tools");
2077 assert_eq!(pending.calls[0].result.as_deref(), Some("ran"));
2078 assert_eq!(pending.calls[1].result.as_deref(), Some("inline"));
2079 assert_eq!(pending.calls[2].result, None);
2080 assert_eq!(folded.tool_call_count, 3);
2081 }
2082
2083 #[test]
2084 fn fold_keeps_only_the_latest_batch_and_ignores_stale_done_records() {
2085 let mut advanced = meta();
2088 advanced.iteration = 1;
2089 let records = vec![
2090 header(),
2091 batch(0, vec![call("c1", None)]),
2092 RunRecord::Progress {
2093 meta: Box::new(advanced),
2094 delta: ContextDelta {
2095 stage_name: "plan".to_string(),
2096 total_tokens: 0,
2097 max_tokens: 10_000,
2098 regions: vec![],
2099 },
2100 at: 11,
2101 },
2102 batch(1, vec![call("c2", None)]),
2103 RunRecord::ToolCallDone {
2104 iteration: 0,
2105 call_id: "c1".to_string(),
2106 result: "stale".to_string(),
2107 at: 12,
2108 },
2109 RunRecord::ToolCallDone {
2110 iteration: 1,
2111 call_id: "unknown".to_string(),
2112 result: "nowhere to land".to_string(),
2113 at: 13,
2114 },
2115 ];
2116 let folded = fold(&records).unwrap();
2117 let pending = folded.pending_batch.expect("latest batch is pending");
2118 assert_eq!(pending.iteration, 1);
2119 assert_eq!(pending.calls.len(), 1);
2120 assert_eq!(pending.calls[0].id, "c2");
2121 assert_eq!(pending.calls[0].result, None, "stale/unknown dones ignored");
2122 }
2123
2124 #[test]
2125 fn fold_clears_a_batch_once_the_iteration_moves_on() {
2126 let mut advanced = meta();
2129 advanced.iteration = 1;
2130 let records = vec![
2131 header(),
2132 batch(0, vec![call("c1", Some("done"))]),
2133 RunRecord::Progress {
2134 meta: Box::new(advanced),
2135 delta: ContextDelta {
2136 stage_name: "plan".to_string(),
2137 total_tokens: 0,
2138 max_tokens: 10_000,
2139 regions: vec![],
2140 },
2141 at: 11,
2142 },
2143 ];
2144 assert_eq!(fold(&records).unwrap().pending_batch, None);
2145 }
2146
2147 #[test]
2148 fn fold_clears_a_batch_whose_turn_already_landed_in_the_window() {
2149 let records = vec![
2152 header(),
2153 batch(0, vec![call("c1", Some("done"))]),
2154 RunRecord::ContextCheckpoint {
2155 snapshot: snapshot("plan", vec![region("conv", vec![turn_entry(&["c1"])])]),
2156 at: 11,
2157 },
2158 ];
2159 assert_eq!(fold(&records).unwrap().pending_batch, None);
2160 }
2161
2162 #[test]
2163 fn context_contains_batch_matches_only_the_batch_turn() {
2164 let pending = PendingToolBatch {
2165 stage_index: 0,
2166 iteration: 0,
2167 response: String::new(),
2168 calls: vec![call("c1", None)],
2169 };
2170 let other = snapshot("plan", vec![region("conv", vec![turn_entry(&["zz"])])]);
2172 assert!(!context_contains_batch(&other, &pending));
2173 let own = snapshot(
2175 "plan",
2176 vec![region("conv", vec![turn_entry(&["c1", "c2"])])],
2177 );
2178 assert!(context_contains_batch(&own, &pending));
2179 let empty = PendingToolBatch {
2181 calls: vec![],
2182 ..pending
2183 };
2184 assert!(!context_contains_batch(&own, &empty));
2185 }
2186
2187 #[test]
2188 fn old_shape_tool_batch_json_still_parses() {
2189 let json = br#"{"ToolBatch":{"calls":[{"id":"c1","name":"shell","arguments":"{}","result":"ok"}],"at":9}}"#;
2193 let mut buf = Vec::new();
2194 buf.extend_from_slice(&(json.len() as u64).to_be_bytes());
2195 buf.extend_from_slice(json);
2196 let record = read_record(&mut buf.as_slice()).unwrap().unwrap();
2197 assert_eq!(
2198 record,
2199 RunRecord::ToolBatch {
2200 calls: vec![call("c1", Some("ok"))],
2201 at: 9,
2202 stage_index: 0,
2203 iteration: 0,
2204 response: String::new(),
2205 }
2206 );
2207 }
2208
2209 fn three_point_records() -> Vec<RunRecord> {
2213 let mut running = meta();
2214 running.status = RunStatus::Running;
2215 vec![
2216 header(),
2217 RunRecord::ContextCheckpoint {
2218 snapshot: snapshot("plan", vec![region("conv", vec![entry("first", 1)])]),
2219 at: 10,
2220 },
2221 RunRecord::ContextDiff {
2222 delta: ContextDelta {
2223 stage_name: "plan".to_string(),
2224 total_tokens: 2,
2225 max_tokens: 10_000,
2226 regions: vec![RegionDelta::Append {
2227 name: "conv".to_string(),
2228 entries: vec![entry("second", 1)],
2229 current_tokens: 2,
2230 }],
2231 },
2232 at: 20,
2233 },
2234 RunRecord::Progress {
2235 meta: Box::new(running),
2236 delta: ContextDelta {
2237 stage_name: "code".to_string(),
2238 total_tokens: 3,
2239 max_tokens: 10_000,
2240 regions: vec![RegionDelta::Append {
2241 name: "conv".to_string(),
2242 entries: vec![entry("third", 1)],
2243 current_tokens: 3,
2244 }],
2245 },
2246 at: 30,
2247 },
2248 ]
2249 }
2250
2251 #[test]
2252 fn visit_points_indexes_points_in_order_and_carries_the_running_window() {
2253 let records = three_point_records();
2254 let mut seen: Vec<(usize, i64, usize)> = Vec::new();
2255 visit_points(&records, &mut |point| {
2256 seen.push((
2257 point.index,
2258 point.at,
2259 point.context.regions[0].entries.len(),
2260 ));
2261 ControlFlow::Continue(())
2262 });
2263 assert_eq!(seen, vec![(0, 10, 1), (1, 20, 2), (2, 30, 3)]);
2265 }
2266
2267 #[test]
2270 fn visit_points_stops_at_the_first_break() {
2271 let records = three_point_records();
2272 let mut visits = 0;
2273 visit_points(&records, &mut |point| {
2274 visits += 1;
2275 if point.index == 1 {
2276 ControlFlow::Break(())
2277 } else {
2278 ControlFlow::Continue(())
2279 }
2280 });
2281 assert_eq!(
2282 visits, 2,
2283 "stopped at the breaking point, did not run the third"
2284 );
2285 }
2286
2287 #[test]
2288 fn visit_points_without_a_header_visits_nothing() {
2289 let mut visits = 0;
2290 {
2291 let mut count = |_: PointRef<'_>| {
2292 visits += 1;
2293 ControlFlow::Continue(())
2294 };
2295
2296 visit_points(&three_point_records(), &mut count);
2300 visit_points(&[], &mut count);
2303 visit_points(
2304 &[RunRecord::ContextCheckpoint {
2305 snapshot: snapshot("plan", vec![]),
2306 at: 1,
2307 }],
2308 &mut count,
2309 );
2310 }
2311 assert_eq!(visits, 3, "only the well-formed journal produced points");
2312 }
2313
2314 #[test]
2318 fn visit_points_and_replay_points_agree() {
2319 for records in [
2320 three_point_records(),
2321 vec![header()],
2322 vec![],
2323 vec![RunRecord::Message {
2324 message: MessageRecord {
2325 role: "user".to_string(),
2326 content: "x".to_string(),
2327 },
2328 at: 1,
2329 }],
2330 ] {
2331 let collected: Vec<RunPoint> = {
2332 let mut out = Vec::new();
2333 visit_points(&records, &mut |point| {
2334 out.push(RunPoint {
2335 meta: point.meta.clone(),
2336 context: point.context.clone(),
2337 at: point.at,
2338 });
2339 ControlFlow::Continue(())
2340 });
2341 out
2342 };
2343 assert_eq!(collected, replay_points(&records));
2344 }
2345 }
2346
2347 #[test]
2348 fn replay_points_requires_a_header() {
2349 assert!(replay_points(&[]).is_empty());
2350 assert!(
2351 replay_points(&[RunRecord::Message {
2352 message: MessageRecord {
2353 role: "user".to_string(),
2354 content: "x".to_string(),
2355 },
2356 at: 1,
2357 }])
2358 .is_empty()
2359 );
2360 }
2361
2362 #[test]
2363 fn replay_points_emits_a_snapshot_per_context_change() {
2364 let mut running = meta();
2367 running.status = RunStatus::Running;
2368 let records = vec![
2369 header(),
2370 RunRecord::Inference {
2371 stage: "plan".to_string(),
2372 iteration: 0,
2373 request: InferenceRequestRecord {
2374 model: "m".to_string(),
2375 system: vec![],
2376 messages: vec![],
2377 tool_names: vec![],
2378 temperature: 0.7,
2379 max_tokens: 10,
2380 },
2381 response: InferenceResponseRecord {
2382 content: "ok".to_string(),
2383 tool_calls: vec![],
2384 prompt_tokens: 1,
2385 completion_tokens: 1,
2386 cached_tokens: 0,
2387 cache_write_tokens: 0,
2388 },
2389 at: 1,
2390 },
2391 batch(0, vec![call("c1", None)]),
2392 RunRecord::ToolCallDone {
2393 iteration: 0,
2394 call_id: "c1".to_string(),
2395 result: "ran".to_string(),
2396 at: 1,
2397 },
2398 RunRecord::ContextCheckpoint {
2399 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
2400 at: 2,
2401 },
2402 RunRecord::StatusChanged {
2403 status: RunStatus::Running,
2404 at: 3,
2405 },
2406 RunRecord::Progress {
2407 meta: Box::new(running),
2408 delta: ContextDelta {
2409 stage_name: "implement".to_string(),
2410 total_tokens: 3,
2411 max_tokens: 10_000,
2412 regions: vec![RegionDelta::Append {
2413 name: "conv".to_string(),
2414 entries: vec![entry("more", 2)],
2415 current_tokens: 3,
2416 }],
2417 },
2418 at: 4,
2419 },
2420 ];
2421 let points = replay_points(&records);
2422 assert_eq!(points.len(), 2, "one point per context change");
2423 assert_eq!(points[0].at, 2);
2425 assert_eq!(points[0].context.regions[0].entries.len(), 1);
2426 assert_eq!(points[1].at, 4);
2429 assert_eq!(points[1].context.regions[0].entries.len(), 2);
2430 assert_eq!(points[1].context.stage_name, "implement");
2431 assert_eq!(points[1].meta.status, RunStatus::Running);
2432 }
2433
2434 #[test]
2435 fn replay_points_handles_context_diff_and_a_later_header() {
2436 let mut relabeled = meta();
2439 relabeled.agent_name = "renamed".to_string();
2440 let records = vec![
2441 header(),
2442 RunRecord::ContextCheckpoint {
2443 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
2444 at: 1,
2445 },
2446 RunRecord::Header {
2447 identity: identity(),
2448 meta: Box::new(relabeled),
2449 },
2450 RunRecord::ContextDiff {
2451 delta: ContextDelta {
2452 stage_name: "plan".to_string(),
2453 total_tokens: 3,
2454 max_tokens: 10_000,
2455 regions: vec![RegionDelta::Append {
2456 name: "conv".to_string(),
2457 entries: vec![entry("more", 2)],
2458 current_tokens: 3,
2459 }],
2460 },
2461 at: 2,
2462 },
2463 ];
2464 let points = replay_points(&records);
2465 assert_eq!(points.len(), 2); assert_eq!(points[1].context.regions[0].entries.len(), 2);
2467 assert_eq!(points[1].meta.agent_name, "renamed");
2469 }
2470
2471 #[test]
2472 fn replay_points_over_a_full_checkpoint() {
2473 let records = vec![
2475 header(),
2476 RunRecord::Checkpoint {
2477 meta: Box::new(meta()),
2478 context: snapshot("review", vec![region("conv", vec![entry("x", 4)])]),
2479 at: 9,
2480 },
2481 ];
2482 let points = replay_points(&records);
2483 assert_eq!(points.len(), 1);
2484 assert_eq!(points[0].context.stage_name, "review");
2485 assert_eq!(points[0].context.regions[0].entries[0].tokens, 4);
2486 }
2487
2488 #[test]
2492 fn every_inference_kind_has_a_distinct_label_and_serialized_name() {
2493 let all = [
2494 (InferenceKind::Stage, "stage"),
2495 (InferenceKind::Compaction, "compaction"),
2496 (InferenceKind::Title, "title"),
2497 (InferenceKind::Routing, "routing"),
2498 ];
2499 for (kind, label) in all {
2500 assert_eq!(kind.label(), label);
2501 assert_eq!(serde_json::to_value(kind).unwrap(), label);
2502 }
2503 let labels: std::collections::HashSet<_> = all.iter().map(|(k, _)| k.label()).collect();
2504 assert_eq!(labels.len(), all.len(), "labels must not collide");
2505 }
2506
2507 #[test]
2511 fn only_a_stage_turn_counts_as_stage_work() {
2512 assert!(InferenceKind::Stage.is_stage_work());
2513 for kind in [
2514 InferenceKind::Compaction,
2515 InferenceKind::Title,
2516 InferenceKind::Routing,
2517 ] {
2518 assert!(
2519 !kind.is_stage_work(),
2520 "{kind:?} is machinery, not stage work"
2521 );
2522 }
2523 }
2524
2525 #[test]
2529 fn a_usage_record_without_a_kind_reads_back_as_stage_work() {
2530 let json = serde_json::json!({
2531 "InferenceUsage": {
2532 "stage": "plan",
2533 "iteration": 1,
2534 "provider": "anthropic",
2535 "model": "claude-sonnet-5",
2536 "prompt_tokens": 10,
2537 "completion_tokens": 2,
2538 "cached_tokens": 0,
2539 "cache_write_tokens": 0,
2540 "at": 5,
2541 }
2542 });
2543 let record: RunRecord = serde_json::from_value(json).unwrap();
2548 assert_eq!(
2549 record,
2550 RunRecord::InferenceUsage {
2551 kind: InferenceKind::Stage,
2552 stage: "plan".to_string(),
2553 iteration: 1,
2554 provider: "anthropic".to_string(),
2555 model: "claude-sonnet-5".to_string(),
2556 prompt_tokens: 10,
2557 completion_tokens: 2,
2558 cached_tokens: 0,
2559 cache_write_tokens: 0,
2560 at: 5,
2561 }
2562 );
2563 }
2564
2565 #[test]
2570 fn folding_keeps_each_call_separate_instead_of_summing_them() {
2571 let usage = |kind, prompt, at| RunRecord::InferenceUsage {
2572 kind,
2573 stage: "plan".to_string(),
2574 iteration: 1,
2575 provider: "anthropic".to_string(),
2576 model: "claude-sonnet-5".to_string(),
2577 prompt_tokens: prompt,
2578 completion_tokens: 1,
2579 cached_tokens: 0,
2580 cache_write_tokens: 0,
2581 at,
2582 };
2583 let folded = fold(&[
2586 header(),
2587 usage(InferenceKind::Compaction, 7000, 1),
2588 usage(InferenceKind::Stage, 21_000, 2),
2589 ])
2590 .unwrap();
2591
2592 assert_eq!(folded.inference_count, 2);
2593 let seen: Vec<_> = folded
2594 .inference_usage
2595 .iter()
2596 .map(|u| (u.kind, u.prompt_tokens))
2597 .collect();
2598 assert_eq!(
2599 seen,
2600 vec![
2601 (InferenceKind::Compaction, 7000),
2602 (InferenceKind::Stage, 21_000)
2603 ]
2604 );
2605 assert!(
2609 folded
2610 .inference_usage
2611 .iter()
2612 .all(|u| u.prompt_tokens < 32_000),
2613 "no single call exceeded the window, and the journal can now prove it"
2614 );
2615 }
2616
2617 #[test]
2620 fn both_inference_record_kinds_count_as_one_call_each() {
2621 let records = all_record_kinds();
2622 let folded = fold(&records).unwrap();
2623 let written = records
2624 .iter()
2625 .filter(|r| {
2626 matches!(
2627 r,
2628 RunRecord::Inference { .. } | RunRecord::InferenceUsage { .. }
2629 )
2630 })
2631 .count();
2632 assert_eq!(folded.inference_count, written);
2633 assert_eq!(folded.inference_usage.len(), 1);
2634 }
2635
2636 #[test]
2645 fn probe_replay_matches_every_step() {
2646 fn region(name: &str, entries: &[(&str, usize)]) -> RegionSnapshot {
2647 RegionSnapshot {
2648 name: name.to_string(),
2649 kind: "temporary".to_string(),
2650 current_tokens: entries.iter().map(|(_, t)| *t).sum(),
2651 max_tokens: 1000,
2652 entries: entries
2653 .iter()
2654 .map(|(c, t)| RegionEntrySnapshot {
2655 content: c.to_string(),
2656 tokens: *t,
2657 key: None,
2658 kind: crate::region::EntryKind::Text,
2659 metadata: None,
2660 taint: crate::taint::TaintLevel::Public,
2661 })
2662 .collect(),
2663 description: None,
2664 }
2665 }
2666
2667 fn snap(entries: &[(&str, usize)]) -> ContextSnapshot {
2668 let r = region("logs", entries);
2669 ContextSnapshot {
2670 stage_name: "s".to_string(),
2671 total_tokens: r.current_tokens,
2672 max_tokens: 1000,
2673 regions: vec![r],
2674 }
2675 }
2676
2677 let steps: Vec<ContextSnapshot> = vec![
2678 snap(&[]),
2679 snap(&[("a", 10)]),
2680 snap(&[("a", 10), ("b", 20)]),
2681 snap(&[("a", 10), ("b", 20), ("c", 30)]),
2682 snap(&[("b", 20), ("c", 30)]),
2684 snap(&[("c", 30), ("d", 40)]),
2687 snap(&[("d", 40)]),
2689 snap(&[("d", 40), ("x", 5)]),
2692 snap(&[("x", 5), ("x", 5)]),
2693 snap(&[("x", 5)]),
2694 snap(&[]),
2695 ];
2696
2697 let mut base = steps[0].clone();
2700 let mut digest = digest_context(&steps[0]);
2701 for (i, next) in steps.iter().enumerate().skip(1) {
2702 let delta = diff_context_digest(&digest, next);
2703 apply_delta(&mut base, &delta);
2704 assert_eq!(
2705 base, *next,
2706 "step {i}: replay drifted from the live state\n delta was {:?}",
2707 delta.regions
2708 );
2709 digest = digest_context(next);
2710 }
2711 }
2712
2713 fn trailing_message() -> RunRecord {
2716 RunRecord::Message {
2717 message: MessageRecord {
2718 role: "user".to_string(),
2719 content: "after the unknown".to_string(),
2720 },
2721 at: 9,
2722 }
2723 }
2724
2725 fn archive_with_an_unknown_record() -> (Vec<u8>, usize) {
2728 let mut buf = Vec::new();
2729 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
2730 write_record(&mut buf, &header()).unwrap();
2731 let payload =
2732 serde_json::to_vec(&serde_json::json!({ "SomethingNew": { "whatever": 1 } })).unwrap();
2733 buf.extend_from_slice(&(payload.len() as u64).to_be_bytes());
2734 buf.extend_from_slice(&payload);
2735 write_record(&mut buf, &trailing_message()).unwrap();
2736 (buf, payload.len())
2737 }
2738
2739 fn record_kind(record: &RunRecord) -> String {
2742 serde_json::to_value(record)
2743 .expect("a RunRecord always serializes")
2744 .as_object()
2745 .expect("externally tagged, so an object")
2746 .keys()
2747 .next()
2748 .expect("with exactly one key")
2749 .clone()
2750 }
2751
2752 #[test]
2758 fn an_unknown_record_kind_is_stepped_over_not_treated_as_the_end() {
2759 let (buf, _) = archive_with_an_unknown_record();
2760 let (version, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
2761 assert_eq!(version, RUN_ARCHIVE_VERSION);
2762 let kinds: Vec<String> = records.iter().map(record_kind).collect();
2763 assert_eq!(
2764 kinds,
2765 vec!["Header".to_string(), "Message".to_string()],
2766 "the header, and the readable record after the gap"
2767 );
2768 assert_eq!(records[1], trailing_message(), "intact, not just present");
2769 }
2770
2771 #[test]
2775 fn the_streaming_reader_also_steps_over_an_unknown_record() {
2776 let (mut buf, _) = archive_with_an_unknown_record();
2779 write_record(
2780 &mut buf,
2781 &RunRecord::ContextCheckpoint {
2782 snapshot: ContextSnapshot {
2783 stage_name: "s".to_string(),
2784 total_tokens: 1,
2785 max_tokens: 10,
2786 regions: vec![],
2787 },
2788 at: 11,
2789 },
2790 )
2791 .unwrap();
2792 let mut points = 0usize;
2793 visit_archive_points(&mut buf.as_slice(), &mut |_point| {
2794 points += 1;
2795 ControlFlow::Continue(())
2796 })
2797 .expect("a valid preamble");
2798 assert_eq!(points, 1, "the walk got past the unknown frame");
2799 }
2800
2801 #[test]
2805 fn a_frame_reports_whether_its_payload_was_readable() {
2806 let (buf, unknown_bytes) = archive_with_an_unknown_record();
2807 let mut r = buf.as_slice();
2808 read_archive_start(&mut r).unwrap();
2809
2810 let mut frames = Vec::new();
2811 while let Some(frame) = read_frame(&mut r).expect("no torn frames here") {
2812 frames.push(frame);
2813 }
2814 assert_eq!(
2815 frames,
2816 vec![
2817 Frame::Record(Box::new(header())),
2818 Frame::Unreadable {
2820 bytes: unknown_bytes
2821 },
2822 Frame::Record(Box::new(trailing_message())),
2823 ],
2824 "one frame per record, with the unreadable one accounted for rather than ending the read"
2825 );
2826 }
2827
2828 #[test]
2831 fn a_torn_frame_still_ends_the_read() {
2832 let mut buf = Vec::new();
2833 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
2834 write_record(&mut buf, &header()).unwrap();
2835 buf.extend_from_slice(&999u64.to_be_bytes());
2837 buf.extend_from_slice(b"not enough");
2838
2839 let (_, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
2840 assert_eq!(records.len(), 1, "everything intact before the tear");
2841 }
2842
2843 #[test]
2847 fn an_archive_from_a_newer_format_is_refused_with_both_versions_named() {
2848 let mut buf = Vec::new();
2849 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION + 1).unwrap();
2850 write_record(&mut buf, &header()).unwrap();
2851
2852 let err = read_archive_lenient(&mut buf.as_slice()).unwrap_err();
2853 let message = err.to_string();
2854 assert!(
2855 message.contains(&(RUN_ARCHIVE_VERSION + 1).to_string()),
2856 "{message}"
2857 );
2858 assert!(message.contains("upgrade leviath"), "{message}");
2859 }
2860
2861 #[test]
2864 fn an_archive_from_an_older_format_still_reads() {
2865 let mut buf = Vec::new();
2866 write_archive_start(&mut buf, 0).unwrap();
2867 write_record(&mut buf, &header()).unwrap();
2868 let (version, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
2869 assert_eq!(version, 0);
2870 assert_eq!(records.len(), 1);
2871 }
2872}