1use std::io::{self, Read, Write};
42use std::ops::ControlFlow;
43
44use serde::{Deserialize, Serialize};
45
46use crate::run_meta::{ContextSnapshot, RegionEntrySnapshot, RegionSnapshot, RunMeta, RunStatus};
47
48pub const RUN_ARCHIVE_MAGIC: &[u8; 4] = b"LVR1";
50
51pub const RUN_ARCHIVE_VERSION: u16 = 1;
53
54#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
61pub struct RunIdentity {
62 pub run_id: String,
64 pub machine_id: String,
66 pub world_id: String,
68 pub created_at: i64,
70}
71
72#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
74pub struct MessageRecord {
75 pub role: String,
77 pub content: String,
79}
80
81#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
83pub struct ToolCallRecord {
84 pub id: String,
86 pub name: String,
88 pub arguments: String,
90 pub result: Option<String>,
92 #[serde(default, skip_serializing_if = "Option::is_none")]
96 pub thought_signature: Option<String>,
97}
98
99#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
102pub struct InferenceRequestRecord {
103 pub model: String,
105 pub system: Vec<String>,
107 pub messages: Vec<MessageRecord>,
109 pub tool_names: Vec<String>,
111 pub temperature: f32,
113 pub max_tokens: usize,
115}
116
117#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
119pub struct InferenceResponseRecord {
120 pub content: String,
122 pub tool_calls: Vec<ToolCallRecord>,
124 pub prompt_tokens: usize,
126 pub completion_tokens: usize,
128 pub cached_tokens: usize,
130 pub cache_write_tokens: usize,
132}
133
134#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
136pub enum RegionDelta {
137 Set(RegionSnapshot),
140 Append {
143 name: String,
145 entries: Vec<RegionEntrySnapshot>,
147 current_tokens: usize,
149 },
150 Clear {
152 name: String,
154 },
155 Remove {
157 name: String,
159 },
160}
161
162#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
164pub struct ContextDelta {
165 pub stage_name: String,
167 pub total_tokens: usize,
169 pub max_tokens: usize,
171 pub regions: Vec<RegionDelta>,
173}
174
175#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
177pub enum RunRecord {
178 Header {
180 identity: RunIdentity,
182 meta: Box<RunMeta>,
184 },
185 OwnershipChanged {
187 machine_id: String,
189 world_id: String,
191 at: i64,
193 },
194 Inference {
196 stage: String,
198 iteration: usize,
200 request: InferenceRequestRecord,
202 response: InferenceResponseRecord,
204 at: i64,
206 },
207 ToolBatch {
215 calls: Vec<ToolCallRecord>,
217 at: i64,
219 #[serde(default)]
221 stage_index: usize,
222 #[serde(default)]
225 iteration: usize,
226 #[serde(default)]
228 response: String,
229 },
230 ToolCallDone {
232 iteration: usize,
234 call_id: String,
236 result: String,
238 at: i64,
240 },
241 ContextCheckpoint {
243 snapshot: ContextSnapshot,
245 at: i64,
247 },
248 ContextDiff {
250 delta: ContextDelta,
252 at: i64,
254 },
255 Message {
257 message: MessageRecord,
259 at: i64,
261 },
262 StatusChanged {
264 status: RunStatus,
266 at: i64,
268 },
269 Checkpoint {
272 meta: Box<RunMeta>,
274 context: ContextSnapshot,
276 at: i64,
278 },
279 Progress {
284 meta: Box<RunMeta>,
286 delta: ContextDelta,
288 at: i64,
290 },
291}
292
293fn is_prefix(prev: &[RegionEntrySnapshot], next: &[RegionEntrySnapshot]) -> bool {
297 prev.len() <= next.len() && next[..prev.len()] == *prev
298}
299
300pub fn diff_context(prev: &ContextSnapshot, next: &ContextSnapshot) -> ContextDelta {
304 let mut regions = Vec::new();
305 for nr in &next.regions {
306 match prev.regions.iter().find(|r| r.name == nr.name) {
307 None => regions.push(RegionDelta::Set(nr.clone())),
308 Some(pr) => {
309 if pr == nr {
310 } else if nr.entries.is_empty() && !pr.entries.is_empty() {
312 regions.push(RegionDelta::Clear {
313 name: nr.name.clone(),
314 });
315 } else if pr.kind == nr.kind
316 && pr.max_tokens == nr.max_tokens
317 && is_prefix(&pr.entries, &nr.entries)
318 {
319 regions.push(RegionDelta::Append {
320 name: nr.name.clone(),
321 entries: nr.entries[pr.entries.len()..].to_vec(),
322 current_tokens: nr.current_tokens,
323 });
324 } else {
325 regions.push(RegionDelta::Set(nr.clone()));
326 }
327 }
328 }
329 }
330 for pr in &prev.regions {
331 if !next.regions.iter().any(|r| r.name == pr.name) {
332 regions.push(RegionDelta::Remove {
333 name: pr.name.clone(),
334 });
335 }
336 }
337 ContextDelta {
338 stage_name: next.stage_name.clone(),
339 total_tokens: next.total_tokens,
340 max_tokens: next.max_tokens,
341 regions,
342 }
343}
344
345#[derive(Debug, Clone, PartialEq)]
356pub struct RegionDigest {
357 pub name: String,
359 pub kind: String,
361 pub current_tokens: usize,
363 pub max_tokens: usize,
365 pub entries: Vec<u64>,
367}
368
369#[derive(Debug, Clone, PartialEq, Default)]
371pub struct ContextDigest {
372 pub regions: Vec<RegionDigest>,
374}
375
376fn entry_digest(entry: &RegionEntrySnapshot) -> u64 {
380 use std::hash::{Hash, Hasher};
381 let mut hasher = std::collections::hash_map::DefaultHasher::new();
382 entry.content.hash(&mut hasher);
383 entry.tokens.hash(&mut hasher);
384 entry.key.hash(&mut hasher);
385 serde_json::to_string(&entry.kind)
388 .expect("EntryKind always serializes")
389 .hash(&mut hasher);
390 serde_json::to_string(&entry.metadata)
391 .expect("entry metadata always serializes")
392 .hash(&mut hasher);
393 serde_json::to_string(&entry.taint)
394 .expect("taint always serializes")
395 .hash(&mut hasher);
396 hasher.finish()
397}
398
399pub fn digest_context(snapshot: &ContextSnapshot) -> ContextDigest {
401 ContextDigest {
402 regions: snapshot
403 .regions
404 .iter()
405 .map(|r| RegionDigest {
406 name: r.name.clone(),
407 kind: r.kind.clone(),
408 current_tokens: r.current_tokens,
409 max_tokens: r.max_tokens,
410 entries: r.entries.iter().map(entry_digest).collect(),
411 })
412 .collect(),
413 }
414}
415
416fn is_prefix_digest(prev: &[u64], next: &[RegionEntrySnapshot]) -> bool {
418 prev.len() <= next.len()
419 && prev
420 .iter()
421 .zip(next)
422 .all(|(hash, entry)| *hash == entry_digest(entry))
423}
424
425pub fn diff_context_digest(prev: &ContextDigest, next: &ContextSnapshot) -> ContextDelta {
430 let mut regions = Vec::new();
431 for nr in &next.regions {
432 match prev.regions.iter().find(|r| r.name == nr.name) {
433 None => regions.push(RegionDelta::Set(nr.clone())),
434 Some(pr) => {
435 let unchanged = pr.kind == nr.kind
436 && pr.max_tokens == nr.max_tokens
437 && pr.current_tokens == nr.current_tokens
438 && pr.entries.len() == nr.entries.len()
439 && is_prefix_digest(&pr.entries, &nr.entries);
440 if unchanged {
441 } else if nr.entries.is_empty() && !pr.entries.is_empty() {
443 regions.push(RegionDelta::Clear {
444 name: nr.name.clone(),
445 });
446 } else if pr.kind == nr.kind
447 && pr.max_tokens == nr.max_tokens
448 && is_prefix_digest(&pr.entries, &nr.entries)
449 {
450 regions.push(RegionDelta::Append {
451 name: nr.name.clone(),
452 entries: nr.entries[pr.entries.len()..].to_vec(),
453 current_tokens: nr.current_tokens,
454 });
455 } else {
456 regions.push(RegionDelta::Set(nr.clone()));
457 }
458 }
459 }
460 }
461 for pr in &prev.regions {
462 if !next.regions.iter().any(|r| r.name == pr.name) {
463 regions.push(RegionDelta::Remove {
464 name: pr.name.clone(),
465 });
466 }
467 }
468 ContextDelta {
469 stage_name: next.stage_name.clone(),
470 total_tokens: next.total_tokens,
471 max_tokens: next.max_tokens,
472 regions,
473 }
474}
475
476pub fn apply_delta(base: &mut ContextSnapshot, delta: &ContextDelta) {
480 base.stage_name = delta.stage_name.clone();
481 base.total_tokens = delta.total_tokens;
482 base.max_tokens = delta.max_tokens;
483 for region_delta in &delta.regions {
484 match region_delta {
485 RegionDelta::Set(snapshot) => {
486 match base.regions.iter_mut().find(|r| r.name == snapshot.name) {
487 Some(existing) => *existing = snapshot.clone(),
488 None => base.regions.push(snapshot.clone()),
489 }
490 }
491 RegionDelta::Append {
492 name,
493 entries,
494 current_tokens,
495 } => {
496 if let Some(region) = base.regions.iter_mut().find(|r| &r.name == name) {
497 region.entries.extend(entries.iter().cloned());
498 region.current_tokens = *current_tokens;
499 }
500 }
501 RegionDelta::Clear { name } => {
502 if let Some(region) = base.regions.iter_mut().find(|r| &r.name == name) {
503 region.entries.clear();
504 region.current_tokens = 0;
505 }
506 }
507 RegionDelta::Remove { name } => {
508 base.regions.retain(|r| &r.name != name);
509 }
510 }
511 }
512}
513
514pub fn write_archive_start(w: &mut dyn Write, version: u16) -> io::Result<()> {
518 w.write_all(RUN_ARCHIVE_MAGIC)?;
519 w.write_all(&version.to_be_bytes())?;
520 Ok(())
521}
522
523pub fn read_archive_start(r: &mut dyn Read) -> io::Result<u16> {
525 let mut magic = [0u8; 4];
526 r.read_exact(&mut magic)?;
527 if &magic != RUN_ARCHIVE_MAGIC {
528 return Err(io::Error::new(
529 io::ErrorKind::InvalidData,
530 "not a leviath run archive (bad magic)",
531 ));
532 }
533 let mut version = [0u8; 2];
534 r.read_exact(&mut version)?;
535 Ok(u16::from_be_bytes(version))
536}
537
538pub fn write_record(w: &mut dyn Write, record: &RunRecord) -> io::Result<()> {
541 let payload = serde_json::to_vec(record).expect("a RunRecord always serializes to JSON");
542 let len = payload.len() as u64;
543 w.write_all(&len.to_be_bytes())?;
544 w.write_all(&payload)?;
545 Ok(())
546}
547
548fn read_exact_or_eof(r: &mut dyn Read, buf: &mut [u8]) -> io::Result<bool> {
551 let mut filled = 0;
552 while filled < buf.len() {
553 match r.read(&mut buf[filled..])? {
554 0 => {
555 if filled == 0 {
556 return Ok(false); }
558 return Err(io::Error::new(
559 io::ErrorKind::UnexpectedEof,
560 "truncated run-archive frame",
561 ));
562 }
563 n => filled += n,
564 }
565 }
566 Ok(true)
567}
568
569const MAX_RECORD_BYTES: u64 = 256 * 1024 * 1024;
576
577pub fn read_record(r: &mut dyn Read) -> io::Result<Option<RunRecord>> {
579 let mut len_bytes = [0u8; 8];
580 if !read_exact_or_eof(r, &mut len_bytes)? {
581 return Ok(None);
582 }
583 let len = u64::from_be_bytes(len_bytes);
584 if len > MAX_RECORD_BYTES {
591 return Err(io::Error::new(
592 io::ErrorKind::InvalidData,
593 format!("run-archive frame claims {len} bytes, over the {MAX_RECORD_BYTES} cap"),
594 ));
595 }
596 let mut payload = vec![0u8; len as usize];
597 if !read_exact_or_eof(r, &mut payload)? {
598 return Err(io::Error::new(
599 io::ErrorKind::UnexpectedEof,
600 "truncated run-archive frame",
601 ));
602 }
603 let record = serde_json::from_slice(&payload)
604 .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
605 Ok(Some(record))
606}
607
608pub fn read_archive(r: &mut dyn Read) -> io::Result<(u16, Vec<RunRecord>)> {
610 let version = read_archive_start(r)?;
611 let mut records = Vec::new();
612 while let Some(record) = read_record(r)? {
613 records.push(record);
614 }
615 Ok((version, records))
616}
617
618pub fn read_archive_lenient(r: &mut dyn Read) -> io::Result<(u16, Vec<RunRecord>)> {
631 let version = read_archive_start(r)?;
632 let mut records = Vec::new();
633 while let Ok(Some(record)) = read_record(r) {
636 records.push(record);
637 }
638 Ok((version, records))
639}
640
641#[derive(Debug, Clone, PartialEq)]
648pub struct PendingToolBatch {
649 pub stage_index: usize,
651 pub iteration: usize,
653 pub response: String,
655 pub calls: Vec<ToolCallRecord>,
657}
658
659#[derive(Debug, Clone, PartialEq)]
662pub struct FoldedRun {
663 pub identity: RunIdentity,
665 pub meta: RunMeta,
667 pub context: ContextSnapshot,
669 pub messages: Vec<MessageRecord>,
671 pub inference_count: usize,
673 pub tool_call_count: usize,
675 pub pending_batch: Option<PendingToolBatch>,
679}
680
681pub fn context_contains_batch(context: &ContextSnapshot, batch: &PendingToolBatch) -> bool {
685 let Some(first_id) = batch.calls.first().map(|c| c.id.as_str()) else {
686 return false;
687 };
688 context.regions.iter().any(|region| {
689 region.entries.iter().any(|entry| {
690 matches!(
691 &entry.kind,
692 crate::region::EntryKind::AssistantTurn { tool_calls }
693 if tool_calls.iter().any(|tc| tc.id == first_id)
694 )
695 })
696 })
697}
698
699pub fn fold(records: &[RunRecord]) -> Option<FoldedRun> {
702 let mut iter = records.iter();
703 let (identity, meta) = match iter.next() {
704 Some(RunRecord::Header { identity, meta }) => (identity.clone(), (**meta).clone()),
705 _ => return None,
706 };
707 let mut folded = FoldedRun {
708 identity,
709 meta,
710 context: ContextSnapshot {
711 stage_name: String::new(),
712 total_tokens: 0,
713 max_tokens: 0,
714 regions: Vec::new(),
715 },
716 messages: Vec::new(),
717 inference_count: 0,
718 tool_call_count: 0,
719 pending_batch: None,
720 };
721 for record in iter {
722 match record {
723 RunRecord::Header { identity, meta } => {
724 folded.identity = identity.clone();
725 folded.meta = (**meta).clone();
726 }
727 RunRecord::OwnershipChanged {
728 machine_id,
729 world_id,
730 ..
731 } => {
732 folded.identity.machine_id = machine_id.clone();
733 folded.identity.world_id = world_id.clone();
734 }
735 RunRecord::Inference { .. } => folded.inference_count += 1,
736 RunRecord::ToolBatch {
737 calls,
738 stage_index,
739 iteration,
740 response,
741 ..
742 } => {
743 folded.tool_call_count += calls.len();
744 folded.pending_batch = Some(PendingToolBatch {
747 stage_index: *stage_index,
748 iteration: *iteration,
749 response: response.clone(),
750 calls: calls.clone(),
751 });
752 }
753 RunRecord::ToolCallDone {
754 iteration,
755 call_id,
756 result,
757 ..
758 } => {
759 if let Some(batch) = folded
762 .pending_batch
763 .as_mut()
764 .filter(|b| b.iteration == *iteration)
765 && let Some(call) = batch.calls.iter_mut().find(|c| c.id == *call_id)
766 {
767 call.result = Some(result.clone());
768 }
769 }
770 RunRecord::ContextCheckpoint { snapshot, .. } => folded.context = snapshot.clone(),
771 RunRecord::ContextDiff { delta, .. } => apply_delta(&mut folded.context, delta),
772 RunRecord::Message { message, .. } => folded.messages.push(message.clone()),
773 RunRecord::StatusChanged { status, .. } => folded.meta.status = status.clone(),
774 RunRecord::Checkpoint { meta, context, .. } => {
775 folded.meta = (**meta).clone();
776 folded.context = context.clone();
777 }
778 RunRecord::Progress { meta, delta, .. } => {
779 folded.meta = (**meta).clone();
780 apply_delta(&mut folded.context, delta);
781 }
782 }
783 }
784 if let Some(batch) = &folded.pending_batch
789 && (folded.meta.iteration != batch.iteration
790 || context_contains_batch(&folded.context, batch))
791 {
792 folded.pending_batch = None;
793 }
794 Some(folded)
795}
796
797#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
800pub struct RunPoint {
801 pub meta: RunMeta,
803 pub context: ContextSnapshot,
805 pub at: i64,
807}
808
809#[derive(Debug)]
812pub struct PointRef<'a> {
813 pub index: usize,
817 pub at: i64,
819 pub meta: &'a RunMeta,
821 pub context: &'a ContextSnapshot,
823}
824
825pub fn visit_points(records: &[RunRecord], visit: &mut dyn FnMut(PointRef<'_>) -> ControlFlow<()>) {
844 let mut iter = records.iter();
845 let Some(mut folder) = (match iter.next() {
846 Some(first) => PointFolder::start(first),
847 None => None,
848 }) else {
849 return;
850 };
851 for record in iter {
852 if folder.push(record, visit).is_break() {
853 return;
854 }
855 }
856}
857
858pub fn visit_archive_points(
868 r: &mut dyn Read,
869 visit: &mut dyn FnMut(PointRef<'_>) -> ControlFlow<()>,
870) -> io::Result<()> {
871 read_archive_start(r)?;
872 let mut folder = match read_record(r) {
873 Ok(Some(first)) => match PointFolder::start(&first) {
874 Some(folder) => folder,
875 None => return Ok(()),
876 },
877 _ => return Ok(()),
878 };
879 while let Ok(Some(record)) = read_record(r) {
880 if folder.push(&record, visit).is_break() {
881 return Ok(());
882 }
883 }
884 Ok(())
885}
886
887struct PointFolder {
892 meta: RunMeta,
893 context: ContextSnapshot,
894 index: usize,
895}
896
897impl PointFolder {
898 fn start(first: &RunRecord) -> Option<Self> {
902 match first {
903 RunRecord::Header { meta, .. } => Some(Self {
904 meta: (**meta).clone(),
905 context: ContextSnapshot {
906 stage_name: String::new(),
907 total_tokens: 0,
908 max_tokens: 0,
909 regions: Vec::new(),
910 },
911 index: 0,
912 }),
913 _ => None,
914 }
915 }
916
917 fn push(
919 &mut self,
920 record: &RunRecord,
921 visit: &mut dyn FnMut(PointRef<'_>) -> ControlFlow<()>,
922 ) -> ControlFlow<()> {
923 let at = match record {
924 RunRecord::Header { meta: m, .. } => {
925 self.meta = (**m).clone();
926 return ControlFlow::Continue(());
927 }
928 RunRecord::StatusChanged { status, .. } => {
929 self.meta.status = status.clone();
930 return ControlFlow::Continue(());
931 }
932 RunRecord::ContextCheckpoint { snapshot, at } => {
933 self.context = snapshot.clone();
934 *at
935 }
936 RunRecord::ContextDiff { delta, at } => {
937 apply_delta(&mut self.context, delta);
938 *at
939 }
940 RunRecord::Checkpoint {
941 meta: m,
942 context: c,
943 at,
944 } => {
945 self.meta = (**m).clone();
946 self.context = c.clone();
947 *at
948 }
949 RunRecord::Progress { meta: m, delta, at } => {
950 self.meta = (**m).clone();
951 apply_delta(&mut self.context, delta);
952 *at
953 }
954 RunRecord::OwnershipChanged { .. }
956 | RunRecord::Inference { .. }
957 | RunRecord::ToolBatch { .. }
958 | RunRecord::ToolCallDone { .. }
959 | RunRecord::Message { .. } => return ControlFlow::Continue(()),
960 };
961 let flow = visit(PointRef {
962 index: self.index,
963 at,
964 meta: &self.meta,
965 context: &self.context,
966 });
967 self.index += 1;
968 flow
969 }
970}
971
972pub fn replay_points(records: &[RunRecord]) -> Vec<RunPoint> {
982 let mut points = Vec::new();
983 visit_points(records, &mut |point| {
984 points.push(RunPoint {
985 meta: point.meta.clone(),
986 context: point.context.clone(),
987 at: point.at,
988 });
989 ControlFlow::Continue(())
990 });
991 points
992}
993
994#[cfg(test)]
995mod tests {
996 use super::*;
997 use crate::run_meta::RunStatus;
998
999 fn identity() -> RunIdentity {
1000 RunIdentity {
1001 run_id: "run-1".to_string(),
1002 machine_id: "machine-a".to_string(),
1003 world_id: "world-x".to_string(),
1004 created_at: 100,
1005 }
1006 }
1007
1008 fn meta() -> RunMeta {
1009 RunMeta::new(
1010 "run-1".to_string(),
1011 "coder".to_string(),
1012 "/agents/coder".to_string(),
1013 "do it".to_string(),
1014 Some("anthropic/claude".to_string()),
1015 "/work".to_string(),
1016 2,
1017 )
1018 }
1019
1020 fn entry(content: &str, tokens: usize) -> RegionEntrySnapshot {
1021 RegionEntrySnapshot {
1022 content: content.to_string(),
1023 tokens,
1024 kind: crate::region::EntryKind::Text,
1025 metadata: None,
1026 key: None,
1027 taint: Default::default(),
1028 }
1029 }
1030
1031 fn region(name: &str, entries: Vec<RegionEntrySnapshot>) -> RegionSnapshot {
1032 let current = entries.iter().map(|e| e.tokens).sum();
1033 RegionSnapshot {
1034 name: name.to_string(),
1035 kind: "clearable".to_string(),
1036 current_tokens: current,
1037 max_tokens: 1000,
1038 entries,
1039 }
1040 }
1041
1042 fn snapshot(stage: &str, regions: Vec<RegionSnapshot>) -> ContextSnapshot {
1043 let total = regions.iter().map(|r| r.current_tokens).sum();
1044 ContextSnapshot {
1045 stage_name: stage.to_string(),
1046 total_tokens: total,
1047 max_tokens: 10_000,
1048 regions,
1049 }
1050 }
1051
1052 fn header() -> RunRecord {
1053 RunRecord::Header {
1054 identity: identity(),
1055 meta: Box::new(meta()),
1056 }
1057 }
1058
1059 fn region_delta_kind(d: &RegionDelta) -> &'static str {
1063 match d {
1064 RegionDelta::Set(_) => "set",
1065 RegionDelta::Append { .. } => "append",
1066 RegionDelta::Clear { .. } => "clear",
1067 RegionDelta::Remove { .. } => "remove",
1068 }
1069 }
1070
1071 fn assert_diff_roundtrip(a: &ContextSnapshot, b: &ContextSnapshot) {
1076 let delta = diff_context(a, b);
1077 let mut base = a.clone();
1078 apply_delta(&mut base, &delta);
1079 assert_eq!(&base, b);
1080 }
1081
1082 #[test]
1083 fn diff_append_only_growth_is_compact() {
1084 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1085 let b = snapshot(
1086 "s1",
1087 vec![region("conv", vec![entry("hi", 1), entry("there", 2)])],
1088 );
1089 let delta = diff_context(&a, &b);
1090 assert_eq!(region_delta_kind(&delta.regions[0]), "append");
1091 assert_diff_roundtrip(&a, &b);
1092 }
1093
1094 #[test]
1095 fn diff_new_region_is_set() {
1096 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1097 let b = snapshot(
1098 "s1",
1099 vec![
1100 region("conv", vec![entry("hi", 1)]),
1101 region("plan", vec![entry("p", 3)]),
1102 ],
1103 );
1104 let delta = diff_context(&a, &b);
1105 assert!(delta.regions.iter().any(|d| region_delta_kind(d) == "set"));
1106 assert_diff_roundtrip(&a, &b);
1107 }
1108
1109 #[test]
1110 fn diff_cleared_region() {
1111 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1112 let b = snapshot("s1", vec![region("conv", vec![])]);
1113 let delta = diff_context(&a, &b);
1114 assert_eq!(region_delta_kind(&delta.regions[0]), "clear");
1115 assert_diff_roundtrip(&a, &b);
1116 }
1117
1118 #[test]
1119 fn diff_removed_region() {
1120 let a = snapshot(
1121 "s1",
1122 vec![
1123 region("conv", vec![entry("hi", 1)]),
1124 region("plan", vec![entry("p", 3)]),
1125 ],
1126 );
1127 let b = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1128 let delta = diff_context(&a, &b);
1129 assert!(
1130 delta
1131 .regions
1132 .iter()
1133 .any(|d| region_delta_kind(d) == "remove")
1134 );
1135 assert_diff_roundtrip(&a, &b);
1136 }
1137
1138 #[test]
1139 fn diff_non_prefix_rewrite_is_set() {
1140 let a = snapshot("s1", vec![region("conv", vec![entry("old", 1)])]);
1142 let b = snapshot("s1", vec![region("conv", vec![entry("new", 1)])]);
1143 let delta = diff_context(&a, &b);
1144 assert_eq!(region_delta_kind(&delta.regions[0]), "set");
1145 assert_diff_roundtrip(&a, &b);
1146 }
1147
1148 #[test]
1149 fn diff_kind_change_is_set_not_append() {
1150 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1152 let mut grown = region("conv", vec![entry("hi", 1), entry("more", 1)]);
1153 grown.kind = "sliding".to_string();
1154 let b = snapshot("s1", vec![grown]);
1155 let delta = diff_context(&a, &b);
1156 assert_eq!(region_delta_kind(&delta.regions[0]), "set");
1157 assert_diff_roundtrip(&a, &b);
1158 }
1159
1160 fn framed(records: &[RunRecord]) -> Vec<u8> {
1164 let mut buf = Vec::new();
1165 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1166 for record in records {
1167 write_record(&mut buf, record).unwrap();
1168 }
1169 buf
1170 }
1171
1172 fn try_collect_streamed(bytes: &[u8]) -> io::Result<Vec<(usize, i64, usize)>> {
1176 let mut seen = Vec::new();
1177 visit_archive_points(&mut &bytes[..], &mut |p| {
1178 seen.push((p.index, p.at, p.context.total_tokens));
1179 ControlFlow::Continue(())
1180 })?;
1181 Ok(seen)
1182 }
1183
1184 fn collect_streamed(bytes: &[u8]) -> Vec<(usize, i64, usize)> {
1186 try_collect_streamed(bytes).unwrap()
1187 }
1188
1189 #[test]
1190 fn visit_archive_points_matches_visit_points() {
1191 let records = vec![
1192 header(),
1193 RunRecord::ContextCheckpoint {
1194 snapshot: snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1195 at: 10,
1196 },
1197 RunRecord::StatusChanged {
1198 status: RunStatus::Running,
1199 at: 11,
1200 },
1201 RunRecord::Progress {
1202 meta: Box::new(meta()),
1203 delta: diff_context(
1204 &snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1205 &snapshot(
1206 "s1",
1207 vec![region("conv", vec![entry("hi", 1), entry("more", 2)])],
1208 ),
1209 ),
1210 at: 12,
1211 },
1212 ];
1213 let mut in_memory = Vec::new();
1214 visit_points(&records, &mut |p| {
1215 in_memory.push((p.index, p.at, p.context.total_tokens));
1216 ControlFlow::Continue(())
1217 });
1218 assert_eq!(collect_streamed(&framed(&records)), in_memory);
1219 assert_eq!(in_memory.len(), 2, "checkpoint + progress = two points");
1220 }
1221
1222 #[test]
1223 fn visit_archive_points_rejects_a_bad_preamble() {
1224 assert!(try_collect_streamed(b"not an archive at all").is_err());
1225 }
1226
1227 #[test]
1228 fn visit_archive_points_is_lenient_about_a_torn_tail() {
1229 let records = vec![
1230 header(),
1231 RunRecord::ContextCheckpoint {
1232 snapshot: snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1233 at: 10,
1234 },
1235 ];
1236 let mut bytes = framed(&records);
1237 bytes.extend_from_slice(&1000u64.to_be_bytes());
1239 bytes.extend_from_slice(b"partial");
1240 assert_eq!(collect_streamed(&bytes).len(), 1, "points before the tear");
1241 }
1242
1243 #[test]
1244 fn visit_archive_points_visits_nothing_without_a_header() {
1245 let records = vec![RunRecord::ContextCheckpoint {
1246 snapshot: snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1247 at: 10,
1248 }];
1249 assert!(collect_streamed(&framed(&records)).is_empty());
1250 assert!(collect_streamed(&framed(&[])).is_empty());
1252 }
1253
1254 #[test]
1255 fn visit_archive_points_stops_on_break() {
1256 let records = vec![
1257 header(),
1258 RunRecord::ContextCheckpoint {
1259 snapshot: snapshot("s1", vec![region("conv", vec![entry("a", 1)])]),
1260 at: 10,
1261 },
1262 RunRecord::ContextCheckpoint {
1263 snapshot: snapshot("s1", vec![region("conv", vec![entry("b", 2)])]),
1264 at: 11,
1265 },
1266 ];
1267 let bytes = framed(&records);
1268 let mut seen = 0;
1269 visit_archive_points(&mut &bytes[..], &mut |_| {
1270 seen += 1;
1271 ControlFlow::Break(())
1272 })
1273 .unwrap();
1274 assert_eq!(seen, 1);
1275 }
1276
1277 fn assert_digest_matches_full_diff(a: &ContextSnapshot, b: &ContextSnapshot) {
1284 let via_digest = diff_context_digest(&digest_context(a), b);
1285 assert_eq!(via_digest, diff_context(a, b));
1286 let mut base = a.clone();
1288 apply_delta(&mut base, &via_digest);
1289 assert_eq!(&base, b);
1290 }
1291
1292 #[test]
1293 fn digest_diff_append_only_growth_is_compact() {
1294 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1295 let b = snapshot(
1296 "s1",
1297 vec![region("conv", vec![entry("hi", 1), entry("there", 2)])],
1298 );
1299 let delta = diff_context_digest(&digest_context(&a), &b);
1300 assert_eq!(region_delta_kind(&delta.regions[0]), "append");
1301 assert_digest_matches_full_diff(&a, &b);
1302 }
1303
1304 #[test]
1305 fn digest_diff_new_cleared_removed_and_rewritten_regions() {
1306 let a = snapshot(
1307 "s1",
1308 vec![
1309 region("conv", vec![entry("hi", 1)]),
1310 region("gone", vec![entry("bye", 1)]),
1311 region("wiped", vec![entry("w", 1)]),
1312 region("rewritten", vec![entry("old", 1)]),
1313 ],
1314 );
1315 let b = snapshot(
1316 "s1",
1317 vec![
1318 region("conv", vec![entry("hi", 1)]),
1319 region("wiped", vec![]),
1320 region("rewritten", vec![entry("new", 1)]),
1321 region("fresh", vec![entry("f", 2)]),
1322 ],
1323 );
1324 let delta = diff_context_digest(&digest_context(&a), &b);
1325 let kinds: Vec<_> = delta.regions.iter().map(region_delta_kind).collect();
1326 assert_eq!(kinds, vec!["clear", "set", "set", "remove"]);
1327 assert_digest_matches_full_diff(&a, &b);
1328 }
1329
1330 #[test]
1331 fn digest_diff_unchanged_region_emits_nothing() {
1332 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1333 let delta = diff_context_digest(&digest_context(&a), &a.clone());
1334 assert!(delta.regions.is_empty());
1335 assert_digest_matches_full_diff(&a, &a.clone());
1336 }
1337
1338 #[test]
1339 fn digest_diff_kind_change_is_set_not_append() {
1340 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1341 let mut grown = region("conv", vec![entry("hi", 1), entry("more", 1)]);
1342 grown.kind = "sliding".to_string();
1343 let b = snapshot("s1", vec![grown]);
1344 let delta = diff_context_digest(&digest_context(&a), &b);
1345 assert_eq!(region_delta_kind(&delta.regions[0]), "set");
1346 assert_digest_matches_full_diff(&a, &b);
1347 }
1348
1349 #[test]
1352 fn digest_diff_token_recount_is_an_empty_append() {
1353 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1354 let mut recounted = region("conv", vec![entry("hi", 1)]);
1355 recounted.current_tokens = 42;
1356 let b = snapshot("s1", vec![recounted]);
1357 let delta = diff_context_digest(&digest_context(&a), &b);
1358 assert_eq!(region_delta_kind(&delta.regions[0]), "append");
1359 assert_digest_matches_full_diff(&a, &b);
1360 }
1361
1362 #[test]
1365 fn entry_digest_covers_every_field() {
1366 let base = entry("text", 1);
1367 let variants = [
1368 entry("other", 1),
1369 entry("text", 2),
1370 RegionEntrySnapshot {
1371 key: Some("k".to_string()),
1372 ..entry("text", 1)
1373 },
1374 RegionEntrySnapshot {
1375 metadata: Some(serde_json::json!({"a": 1})),
1376 ..entry("text", 1)
1377 },
1378 RegionEntrySnapshot {
1379 kind: crate::region::EntryKind::ToolResult {
1380 tool_call_id: "c1".to_string(),
1381 tool_name: "shell".to_string(),
1382 is_error: false,
1383 },
1384 ..entry("text", 1)
1385 },
1386 ];
1387 let base_hash = entry_digest(&base);
1388 for variant in &variants {
1389 assert_ne!(
1390 entry_digest(variant),
1391 base_hash,
1392 "field change must change the digest: {variant:?}"
1393 );
1394 }
1395 assert_eq!(entry_digest(&base), entry_digest(&entry("text", 1)));
1397 }
1398
1399 #[test]
1400 fn diff_unchanged_region_emits_nothing() {
1401 let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1402 let b = a.clone();
1403 let delta = diff_context(&a, &b);
1404 assert!(delta.regions.is_empty());
1405 assert_diff_roundtrip(&a, &b);
1406 }
1407
1408 #[test]
1409 fn apply_delta_skips_unknown_regions_leniently() {
1410 let mut base = snapshot("s1", vec![]);
1412 let delta = ContextDelta {
1413 stage_name: "s1".to_string(),
1414 total_tokens: 0,
1415 max_tokens: 10_000,
1416 regions: vec![
1417 RegionDelta::Append {
1418 name: "ghost".to_string(),
1419 entries: vec![entry("x", 1)],
1420 current_tokens: 1,
1421 },
1422 RegionDelta::Clear {
1423 name: "ghost".to_string(),
1424 },
1425 RegionDelta::Remove {
1426 name: "ghost".to_string(),
1427 },
1428 ],
1429 };
1430 apply_delta(&mut base, &delta);
1431 assert!(base.regions.is_empty());
1432 }
1433
1434 fn all_record_kinds() -> Vec<RunRecord> {
1437 vec![
1438 header(),
1439 RunRecord::OwnershipChanged {
1440 machine_id: "machine-b".to_string(),
1441 world_id: "world-y".to_string(),
1442 at: 101,
1443 },
1444 RunRecord::Inference {
1445 stage: "plan".to_string(),
1446 iteration: 0,
1447 request: InferenceRequestRecord {
1448 model: "m".to_string(),
1449 system: vec!["sys".to_string()],
1450 messages: vec![MessageRecord {
1451 role: "user".to_string(),
1452 content: "hi".to_string(),
1453 }],
1454 tool_names: vec!["read_file".to_string()],
1455 temperature: 0.7,
1456 max_tokens: 1024,
1457 },
1458 response: InferenceResponseRecord {
1459 content: "ok".to_string(),
1460 tool_calls: vec![],
1461 prompt_tokens: 10,
1462 completion_tokens: 5,
1463 cached_tokens: 0,
1464 cache_write_tokens: 0,
1465 },
1466 at: 102,
1467 },
1468 RunRecord::ToolBatch {
1469 calls: vec![ToolCallRecord {
1470 id: "c1".to_string(),
1471 name: "read_file".to_string(),
1472 arguments: "{}".to_string(),
1473 result: Some("body".to_string()),
1474 thought_signature: Some("sig".to_string()),
1475 }],
1476 at: 103,
1477 stage_index: 0,
1478 iteration: 0,
1479 response: "reading".to_string(),
1480 },
1481 RunRecord::ToolCallDone {
1482 iteration: 0,
1483 call_id: "c1".to_string(),
1484 result: "body".to_string(),
1485 at: 103,
1486 },
1487 RunRecord::ContextCheckpoint {
1488 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1489 at: 104,
1490 },
1491 RunRecord::ContextDiff {
1492 delta: ContextDelta {
1493 stage_name: "plan".to_string(),
1494 total_tokens: 3,
1495 max_tokens: 10_000,
1496 regions: vec![RegionDelta::Append {
1497 name: "conv".to_string(),
1498 entries: vec![entry("more", 2)],
1499 current_tokens: 3,
1500 }],
1501 },
1502 at: 105,
1503 },
1504 RunRecord::Message {
1505 message: MessageRecord {
1506 role: "user".to_string(),
1507 content: "another".to_string(),
1508 },
1509 at: 106,
1510 },
1511 RunRecord::StatusChanged {
1512 status: RunStatus::Complete,
1513 at: 107,
1514 },
1515 RunRecord::Checkpoint {
1516 meta: Box::new(meta()),
1517 context: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1518 at: 108,
1519 },
1520 RunRecord::Progress {
1521 meta: Box::new(meta()),
1522 delta: ContextDelta {
1523 stage_name: "plan".to_string(),
1524 total_tokens: 3,
1525 max_tokens: 10_000,
1526 regions: vec![RegionDelta::Append {
1527 name: "conv".to_string(),
1528 entries: vec![entry("step", 2)],
1529 current_tokens: 3,
1530 }],
1531 },
1532 at: 109,
1533 },
1534 ]
1535 }
1536
1537 #[test]
1538 fn archive_write_then_read_roundtrips_every_record_kind() {
1539 let records = all_record_kinds();
1540 let mut buf = Vec::new();
1541 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1542 for r in &records {
1543 write_record(&mut buf, r).unwrap();
1544 }
1545 let (version, read) = read_archive(&mut buf.as_slice()).unwrap();
1546 assert_eq!(version, RUN_ARCHIVE_VERSION);
1547 assert_eq!(read, records);
1548 }
1549
1550 #[test]
1551 fn read_archive_start_rejects_bad_magic() {
1552 let mut bytes: &[u8] = b"XXXX\x00\x01";
1553 let err = read_archive_start(&mut bytes).unwrap_err();
1554 assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1555 }
1556
1557 #[test]
1558 fn read_archive_start_reports_version() {
1559 let mut buf = Vec::new();
1560 write_archive_start(&mut buf, 7).unwrap();
1561 assert_eq!(read_archive_start(&mut buf.as_slice()).unwrap(), 7);
1562 }
1563
1564 #[test]
1565 fn read_record_returns_none_at_clean_eof() {
1566 let empty: &[u8] = &[];
1567 assert!(read_record(&mut { empty }).unwrap().is_none());
1568 }
1569
1570 #[test]
1571 fn read_record_errors_on_truncated_length_prefix() {
1572 let mut bytes: &[u8] = &[0, 0];
1574 let err = read_record(&mut bytes).unwrap_err();
1575 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1576 }
1577
1578 #[test]
1579 fn read_record_errors_on_truncated_payload() {
1580 let mut bytes: &[u8] = &[0, 0, 0, 0, 0, 0, 0, 10, 1, 2];
1582 let err = read_record(&mut bytes).unwrap_err();
1583 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1584 }
1585
1586 #[test]
1587 fn read_record_errors_on_empty_payload_at_boundary() {
1588 let mut bytes: &[u8] = &[0, 0, 0, 0, 0, 0, 0, 10];
1591 let err = read_record(&mut bytes).unwrap_err();
1592 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1593 }
1594
1595 #[test]
1596 fn read_record_errors_on_invalid_json_payload() {
1597 let mut buf = Vec::new();
1599 let bad = b"not json";
1600 buf.extend_from_slice(&(bad.len() as u64).to_be_bytes());
1601 buf.extend_from_slice(bad);
1602 let err = read_record(&mut buf.as_slice()).unwrap_err();
1603 assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1604 }
1605
1606 struct FailingReader;
1609 impl Read for FailingReader {
1610 fn read(&mut self, _buf: &mut [u8]) -> io::Result<usize> {
1611 Err(io::Error::other("device error"))
1612 }
1613 }
1614
1615 #[test]
1616 fn read_record_propagates_reader_errors() {
1617 let err = read_record(&mut FailingReader).unwrap_err();
1618 assert_eq!(err.kind(), io::ErrorKind::Other);
1619 }
1620
1621 #[test]
1622 fn read_archive_propagates_a_bad_preamble() {
1623 let mut bytes: &[u8] = b"LV";
1625 assert!(read_archive(&mut bytes).is_err());
1626 }
1627
1628 #[test]
1629 fn read_archive_propagates_a_bad_frame() {
1630 let mut buf = Vec::new();
1632 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1633 buf.extend_from_slice(&[0, 0, 0, 0, 0, 0, 0, 5, 1, 2]); let err = read_archive(&mut buf.as_slice()).unwrap_err();
1635 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1636 }
1637
1638 struct FailAfter {
1640 remaining: usize,
1641 }
1642 impl Write for FailAfter {
1643 fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
1644 if self.remaining == 0 {
1645 return Err(io::Error::other("disk full"));
1646 }
1647 let n = buf.len().min(self.remaining);
1648 self.remaining -= n;
1649 Ok(n)
1650 }
1651 fn flush(&mut self) -> io::Result<()> {
1652 Ok(())
1653 }
1654 }
1655
1656 #[test]
1657 fn fail_after_writer_flush_is_a_noop() {
1658 assert!(FailAfter { remaining: 1 }.flush().is_ok());
1659 }
1660
1661 #[test]
1662 fn write_archive_start_propagates_write_errors() {
1663 assert!(write_archive_start(&mut FailAfter { remaining: 0 }, 1).is_err());
1665 assert!(write_archive_start(&mut FailAfter { remaining: 4 }, 1).is_err());
1666 }
1667
1668 #[test]
1669 fn write_record_propagates_write_errors() {
1670 let rec = header();
1671 assert!(write_record(&mut FailAfter { remaining: 0 }, &rec).is_err());
1673 assert!(write_record(&mut FailAfter { remaining: 8 }, &rec).is_err());
1674 }
1675
1676 #[test]
1682 fn an_absurd_frame_length_is_an_error_not_an_allocation() {
1683 let mut buf = Vec::new();
1684 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1685 write_record(&mut buf, &header()).unwrap();
1686 buf.extend_from_slice(&u64::MAX.to_be_bytes());
1688
1689 let err = read_archive(&mut buf.as_slice())
1690 .expect_err("the strict reader must refuse an impossible frame");
1691 assert_eq!(err.kind(), io::ErrorKind::InvalidData, "{err}");
1692
1693 let (_, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
1696 assert_eq!(records, vec![header()]);
1697 }
1698
1699 #[test]
1700 fn read_archive_lenient_matches_strict_on_a_clean_archive() {
1701 let records = all_record_kinds();
1704 let mut buf = Vec::new();
1705 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1706 for r in &records {
1707 write_record(&mut buf, r).unwrap();
1708 }
1709 let (version, read) = read_archive_lenient(&mut buf.as_slice()).unwrap();
1710 assert_eq!(version, RUN_ARCHIVE_VERSION);
1711 assert_eq!(read, records);
1712 }
1713
1714 #[test]
1715 fn read_archive_lenient_keeps_valid_prefix_before_a_torn_tail() {
1716 let mut buf = Vec::new();
1720 write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1721 write_record(&mut buf, &header()).unwrap();
1722 write_record(
1723 &mut buf,
1724 &RunRecord::ContextCheckpoint {
1725 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1726 at: 1,
1727 },
1728 )
1729 .unwrap();
1730 buf.extend_from_slice(&[0, 0, 0, 0, 0, 0, 0, 10, 1, 2]);
1732
1733 assert!(read_archive(&mut buf.as_slice()).is_err());
1735 let (version, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
1737 assert_eq!(version, RUN_ARCHIVE_VERSION);
1738 assert_eq!(records.len(), 2);
1739 let folded = fold(&records).expect("prefix starts with a Header");
1740 assert_eq!(folded.context.regions[0].entries.len(), 1);
1741 }
1742
1743 #[test]
1744 fn read_archive_lenient_still_errors_on_a_bad_preamble() {
1745 let mut bad_magic: &[u8] = b"XXXX\x00\x01";
1748 assert!(read_archive_lenient(&mut bad_magic).is_err());
1749 let mut short: &[u8] = b"LVR1";
1751 assert!(read_archive_lenient(&mut short).is_err());
1752 }
1753
1754 #[test]
1755 fn read_archive_start_errors_on_truncated_version() {
1756 let mut bytes: &[u8] = b"LVR1";
1758 let err = read_archive_start(&mut bytes).unwrap_err();
1759 assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1760 }
1761
1762 #[test]
1765 fn fold_requires_a_header_first() {
1766 assert!(fold(&[]).is_none());
1767 assert!(
1768 fold(&[RunRecord::StatusChanged {
1769 status: RunStatus::Complete,
1770 at: 1
1771 }])
1772 .is_none()
1773 );
1774 }
1775
1776 #[test]
1777 fn fold_reconstructs_state_from_the_journal() {
1778 let records = all_record_kinds();
1779 let folded = fold(&records).expect("has header");
1780 assert_eq!(folded.identity.machine_id, "machine-b");
1782 assert_eq!(folded.identity.world_id, "world-y");
1783 assert_eq!(folded.inference_count, 1);
1785 assert_eq!(folded.tool_call_count, 1);
1786 assert_eq!(folded.messages.len(), 1);
1788 assert_eq!(folded.messages[0].content, "another");
1789 assert_eq!(folded.context.regions[0].name, "conv");
1792 assert_eq!(folded.context.regions[0].entries.len(), 2);
1793 assert_eq!(folded.context.total_tokens, 3);
1794 assert_eq!(folded.meta.run_id, "run-1");
1795 let pending = folded.pending_batch.expect("batch never applied");
1798 assert_eq!(pending.calls[0].result.as_deref(), Some("body"));
1799 }
1800
1801 #[test]
1802 fn fold_applies_context_diffs_over_a_checkpoint() {
1803 let records = vec![
1805 header(),
1806 RunRecord::ContextCheckpoint {
1807 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1808 at: 1,
1809 },
1810 RunRecord::ContextDiff {
1811 delta: ContextDelta {
1812 stage_name: "plan".to_string(),
1813 total_tokens: 3,
1814 max_tokens: 10_000,
1815 regions: vec![RegionDelta::Append {
1816 name: "conv".to_string(),
1817 entries: vec![entry("there", 2)],
1818 current_tokens: 3,
1819 }],
1820 },
1821 at: 2,
1822 },
1823 ];
1824 let folded = fold(&records).unwrap();
1825 assert_eq!(folded.context.regions[0].entries.len(), 2);
1826 assert_eq!(folded.context.total_tokens, 3);
1827 }
1828
1829 #[test]
1830 fn fold_later_header_updates_identity_and_meta() {
1831 let mut second_meta = meta();
1833 second_meta.status = RunStatus::Running;
1834 let records = vec![
1835 header(),
1836 RunRecord::Header {
1837 identity: RunIdentity {
1838 run_id: "run-1".to_string(),
1839 machine_id: "machine-c".to_string(),
1840 world_id: "world-z".to_string(),
1841 created_at: 200,
1842 },
1843 meta: Box::new(second_meta),
1844 },
1845 ];
1846 let folded = fold(&records).unwrap();
1847 assert_eq!(folded.identity.machine_id, "machine-c");
1848 assert_eq!(folded.meta.status, RunStatus::Running);
1849 }
1850
1851 #[test]
1852 fn fold_progress_applies_meta_and_context_diff() {
1853 let mut advanced = meta();
1854 advanced.status = RunStatus::Running;
1855 advanced.iteration = 5;
1856 let records = vec![
1857 header(),
1858 RunRecord::ContextCheckpoint {
1859 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1860 at: 1,
1861 },
1862 RunRecord::Progress {
1863 meta: Box::new(advanced),
1864 delta: ContextDelta {
1865 stage_name: "plan".to_string(),
1866 total_tokens: 3,
1867 max_tokens: 10_000,
1868 regions: vec![RegionDelta::Append {
1869 name: "conv".to_string(),
1870 entries: vec![entry("there", 2)],
1871 current_tokens: 3,
1872 }],
1873 },
1874 at: 2,
1875 },
1876 ];
1877 let folded = fold(&records).unwrap();
1878 assert_eq!(folded.meta.iteration, 5);
1879 assert_eq!(folded.meta.status, RunStatus::Running);
1880 assert_eq!(folded.context.regions[0].entries.len(), 2);
1881 }
1882
1883 #[test]
1887 fn fold_carries_a_submitted_final_output_through_progress() {
1888 let mut answered = meta();
1889 answered.final_output = Some(
1890 crate::output::FinalOutput::new(
1891 "renamed two helpers",
1892 Some("markdown".to_string()),
1893 "summary".to_string(),
1894 9,
1895 )
1896 .descriptor(),
1897 );
1898 answered.output_request = Some(crate::output::OutputSpec {
1899 format: Some("a2ui".to_string()),
1900 ..Default::default()
1901 });
1902 let records = vec![
1903 header(),
1904 RunRecord::Progress {
1905 meta: Box::new(answered),
1906 delta: ContextDelta {
1907 stage_name: "summary".to_string(),
1908 total_tokens: 0,
1909 max_tokens: 10_000,
1910 regions: vec![],
1911 },
1912 at: 2,
1913 },
1914 ];
1915 let folded = fold(&records).unwrap();
1916 let output = folded.meta.final_output.expect("the answer folded through");
1917 assert_eq!(output.bytes, "renamed two helpers".len());
1920 assert_eq!(output.stage, "summary");
1921 assert_eq!(
1922 folded.meta.output_request.and_then(|s| s.format).as_deref(),
1923 Some("a2ui")
1924 );
1925 }
1926
1927 fn call(id: &str, result: Option<&str>) -> ToolCallRecord {
1930 ToolCallRecord {
1931 id: id.to_string(),
1932 name: "shell".to_string(),
1933 arguments: "{}".to_string(),
1934 result: result.map(str::to_string),
1935 thought_signature: None,
1936 }
1937 }
1938
1939 fn batch(iteration: usize, calls: Vec<ToolCallRecord>) -> RunRecord {
1940 RunRecord::ToolBatch {
1941 calls,
1942 at: 10,
1943 stage_index: 0,
1944 iteration,
1945 response: "running tools".to_string(),
1946 }
1947 }
1948
1949 fn turn_entry(call_ids: &[&str]) -> RegionEntrySnapshot {
1951 let mut e = entry("turn", 1);
1952 e.kind = crate::region::EntryKind::AssistantTurn {
1953 tool_calls: call_ids
1954 .iter()
1955 .map(|id| crate::region::SerializedToolCall {
1956 id: id.to_string(),
1957 name: "shell".to_string(),
1958 arguments: serde_json::Value::Null,
1959 thought_signature: None,
1960 })
1961 .collect(),
1962 };
1963 e
1964 }
1965
1966 #[test]
1967 fn fold_surfaces_a_pending_batch_with_merged_results() {
1968 let records = vec![
1973 header(),
1974 batch(
1975 0,
1976 vec![
1977 call("c1", None),
1978 call("c2", Some("inline")),
1979 call("c3", None),
1980 ],
1981 ),
1982 RunRecord::ToolCallDone {
1983 iteration: 0,
1984 call_id: "c1".to_string(),
1985 result: "ran".to_string(),
1986 at: 11,
1987 },
1988 ];
1989 let folded = fold(&records).unwrap();
1990 let pending = folded.pending_batch.expect("batch is pending");
1991 assert_eq!(pending.iteration, 0);
1992 assert_eq!(pending.response, "running tools");
1993 assert_eq!(pending.calls[0].result.as_deref(), Some("ran"));
1994 assert_eq!(pending.calls[1].result.as_deref(), Some("inline"));
1995 assert_eq!(pending.calls[2].result, None);
1996 assert_eq!(folded.tool_call_count, 3);
1997 }
1998
1999 #[test]
2000 fn fold_keeps_only_the_latest_batch_and_ignores_stale_done_records() {
2001 let mut advanced = meta();
2004 advanced.iteration = 1;
2005 let records = vec![
2006 header(),
2007 batch(0, vec![call("c1", None)]),
2008 RunRecord::Progress {
2009 meta: Box::new(advanced),
2010 delta: ContextDelta {
2011 stage_name: "plan".to_string(),
2012 total_tokens: 0,
2013 max_tokens: 10_000,
2014 regions: vec![],
2015 },
2016 at: 11,
2017 },
2018 batch(1, vec![call("c2", None)]),
2019 RunRecord::ToolCallDone {
2020 iteration: 0,
2021 call_id: "c1".to_string(),
2022 result: "stale".to_string(),
2023 at: 12,
2024 },
2025 RunRecord::ToolCallDone {
2026 iteration: 1,
2027 call_id: "unknown".to_string(),
2028 result: "nowhere to land".to_string(),
2029 at: 13,
2030 },
2031 ];
2032 let folded = fold(&records).unwrap();
2033 let pending = folded.pending_batch.expect("latest batch is pending");
2034 assert_eq!(pending.iteration, 1);
2035 assert_eq!(pending.calls.len(), 1);
2036 assert_eq!(pending.calls[0].id, "c2");
2037 assert_eq!(pending.calls[0].result, None, "stale/unknown dones ignored");
2038 }
2039
2040 #[test]
2041 fn fold_clears_a_batch_once_the_iteration_moves_on() {
2042 let mut advanced = meta();
2045 advanced.iteration = 1;
2046 let records = vec![
2047 header(),
2048 batch(0, vec![call("c1", Some("done"))]),
2049 RunRecord::Progress {
2050 meta: Box::new(advanced),
2051 delta: ContextDelta {
2052 stage_name: "plan".to_string(),
2053 total_tokens: 0,
2054 max_tokens: 10_000,
2055 regions: vec![],
2056 },
2057 at: 11,
2058 },
2059 ];
2060 assert_eq!(fold(&records).unwrap().pending_batch, None);
2061 }
2062
2063 #[test]
2064 fn fold_clears_a_batch_whose_turn_already_landed_in_the_window() {
2065 let records = vec![
2068 header(),
2069 batch(0, vec![call("c1", Some("done"))]),
2070 RunRecord::ContextCheckpoint {
2071 snapshot: snapshot("plan", vec![region("conv", vec![turn_entry(&["c1"])])]),
2072 at: 11,
2073 },
2074 ];
2075 assert_eq!(fold(&records).unwrap().pending_batch, None);
2076 }
2077
2078 #[test]
2079 fn context_contains_batch_matches_only_the_batch_turn() {
2080 let pending = PendingToolBatch {
2081 stage_index: 0,
2082 iteration: 0,
2083 response: String::new(),
2084 calls: vec![call("c1", None)],
2085 };
2086 let other = snapshot("plan", vec![region("conv", vec![turn_entry(&["zz"])])]);
2088 assert!(!context_contains_batch(&other, &pending));
2089 let own = snapshot(
2091 "plan",
2092 vec![region("conv", vec![turn_entry(&["c1", "c2"])])],
2093 );
2094 assert!(context_contains_batch(&own, &pending));
2095 let empty = PendingToolBatch {
2097 calls: vec![],
2098 ..pending
2099 };
2100 assert!(!context_contains_batch(&own, &empty));
2101 }
2102
2103 #[test]
2104 fn old_shape_tool_batch_json_still_parses() {
2105 let json = br#"{"ToolBatch":{"calls":[{"id":"c1","name":"shell","arguments":"{}","result":"ok"}],"at":9}}"#;
2109 let mut buf = Vec::new();
2110 buf.extend_from_slice(&(json.len() as u64).to_be_bytes());
2111 buf.extend_from_slice(json);
2112 let record = read_record(&mut buf.as_slice()).unwrap().unwrap();
2113 assert_eq!(
2114 record,
2115 RunRecord::ToolBatch {
2116 calls: vec![call("c1", Some("ok"))],
2117 at: 9,
2118 stage_index: 0,
2119 iteration: 0,
2120 response: String::new(),
2121 }
2122 );
2123 }
2124
2125 fn three_point_records() -> Vec<RunRecord> {
2129 let mut running = meta();
2130 running.status = RunStatus::Running;
2131 vec![
2132 header(),
2133 RunRecord::ContextCheckpoint {
2134 snapshot: snapshot("plan", vec![region("conv", vec![entry("first", 1)])]),
2135 at: 10,
2136 },
2137 RunRecord::ContextDiff {
2138 delta: ContextDelta {
2139 stage_name: "plan".to_string(),
2140 total_tokens: 2,
2141 max_tokens: 10_000,
2142 regions: vec![RegionDelta::Append {
2143 name: "conv".to_string(),
2144 entries: vec![entry("second", 1)],
2145 current_tokens: 2,
2146 }],
2147 },
2148 at: 20,
2149 },
2150 RunRecord::Progress {
2151 meta: Box::new(running),
2152 delta: ContextDelta {
2153 stage_name: "code".to_string(),
2154 total_tokens: 3,
2155 max_tokens: 10_000,
2156 regions: vec![RegionDelta::Append {
2157 name: "conv".to_string(),
2158 entries: vec![entry("third", 1)],
2159 current_tokens: 3,
2160 }],
2161 },
2162 at: 30,
2163 },
2164 ]
2165 }
2166
2167 #[test]
2168 fn visit_points_indexes_points_in_order_and_carries_the_running_window() {
2169 let records = three_point_records();
2170 let mut seen: Vec<(usize, i64, usize)> = Vec::new();
2171 visit_points(&records, &mut |point| {
2172 seen.push((
2173 point.index,
2174 point.at,
2175 point.context.regions[0].entries.len(),
2176 ));
2177 ControlFlow::Continue(())
2178 });
2179 assert_eq!(seen, vec![(0, 10, 1), (1, 20, 2), (2, 30, 3)]);
2181 }
2182
2183 #[test]
2186 fn visit_points_stops_at_the_first_break() {
2187 let records = three_point_records();
2188 let mut visits = 0;
2189 visit_points(&records, &mut |point| {
2190 visits += 1;
2191 if point.index == 1 {
2192 ControlFlow::Break(())
2193 } else {
2194 ControlFlow::Continue(())
2195 }
2196 });
2197 assert_eq!(
2198 visits, 2,
2199 "stopped at the breaking point, did not run the third"
2200 );
2201 }
2202
2203 #[test]
2204 fn visit_points_without_a_header_visits_nothing() {
2205 let mut visits = 0;
2206 {
2207 let mut count = |_: PointRef<'_>| {
2208 visits += 1;
2209 ControlFlow::Continue(())
2210 };
2211
2212 visit_points(&three_point_records(), &mut count);
2216 visit_points(&[], &mut count);
2219 visit_points(
2220 &[RunRecord::ContextCheckpoint {
2221 snapshot: snapshot("plan", vec![]),
2222 at: 1,
2223 }],
2224 &mut count,
2225 );
2226 }
2227 assert_eq!(visits, 3, "only the well-formed journal produced points");
2228 }
2229
2230 #[test]
2234 fn visit_points_and_replay_points_agree() {
2235 for records in [
2236 three_point_records(),
2237 vec![header()],
2238 vec![],
2239 vec![RunRecord::Message {
2240 message: MessageRecord {
2241 role: "user".to_string(),
2242 content: "x".to_string(),
2243 },
2244 at: 1,
2245 }],
2246 ] {
2247 let collected: Vec<RunPoint> = {
2248 let mut out = Vec::new();
2249 visit_points(&records, &mut |point| {
2250 out.push(RunPoint {
2251 meta: point.meta.clone(),
2252 context: point.context.clone(),
2253 at: point.at,
2254 });
2255 ControlFlow::Continue(())
2256 });
2257 out
2258 };
2259 assert_eq!(collected, replay_points(&records));
2260 }
2261 }
2262
2263 #[test]
2264 fn replay_points_requires_a_header() {
2265 assert!(replay_points(&[]).is_empty());
2266 assert!(
2267 replay_points(&[RunRecord::Message {
2268 message: MessageRecord {
2269 role: "user".to_string(),
2270 content: "x".to_string(),
2271 },
2272 at: 1,
2273 }])
2274 .is_empty()
2275 );
2276 }
2277
2278 #[test]
2279 fn replay_points_emits_a_snapshot_per_context_change() {
2280 let mut running = meta();
2283 running.status = RunStatus::Running;
2284 let records = vec![
2285 header(),
2286 RunRecord::Inference {
2287 stage: "plan".to_string(),
2288 iteration: 0,
2289 request: InferenceRequestRecord {
2290 model: "m".to_string(),
2291 system: vec![],
2292 messages: vec![],
2293 tool_names: vec![],
2294 temperature: 0.7,
2295 max_tokens: 10,
2296 },
2297 response: InferenceResponseRecord {
2298 content: "ok".to_string(),
2299 tool_calls: vec![],
2300 prompt_tokens: 1,
2301 completion_tokens: 1,
2302 cached_tokens: 0,
2303 cache_write_tokens: 0,
2304 },
2305 at: 1,
2306 },
2307 batch(0, vec![call("c1", None)]),
2308 RunRecord::ToolCallDone {
2309 iteration: 0,
2310 call_id: "c1".to_string(),
2311 result: "ran".to_string(),
2312 at: 1,
2313 },
2314 RunRecord::ContextCheckpoint {
2315 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
2316 at: 2,
2317 },
2318 RunRecord::StatusChanged {
2319 status: RunStatus::Running,
2320 at: 3,
2321 },
2322 RunRecord::Progress {
2323 meta: Box::new(running),
2324 delta: ContextDelta {
2325 stage_name: "implement".to_string(),
2326 total_tokens: 3,
2327 max_tokens: 10_000,
2328 regions: vec![RegionDelta::Append {
2329 name: "conv".to_string(),
2330 entries: vec![entry("more", 2)],
2331 current_tokens: 3,
2332 }],
2333 },
2334 at: 4,
2335 },
2336 ];
2337 let points = replay_points(&records);
2338 assert_eq!(points.len(), 2, "one point per context change");
2339 assert_eq!(points[0].at, 2);
2341 assert_eq!(points[0].context.regions[0].entries.len(), 1);
2342 assert_eq!(points[1].at, 4);
2345 assert_eq!(points[1].context.regions[0].entries.len(), 2);
2346 assert_eq!(points[1].context.stage_name, "implement");
2347 assert_eq!(points[1].meta.status, RunStatus::Running);
2348 }
2349
2350 #[test]
2351 fn replay_points_handles_context_diff_and_a_later_header() {
2352 let mut relabeled = meta();
2355 relabeled.agent_name = "renamed".to_string();
2356 let records = vec![
2357 header(),
2358 RunRecord::ContextCheckpoint {
2359 snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
2360 at: 1,
2361 },
2362 RunRecord::Header {
2363 identity: identity(),
2364 meta: Box::new(relabeled),
2365 },
2366 RunRecord::ContextDiff {
2367 delta: ContextDelta {
2368 stage_name: "plan".to_string(),
2369 total_tokens: 3,
2370 max_tokens: 10_000,
2371 regions: vec![RegionDelta::Append {
2372 name: "conv".to_string(),
2373 entries: vec![entry("more", 2)],
2374 current_tokens: 3,
2375 }],
2376 },
2377 at: 2,
2378 },
2379 ];
2380 let points = replay_points(&records);
2381 assert_eq!(points.len(), 2); assert_eq!(points[1].context.regions[0].entries.len(), 2);
2383 assert_eq!(points[1].meta.agent_name, "renamed");
2385 }
2386
2387 #[test]
2388 fn replay_points_over_a_full_checkpoint() {
2389 let records = vec![
2391 header(),
2392 RunRecord::Checkpoint {
2393 meta: Box::new(meta()),
2394 context: snapshot("review", vec![region("conv", vec![entry("x", 4)])]),
2395 at: 9,
2396 },
2397 ];
2398 let points = replay_points(&records);
2399 assert_eq!(points.len(), 1);
2400 assert_eq!(points[0].context.stage_name, "review");
2401 assert_eq!(points[0].context.regions[0].entries[0].tokens, 4);
2402 }
2403}