Skip to main content

leviath_core/
run_archive.rs

1//! A portable, self-contained record of an entire agent run.
2//!
3//! A run archive is a single append-only file that captures everything about a
4//! run - who owns it (which machine, which world/daemon instance), its metadata,
5//! every inference and tool batch, inbound messages, and the evolving context
6//! window - with enough fidelity that copying the file to another machine lets a
7//! daemon **continue the run where it left off** (LLM non-determinism aside) or
8//! replay it for debugging.
9//!
10//! ## Layout
11//!
12//! ```text
13//! MAGIC ("LVR1") | version (u16 BE) | frame*
14//! frame := len (u64 BE) | JSON-encoded RunRecord
15//! ```
16//!
17//! The framing is binary and codec-agnostic (a future release can swap the JSON
18//! payload for a compact binary codec without changing readers that only seek by
19//! frame length). The first record is always a [`RunRecord::Header`].
20//!
21//! ## Portability / future migration
22//!
23//! [`RunIdentity`] records which machine + world/daemon instance owns a run, and
24//! [`RunRecord::OwnershipChanged`] records a handoff. This is deliberately more
25//! than today needs: the format is meant to eventually let a run start on one
26//! machine, pause, and resume on another - including a machine declining a run
27//! whose tools it lacks and waiting for a capable host. That scheduling logic
28//! isn't built yet; the format simply reserves room for it (ownership handoffs
29//! are first-class, the version field gates changes, and new record variants can
30//! be added without disturbing the frame layout).
31//!
32//! ## Efficiency
33//!
34//! Context windows are the bulk of a run. Rather than snapshot the whole window
35//! on every step, a writer emits an occasional full [`RunRecord::ContextCheckpoint`]
36//! and, between checkpoints, small [`RunRecord::ContextDiff`] records describing
37//! only what changed (the common case between inferences is a pure append to one
38//! region). [`diff_context`]/[`apply_delta`] compute and replay those diffs, and
39//! [`fold`] reconstructs the current state from the whole journal.
40
41use 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/// Identity + ownership of a run.
49///
50/// `machine_id` + `world_id` make a run unambiguously attributable even when
51/// several daemons share a filesystem and might otherwise pick the same
52/// `run_id` - a daemon can read a run's owner before deciding whether to resume
53/// or leave it alone.
54#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
55pub struct RunIdentity {
56    /// The run's id (its directory/file name).
57    pub run_id: String,
58    /// Stable fingerprint of the machine that owns the run.
59    pub machine_id: String,
60    /// Id of the specific world/daemon instance that owns the run.
61    pub world_id: String,
62    /// Unix seconds when the archive was created.
63    pub created_at: i64,
64}
65
66/// A conversation message as recorded in the archive.
67#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
68pub struct MessageRecord {
69    /// `"user"` / `"assistant"` / `"tool"` / `"system"`.
70    pub role: String,
71    /// The message text.
72    pub content: String,
73}
74
75/// A single tool call and (once executed) its result.
76#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
77pub struct ToolCallRecord {
78    /// The tool-call id.
79    pub id: String,
80    /// The tool name.
81    pub name: String,
82    /// The JSON arguments, stringified.
83    pub arguments: String,
84    /// The result text, once the tool has run (`None` while pending).
85    pub result: Option<String>,
86    /// Opaque provider token that must be replayed with this call (Gemini's
87    /// `thought_signature`). Carried so a restored batch can rebuild the exact
88    /// assistant turn.
89    #[serde(default, skip_serializing_if = "Option::is_none")]
90    pub thought_signature: Option<String>,
91}
92
93/// The outbound request of one inference (a provider-agnostic digest - enough to
94/// reproduce/debug the call without depending on `leviath-providers`).
95#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
96pub struct InferenceRequestRecord {
97    /// The model the request targeted.
98    pub model: String,
99    /// System-block texts, in order.
100    pub system: Vec<String>,
101    /// The conversation messages sent.
102    pub messages: Vec<MessageRecord>,
103    /// The tool names offered to the model.
104    pub tool_names: Vec<String>,
105    /// The temperature used.
106    pub temperature: f32,
107    /// The max output tokens requested.
108    pub max_tokens: usize,
109}
110
111/// The response of one inference.
112#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
113pub struct InferenceResponseRecord {
114    /// The assistant's text.
115    pub content: String,
116    /// Any tool calls the model requested.
117    pub tool_calls: Vec<ToolCallRecord>,
118    /// Prompt tokens billed.
119    pub prompt_tokens: usize,
120    /// Completion tokens billed.
121    pub completion_tokens: usize,
122    /// Tokens read from provider cache.
123    pub cached_tokens: usize,
124    /// Tokens written to provider cache.
125    pub cache_write_tokens: usize,
126}
127
128/// A per-region change within a [`ContextDelta`].
129#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
130pub enum RegionDelta {
131    /// A new region, or a region whose kind/max changed or whose entries were
132    /// rewritten in a non-append way - carried in full.
133    Set(RegionSnapshot),
134    /// Entries appended to an existing region (the common between-inference
135    /// case). The region's kind/max are unchanged.
136    Append {
137        /// The region name.
138        name: String,
139        /// The entries appended after the previously-recorded ones.
140        entries: Vec<RegionEntrySnapshot>,
141        /// The region's new token count.
142        current_tokens: usize,
143    },
144    /// An existing region emptied of entries.
145    Clear {
146        /// The region name.
147        name: String,
148    },
149    /// A region that no longer exists.
150    Remove {
151        /// The region name.
152        name: String,
153    },
154}
155
156/// The change to a context window since the previously-recorded snapshot.
157#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
158pub struct ContextDelta {
159    /// The window's stage name at this point.
160    pub stage_name: String,
161    /// The window's total token count at this point.
162    pub total_tokens: usize,
163    /// The window's max token budget at this point.
164    pub max_tokens: usize,
165    /// Per-region changes.
166    pub regions: Vec<RegionDelta>,
167}
168
169/// Which provider call a [`RunRecord::InferenceUsage`] belongs to.
170///
171/// A run bills for more than its stage turns, and the three auxiliary kinds are
172/// invisible in every other surface: they do not appear in the stage ledger and
173/// nothing else names them. Recording which kind spent the tokens is what turns
174/// a total into an explanation - "this run cost double what its stages did
175/// because its edges compact" is a sentence the journal can now support.
176#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
177#[serde(rename_all = "snake_case")]
178pub enum InferenceKind {
179    /// An ordinary stage turn: the agent thinking or calling tools. The default
180    /// so a journal written before this field existed reads back as stage work,
181    /// which is what every record in one is.
182    #[default]
183    Stage,
184    /// A region-summarizing call, from memory pressure or an edge transform.
185    Compaction,
186    /// The one-off call that names the run.
187    Title,
188    /// A call asking the model which stage to move to next.
189    Routing,
190}
191
192impl InferenceKind {
193    /// A short stable label, for logs and wire formats that want a string.
194    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    /// Whether this call is stage work the agent asked for, as opposed to
204    /// machinery the runtime ran on its behalf.
205    pub fn is_stage_work(&self) -> bool {
206        matches!(self, InferenceKind::Stage)
207    }
208}
209
210/// One entry in the run journal. Folding the sequence reconstructs the run.
211#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
212pub enum RunRecord {
213    /// The run's identity + static metadata. Always the first record.
214    Header {
215        /// Ownership/identity.
216        identity: RunIdentity,
217        /// The run metadata at archive-creation time.
218        meta: Box<RunMeta>,
219    },
220    /// Ownership handed to a different world/machine (e.g. resumed elsewhere).
221    OwnershipChanged {
222        /// The new owning machine.
223        machine_id: String,
224        /// The new owning world/daemon instance.
225        world_id: String,
226        /// Unix seconds.
227        at: i64,
228    },
229    /// One inference: what went out and what came back.
230    Inference {
231        /// The stage the agent was in.
232        stage: String,
233        /// The stage-local iteration index.
234        iteration: usize,
235        /// The request digest.
236        request: InferenceRequestRecord,
237        /// The response.
238        response: InferenceResponseRecord,
239        /// Unix seconds.
240        at: i64,
241    },
242    /// What one provider call cost, written as it lands.
243    ///
244    /// [`RunRecord::Progress`] carries cumulative counters, so two calls between
245    /// two ticks are indistinguishable downstream: their sum arrives as one
246    /// number, and a chart of it shows a spike no single call ever made. This
247    /// record is per call, so token-over-time telemetry is exact and "no request
248    /// ever exceeded the window" is provable from the journal rather than
249    /// inferred from region-size checkpoints.
250    ///
251    /// Deliberately lighter than [`RunRecord::Inference`], which carries the
252    /// full request and response: every call re-sends the whole window, so
253    /// journaling those bodies for each of them would multiply the file by the
254    /// context size. The window is already recoverable from the surrounding
255    /// [`RunRecord::ContextCheckpoint`] and [`RunRecord::ContextDiff`] records.
256    InferenceUsage {
257        /// Which kind of call this was.
258        #[serde(default)]
259        kind: InferenceKind,
260        /// The stage the run was in. Empty for a call with no stage of its own
261        /// (the title call, which runs once at spawn).
262        stage: String,
263        /// The stage-local iteration index.
264        iteration: usize,
265        /// The provider that served the call.
266        provider: String,
267        /// The model the call targeted.
268        model: String,
269        /// Prompt tokens billed.
270        prompt_tokens: usize,
271        /// Completion tokens billed.
272        completion_tokens: usize,
273        /// Tokens read from provider cache.
274        cached_tokens: usize,
275        /// Tokens written to provider cache.
276        cache_write_tokens: usize,
277        /// Unix seconds.
278        at: i64,
279    },
280    /// A batch of tool calls, written when the batch is dispatched to the tool
281    /// lane - before anything runs. Calls the dispatcher already resolved inline
282    /// (context tools, refusals, gate denials) carry `result: Some(..)`; lane
283    /// calls start at `result: None` and are completed by matching
284    /// [`RunRecord::ToolCallDone`] records as each call finishes. A batch still
285    /// pending at fold time surfaces as [`FoldedRun::pending_batch`] so a
286    /// crash-resume can replay executed calls instead of re-running them.
287    ToolBatch {
288        /// The calls (inline results pre-filled; lane calls pending).
289        calls: Vec<ToolCallRecord>,
290        /// Unix seconds.
291        at: i64,
292        /// The stage index the batch was dispatched in.
293        #[serde(default)]
294        stage_index: usize,
295        /// The stage-local iteration that produced the batch - the batch key
296        /// (one batch per iteration).
297        #[serde(default)]
298        iteration: usize,
299        /// The assistant text of the turn that issued the calls.
300        #[serde(default)]
301        response: String,
302    },
303    /// One tool call of the pending batch finished; its result.
304    ToolCallDone {
305        /// The iteration of the [`RunRecord::ToolBatch`] this belongs to.
306        iteration: usize,
307        /// The tool-call id.
308        call_id: String,
309        /// The result text.
310        result: String,
311        /// Unix seconds.
312        at: i64,
313    },
314    /// A full context-window snapshot that subsequent diffs rebase on.
315    ContextCheckpoint {
316        /// The full window snapshot.
317        snapshot: ContextSnapshot,
318        /// Unix seconds.
319        at: i64,
320    },
321    /// A context-window change since the previous snapshot/diff.
322    ContextDiff {
323        /// The delta.
324        delta: ContextDelta,
325        /// Unix seconds.
326        at: i64,
327    },
328    /// An inbound message.
329    Message {
330        /// The message.
331        message: MessageRecord,
332        /// Unix seconds.
333        at: i64,
334    },
335    /// A run-status change.
336    StatusChanged {
337        /// The new status.
338        status: RunStatus,
339        /// Unix seconds.
340        at: i64,
341    },
342    /// A full resumable checkpoint: the updated metadata + the full window, so a
343    /// reader can continue without folding the whole journal.
344    Checkpoint {
345        /// The run metadata as of this checkpoint.
346        meta: Box<RunMeta>,
347        /// The full window snapshot as of this checkpoint.
348        context: ContextSnapshot,
349        /// Unix seconds.
350        at: i64,
351    },
352    /// A step forward: the updated metadata plus a *diff* of the context window
353    /// since the previous point. This is the compact per-tick record the writer
354    /// emits between full checkpoints - meta is small, and the context (the bulk)
355    /// is carried as a [`ContextDelta`] rather than a full snapshot.
356    Progress {
357        /// The run metadata as of this step.
358        meta: Box<RunMeta>,
359        /// The context change since the previous recorded point.
360        delta: ContextDelta,
361        /// Unix seconds.
362        at: i64,
363    },
364}
365
366// ─── context diffing ────────────────────────────────────────────────────────
367
368/// Whether `prev` is a prefix of `next` (same entries, in order, at the front).
369fn is_prefix(prev: &[RegionEntrySnapshot], next: &[RegionEntrySnapshot]) -> bool {
370    prev.len() <= next.len() && next[..prev.len()] == *prev
371}
372
373/// Compute the minimal-ish [`ContextDelta`] turning `prev` into `next`. Regions
374/// that only grew at the tail become a compact `Append`; everything else is
375/// carried as a `Set`/`Clear`/`Remove`.
376pub 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                    // unchanged - emit nothing
384                } 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// ─── digest-based diffing ───────────────────────────────────────────────────
419//
420// `diff_context` needs the previous snapshot only to answer two questions per
421// region: "did anything change?" and "did it change by appending at the tail?".
422// A per-entry fingerprint answers both, so the writer can retain this digest
423// instead of a full copy of every live run's context window (which doubled the
424// per-run resident cost of the persistence lane).
425
426/// Fingerprint of one region: everything `diff_context` compares except the
427/// entry contents themselves, which are folded into per-entry hashes.
428#[derive(Debug, Clone, PartialEq)]
429pub struct RegionDigest {
430    /// The region name.
431    pub name: String,
432    /// The region's stringified kind.
433    pub kind: String,
434    /// The region's token count at digest time.
435    pub current_tokens: usize,
436    /// The region's token budget at digest time.
437    pub max_tokens: usize,
438    /// One hash per entry, in order.
439    pub entries: Vec<u64>,
440}
441
442/// Fingerprint of a whole context window, cheap to retain per live run.
443#[derive(Debug, Clone, PartialEq, Default)]
444pub struct ContextDigest {
445    /// Per-region fingerprints, in snapshot order.
446    pub regions: Vec<RegionDigest>,
447}
448
449/// Hash one region entry. Every field participates: two entries that differ
450/// anywhere must digest differently, or a real change would be recorded as
451/// "unchanged" and the folded archive would silently drift from the run.
452fn 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    // kind / metadata / taint are small enums and values without a Hash impl;
459    // their serialized form is tiny next to `content` and hashes faithfully.
460    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
472/// Compute the retained fingerprint of `snapshot`.
473pub 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
489/// Whether `prev`'s entry hashes are a prefix of `next`'s entries.
490fn 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
498/// [`diff_context`] against a retained [`ContextDigest`] instead of a full
499/// previous snapshot. Produces the same delta shapes for the same changes:
500/// unchanged regions emit nothing, tail growth becomes `Append`, everything
501/// else `Set`/`Clear`/`Remove`.
502pub 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                    // emit nothing
515                } 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
549/// Apply a [`ContextDelta`] to `base` in place. Lenient: a delta referencing a
550/// region that isn't present is skipped rather than erroring, so folding never
551/// fails on a malformed diff.
552pub 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// ─── fold ───────────────────────────────────────────────────────────────────
595
596/// A tool batch that was dispatched but whose results never reached the context
597/// window - what a crash-resume must replay instead of re-running. `calls` carry
598/// every result recorded before the crash ([`RunRecord::ToolCallDone`] merged
599/// in); a call still at `result: None` genuinely never finished.
600#[derive(Debug, Clone, PartialEq)]
601pub struct PendingToolBatch {
602    /// The stage index the batch was dispatched in.
603    pub stage_index: usize,
604    /// The stage-local iteration that produced the batch.
605    pub iteration: usize,
606    /// The assistant text of the turn that issued the calls.
607    pub response: String,
608    /// The calls, with every recorded result merged in.
609    pub calls: Vec<ToolCallRecord>,
610}
611
612/// One provider call's cost, as folded out of the journal.
613///
614/// The flattened form of [`RunRecord::InferenceUsage`], so a consumer walking a
615/// folded run does not have to match the record enum to read a number.
616#[derive(Debug, Clone, PartialEq, Eq)]
617pub struct InferenceUsageRecord {
618    /// Which kind of call this was.
619    pub kind: InferenceKind,
620    /// The stage the run was in, empty for the title call.
621    pub stage: String,
622    /// The stage-local iteration index.
623    pub iteration: usize,
624    /// The provider that served the call.
625    pub provider: String,
626    /// The model the call targeted.
627    pub model: String,
628    /// Prompt tokens billed.
629    pub prompt_tokens: usize,
630    /// Completion tokens billed.
631    pub completion_tokens: usize,
632    /// Tokens read from provider cache.
633    pub cached_tokens: usize,
634    /// Tokens written to provider cache.
635    pub cache_write_tokens: usize,
636    /// Unix seconds.
637    pub at: i64,
638}
639
640/// The state reconstructed from a run journal - enough to resume or inspect the
641/// run at its latest recorded point.
642#[derive(Debug, Clone, PartialEq)]
643pub struct FoldedRun {
644    /// The run's current owner/identity.
645    pub identity: RunIdentity,
646    /// The latest run metadata.
647    pub meta: RunMeta,
648    /// The reconstructed current context window.
649    pub context: ContextSnapshot,
650    /// The recorded inbound messages, in order.
651    pub messages: Vec<MessageRecord>,
652    /// Number of inferences recorded.
653    pub inference_count: usize,
654    /// Per-call usage, in the order the calls landed.
655    ///
656    /// The point of keeping every entry rather than a running sum: a sum is
657    /// already available from [`RunMeta`], and what it cannot answer is whether
658    /// any single call exceeded the window, or which kind of call the spend went
659    /// to.
660    pub inference_usage: Vec<InferenceUsageRecord>,
661    /// Number of tool calls recorded.
662    pub tool_call_count: usize,
663    /// A dispatched tool batch whose results never made it into the context
664    /// window (the run crashed mid-batch). `None` when the run has no batch in
665    /// flight or the batch's turn already landed in `context`.
666    pub pending_batch: Option<PendingToolBatch>,
667}
668
669/// Whether `context` already contains the assistant turn of `batch` - i.e. the
670/// batch completed and `apply_tool_results` landed it before the crash, so there
671/// is nothing to replay. Matched by the first call id, which is unique per batch.
672pub 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
687/// Reconstruct a run's current state from its journal. Returns `None` if the
688/// records don't start with a [`RunRecord::Header`].
689pub 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                // Counted alongside the heavy variant: both name one provider
738                // call, and a consumer asking "how many calls" should not have
739                // to know which of the two the writer chose.
740                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                // A later batch replaces an earlier one - only the newest can
763                // still be in flight.
764                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                // Fill the matching pending call; a stale record for a replaced
778                // batch (iteration mismatch) is ignored.
779                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    // The batch is only pending if it was never applied. Two applied signals: a
803    // later inference moved the iteration on (even if a sliding window has since
804    // evicted the turn), or the batch's assistant turn is already in the folded
805    // window (the Progress carrying it landed before the crash).
806    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/// A run's context window at one recorded point in time, with the metadata
816/// (stage, iteration, status, …) in effect then. Produced by [`replay_points`].
817#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
818pub struct RunPoint {
819    /// The run metadata at this point.
820    pub meta: RunMeta,
821    /// The full context window at this point.
822    pub context: ContextSnapshot,
823    /// Unix seconds this point was recorded.
824    pub at: i64,
825}
826
827/// One replayed point, lent to a [`visit_points`] visitor rather than handed
828/// over. Borrowing is the whole purpose: see that function.
829#[derive(Debug)]
830pub struct PointRef<'a> {
831    /// Position in the timeline, counting only records that produce a point.
832    /// Stable for a given journal prefix, because the journal is append-only -
833    /// which is what makes it usable as a pagination cursor.
834    pub index: usize,
835    /// Unix seconds this point was recorded.
836    pub at: i64,
837    /// The run metadata in effect at this point.
838    pub meta: &'a RunMeta,
839    /// The full context window at this point.
840    pub context: &'a ContextSnapshot,
841}
842
843/// Replay a run journal, calling `visit` once per record that changes the
844/// context (a checkpoint, diff, or progress step), in order. Stops early if the
845/// visitor returns [`ControlFlow::Break`]. Does nothing if the records don't
846/// start with a [`RunRecord::Header`].
847///
848/// The point of lending each point instead of collecting them: replaying a run
849/// means carrying one running window and mutating it, so materializing the
850/// timeline costs a **full deep copy of the context window per point** - and a
851/// window holds every region's entry text. On a megabyte-scale journal that is
852/// hundreds of whole-window clones, which is why anything that wants a slice of
853/// the timeline, or just an answer to "does any point contain this text",
854/// should come through here rather than [`replay_points`].
855///
856/// `&mut dyn FnMut` rather than a generic parameter, deliberately: this is
857/// called from a handful of places with unrelated closure types, and one
858/// monomorphization keeps both the compiled size and the coverage instantiation
859/// count at one - the same reasoning `execute_with_shutdown` documents in the
860/// serve module.
861pub 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
876/// Streaming [`visit_points`] over a framed archive: validate the preamble,
877/// then read one record at a time and fold it into the running window - so a
878/// multi-megabyte `run.lvr` is walked holding one record and one window in
879/// memory, instead of the whole parsed journal (`read_archive` materializes
880/// every record first, typically 2-4x the file's bytes as structs).
881///
882/// Errors only on a bad preamble. Like [`read_archive_lenient`], a torn or
883/// unreadable frame ends the walk with the points already visited: the tail of
884/// a live run's journal can legitimately be mid-append.
885pub 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    // The first record has to be a Header, and a Header is a kind every build
891    // knows - so an unreadable frame here means this is not a foldable archive.
892    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    // A record kind from a later build is stepped over rather than ending the
900    // walk: it carries no context change this build can apply, and everything
901    // after it still does.
902    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
913/// The running state of a point replay: the metadata and window in effect,
914/// folded record by record. Shared by [`visit_points`] (in-memory records) and
915/// [`visit_archive_points`] (streamed records) so the two can never disagree
916/// about what a record means.
917struct PointFolder {
918    meta: RunMeta,
919    context: ContextSnapshot,
920    index: usize,
921}
922
923impl PointFolder {
924    /// Start a replay from the first record, which must be the Header -
925    /// anything else means this isn't a run journal, and the replay visits
926    /// nothing (`None`).
927    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    /// Fold one record; when it produces a timeline point, lend it to `visit`.
944    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            // Non-context records don't add a timeline point. Usage included:
981            // it says what a call cost, not what the window then held, and
982            // emitting a point per call would double the timeline with entries
983            // whose context is identical to their neighbour's.
984            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
1002/// Replay a run journal into the sequence of context-window snapshots over time,
1003/// one [`RunPoint`] per record that changes the context (a checkpoint, diff, or
1004/// progress step). This is what the context-history views (TUI/CLI/API) consume
1005/// to show the window "at each stage and point". Returns an empty vec if the
1006/// records don't start with a [`RunRecord::Header`].
1007///
1008/// Materializes every point, so it deep-copies the whole context window once per
1009/// point. Prefer [`visit_points`] when only part of the timeline is wanted, or
1010/// when the answer is a predicate rather than the points themselves.
1011pub 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    // Only the tests write through the trait now; the codec moved out.
1028    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    /// The instant the fixture pretends it is, on every construction.
1041    ///
1042    /// Arbitrary, and deliberately not the wall clock: `RunMeta::new` stamps
1043    /// `started_at`/`updated_at` from it, and these tests build the fixture
1044    /// once to write and again to compare against. Two reads straddling a
1045    /// second boundary produced two unequal `RunMeta`s, which failed whichever
1046    /// round-trip assertion happened to span the tick.
1047    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    /// Two constructions of the fixture are equal however much time passes
1065    /// between them. This is the property every round-trip assertion in this
1066    /// module rests on, and the one a wall-clock stamp quietly broke.
1067    #[test]
1068    fn the_fixture_does_not_move_with_the_clock() {
1069        let first = meta();
1070        let mut later = meta();
1071        // Rather than sleeping across a real second boundary, move the clock
1072        // the way it would have moved: an unpinned fixture differs by exactly
1073        // this, and a pinned one is rebuilt identically.
1074        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    /// A stable tag per region-delta shape - asserting on this avoids the
1123    /// uncovered `false` arm a `matches!` leaves when the assertion passes.
1124    /// Every arm is exercised across the diff tests below.
1125    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    // ── diff / apply round-trips ──
1135
1136    /// Applying `diff(a, b)` to a clone of `a` must reproduce `b`, for every
1137    /// region-delta shape (new, append, clear, remove, full-replace, unchanged).
1138    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        // Entries changed at the front (not an append) → full Set.
1204        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        // Same prefix entries but the region's kind changed → Set, not Append.
1214        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    // ── streaming point replay ──
1224
1225    /// Frame `records` exactly as `run.lvr` stores them.
1226    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    /// Collect `(index, at, total_tokens)` per visited point, or the stream
1236    /// error. One closure shared by every streamed test, including the
1237    /// bad-preamble one whose visitor never runs.
1238    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    /// Collect `(index, at, total_tokens)` per visited point.
1248    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        // A torn frame: a length prefix promising more than exists.
1301        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        // And an archive with no records at all visits nothing.
1314        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    // ── digest-based diffing ──
1341    //
1342    // `diff_context_digest(digest(a), b)` must produce the same delta as
1343    // `diff_context(a, b)` for every shape: the persistence lane retains only
1344    // the digest, and any divergence would silently corrupt the archive.
1345
1346    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        // And the digest-produced delta still round-trips.
1350        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    /// A token-count change with identical entries is still an (empty) Append
1413    /// carrying the new count, exactly as the full diff records it.
1414    #[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    /// Every field of an entry participates in its digest: a change anywhere
1426    /// must change the hash, or a real edit would fold as "unchanged".
1427    #[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        // And digesting the same entry twice is stable.
1459        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        // Append/Clear targeting a region not present are no-ops (not errors).
1474        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    // ── codec round-trips ──
1498
1499    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    /// The preamble round-trips the version it was written with. Previously
1633    /// checked with an arbitrary 7; that now names a framing generation this
1634    /// build cannot read, and is refused - see
1635    /// `an_archive_from_a_newer_format_is_refused_with_both_versions_named`.
1636    #[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        // Two bytes where an 8-byte length is expected → partial read → error.
1655        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        // A frame claiming 10 bytes but only 2 present after the 8-byte length.
1663        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        // A non-zero length with zero payload bytes → clean EOF at the payload
1671        // start is still a truncation (the frame promised bytes).
1672        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        // A well-framed payload that isn't a valid RunRecord.
1680        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    /// A reader whose `read` always errors, to exercise the read error path
1689    /// inside `read_exact_or_eof` (distinct from a clean EOF).
1690    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        // Too short to even hold the magic → the preamble read errors.
1706        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        // Valid preamble, then a truncated frame → the record read errors.
1713        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]); // len 5, 2 present
1716        let err = read_archive(&mut buf.as_slice()).unwrap_err();
1717        assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1718    }
1719
1720    /// A writer that fails after `ok_bytes` bytes, to exercise write error paths.
1721    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        // Fail on the magic write (0 bytes allowed) and on the version write.
1746        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        // Fail on the 8-byte length prefix, and (after it) on the payload.
1754        assert!(write_record(&mut FailAfter { remaining: 0 }, &rec).is_err());
1755        assert!(write_record(&mut FailAfter { remaining: 8 }, &rec).is_err());
1756    }
1757
1758    /// A torn *length prefix* is where a nonsense `u64` comes from, and the
1759    /// lenient reader exists precisely to survive a torn tail. Taking the
1760    /// length at its word would turn a crash-truncated archive into an
1761    /// allocation of that size - during daemon recovery, the one moment this
1762    /// reader is there to keep working.
1763    #[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        // A crash mid-append that left a garbage length behind.
1769        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        // And the lenient reader folds back to the intact record before it,
1776        // which is the behaviour recovery depends on.
1777        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        // With no torn tail, the lenient reader returns exactly what the strict
1784        // reader does.
1785        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        // A valid preamble + two full records, then a truncated frame (a crash
1799        // mid-append). The strict reader would reject the whole file; the lenient
1800        // reader returns the two intact records and stops at the torn tail.
1801        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        // A frame claiming 10 payload bytes but only 2 present → torn tail.
1813        buf.extend_from_slice(&[0, 0, 0, 0, 0, 0, 0, 10, 1, 2]);
1814
1815        // Strict rejects the whole archive.
1816        assert!(read_archive(&mut buf.as_slice()).is_err());
1817        // Lenient keeps the valid prefix and folds cleanly.
1818        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        // The preamble is validated strictly: a file that isn't a run archive at
1828        // all errors rather than folding to nothing.
1829        let mut bad_magic: &[u8] = b"XXXX\x00\x01";
1830        assert!(read_archive_lenient(&mut bad_magic).is_err());
1831        // A truncated version (valid magic, no version bytes) also errors.
1832        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        // Valid 4-byte magic but no version bytes → the version read errors.
1839        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    // ── fold ──
1845
1846    #[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        // Ownership was reassigned mid-journal.
1863        assert_eq!(folded.identity.machine_id, "machine-b");
1864        assert_eq!(folded.identity.world_id, "world-y");
1865        // Counters. Two inferences: the fixture carries one record of each
1866        // kind, and both name one provider call.
1867        assert_eq!(folded.inference_count, 2);
1868        assert_eq!(folded.inference_usage.len(), 1);
1869        assert_eq!(folded.tool_call_count, 1);
1870        // One inbound message recorded.
1871        assert_eq!(folded.messages.len(), 1);
1872        assert_eq!(folded.messages[0].content, "another");
1873        // The Progress step is the last context-affecting record: it layers its
1874        // append diff onto the preceding Checkpoint's window (hi + step).
1875        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        // The batch shares the meta's iteration and its turn never reached the
1880        // window, so it folds as pending (with the ToolCallDone merged in).
1881        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        // Header → checkpoint → diff (append). The diff must layer on the checkpoint.
1888        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        // A second Header (unusual, but tolerated) refreshes identity + meta.
1916        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    /// A submitted answer needs no record type of its own: `Progress` and
1968    /// `Checkpoint` both replace the whole `RunMeta`, so it folds along with
1969    /// everything else and a crash-resume finds the answer already there.
1970    #[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        // The descriptor, not the bytes: the answer itself is a sidecar file,
2002        // so what folds is the record of it.
2003        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    // ── pending tool batch (fold) ──
2012
2013    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    /// An entry whose kind is the assistant turn that issued `call_ids`.
2034    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        // meta().iteration is 0, matching the batch, and the context has no
2053        // assistant turn for it - so the batch is genuinely pending. c1's
2054        // ToolCallDone merges in; c2 keeps its dispatch-time inline result; c3
2055        // stays pending.
2056        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        // The second batch replaces the first; a ToolCallDone for the replaced
2086        // iteration is ignored, as is one naming a call the batch doesn't have.
2087        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        // A later inference bumped meta.iteration past the batch: the batch was
2127        // applied (even if a sliding window evicted the turn), nothing to replay.
2128        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        // Same iteration, but the context already holds the batch's assistant
2150        // turn: apply_tool_results ran before the crash, nothing to replay.
2151        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        // A window with an unrelated turn does not match.
2171        let other = snapshot("plan", vec![region("conv", vec![turn_entry(&["zz"])])]);
2172        assert!(!context_contains_batch(&other, &pending));
2173        // The batch's own turn matches by its first call id.
2174        let own = snapshot(
2175            "plan",
2176            vec![region("conv", vec![turn_entry(&["c1", "c2"])])],
2177        );
2178        assert!(context_contains_batch(&own, &pending));
2179        // A batch with no calls can never match.
2180        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        // Archives written before the batch-journal fields existed carry
2190        // ToolBatch records without stage_index/iteration/response (and calls
2191        // without thought_signature); serde defaults fill them in.
2192        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    // ── replay_points (context-window history) ──
2210
2211    /// Three context changes, so a windowing caller has something to page over.
2212    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        // Index counts points, not records - the Header produces none.
2264        assert_eq!(seen, vec![(0, 10, 1), (1, 20, 2), (2, 30, 3)]);
2265    }
2266
2267    /// The reason this function exists: a caller wanting one window, or an
2268    /// answer to "does any point match", must be able to stop.
2269    #[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            // A well-formed journal first, with the *same* visitor. Without
2297            // this the test would pass against a visitor that can never run at
2298            // all, which is exactly the reassurance it is not meant to give.
2299            visit_points(&three_point_records(), &mut count);
2300            // Neither of these starts with a Header, so neither is a replayable
2301            // journal and neither may produce a point.
2302            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    /// `replay_points` is now a thin collector over `visit_points`, so this
2315    /// pins the two together: if the reimplementation ever drifts, the borrowed
2316    /// walk and the materialized one stop agreeing here first.
2317    #[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        // Header (no point) → checkpoint (point 1) → status (no point, but tracked)
2365        // → progress diff (point 2). Non-context records don't add points.
2366        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        // First point: the checkpoint window.
2424        assert_eq!(points[0].at, 2);
2425        assert_eq!(points[0].context.regions[0].entries.len(), 1);
2426        // Second point: the progress diff layered on, with the running status
2427        // carried from the StatusChanged + the progress meta.
2428        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        // A standalone ContextDiff is a point; a second Header refreshes meta
2437        // without adding a point.
2438        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); // checkpoint + diff (header adds no point)
2466        assert_eq!(points[1].context.regions[0].entries.len(), 2);
2467        // The later Header's meta is in effect at the diff point.
2468        assert_eq!(points[1].meta.agent_name, "renamed");
2469    }
2470
2471    #[test]
2472    fn replay_points_over_a_full_checkpoint() {
2473        // A `Checkpoint` (full meta+context) is also a point.
2474        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    /// Each kind has to survive the wire under its own name: the label is what
2489    /// a consumer groups a token chart by, so a rename that silently reordered
2490    /// the enum would re-attribute somebody's spend.
2491    #[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    /// Only stage turns are work the agent asked for. The other three are
2508    /// machinery the runtime ran on its behalf, which is the split anything
2509    /// reporting "what did my agent actually do" needs.
2510    #[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    /// A journal written before the field existed has to read back as stage
2526    /// work rather than failing to parse - every record in one is a stage turn,
2527    /// because nothing else was ever written.
2528    #[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        // Compared whole rather than destructured: the point is that the
2544        // missing field defaults and every present one still lands, and a
2545        // destructure that pulled out `kind` alone would pass even if the rest
2546        // had been dropped.
2547        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    /// The point of the record. Folding keeps every call in order, so a
2566    /// consumer can ask what any single one cost - which the cumulative
2567    /// counters cannot answer, because two calls between two ticks arrive as
2568    /// their sum.
2569    #[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        // The shape the issue reported: a compaction call and a stage call
2584        // landing between the same two progress ticks.
2585        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        // The whole reason this exists: their sum is 28k, and a reader of the
2606        // cumulative counter alone would see one 28k request and reasonably ask
2607        // whether a 32k window had been violated. Neither call came close.
2608        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    /// The heavy variant and the light one both name one provider call, so a
2618    /// consumer counting calls should not have to know which the writer chose.
2619    #[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    /// Replaying a journal has to land on exactly the state the run was in.
2637    ///
2638    /// Issue #455 reports evictable regions drifting - a folded `logs` region
2639    /// holding entries the live agent had already lost, and elsewhere fewer
2640    /// than it held. This walks a region through the mutations a temporary
2641    /// region actually performs (append, evict-oldest, evict-and-append in one
2642    /// step, clear) and checks the digest -> delta -> apply chain reproduces
2643    /// every intermediate state exactly.
2644    #[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            // Evict oldest.
2683            snap(&[("b", 20), ("c", 30)]),
2684            // Evict and append in one step - what a temporary region does when
2685            // an add pushes it over its bound.
2686            snap(&[("c", 30), ("d", 40)]),
2687            // Several evictions at once.
2688            snap(&[("d", 40)]),
2689            // Repeated identical content, the case a content hash cannot tell
2690            // apart by value alone.
2691            snap(&[("d", 40), ("x", 5)]),
2692            snap(&[("x", 5), ("x", 5)]),
2693            snap(&[("x", 5)]),
2694            snap(&[]),
2695        ];
2696
2697        // Fold exactly as the lane writes and the reader replays: retain a
2698        // digest, diff the next snapshot against it, apply to the running base.
2699        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    /// One frame from a later build, written the way a later build would write
2714    /// it: correctly framed, with a variant name this one has never heard of.
2715    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    /// Returns the archive and the size of the unknown frame's payload, so a
2726    /// caller can assert on the exact `Frame` it expects back.
2727    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    /// A record's variant name, read off its serialized form. Avoids a match
2740    /// whose unreached arms would be uncovered, and pins the wire spelling.
2741    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    /// The forward-compatibility guarantee, and the reason it is worth having:
2753    /// adding a record kind used to truncate the journal for every older
2754    /// reader. The lenient reader stopped at the first unknown record and
2755    /// returned the prefix, so a build predating `InferenceUsage` would have
2756    /// read a 0.3.10 journal as "header, then nothing" - with no error.
2757    #[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    /// The streaming reader skips the same way the buffering one does. It has
2772    /// its own loop, so "both apply the rule" is a claim that needs checking
2773    /// rather than assuming.
2774    #[test]
2775    fn the_streaming_reader_also_steps_over_an_unknown_record() {
2776        // A context record *after* the unknown frame: the only way to reach it
2777        // is to step over that frame, so the point count is the proof.
2778        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    /// The frame reader distinguishes "cannot parse this" from "cannot find
2802    /// the end of this". Only the second is fatal, and the difference is what
2803    /// makes stepping over the first safe.
2804    #[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                // Stepped over, and it says how far.
2819                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    /// A torn tail still ends the read. Skipping is for frames whose bytes are
2829    /// all present; a truncated one has no known length to step over.
2830    #[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        // A length prefix promising more than follows.
2836        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    /// An archive from a future framing generation is refused rather than
2844    /// misread. Reading it anyway would not fail cleanly - it would take
2845    /// whatever the length prefixes happened to say and produce nonsense.
2846    #[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    /// An older archive is read normally: framing has not changed under it, so
2862    /// the only difference is which record kinds it happens to contain.
2863    #[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}