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, once the tool has run (`None` while pending). Text and
85    /// any stored parts the tool produced; a plain string on the wire when
86    /// it is text alone, which is what every journal written before parts
87    /// existed holds.
88    pub result: Option<crate::region::EntryContent>,
89    /// Opaque provider token that must be replayed with this call (Gemini's
90    /// `thought_signature`). Carried so a restored batch can rebuild the exact
91    /// assistant turn.
92    #[serde(default, skip_serializing_if = "Option::is_none")]
93    pub thought_signature: Option<String>,
94}
95
96/// The outbound request of one inference (a provider-agnostic digest - enough to
97/// reproduce/debug the call without depending on `leviath-providers`).
98#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
99pub struct InferenceRequestRecord {
100    /// The model the request targeted.
101    pub model: String,
102    /// System-block texts, in order.
103    pub system: Vec<String>,
104    /// The conversation messages sent.
105    pub messages: Vec<MessageRecord>,
106    /// The tool names offered to the model.
107    pub tool_names: Vec<String>,
108    /// The temperature used.
109    pub temperature: f32,
110    /// The max output tokens requested.
111    pub max_tokens: usize,
112}
113
114/// The response of one inference.
115#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
116pub struct InferenceResponseRecord {
117    /// The assistant's text.
118    pub content: String,
119    /// Any tool calls the model requested.
120    pub tool_calls: Vec<ToolCallRecord>,
121    /// Prompt tokens billed.
122    pub prompt_tokens: usize,
123    /// Completion tokens billed.
124    pub completion_tokens: usize,
125    /// Tokens read from provider cache.
126    pub cached_tokens: usize,
127    /// Tokens written to provider cache.
128    pub cache_write_tokens: usize,
129}
130
131/// A per-region change within a [`ContextDelta`].
132#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
133pub enum RegionDelta {
134    /// A new region, or a region whose kind/max changed or whose entries were
135    /// rewritten in a non-append way - carried in full.
136    Set(RegionSnapshot),
137    /// Entries appended to an existing region (the common between-inference
138    /// case). The region's kind/max are unchanged.
139    Append {
140        /// The region name.
141        name: String,
142        /// The entries appended after the previously-recorded ones.
143        entries: Vec<RegionEntrySnapshot>,
144        /// The region's new token count.
145        current_tokens: usize,
146    },
147    /// An existing region emptied of entries.
148    Clear {
149        /// The region name.
150        name: String,
151    },
152    /// A region that no longer exists.
153    Remove {
154        /// The region name.
155        name: String,
156    },
157}
158
159/// The change to a context window since the previously-recorded snapshot.
160#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
161pub struct ContextDelta {
162    /// The window's stage name at this point.
163    pub stage_name: String,
164    /// The window's total token count at this point.
165    pub total_tokens: usize,
166    /// The window's max token budget at this point.
167    pub max_tokens: usize,
168    /// Per-region changes.
169    pub regions: Vec<RegionDelta>,
170}
171
172/// Which provider call a [`RunRecord::InferenceUsage`] belongs to.
173///
174/// A run bills for more than its stage turns, and the three auxiliary kinds are
175/// invisible in every other surface: they do not appear in the stage ledger and
176/// nothing else names them. Recording which kind spent the tokens is what turns
177/// a total into an explanation - "this run cost double what its stages did
178/// because its edges compact" is a sentence the journal can now support.
179#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
180#[serde(rename_all = "snake_case")]
181pub enum InferenceKind {
182    /// An ordinary stage turn: the agent thinking or calling tools. The default
183    /// so a journal written before this field existed reads back as stage work,
184    /// which is what every record in one is.
185    #[default]
186    Stage,
187    /// A region-summarizing call, from memory pressure or an edge transform.
188    Compaction,
189    /// The one-off call that names the run.
190    Title,
191    /// A call asking the model which stage to move to next.
192    Routing,
193}
194
195impl InferenceKind {
196    /// A short stable label, for logs and wire formats that want a string.
197    pub fn label(&self) -> &'static str {
198        match self {
199            InferenceKind::Stage => "stage",
200            InferenceKind::Compaction => "compaction",
201            InferenceKind::Title => "title",
202            InferenceKind::Routing => "routing",
203        }
204    }
205
206    /// Whether this call is stage work the agent asked for, as opposed to
207    /// machinery the runtime ran on its behalf.
208    pub fn is_stage_work(&self) -> bool {
209        matches!(self, InferenceKind::Stage)
210    }
211}
212
213/// One entry in the run journal. Folding the sequence reconstructs the run.
214#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
215pub enum RunRecord {
216    /// The run's identity + static metadata. Always the first record.
217    Header {
218        /// Ownership/identity.
219        identity: RunIdentity,
220        /// The run metadata at archive-creation time.
221        meta: Box<RunMeta>,
222    },
223    /// Ownership handed to a different world/machine (e.g. resumed elsewhere).
224    OwnershipChanged {
225        /// The new owning machine.
226        machine_id: String,
227        /// The new owning world/daemon instance.
228        world_id: String,
229        /// Unix seconds.
230        at: i64,
231    },
232    /// One inference: what went out and what came back.
233    Inference {
234        /// The stage the agent was in.
235        stage: String,
236        /// The stage-local iteration index.
237        iteration: usize,
238        /// The request digest.
239        request: InferenceRequestRecord,
240        /// The response.
241        response: InferenceResponseRecord,
242        /// Unix seconds.
243        at: i64,
244    },
245    /// What one provider call cost, written as it lands.
246    ///
247    /// [`RunRecord::Progress`] carries cumulative counters, so two calls between
248    /// two ticks are indistinguishable downstream: their sum arrives as one
249    /// number, and a chart of it shows a spike no single call ever made. This
250    /// record is per call, so token-over-time telemetry is exact and "no request
251    /// ever exceeded the window" is provable from the journal rather than
252    /// inferred from region-size checkpoints.
253    ///
254    /// Deliberately lighter than [`RunRecord::Inference`], which carries the
255    /// full request and response: every call re-sends the whole window, so
256    /// journaling those bodies for each of them would multiply the file by the
257    /// context size. The window is already recoverable from the surrounding
258    /// [`RunRecord::ContextCheckpoint`] and [`RunRecord::ContextDiff`] records.
259    InferenceUsage {
260        /// Which kind of call this was.
261        #[serde(default)]
262        kind: InferenceKind,
263        /// The stage the run was in. Empty for a call with no stage of its own
264        /// (the title call, which runs once at spawn).
265        stage: String,
266        /// The stage-local iteration index.
267        iteration: usize,
268        /// The provider that served the call.
269        provider: String,
270        /// The model the call targeted.
271        model: String,
272        /// Prompt tokens billed.
273        prompt_tokens: usize,
274        /// Completion tokens billed.
275        completion_tokens: usize,
276        /// Tokens read from provider cache.
277        cached_tokens: usize,
278        /// Tokens written to provider cache.
279        cache_write_tokens: usize,
280        /// What this one call cost in USD, when it could be established.
281        ///
282        /// Per call rather than only per run, because the model can change
283        /// mid-run: a stage that fails over to a second provider has spent at
284        /// two different rates, and a run-level total cannot be re-derived from
285        /// `meta.model` afterwards. Recorded here, the journal answers "which
286        /// call cost that" without re-pricing anything.
287        ///
288        /// `None` means unpriced - the provider reported no cost and no rates
289        /// were known - never that the call was free.
290        #[serde(default, skip_serializing_if = "Option::is_none")]
291        cost_usd: Option<f64>,
292        /// Whether `cost_usd` is the provider's own figure rather than one
293        /// computed from published rates. Absent when there is no cost.
294        #[serde(default, skip_serializing_if = "Option::is_none")]
295        cost_reported_by_provider: Option<bool>,
296        /// Unix seconds.
297        at: i64,
298    },
299    /// A batch of tool calls, written when the batch is dispatched to the tool
300    /// lane - before anything runs. Calls the dispatcher already resolved inline
301    /// (context tools, refusals, gate denials) carry `result: Some(..)`; lane
302    /// calls start at `result: None` and are completed by matching
303    /// [`RunRecord::ToolCallDone`] records as each call finishes. A batch still
304    /// pending at fold time surfaces as [`FoldedRun::pending_batch`] so a
305    /// crash-resume can replay executed calls instead of re-running them.
306    ToolBatch {
307        /// The calls (inline results pre-filled; lane calls pending).
308        calls: Vec<ToolCallRecord>,
309        /// Unix seconds.
310        at: i64,
311        /// The stage index the batch was dispatched in.
312        #[serde(default)]
313        stage_index: usize,
314        /// The stage-local iteration that produced the batch - the batch key
315        /// (one batch per iteration).
316        #[serde(default)]
317        iteration: usize,
318        /// The assistant text of the turn that issued the calls.
319        #[serde(default)]
320        response: String,
321    },
322    /// One tool call of the pending batch finished; its result.
323    ToolCallDone {
324        /// The iteration of the [`RunRecord::ToolBatch`] this belongs to.
325        iteration: usize,
326        /// The tool-call id.
327        call_id: String,
328        /// The result: text and any stored parts.
329        result: crate::region::EntryContent,
330        /// Unix seconds.
331        at: i64,
332    },
333    /// A full context-window snapshot that subsequent diffs rebase on.
334    ContextCheckpoint {
335        /// The full window snapshot.
336        snapshot: ContextSnapshot,
337        /// Unix seconds.
338        at: i64,
339    },
340    /// A context-window change since the previous snapshot/diff.
341    ContextDiff {
342        /// The delta.
343        delta: ContextDelta,
344        /// Unix seconds.
345        at: i64,
346    },
347    /// An inbound message.
348    Message {
349        /// The message.
350        message: MessageRecord,
351        /// Unix seconds.
352        at: i64,
353    },
354    /// A run-status change.
355    StatusChanged {
356        /// The new status.
357        status: RunStatus,
358        /// Unix seconds.
359        at: i64,
360    },
361    /// A full resumable checkpoint: the updated metadata + the full window, so a
362    /// reader can continue without folding the whole journal.
363    Checkpoint {
364        /// The run metadata as of this checkpoint.
365        meta: Box<RunMeta>,
366        /// The full window snapshot as of this checkpoint.
367        context: ContextSnapshot,
368        /// Unix seconds.
369        at: i64,
370    },
371    /// A step forward: the updated metadata plus a *diff* of the context window
372    /// since the previous point. This is the compact per-tick record the writer
373    /// emits between full checkpoints - meta is small, and the context (the bulk)
374    /// is carried as a [`ContextDelta`] rather than a full snapshot.
375    Progress {
376        /// The run metadata as of this step.
377        meta: Box<RunMeta>,
378        /// The context change since the previous recorded point.
379        delta: ContextDelta,
380        /// Unix seconds.
381        at: i64,
382    },
383}
384
385// ─── context diffing ────────────────────────────────────────────────────────
386
387/// Whether `prev` is a prefix of `next` (same entries, in order, at the front).
388fn is_prefix(prev: &[RegionEntrySnapshot], next: &[RegionEntrySnapshot]) -> bool {
389    prev.len() <= next.len() && next[..prev.len()] == *prev
390}
391
392/// What the diff needs to know about a region as it was.
393///
394/// Two questions per region: did anything change, and if so, did it change
395/// only by appending at the tail? A full previous snapshot answers them by
396/// comparing entries; a retained [`RegionDigest`] answers them from per-entry
397/// hashes. Everything else about the diff - which regions are new, which
398/// went away, which were cleared - is the same algorithm over either.
399trait PriorRegion {
400    /// The region's name, which is what pairs it with its successor.
401    fn name(&self) -> &str;
402    /// Whether `next` is this region, unchanged.
403    fn unchanged(&self, next: &RegionSnapshot) -> bool;
404    /// Whether `next` kept this region's kind and budget and only appended.
405    fn appended_to(&self, next: &RegionSnapshot) -> bool;
406    /// How many entries this region held.
407    fn entry_count(&self) -> usize;
408}
409
410impl PriorRegion for RegionSnapshot {
411    fn name(&self) -> &str {
412        &self.name
413    }
414
415    fn unchanged(&self, next: &RegionSnapshot) -> bool {
416        self == next
417    }
418
419    fn appended_to(&self, next: &RegionSnapshot) -> bool {
420        self.kind == next.kind
421            && self.max_tokens == next.max_tokens
422            && is_prefix(&self.entries, &next.entries)
423    }
424
425    fn entry_count(&self) -> usize {
426        self.entries.len()
427    }
428}
429
430/// The delta turning the regions `prev` describes into `next`.
431///
432/// Regions that only grew at the tail become a compact `Append`; everything
433/// else is carried as a `Set`/`Clear`/`Remove`.
434fn diff_regions<P: PriorRegion>(prev: &[P], next: &ContextSnapshot) -> ContextDelta {
435    let mut regions = Vec::new();
436    for nr in &next.regions {
437        match prev.iter().find(|r| r.name() == nr.name) {
438            None => regions.push(RegionDelta::Set(nr.clone())),
439            Some(pr) => {
440                if pr.unchanged(nr) {
441                    // unchanged - emit nothing
442                } else if nr.entries.is_empty() && pr.entry_count() > 0 {
443                    regions.push(RegionDelta::Clear {
444                        name: nr.name.clone(),
445                    });
446                } else if pr.appended_to(nr) {
447                    regions.push(RegionDelta::Append {
448                        name: nr.name.clone(),
449                        entries: nr.entries[pr.entry_count()..].to_vec(),
450                        current_tokens: nr.current_tokens,
451                    });
452                } else {
453                    regions.push(RegionDelta::Set(nr.clone()));
454                }
455            }
456        }
457    }
458    for pr in prev {
459        if !next.regions.iter().any(|r| r.name == pr.name()) {
460            regions.push(RegionDelta::Remove {
461                name: pr.name().to_string(),
462            });
463        }
464    }
465    ContextDelta {
466        stage_name: next.stage_name.clone(),
467        total_tokens: next.total_tokens,
468        max_tokens: next.max_tokens,
469        regions,
470    }
471}
472
473/// Compute the minimal-ish [`ContextDelta`] turning `prev` into `next`. Regions
474/// that only grew at the tail become a compact `Append`; everything else is
475/// carried as a `Set`/`Clear`/`Remove`.
476pub fn diff_context(prev: &ContextSnapshot, next: &ContextSnapshot) -> ContextDelta {
477    diff_regions(&prev.regions, next)
478}
479
480// ─── digest-based diffing ───────────────────────────────────────────────────
481//
482// `diff_context` needs the previous snapshot only to answer two questions per
483// region: "did anything change?" and "did it change by appending at the tail?".
484// A per-entry fingerprint answers both, so the writer can retain this digest
485// instead of a full copy of every live run's context window (which doubled the
486// per-run resident cost of the persistence lane).
487
488/// Fingerprint of one region: everything `diff_context` compares except the
489/// entry contents themselves, which are folded into per-entry hashes.
490#[derive(Debug, Clone, PartialEq)]
491pub struct RegionDigest {
492    /// The region name.
493    pub name: String,
494    /// The region's stringified kind.
495    pub kind: String,
496    /// The region's token count at digest time.
497    pub current_tokens: usize,
498    /// The region's token budget at digest time.
499    pub max_tokens: usize,
500    /// One hash per entry, in order.
501    pub entries: Vec<u64>,
502}
503
504/// Fingerprint of a whole context window, cheap to retain per live run.
505#[derive(Debug, Clone, PartialEq, Default)]
506pub struct ContextDigest {
507    /// Per-region fingerprints, in snapshot order.
508    pub regions: Vec<RegionDigest>,
509}
510
511/// Hash one region entry. Every field participates: two entries that differ
512/// anywhere must digest differently, or a real change would be recorded as
513/// "unchanged" and the folded archive would silently drift from the run.
514fn entry_digest(entry: &RegionEntrySnapshot) -> u64 {
515    use std::hash::{Hash, Hasher};
516    let mut hasher = std::collections::hash_map::DefaultHasher::new();
517    entry.content.hash(&mut hasher);
518    entry.tokens.hash(&mut hasher);
519    entry.key.hash(&mut hasher);
520    // kind / metadata / taint are small enums and values without a Hash impl;
521    // their serialized form is tiny next to `content` and hashes faithfully.
522    serde_json::to_string(&entry.kind)
523        .expect("EntryKind always serializes")
524        .hash(&mut hasher);
525    serde_json::to_string(&entry.metadata)
526        .expect("entry metadata always serializes")
527        .hash(&mut hasher);
528    serde_json::to_string(&entry.taint)
529        .expect("taint always serializes")
530        .hash(&mut hasher);
531    hasher.finish()
532}
533
534/// Compute the retained fingerprint of `snapshot`.
535pub fn digest_context(snapshot: &ContextSnapshot) -> ContextDigest {
536    ContextDigest {
537        regions: snapshot
538            .regions
539            .iter()
540            .map(|r| RegionDigest {
541                name: r.name.clone(),
542                kind: r.kind.clone(),
543                current_tokens: r.current_tokens,
544                max_tokens: r.max_tokens,
545                entries: r.entries.iter().map(entry_digest).collect(),
546            })
547            .collect(),
548    }
549}
550
551/// Whether `prev`'s entry hashes are a prefix of `next`'s entries.
552fn is_prefix_digest(prev: &[u64], next: &[RegionEntrySnapshot]) -> bool {
553    prev.len() <= next.len()
554        && prev
555            .iter()
556            .zip(next)
557            .all(|(hash, entry)| *hash == entry_digest(entry))
558}
559
560impl PriorRegion for RegionDigest {
561    fn name(&self) -> &str {
562        &self.name
563    }
564
565    fn unchanged(&self, next: &RegionSnapshot) -> bool {
566        self.kind == next.kind
567            && self.max_tokens == next.max_tokens
568            && self.current_tokens == next.current_tokens
569            && self.entries.len() == next.entries.len()
570            && is_prefix_digest(&self.entries, &next.entries)
571    }
572
573    fn appended_to(&self, next: &RegionSnapshot) -> bool {
574        self.kind == next.kind
575            && self.max_tokens == next.max_tokens
576            && is_prefix_digest(&self.entries, &next.entries)
577    }
578
579    fn entry_count(&self) -> usize {
580        self.entries.len()
581    }
582}
583
584/// [`diff_context`] against a retained [`ContextDigest`] instead of a full
585/// previous snapshot. The same algorithm over a different idea of "before":
586/// unchanged regions emit nothing, tail growth becomes `Append`, everything
587/// else `Set`/`Clear`/`Remove`.
588pub fn diff_context_digest(prev: &ContextDigest, next: &ContextSnapshot) -> ContextDelta {
589    diff_regions(&prev.regions, next)
590}
591
592/// Apply a [`ContextDelta`] to `base` in place. Lenient: a delta referencing a
593/// region that isn't present is skipped rather than erroring, so folding never
594/// fails on a malformed diff.
595pub fn apply_delta(base: &mut ContextSnapshot, delta: &ContextDelta) {
596    base.stage_name = delta.stage_name.clone();
597    base.total_tokens = delta.total_tokens;
598    base.max_tokens = delta.max_tokens;
599    for region_delta in &delta.regions {
600        match region_delta {
601            RegionDelta::Set(snapshot) => {
602                match base.regions.iter_mut().find(|r| r.name == snapshot.name) {
603                    Some(existing) => *existing = snapshot.clone(),
604                    None => base.regions.push(snapshot.clone()),
605                }
606            }
607            RegionDelta::Append {
608                name,
609                entries,
610                current_tokens,
611            } => {
612                if let Some(region) = base.regions.iter_mut().find(|r| &r.name == name) {
613                    region.entries.extend(entries.iter().cloned());
614                    region.current_tokens = *current_tokens;
615                }
616            }
617            RegionDelta::Clear { name } => {
618                if let Some(region) = base.regions.iter_mut().find(|r| &r.name == name) {
619                    region.entries.clear();
620                    region.current_tokens = 0;
621                }
622            }
623            RegionDelta::Remove { name } => {
624                base.regions.retain(|r| &r.name != name);
625            }
626        }
627    }
628}
629
630mod codec;
631
632pub use codec::{
633    Frame, RUN_ARCHIVE_MAGIC, RUN_ARCHIVE_VERSION, read_archive, read_archive_lenient,
634    read_archive_start, read_frame, read_record, write_archive_start, write_record,
635};
636
637// ─── fold ───────────────────────────────────────────────────────────────────
638
639/// A tool batch that was dispatched but whose results never reached the context
640/// window - what a crash-resume must replay instead of re-running. `calls` carry
641/// every result recorded before the crash ([`RunRecord::ToolCallDone`] merged
642/// in); a call still at `result: None` genuinely never finished.
643#[derive(Debug, Clone, PartialEq)]
644pub struct PendingToolBatch {
645    /// The stage index the batch was dispatched in.
646    pub stage_index: usize,
647    /// The stage-local iteration that produced the batch.
648    pub iteration: usize,
649    /// The assistant text of the turn that issued the calls.
650    pub response: String,
651    /// The calls, with every recorded result merged in.
652    pub calls: Vec<ToolCallRecord>,
653}
654
655/// One provider call's cost, as folded out of the journal.
656///
657/// The flattened form of [`RunRecord::InferenceUsage`], so a consumer walking a
658/// folded run does not have to match the record enum to read a number.
659// `Eq` is not derivable once a cost is present: `f64` has no total equality.
660// `PartialEq` is what the tests compare with anyway.
661#[derive(Debug, Clone, PartialEq)]
662pub struct InferenceUsageRecord {
663    /// Which kind of call this was.
664    pub kind: InferenceKind,
665    /// The stage the run was in, empty for the title call.
666    pub stage: String,
667    /// The stage-local iteration index.
668    pub iteration: usize,
669    /// The provider that served the call.
670    pub provider: String,
671    /// The model the call targeted.
672    pub model: String,
673    /// Prompt tokens billed.
674    pub prompt_tokens: usize,
675    /// Completion tokens billed.
676    pub completion_tokens: usize,
677    /// Tokens read from provider cache.
678    pub cached_tokens: usize,
679    /// Tokens written to provider cache.
680    pub cache_write_tokens: usize,
681    /// What this call cost in USD, when it could be established at all.
682    /// `None` is unpriced, never free.
683    pub cost_usd: Option<f64>,
684    /// Whether `cost_usd` came from the provider rather than from rates.
685    pub cost_reported_by_provider: Option<bool>,
686    /// Unix seconds.
687    pub at: i64,
688}
689
690/// The state reconstructed from a run journal - enough to resume or inspect the
691/// run at its latest recorded point.
692#[derive(Debug, Clone, PartialEq)]
693pub struct FoldedRun {
694    /// The run's current owner/identity.
695    pub identity: RunIdentity,
696    /// The latest run metadata.
697    pub meta: RunMeta,
698    /// The reconstructed current context window.
699    pub context: ContextSnapshot,
700    /// The recorded inbound messages, in order.
701    pub messages: Vec<MessageRecord>,
702    /// Number of inferences recorded.
703    pub inference_count: usize,
704    /// Per-call usage, in the order the calls landed.
705    ///
706    /// The point of keeping every entry rather than a running sum: a sum is
707    /// already available from [`RunMeta`], and what it cannot answer is whether
708    /// any single call exceeded the window, or which kind of call the spend went
709    /// to.
710    pub inference_usage: Vec<InferenceUsageRecord>,
711    /// Number of tool calls recorded.
712    pub tool_call_count: usize,
713    /// A dispatched tool batch whose results never made it into the context
714    /// window (the run crashed mid-batch). `None` when the run has no batch in
715    /// flight or the batch's turn already landed in `context`.
716    pub pending_batch: Option<PendingToolBatch>,
717}
718
719/// Whether `context` already contains the assistant turn of `batch` - i.e. the
720/// batch completed and `apply_tool_results` landed it before the crash, so there
721/// is nothing to replay. Matched by the first call id, which is unique per batch.
722pub fn context_contains_batch(context: &ContextSnapshot, batch: &PendingToolBatch) -> bool {
723    let Some(first_id) = batch.calls.first().map(|c| c.id.as_str()) else {
724        return false;
725    };
726    context.regions.iter().any(|region| {
727        region.entries.iter().any(|entry| {
728            matches!(
729                &entry.kind,
730                crate::region::EntryKind::AssistantTurn { tool_calls }
731                    if tool_calls.iter().any(|tc| tc.id == first_id)
732            )
733        })
734    })
735}
736
737/// Reconstruct a run's current state from its journal. Returns `None` if the
738/// records don't start with a [`RunRecord::Header`].
739pub fn fold(records: &[RunRecord]) -> Option<FoldedRun> {
740    let mut iter = records.iter();
741    let (identity, meta) = match iter.next() {
742        Some(RunRecord::Header { identity, meta }) => (identity.clone(), (**meta).clone()),
743        _ => return None,
744    };
745    let mut folded = FoldedRun {
746        identity,
747        meta,
748        context: ContextSnapshot {
749            stage_name: String::new(),
750            total_tokens: 0,
751            max_tokens: 0,
752            regions: Vec::new(),
753        },
754        messages: Vec::new(),
755        inference_count: 0,
756        inference_usage: Vec::new(),
757        tool_call_count: 0,
758        pending_batch: None,
759    };
760    for record in iter {
761        match record {
762            RunRecord::Header { identity, meta } => {
763                folded.identity = identity.clone();
764                folded.meta = (**meta).clone();
765            }
766            RunRecord::OwnershipChanged {
767                machine_id,
768                world_id,
769                ..
770            } => {
771                folded.identity.machine_id = machine_id.clone();
772                folded.identity.world_id = world_id.clone();
773            }
774            RunRecord::Inference { .. } => folded.inference_count += 1,
775            RunRecord::InferenceUsage {
776                kind,
777                stage,
778                iteration,
779                provider,
780                model,
781                prompt_tokens,
782                completion_tokens,
783                cached_tokens,
784                cache_write_tokens,
785                cost_usd,
786                cost_reported_by_provider,
787                at,
788            } => {
789                // Counted alongside the heavy variant: both name one provider
790                // call, and a consumer asking "how many calls" should not have
791                // to know which of the two the writer chose.
792                folded.inference_count += 1;
793                folded.inference_usage.push(InferenceUsageRecord {
794                    kind: *kind,
795                    stage: stage.clone(),
796                    iteration: *iteration,
797                    provider: provider.clone(),
798                    model: model.clone(),
799                    prompt_tokens: *prompt_tokens,
800                    completion_tokens: *completion_tokens,
801                    cached_tokens: *cached_tokens,
802                    cache_write_tokens: *cache_write_tokens,
803                    cost_usd: *cost_usd,
804                    cost_reported_by_provider: *cost_reported_by_provider,
805                    at: *at,
806                });
807            }
808            RunRecord::ToolBatch {
809                calls,
810                stage_index,
811                iteration,
812                response,
813                ..
814            } => {
815                folded.tool_call_count += calls.len();
816                // A later batch replaces an earlier one - only the newest can
817                // still be in flight.
818                folded.pending_batch = Some(PendingToolBatch {
819                    stage_index: *stage_index,
820                    iteration: *iteration,
821                    response: response.clone(),
822                    calls: calls.clone(),
823                });
824            }
825            RunRecord::ToolCallDone {
826                iteration,
827                call_id,
828                result,
829                ..
830            } => {
831                // Fill the matching pending call; a stale record for a replaced
832                // batch (iteration mismatch) is ignored.
833                if let Some(batch) = folded
834                    .pending_batch
835                    .as_mut()
836                    .filter(|b| b.iteration == *iteration)
837                    && let Some(call) = batch.calls.iter_mut().find(|c| c.id == *call_id)
838                {
839                    call.result = Some(result.clone());
840                }
841            }
842            RunRecord::ContextCheckpoint { snapshot, .. } => folded.context = snapshot.clone(),
843            RunRecord::ContextDiff { delta, .. } => apply_delta(&mut folded.context, delta),
844            RunRecord::Message { message, .. } => folded.messages.push(message.clone()),
845            RunRecord::StatusChanged { status, .. } => folded.meta.status = status.clone(),
846            RunRecord::Checkpoint { meta, context, .. } => {
847                folded.meta = (**meta).clone();
848                folded.context = context.clone();
849            }
850            RunRecord::Progress { meta, delta, .. } => {
851                folded.meta = (**meta).clone();
852                apply_delta(&mut folded.context, delta);
853            }
854        }
855    }
856    // The batch is only pending if it was never applied. Two applied signals: a
857    // later inference moved the iteration on (even if a sliding window has since
858    // evicted the turn), or the batch's assistant turn is already in the folded
859    // window (the Progress carrying it landed before the crash).
860    if let Some(batch) = &folded.pending_batch
861        && (folded.meta.iteration != batch.iteration
862            || context_contains_batch(&folded.context, batch))
863    {
864        folded.pending_batch = None;
865    }
866    Some(folded)
867}
868
869/// A run's context window at one recorded point in time, with the metadata
870/// (stage, iteration, status, …) in effect then. Produced by [`replay_points`].
871#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
872pub struct RunPoint {
873    /// The run metadata at this point.
874    pub meta: RunMeta,
875    /// The full context window at this point.
876    pub context: ContextSnapshot,
877    /// Unix seconds this point was recorded.
878    pub at: i64,
879}
880
881/// One replayed point, lent to a [`visit_points`] visitor rather than handed
882/// over. Borrowing is the whole purpose: see that function.
883#[derive(Debug)]
884pub struct PointRef<'a> {
885    /// Position in the timeline, counting only records that produce a point.
886    /// Stable for a given journal prefix, because the journal is append-only -
887    /// which is what makes it usable as a pagination cursor.
888    pub index: usize,
889    /// Unix seconds this point was recorded.
890    pub at: i64,
891    /// The run metadata in effect at this point.
892    pub meta: &'a RunMeta,
893    /// The full context window at this point.
894    pub context: &'a ContextSnapshot,
895}
896
897/// Replay a run journal, calling `visit` once per record that changes the
898/// context (a checkpoint, diff, or progress step), in order. Stops early if the
899/// visitor returns [`ControlFlow::Break`]. Does nothing if the records don't
900/// start with a [`RunRecord::Header`].
901///
902/// The point of lending each point instead of collecting them: replaying a run
903/// means carrying one running window and mutating it, so materializing the
904/// timeline costs a **full deep copy of the context window per point** - and a
905/// window holds every region's entry text. On a megabyte-scale journal that is
906/// hundreds of whole-window clones, which is why anything that wants a slice of
907/// the timeline, or just an answer to "does any point contain this text",
908/// should come through here rather than [`replay_points`].
909///
910/// `&mut dyn FnMut` rather than a generic parameter, deliberately: this is
911/// called from a handful of places with unrelated closure types, and one
912/// monomorphization keeps both the compiled size and the coverage instantiation
913/// count at one - the same reasoning `execute_with_shutdown` documents in the
914/// serve module.
915pub fn visit_points(records: &[RunRecord], visit: &mut dyn FnMut(PointRef<'_>) -> ControlFlow<()>) {
916    let mut iter = records.iter();
917    let Some(mut folder) = (match iter.next() {
918        Some(first) => PointFolder::start(first),
919        None => None,
920    }) else {
921        return;
922    };
923    for record in iter {
924        if folder.push(record, visit).is_break() {
925            return;
926        }
927    }
928}
929
930/// Streaming [`visit_points`] over a framed archive: validate the preamble,
931/// then read one record at a time and fold it into the running window - so a
932/// multi-megabyte `run.lvr` is walked holding one record and one window in
933/// memory, instead of the whole parsed journal (`read_archive` materializes
934/// every record first, typically 2-4x the file's bytes as structs).
935///
936/// Errors only on a bad preamble. Like [`read_archive_lenient`], a torn or
937/// unreadable frame ends the walk with the points already visited: the tail of
938/// a live run's journal can legitimately be mid-append.
939pub fn visit_archive_points(
940    r: &mut dyn Read,
941    visit: &mut dyn FnMut(PointRef<'_>) -> ControlFlow<()>,
942) -> io::Result<()> {
943    read_archive_start(r)?;
944    // The first record has to be a Header, and a Header is a kind every build
945    // knows - so an unreadable frame here means this is not a foldable archive.
946    let mut folder = match read_frame(r) {
947        Ok(Some(Frame::Record(first))) => match PointFolder::start(&first) {
948            Some(folder) => folder,
949            None => return Ok(()),
950        },
951        _ => return Ok(()),
952    };
953    // A record kind from a later build is stepped over rather than ending the
954    // walk: it carries no context change this build can apply, and everything
955    // after it still does.
956    while let Ok(Some(frame)) = read_frame(r) {
957        let Frame::Record(record) = frame else {
958            continue;
959        };
960        if folder.push(&record, visit).is_break() {
961            return Ok(());
962        }
963    }
964    Ok(())
965}
966
967/// The running state of a point replay: the metadata and window in effect,
968/// folded record by record. Shared by [`visit_points`] (in-memory records) and
969/// [`visit_archive_points`] (streamed records) so the two can never disagree
970/// about what a record means.
971struct PointFolder {
972    meta: RunMeta,
973    context: ContextSnapshot,
974    index: usize,
975}
976
977impl PointFolder {
978    /// Start a replay from the first record, which must be the Header -
979    /// anything else means this isn't a run journal, and the replay visits
980    /// nothing (`None`).
981    fn start(first: &RunRecord) -> Option<Self> {
982        match first {
983            RunRecord::Header { meta, .. } => Some(Self {
984                meta: (**meta).clone(),
985                context: ContextSnapshot {
986                    stage_name: String::new(),
987                    total_tokens: 0,
988                    max_tokens: 0,
989                    regions: Vec::new(),
990                },
991                index: 0,
992            }),
993            _ => None,
994        }
995    }
996
997    /// Fold one record; when it produces a timeline point, lend it to `visit`.
998    fn push(
999        &mut self,
1000        record: &RunRecord,
1001        visit: &mut dyn FnMut(PointRef<'_>) -> ControlFlow<()>,
1002    ) -> ControlFlow<()> {
1003        let at = match record {
1004            RunRecord::Header { meta: m, .. } => {
1005                self.meta = (**m).clone();
1006                return ControlFlow::Continue(());
1007            }
1008            RunRecord::StatusChanged { status, .. } => {
1009                self.meta.status = status.clone();
1010                return ControlFlow::Continue(());
1011            }
1012            RunRecord::ContextCheckpoint { snapshot, at } => {
1013                self.context = snapshot.clone();
1014                *at
1015            }
1016            RunRecord::ContextDiff { delta, at } => {
1017                apply_delta(&mut self.context, delta);
1018                *at
1019            }
1020            RunRecord::Checkpoint {
1021                meta: m,
1022                context: c,
1023                at,
1024            } => {
1025                self.meta = (**m).clone();
1026                self.context = c.clone();
1027                *at
1028            }
1029            RunRecord::Progress { meta: m, delta, at } => {
1030                self.meta = (**m).clone();
1031                apply_delta(&mut self.context, delta);
1032                *at
1033            }
1034            // Non-context records don't add a timeline point. Usage included:
1035            // it says what a call cost, not what the window then held, and
1036            // emitting a point per call would double the timeline with entries
1037            // whose context is identical to their neighbour's.
1038            RunRecord::OwnershipChanged { .. }
1039            | RunRecord::Inference { .. }
1040            | RunRecord::InferenceUsage { .. }
1041            | RunRecord::ToolBatch { .. }
1042            | RunRecord::ToolCallDone { .. }
1043            | RunRecord::Message { .. } => return ControlFlow::Continue(()),
1044        };
1045        let flow = visit(PointRef {
1046            index: self.index,
1047            at,
1048            meta: &self.meta,
1049            context: &self.context,
1050        });
1051        self.index += 1;
1052        flow
1053    }
1054}
1055
1056/// Replay a run journal into the sequence of context-window snapshots over time,
1057/// one [`RunPoint`] per record that changes the context (a checkpoint, diff, or
1058/// progress step). This is what the context-history views (TUI/CLI/API) consume
1059/// to show the window "at each stage and point". Returns an empty vec if the
1060/// records don't start with a [`RunRecord::Header`].
1061///
1062/// Materializes every point, so it deep-copies the whole context window once per
1063/// point. Prefer [`visit_points`] when only part of the timeline is wanted, or
1064/// when the answer is a predicate rather than the points themselves.
1065pub fn replay_points(records: &[RunRecord]) -> Vec<RunPoint> {
1066    let mut points = Vec::new();
1067    visit_points(records, &mut |point| {
1068        points.push(RunPoint {
1069            meta: point.meta.clone(),
1070            context: point.context.clone(),
1071            at: point.at,
1072        });
1073        ControlFlow::Continue(())
1074    });
1075    points
1076}
1077
1078#[cfg(test)]
1079mod tests {
1080    use super::*;
1081    // Only the tests write through the trait now; the codec moved out.
1082    use crate::run_meta::RunStatus;
1083    use std::io::Write;
1084
1085    fn identity() -> RunIdentity {
1086        RunIdentity {
1087            run_id: "run-1".to_string(),
1088            machine_id: "machine-a".to_string(),
1089            world_id: "world-x".to_string(),
1090            created_at: 100,
1091        }
1092    }
1093
1094    /// The instant the fixture pretends it is, on every construction.
1095    ///
1096    /// Arbitrary, and deliberately not the wall clock: `RunMeta::new` stamps
1097    /// `started_at`/`updated_at` from it, and these tests build the fixture
1098    /// once to write and again to compare against. Two reads straddling a
1099    /// second boundary produced two unequal `RunMeta`s, which failed whichever
1100    /// round-trip assertion happened to span the tick.
1101    const FIXTURE_NOW: i64 = 1_700_000_000;
1102
1103    fn meta() -> RunMeta {
1104        let mut meta = RunMeta::new(
1105            "run-1".to_string(),
1106            "coder".to_string(),
1107            "/agents/coder".to_string(),
1108            "do it".to_string(),
1109            Some("anthropic/claude".to_string()),
1110            "/work".to_string(),
1111            2,
1112        );
1113        meta.started_at = FIXTURE_NOW;
1114        meta.updated_at = FIXTURE_NOW;
1115        meta
1116    }
1117
1118    /// Two constructions of the fixture are equal however much time passes
1119    /// between them. This is the property every round-trip assertion in this
1120    /// module rests on, and the one a wall-clock stamp quietly broke.
1121    #[test]
1122    fn the_fixture_does_not_move_with_the_clock() {
1123        let first = meta();
1124        let mut later = meta();
1125        // Rather than sleeping across a real second boundary, move the clock
1126        // the way it would have moved: an unpinned fixture differs by exactly
1127        // this, and a pinned one is rebuilt identically.
1128        assert_eq!(first, later, "the fixture is rebuilt identically");
1129        later.started_at += 1;
1130        assert_ne!(
1131            first, later,
1132            "and the comparison is sensitive to the field that used to drift"
1133        );
1134    }
1135
1136    fn entry(content: &str, tokens: usize) -> RegionEntrySnapshot {
1137        RegionEntrySnapshot {
1138            content: content.into(),
1139            tokens,
1140            kind: crate::region::EntryKind::Text,
1141            metadata: None,
1142            key: None,
1143            taint: Default::default(),
1144            reasoning: None,
1145        }
1146    }
1147
1148    fn region(name: &str, entries: Vec<RegionEntrySnapshot>) -> RegionSnapshot {
1149        let current = entries.iter().map(|e| e.tokens).sum();
1150        RegionSnapshot {
1151            name: name.to_string(),
1152            kind: "clearable".to_string(),
1153            current_tokens: current,
1154            max_tokens: 1000,
1155            entries,
1156            description: None,
1157        }
1158    }
1159
1160    fn snapshot(stage: &str, regions: Vec<RegionSnapshot>) -> ContextSnapshot {
1161        let total = regions.iter().map(|r| r.current_tokens).sum();
1162        ContextSnapshot {
1163            stage_name: stage.to_string(),
1164            total_tokens: total,
1165            max_tokens: 10_000,
1166            regions,
1167        }
1168    }
1169
1170    fn header() -> RunRecord {
1171        RunRecord::Header {
1172            identity: identity(),
1173            meta: Box::new(meta()),
1174        }
1175    }
1176
1177    /// A stable tag per region-delta shape - asserting on this avoids the
1178    /// uncovered `false` arm a `matches!` leaves when the assertion passes.
1179    /// Every arm is exercised across the diff tests below.
1180    fn region_delta_kind(d: &RegionDelta) -> &'static str {
1181        match d {
1182            RegionDelta::Set(_) => "set",
1183            RegionDelta::Append { .. } => "append",
1184            RegionDelta::Clear { .. } => "clear",
1185            RegionDelta::Remove { .. } => "remove",
1186        }
1187    }
1188
1189    // ── diff / apply round-trips ──
1190
1191    /// Applying `diff(a, b)` to a clone of `a` must reproduce `b`, for every
1192    /// region-delta shape (new, append, clear, remove, full-replace, unchanged).
1193    fn assert_diff_roundtrip(a: &ContextSnapshot, b: &ContextSnapshot) {
1194        let delta = diff_context(a, b);
1195        let mut base = a.clone();
1196        apply_delta(&mut base, &delta);
1197        assert_eq!(&base, b);
1198    }
1199
1200    #[test]
1201    fn diff_append_only_growth_is_compact() {
1202        let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1203        let b = snapshot(
1204            "s1",
1205            vec![region("conv", vec![entry("hi", 1), entry("there", 2)])],
1206        );
1207        let delta = diff_context(&a, &b);
1208        assert_eq!(region_delta_kind(&delta.regions[0]), "append");
1209        assert_diff_roundtrip(&a, &b);
1210    }
1211
1212    #[test]
1213    fn diff_new_region_is_set() {
1214        let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1215        let b = snapshot(
1216            "s1",
1217            vec![
1218                region("conv", vec![entry("hi", 1)]),
1219                region("plan", vec![entry("p", 3)]),
1220            ],
1221        );
1222        let delta = diff_context(&a, &b);
1223        assert!(delta.regions.iter().any(|d| region_delta_kind(d) == "set"));
1224        assert_diff_roundtrip(&a, &b);
1225    }
1226
1227    #[test]
1228    fn diff_cleared_region() {
1229        let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1230        let b = snapshot("s1", vec![region("conv", vec![])]);
1231        let delta = diff_context(&a, &b);
1232        assert_eq!(region_delta_kind(&delta.regions[0]), "clear");
1233        assert_diff_roundtrip(&a, &b);
1234    }
1235
1236    #[test]
1237    fn diff_removed_region() {
1238        let a = snapshot(
1239            "s1",
1240            vec![
1241                region("conv", vec![entry("hi", 1)]),
1242                region("plan", vec![entry("p", 3)]),
1243            ],
1244        );
1245        let b = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1246        let delta = diff_context(&a, &b);
1247        assert!(
1248            delta
1249                .regions
1250                .iter()
1251                .any(|d| region_delta_kind(d) == "remove")
1252        );
1253        assert_diff_roundtrip(&a, &b);
1254    }
1255
1256    #[test]
1257    fn diff_non_prefix_rewrite_is_set() {
1258        // Entries changed at the front (not an append) → full Set.
1259        let a = snapshot("s1", vec![region("conv", vec![entry("old", 1)])]);
1260        let b = snapshot("s1", vec![region("conv", vec![entry("new", 1)])]);
1261        let delta = diff_context(&a, &b);
1262        assert_eq!(region_delta_kind(&delta.regions[0]), "set");
1263        assert_diff_roundtrip(&a, &b);
1264    }
1265
1266    #[test]
1267    fn diff_kind_change_is_set_not_append() {
1268        // Same prefix entries but the region's kind changed → Set, not Append.
1269        let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1270        let mut grown = region("conv", vec![entry("hi", 1), entry("more", 1)]);
1271        grown.kind = "sliding".to_string();
1272        let b = snapshot("s1", vec![grown]);
1273        let delta = diff_context(&a, &b);
1274        assert_eq!(region_delta_kind(&delta.regions[0]), "set");
1275        assert_diff_roundtrip(&a, &b);
1276    }
1277
1278    // ── streaming point replay ──
1279
1280    /// Frame `records` exactly as `run.lvr` stores them.
1281    fn framed(records: &[RunRecord]) -> Vec<u8> {
1282        let mut buf = Vec::new();
1283        write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1284        for record in records {
1285            write_record(&mut buf, record).unwrap();
1286        }
1287        buf
1288    }
1289
1290    /// Collect `(index, at, total_tokens)` per visited point, or the stream
1291    /// error. One closure shared by every streamed test, including the
1292    /// bad-preamble one whose visitor never runs.
1293    fn try_collect_streamed(bytes: &[u8]) -> io::Result<Vec<(usize, i64, usize)>> {
1294        let mut seen = Vec::new();
1295        visit_archive_points(&mut &bytes[..], &mut |p| {
1296            seen.push((p.index, p.at, p.context.total_tokens));
1297            ControlFlow::Continue(())
1298        })?;
1299        Ok(seen)
1300    }
1301
1302    /// Collect `(index, at, total_tokens)` per visited point.
1303    fn collect_streamed(bytes: &[u8]) -> Vec<(usize, i64, usize)> {
1304        try_collect_streamed(bytes).unwrap()
1305    }
1306
1307    #[test]
1308    fn visit_archive_points_matches_visit_points() {
1309        let records = vec![
1310            header(),
1311            RunRecord::ContextCheckpoint {
1312                snapshot: snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1313                at: 10,
1314            },
1315            RunRecord::StatusChanged {
1316                status: RunStatus::Running,
1317                at: 11,
1318            },
1319            RunRecord::Progress {
1320                meta: Box::new(meta()),
1321                delta: diff_context(
1322                    &snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1323                    &snapshot(
1324                        "s1",
1325                        vec![region("conv", vec![entry("hi", 1), entry("more", 2)])],
1326                    ),
1327                ),
1328                at: 12,
1329            },
1330        ];
1331        let mut in_memory = Vec::new();
1332        visit_points(&records, &mut |p| {
1333            in_memory.push((p.index, p.at, p.context.total_tokens));
1334            ControlFlow::Continue(())
1335        });
1336        assert_eq!(collect_streamed(&framed(&records)), in_memory);
1337        assert_eq!(in_memory.len(), 2, "checkpoint + progress = two points");
1338    }
1339
1340    #[test]
1341    fn visit_archive_points_rejects_a_bad_preamble() {
1342        assert!(try_collect_streamed(b"not an archive at all").is_err());
1343    }
1344
1345    #[test]
1346    fn visit_archive_points_is_lenient_about_a_torn_tail() {
1347        let records = vec![
1348            header(),
1349            RunRecord::ContextCheckpoint {
1350                snapshot: snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1351                at: 10,
1352            },
1353        ];
1354        let mut bytes = framed(&records);
1355        // A torn frame: a length prefix promising more than exists.
1356        bytes.extend_from_slice(&1000u64.to_be_bytes());
1357        bytes.extend_from_slice(b"partial");
1358        assert_eq!(collect_streamed(&bytes).len(), 1, "points before the tear");
1359    }
1360
1361    #[test]
1362    fn visit_archive_points_visits_nothing_without_a_header() {
1363        let records = vec![RunRecord::ContextCheckpoint {
1364            snapshot: snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]),
1365            at: 10,
1366        }];
1367        assert!(collect_streamed(&framed(&records)).is_empty());
1368        // And an archive with no records at all visits nothing.
1369        assert!(collect_streamed(&framed(&[])).is_empty());
1370    }
1371
1372    #[test]
1373    fn visit_archive_points_stops_on_break() {
1374        let records = vec![
1375            header(),
1376            RunRecord::ContextCheckpoint {
1377                snapshot: snapshot("s1", vec![region("conv", vec![entry("a", 1)])]),
1378                at: 10,
1379            },
1380            RunRecord::ContextCheckpoint {
1381                snapshot: snapshot("s1", vec![region("conv", vec![entry("b", 2)])]),
1382                at: 11,
1383            },
1384        ];
1385        let bytes = framed(&records);
1386        let mut seen = 0;
1387        visit_archive_points(&mut &bytes[..], &mut |_| {
1388            seen += 1;
1389            ControlFlow::Break(())
1390        })
1391        .unwrap();
1392        assert_eq!(seen, 1);
1393    }
1394
1395    // ── digest-based diffing ──
1396    //
1397    // `diff_context_digest(digest(a), b)` must produce the same delta as
1398    // `diff_context(a, b)` for every shape: the persistence lane retains only
1399    // the digest, and any divergence would silently corrupt the archive.
1400
1401    fn assert_digest_matches_full_diff(a: &ContextSnapshot, b: &ContextSnapshot) {
1402        let via_digest = diff_context_digest(&digest_context(a), b);
1403        assert_eq!(via_digest, diff_context(a, b));
1404        // And the digest-produced delta still round-trips.
1405        let mut base = a.clone();
1406        apply_delta(&mut base, &via_digest);
1407        assert_eq!(&base, b);
1408    }
1409
1410    #[test]
1411    fn digest_diff_append_only_growth_is_compact() {
1412        let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1413        let b = snapshot(
1414            "s1",
1415            vec![region("conv", vec![entry("hi", 1), entry("there", 2)])],
1416        );
1417        let delta = diff_context_digest(&digest_context(&a), &b);
1418        assert_eq!(region_delta_kind(&delta.regions[0]), "append");
1419        assert_digest_matches_full_diff(&a, &b);
1420    }
1421
1422    #[test]
1423    fn digest_diff_new_cleared_removed_and_rewritten_regions() {
1424        let a = snapshot(
1425            "s1",
1426            vec![
1427                region("conv", vec![entry("hi", 1)]),
1428                region("gone", vec![entry("bye", 1)]),
1429                region("wiped", vec![entry("w", 1)]),
1430                region("rewritten", vec![entry("old", 1)]),
1431            ],
1432        );
1433        let b = snapshot(
1434            "s1",
1435            vec![
1436                region("conv", vec![entry("hi", 1)]),
1437                region("wiped", vec![]),
1438                region("rewritten", vec![entry("new", 1)]),
1439                region("fresh", vec![entry("f", 2)]),
1440            ],
1441        );
1442        let delta = diff_context_digest(&digest_context(&a), &b);
1443        let kinds: Vec<_> = delta.regions.iter().map(region_delta_kind).collect();
1444        assert_eq!(kinds, vec!["clear", "set", "set", "remove"]);
1445        assert_digest_matches_full_diff(&a, &b);
1446    }
1447
1448    #[test]
1449    fn digest_diff_unchanged_region_emits_nothing() {
1450        let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1451        let delta = diff_context_digest(&digest_context(&a), &a.clone());
1452        assert!(delta.regions.is_empty());
1453        assert_digest_matches_full_diff(&a, &a.clone());
1454    }
1455
1456    #[test]
1457    fn digest_diff_kind_change_is_set_not_append() {
1458        let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1459        let mut grown = region("conv", vec![entry("hi", 1), entry("more", 1)]);
1460        grown.kind = "sliding".to_string();
1461        let b = snapshot("s1", vec![grown]);
1462        let delta = diff_context_digest(&digest_context(&a), &b);
1463        assert_eq!(region_delta_kind(&delta.regions[0]), "set");
1464        assert_digest_matches_full_diff(&a, &b);
1465    }
1466
1467    /// A token-count change with identical entries is still an (empty) Append
1468    /// carrying the new count, exactly as the full diff records it.
1469    #[test]
1470    fn digest_diff_token_recount_is_an_empty_append() {
1471        let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1472        let mut recounted = region("conv", vec![entry("hi", 1)]);
1473        recounted.current_tokens = 42;
1474        let b = snapshot("s1", vec![recounted]);
1475        let delta = diff_context_digest(&digest_context(&a), &b);
1476        assert_eq!(region_delta_kind(&delta.regions[0]), "append");
1477        assert_digest_matches_full_diff(&a, &b);
1478    }
1479
1480    /// Every field of an entry participates in its digest: a change anywhere
1481    /// must change the hash, or a real edit would fold as "unchanged".
1482    #[test]
1483    fn entry_digest_covers_every_field() {
1484        let base = entry("text", 1);
1485        let variants = [
1486            entry("other", 1),
1487            entry("text", 2),
1488            RegionEntrySnapshot {
1489                key: Some("k".to_string()),
1490                ..entry("text", 1)
1491            },
1492            RegionEntrySnapshot {
1493                metadata: Some(serde_json::json!({"a": 1})),
1494                ..entry("text", 1)
1495            },
1496            RegionEntrySnapshot {
1497                kind: crate::region::EntryKind::ToolResult {
1498                    tool_call_id: "c1".to_string(),
1499                    tool_name: "shell".to_string(),
1500                    is_error: false,
1501                },
1502                ..entry("text", 1)
1503            },
1504        ];
1505        let base_hash = entry_digest(&base);
1506        for variant in &variants {
1507            assert_ne!(
1508                entry_digest(variant),
1509                base_hash,
1510                "field change must change the digest: {variant:?}"
1511            );
1512        }
1513        // And digesting the same entry twice is stable.
1514        assert_eq!(entry_digest(&base), entry_digest(&entry("text", 1)));
1515    }
1516
1517    #[test]
1518    fn diff_unchanged_region_emits_nothing() {
1519        let a = snapshot("s1", vec![region("conv", vec![entry("hi", 1)])]);
1520        let b = a.clone();
1521        let delta = diff_context(&a, &b);
1522        assert!(delta.regions.is_empty());
1523        assert_diff_roundtrip(&a, &b);
1524    }
1525
1526    #[test]
1527    fn apply_delta_skips_unknown_regions_leniently() {
1528        // Append/Clear targeting a region not present are no-ops (not errors).
1529        let mut base = snapshot("s1", vec![]);
1530        let delta = ContextDelta {
1531            stage_name: "s1".to_string(),
1532            total_tokens: 0,
1533            max_tokens: 10_000,
1534            regions: vec![
1535                RegionDelta::Append {
1536                    name: "ghost".to_string(),
1537                    entries: vec![entry("x", 1)],
1538                    current_tokens: 1,
1539                },
1540                RegionDelta::Clear {
1541                    name: "ghost".to_string(),
1542                },
1543                RegionDelta::Remove {
1544                    name: "ghost".to_string(),
1545                },
1546            ],
1547        };
1548        apply_delta(&mut base, &delta);
1549        assert!(base.regions.is_empty());
1550    }
1551
1552    // ── codec round-trips ──
1553
1554    fn all_record_kinds() -> Vec<RunRecord> {
1555        vec![
1556            header(),
1557            RunRecord::OwnershipChanged {
1558                machine_id: "machine-b".to_string(),
1559                world_id: "world-y".to_string(),
1560                at: 101,
1561            },
1562            RunRecord::Inference {
1563                stage: "plan".to_string(),
1564                iteration: 0,
1565                request: InferenceRequestRecord {
1566                    model: "m".to_string(),
1567                    system: vec!["sys".to_string()],
1568                    messages: vec![MessageRecord {
1569                        role: "user".to_string(),
1570                        content: "hi".to_string(),
1571                    }],
1572                    tool_names: vec!["read_file".to_string()],
1573                    temperature: 0.7,
1574                    max_tokens: 1024,
1575                },
1576                response: InferenceResponseRecord {
1577                    content: "ok".to_string(),
1578                    tool_calls: vec![],
1579                    prompt_tokens: 10,
1580                    completion_tokens: 5,
1581                    cached_tokens: 0,
1582                    cache_write_tokens: 0,
1583                },
1584                at: 102,
1585            },
1586            RunRecord::InferenceUsage {
1587                kind: InferenceKind::Compaction,
1588                stage: "plan".to_string(),
1589                iteration: 2,
1590                provider: "anthropic".to_string(),
1591                model: "claude-sonnet-5".to_string(),
1592                prompt_tokens: 7000,
1593                completion_tokens: 70,
1594                cached_tokens: 12,
1595                cache_write_tokens: 34,
1596                cost_usd: None,
1597                cost_reported_by_provider: None,
1598                at: 102,
1599            },
1600            RunRecord::ToolBatch {
1601                calls: vec![ToolCallRecord {
1602                    id: "c1".to_string(),
1603                    name: "read_file".to_string(),
1604                    arguments: "{}".to_string(),
1605                    result: Some("body".to_string().into()),
1606                    thought_signature: Some("sig".to_string()),
1607                }],
1608                at: 103,
1609                stage_index: 0,
1610                iteration: 0,
1611                response: "reading".to_string(),
1612            },
1613            RunRecord::ToolCallDone {
1614                iteration: 0,
1615                call_id: "c1".to_string(),
1616                result: "body".to_string().into(),
1617                at: 103,
1618            },
1619            RunRecord::ContextCheckpoint {
1620                snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1621                at: 104,
1622            },
1623            RunRecord::ContextDiff {
1624                delta: ContextDelta {
1625                    stage_name: "plan".to_string(),
1626                    total_tokens: 3,
1627                    max_tokens: 10_000,
1628                    regions: vec![RegionDelta::Append {
1629                        name: "conv".to_string(),
1630                        entries: vec![entry("more", 2)],
1631                        current_tokens: 3,
1632                    }],
1633                },
1634                at: 105,
1635            },
1636            RunRecord::Message {
1637                message: MessageRecord {
1638                    role: "user".to_string(),
1639                    content: "another".to_string(),
1640                },
1641                at: 106,
1642            },
1643            RunRecord::StatusChanged {
1644                status: RunStatus::Complete,
1645                at: 107,
1646            },
1647            RunRecord::Checkpoint {
1648                meta: Box::new(meta()),
1649                context: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1650                at: 108,
1651            },
1652            RunRecord::Progress {
1653                meta: Box::new(meta()),
1654                delta: ContextDelta {
1655                    stage_name: "plan".to_string(),
1656                    total_tokens: 3,
1657                    max_tokens: 10_000,
1658                    regions: vec![RegionDelta::Append {
1659                        name: "conv".to_string(),
1660                        entries: vec![entry("step", 2)],
1661                        current_tokens: 3,
1662                    }],
1663                },
1664                at: 109,
1665            },
1666        ]
1667    }
1668
1669    #[test]
1670    fn archive_write_then_read_roundtrips_every_record_kind() {
1671        let records = all_record_kinds();
1672        let mut buf = Vec::new();
1673        write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1674        for r in &records {
1675            write_record(&mut buf, r).unwrap();
1676        }
1677        let (version, read) = read_archive(&mut buf.as_slice()).unwrap();
1678        assert_eq!(version, RUN_ARCHIVE_VERSION);
1679        assert_eq!(read, records);
1680    }
1681
1682    #[test]
1683    fn read_archive_start_rejects_bad_magic() {
1684        let mut bytes: &[u8] = b"XXXX\x00\x01";
1685        let err = read_archive_start(&mut bytes).unwrap_err();
1686        assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1687    }
1688
1689    /// The preamble round-trips the version it was written with. Previously
1690    /// checked with an arbitrary 7; that now names a framing generation this
1691    /// build cannot read, and is refused - see
1692    /// `an_archive_from_a_newer_format_is_refused_with_both_versions_named`.
1693    #[test]
1694    fn read_archive_start_reports_version() {
1695        let mut buf = Vec::new();
1696        write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1697        assert_eq!(
1698            read_archive_start(&mut buf.as_slice()).unwrap(),
1699            RUN_ARCHIVE_VERSION
1700        );
1701    }
1702
1703    #[test]
1704    fn read_record_returns_none_at_clean_eof() {
1705        let empty: &[u8] = &[];
1706        assert!(read_record(&mut { empty }).unwrap().is_none());
1707    }
1708
1709    #[test]
1710    fn read_record_errors_on_truncated_length_prefix() {
1711        // Two bytes where an 8-byte length is expected → partial read → error.
1712        let mut bytes: &[u8] = &[0, 0];
1713        let err = read_record(&mut bytes).unwrap_err();
1714        assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1715    }
1716
1717    #[test]
1718    fn read_record_errors_on_truncated_payload() {
1719        // A frame claiming 10 bytes but only 2 present after the 8-byte length.
1720        let mut bytes: &[u8] = &[0, 0, 0, 0, 0, 0, 0, 10, 1, 2];
1721        let err = read_record(&mut bytes).unwrap_err();
1722        assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1723    }
1724
1725    #[test]
1726    fn read_record_errors_on_empty_payload_at_boundary() {
1727        // A non-zero length with zero payload bytes → clean EOF at the payload
1728        // start is still a truncation (the frame promised bytes).
1729        let mut bytes: &[u8] = &[0, 0, 0, 0, 0, 0, 0, 10];
1730        let err = read_record(&mut bytes).unwrap_err();
1731        assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1732    }
1733
1734    #[test]
1735    fn read_record_errors_on_invalid_json_payload() {
1736        // A well-framed payload that isn't a valid RunRecord.
1737        let mut buf = Vec::new();
1738        let bad = b"not json";
1739        buf.extend_from_slice(&(bad.len() as u64).to_be_bytes());
1740        buf.extend_from_slice(bad);
1741        let err = read_record(&mut buf.as_slice()).unwrap_err();
1742        assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1743    }
1744
1745    /// A reader whose `read` always errors, to exercise the read error path
1746    /// inside `read_exact_or_eof` (distinct from a clean EOF).
1747    struct FailingReader;
1748    impl Read for FailingReader {
1749        fn read(&mut self, _buf: &mut [u8]) -> io::Result<usize> {
1750            Err(io::Error::other("device error"))
1751        }
1752    }
1753
1754    #[test]
1755    fn read_record_propagates_reader_errors() {
1756        let err = read_record(&mut FailingReader).unwrap_err();
1757        assert_eq!(err.kind(), io::ErrorKind::Other);
1758    }
1759
1760    #[test]
1761    fn read_archive_propagates_a_bad_preamble() {
1762        // Too short to even hold the magic → the preamble read errors.
1763        let mut bytes: &[u8] = b"LV";
1764        assert!(read_archive(&mut bytes).is_err());
1765    }
1766
1767    #[test]
1768    fn read_archive_propagates_a_bad_frame() {
1769        // Valid preamble, then a truncated frame → the record read errors.
1770        let mut buf = Vec::new();
1771        write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1772        buf.extend_from_slice(&[0, 0, 0, 0, 0, 0, 0, 5, 1, 2]); // len 5, 2 present
1773        let err = read_archive(&mut buf.as_slice()).unwrap_err();
1774        assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1775    }
1776
1777    /// A writer that fails after `ok_bytes` bytes, to exercise write error paths.
1778    struct FailAfter {
1779        remaining: usize,
1780    }
1781    impl Write for FailAfter {
1782        fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
1783            if self.remaining == 0 {
1784                return Err(io::Error::other("disk full"));
1785            }
1786            let n = buf.len().min(self.remaining);
1787            self.remaining -= n;
1788            Ok(n)
1789        }
1790        fn flush(&mut self) -> io::Result<()> {
1791            Ok(())
1792        }
1793    }
1794
1795    #[test]
1796    fn fail_after_writer_flush_is_a_noop() {
1797        assert!(FailAfter { remaining: 1 }.flush().is_ok());
1798    }
1799
1800    #[test]
1801    fn write_archive_start_propagates_write_errors() {
1802        // Fail on the magic write (0 bytes allowed) and on the version write.
1803        assert!(write_archive_start(&mut FailAfter { remaining: 0 }, 1).is_err());
1804        assert!(write_archive_start(&mut FailAfter { remaining: 4 }, 1).is_err());
1805    }
1806
1807    #[test]
1808    fn write_record_propagates_write_errors() {
1809        let rec = header();
1810        // Fail on the 8-byte length prefix, and (after it) on the payload.
1811        assert!(write_record(&mut FailAfter { remaining: 0 }, &rec).is_err());
1812        assert!(write_record(&mut FailAfter { remaining: 8 }, &rec).is_err());
1813    }
1814
1815    /// A torn *length prefix* is where a nonsense `u64` comes from, and the
1816    /// lenient reader exists precisely to survive a torn tail. Taking the
1817    /// length at its word would turn a crash-truncated archive into an
1818    /// allocation of that size - during daemon recovery, the one moment this
1819    /// reader is there to keep working.
1820    #[test]
1821    fn an_absurd_frame_length_is_an_error_not_an_allocation() {
1822        let mut buf = Vec::new();
1823        write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1824        write_record(&mut buf, &header()).unwrap();
1825        // A crash mid-append that left a garbage length behind.
1826        buf.extend_from_slice(&u64::MAX.to_be_bytes());
1827
1828        let err = read_archive(&mut buf.as_slice())
1829            .expect_err("the strict reader must refuse an impossible frame");
1830        assert_eq!(err.kind(), io::ErrorKind::InvalidData, "{err}");
1831
1832        // And the lenient reader folds back to the intact record before it,
1833        // which is the behaviour recovery depends on.
1834        let (_, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
1835        assert_eq!(records, vec![header()]);
1836    }
1837
1838    #[test]
1839    fn read_archive_lenient_matches_strict_on_a_clean_archive() {
1840        // With no torn tail, the lenient reader returns exactly what the strict
1841        // reader does.
1842        let records = all_record_kinds();
1843        let mut buf = Vec::new();
1844        write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1845        for r in &records {
1846            write_record(&mut buf, r).unwrap();
1847        }
1848        let (version, read) = read_archive_lenient(&mut buf.as_slice()).unwrap();
1849        assert_eq!(version, RUN_ARCHIVE_VERSION);
1850        assert_eq!(read, records);
1851    }
1852
1853    #[test]
1854    fn read_archive_lenient_keeps_valid_prefix_before_a_torn_tail() {
1855        // A valid preamble + two full records, then a truncated frame (a crash
1856        // mid-append). The strict reader would reject the whole file; the lenient
1857        // reader returns the two intact records and stops at the torn tail.
1858        let mut buf = Vec::new();
1859        write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
1860        write_record(&mut buf, &header()).unwrap();
1861        write_record(
1862            &mut buf,
1863            &RunRecord::ContextCheckpoint {
1864                snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1865                at: 1,
1866            },
1867        )
1868        .unwrap();
1869        // A frame claiming 10 payload bytes but only 2 present → torn tail.
1870        buf.extend_from_slice(&[0, 0, 0, 0, 0, 0, 0, 10, 1, 2]);
1871
1872        // Strict rejects the whole archive.
1873        assert!(read_archive(&mut buf.as_slice()).is_err());
1874        // Lenient keeps the valid prefix and folds cleanly.
1875        let (version, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
1876        assert_eq!(version, RUN_ARCHIVE_VERSION);
1877        assert_eq!(records.len(), 2);
1878        let folded = fold(&records).expect("prefix starts with a Header");
1879        assert_eq!(folded.context.regions[0].entries.len(), 1);
1880    }
1881
1882    #[test]
1883    fn read_archive_lenient_still_errors_on_a_bad_preamble() {
1884        // The preamble is validated strictly: a file that isn't a run archive at
1885        // all errors rather than folding to nothing.
1886        let mut bad_magic: &[u8] = b"XXXX\x00\x01";
1887        assert!(read_archive_lenient(&mut bad_magic).is_err());
1888        // A truncated version (valid magic, no version bytes) also errors.
1889        let mut short: &[u8] = b"LVR1";
1890        assert!(read_archive_lenient(&mut short).is_err());
1891    }
1892
1893    #[test]
1894    fn read_archive_start_errors_on_truncated_version() {
1895        // Valid 4-byte magic but no version bytes → the version read errors.
1896        let mut bytes: &[u8] = b"LVR1";
1897        let err = read_archive_start(&mut bytes).unwrap_err();
1898        assert_eq!(err.kind(), io::ErrorKind::UnexpectedEof);
1899    }
1900
1901    // ── fold ──
1902
1903    #[test]
1904    fn fold_requires_a_header_first() {
1905        assert!(fold(&[]).is_none());
1906        assert!(
1907            fold(&[RunRecord::StatusChanged {
1908                status: RunStatus::Complete,
1909                at: 1
1910            }])
1911            .is_none()
1912        );
1913    }
1914
1915    #[test]
1916    fn fold_reconstructs_state_from_the_journal() {
1917        let records = all_record_kinds();
1918        let folded = fold(&records).expect("has header");
1919        // Ownership was reassigned mid-journal.
1920        assert_eq!(folded.identity.machine_id, "machine-b");
1921        assert_eq!(folded.identity.world_id, "world-y");
1922        // Counters. Two inferences: the fixture carries one record of each
1923        // kind, and both name one provider call.
1924        assert_eq!(folded.inference_count, 2);
1925        assert_eq!(folded.inference_usage.len(), 1);
1926        assert_eq!(folded.tool_call_count, 1);
1927        // One inbound message recorded.
1928        assert_eq!(folded.messages.len(), 1);
1929        assert_eq!(folded.messages[0].content, "another");
1930        // The Progress step is the last context-affecting record: it layers its
1931        // append diff onto the preceding Checkpoint's window (hi + step).
1932        assert_eq!(folded.context.regions[0].name, "conv");
1933        assert_eq!(folded.context.regions[0].entries.len(), 2);
1934        assert_eq!(folded.context.total_tokens, 3);
1935        assert_eq!(folded.meta.run_id, "run-1");
1936        // The batch shares the meta's iteration and its turn never reached the
1937        // window, so it folds as pending (with the ToolCallDone merged in).
1938        let pending = folded.pending_batch.expect("batch never applied");
1939        assert_eq!(pending.calls[0].result.as_deref(), Some("body"));
1940    }
1941
1942    #[test]
1943    fn fold_applies_context_diffs_over_a_checkpoint() {
1944        // Header → checkpoint → diff (append). The diff must layer on the checkpoint.
1945        let records = vec![
1946            header(),
1947            RunRecord::ContextCheckpoint {
1948                snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
1949                at: 1,
1950            },
1951            RunRecord::ContextDiff {
1952                delta: ContextDelta {
1953                    stage_name: "plan".to_string(),
1954                    total_tokens: 3,
1955                    max_tokens: 10_000,
1956                    regions: vec![RegionDelta::Append {
1957                        name: "conv".to_string(),
1958                        entries: vec![entry("there", 2)],
1959                        current_tokens: 3,
1960                    }],
1961                },
1962                at: 2,
1963            },
1964        ];
1965        let folded = fold(&records).unwrap();
1966        assert_eq!(folded.context.regions[0].entries.len(), 2);
1967        assert_eq!(folded.context.total_tokens, 3);
1968    }
1969
1970    #[test]
1971    fn fold_later_header_updates_identity_and_meta() {
1972        // A second Header (unusual, but tolerated) refreshes identity + meta.
1973        let mut second_meta = meta();
1974        second_meta.status = RunStatus::Running;
1975        let records = vec![
1976            header(),
1977            RunRecord::Header {
1978                identity: RunIdentity {
1979                    run_id: "run-1".to_string(),
1980                    machine_id: "machine-c".to_string(),
1981                    world_id: "world-z".to_string(),
1982                    created_at: 200,
1983                },
1984                meta: Box::new(second_meta),
1985            },
1986        ];
1987        let folded = fold(&records).unwrap();
1988        assert_eq!(folded.identity.machine_id, "machine-c");
1989        assert_eq!(folded.meta.status, RunStatus::Running);
1990    }
1991
1992    #[test]
1993    fn fold_progress_applies_meta_and_context_diff() {
1994        let mut advanced = meta();
1995        advanced.status = RunStatus::Running;
1996        advanced.iteration = 5;
1997        let records = vec![
1998            header(),
1999            RunRecord::ContextCheckpoint {
2000                snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
2001                at: 1,
2002            },
2003            RunRecord::Progress {
2004                meta: Box::new(advanced),
2005                delta: ContextDelta {
2006                    stage_name: "plan".to_string(),
2007                    total_tokens: 3,
2008                    max_tokens: 10_000,
2009                    regions: vec![RegionDelta::Append {
2010                        name: "conv".to_string(),
2011                        entries: vec![entry("there", 2)],
2012                        current_tokens: 3,
2013                    }],
2014                },
2015                at: 2,
2016            },
2017        ];
2018        let folded = fold(&records).unwrap();
2019        assert_eq!(folded.meta.iteration, 5);
2020        assert_eq!(folded.meta.status, RunStatus::Running);
2021        assert_eq!(folded.context.regions[0].entries.len(), 2);
2022    }
2023
2024    /// A submitted answer needs no record type of its own: `Progress` and
2025    /// `Checkpoint` both replace the whole `RunMeta`, so it folds along with
2026    /// everything else and a crash-resume finds the answer already there.
2027    #[test]
2028    fn fold_carries_a_submitted_final_output_through_progress() {
2029        let mut answered = meta();
2030        answered.final_output = Some(
2031            crate::output::FinalOutput::new(
2032                "renamed two helpers",
2033                Some("markdown".to_string()),
2034                "summary".to_string(),
2035                9,
2036            )
2037            .descriptor(),
2038        );
2039        answered.output_request = Some(crate::output::OutputSpec {
2040            format: Some("a2ui".to_string()),
2041            ..Default::default()
2042        });
2043        let records = vec![
2044            header(),
2045            RunRecord::Progress {
2046                meta: Box::new(answered),
2047                delta: ContextDelta {
2048                    stage_name: "summary".to_string(),
2049                    total_tokens: 0,
2050                    max_tokens: 10_000,
2051                    regions: vec![],
2052                },
2053                at: 2,
2054            },
2055        ];
2056        let folded = fold(&records).unwrap();
2057        let output = folded.meta.final_output.expect("the answer folded through");
2058        // The descriptor, not the bytes: the answer itself is a sidecar file,
2059        // so what folds is the record of it.
2060        assert_eq!(output.bytes, "renamed two helpers".len());
2061        assert_eq!(output.stage, "summary");
2062        assert_eq!(
2063            folded.meta.output_request.and_then(|s| s.format).as_deref(),
2064            Some("a2ui")
2065        );
2066    }
2067
2068    // ── pending tool batch (fold) ──
2069
2070    fn call(id: &str, result: Option<&str>) -> ToolCallRecord {
2071        ToolCallRecord {
2072            id: id.to_string(),
2073            name: "shell".to_string(),
2074            arguments: "{}".to_string(),
2075            result: result.map(Into::into),
2076            thought_signature: None,
2077        }
2078    }
2079
2080    fn batch(iteration: usize, calls: Vec<ToolCallRecord>) -> RunRecord {
2081        RunRecord::ToolBatch {
2082            calls,
2083            at: 10,
2084            stage_index: 0,
2085            iteration,
2086            response: "running tools".to_string(),
2087        }
2088    }
2089
2090    /// An entry whose kind is the assistant turn that issued `call_ids`.
2091    fn turn_entry(call_ids: &[&str]) -> RegionEntrySnapshot {
2092        let mut e = entry("turn", 1);
2093        e.kind = crate::region::EntryKind::AssistantTurn {
2094            tool_calls: call_ids
2095                .iter()
2096                .map(|id| crate::region::SerializedToolCall {
2097                    id: id.to_string(),
2098                    name: "shell".to_string(),
2099                    arguments: serde_json::Value::Null,
2100                    thought_signature: None,
2101                })
2102                .collect(),
2103        };
2104        e
2105    }
2106
2107    #[test]
2108    fn fold_surfaces_a_pending_batch_with_merged_results() {
2109        // meta().iteration is 0, matching the batch, and the context has no
2110        // assistant turn for it - so the batch is genuinely pending. c1's
2111        // ToolCallDone merges in; c2 keeps its dispatch-time inline result; c3
2112        // stays pending.
2113        let records = vec![
2114            header(),
2115            batch(
2116                0,
2117                vec![
2118                    call("c1", None),
2119                    call("c2", Some("inline")),
2120                    call("c3", None),
2121                ],
2122            ),
2123            RunRecord::ToolCallDone {
2124                iteration: 0,
2125                call_id: "c1".to_string(),
2126                result: "ran".to_string().into(),
2127                at: 11,
2128            },
2129        ];
2130        let folded = fold(&records).unwrap();
2131        let pending = folded.pending_batch.expect("batch is pending");
2132        assert_eq!(pending.iteration, 0);
2133        assert_eq!(pending.response, "running tools");
2134        assert_eq!(pending.calls[0].result.as_deref(), Some("ran"));
2135        assert_eq!(pending.calls[1].result.as_deref(), Some("inline"));
2136        assert_eq!(pending.calls[2].result, None);
2137        assert_eq!(folded.tool_call_count, 3);
2138    }
2139
2140    #[test]
2141    fn fold_keeps_only_the_latest_batch_and_ignores_stale_done_records() {
2142        // The second batch replaces the first; a ToolCallDone for the replaced
2143        // iteration is ignored, as is one naming a call the batch doesn't have.
2144        let mut advanced = meta();
2145        advanced.iteration = 1;
2146        let records = vec![
2147            header(),
2148            batch(0, vec![call("c1", None)]),
2149            RunRecord::Progress {
2150                meta: Box::new(advanced),
2151                delta: ContextDelta {
2152                    stage_name: "plan".to_string(),
2153                    total_tokens: 0,
2154                    max_tokens: 10_000,
2155                    regions: vec![],
2156                },
2157                at: 11,
2158            },
2159            batch(1, vec![call("c2", None)]),
2160            RunRecord::ToolCallDone {
2161                iteration: 0,
2162                call_id: "c1".to_string(),
2163                result: "stale".to_string().into(),
2164                at: 12,
2165            },
2166            RunRecord::ToolCallDone {
2167                iteration: 1,
2168                call_id: "unknown".to_string(),
2169                result: "nowhere to land".to_string().into(),
2170                at: 13,
2171            },
2172        ];
2173        let folded = fold(&records).unwrap();
2174        let pending = folded.pending_batch.expect("latest batch is pending");
2175        assert_eq!(pending.iteration, 1);
2176        assert_eq!(pending.calls.len(), 1);
2177        assert_eq!(pending.calls[0].id, "c2");
2178        assert_eq!(pending.calls[0].result, None, "stale/unknown dones ignored");
2179    }
2180
2181    #[test]
2182    fn fold_clears_a_batch_once_the_iteration_moves_on() {
2183        // A later inference bumped meta.iteration past the batch: the batch was
2184        // applied (even if a sliding window evicted the turn), nothing to replay.
2185        let mut advanced = meta();
2186        advanced.iteration = 1;
2187        let records = vec![
2188            header(),
2189            batch(0, vec![call("c1", Some("done"))]),
2190            RunRecord::Progress {
2191                meta: Box::new(advanced),
2192                delta: ContextDelta {
2193                    stage_name: "plan".to_string(),
2194                    total_tokens: 0,
2195                    max_tokens: 10_000,
2196                    regions: vec![],
2197                },
2198                at: 11,
2199            },
2200        ];
2201        assert_eq!(fold(&records).unwrap().pending_batch, None);
2202    }
2203
2204    #[test]
2205    fn fold_clears_a_batch_whose_turn_already_landed_in_the_window() {
2206        // Same iteration, but the context already holds the batch's assistant
2207        // turn: apply_tool_results ran before the crash, nothing to replay.
2208        let records = vec![
2209            header(),
2210            batch(0, vec![call("c1", Some("done"))]),
2211            RunRecord::ContextCheckpoint {
2212                snapshot: snapshot("plan", vec![region("conv", vec![turn_entry(&["c1"])])]),
2213                at: 11,
2214            },
2215        ];
2216        assert_eq!(fold(&records).unwrap().pending_batch, None);
2217    }
2218
2219    #[test]
2220    fn context_contains_batch_matches_only_the_batch_turn() {
2221        let pending = PendingToolBatch {
2222            stage_index: 0,
2223            iteration: 0,
2224            response: String::new(),
2225            calls: vec![call("c1", None)],
2226        };
2227        // A window with an unrelated turn does not match.
2228        let other = snapshot("plan", vec![region("conv", vec![turn_entry(&["zz"])])]);
2229        assert!(!context_contains_batch(&other, &pending));
2230        // The batch's own turn matches by its first call id.
2231        let own = snapshot(
2232            "plan",
2233            vec![region("conv", vec![turn_entry(&["c1", "c2"])])],
2234        );
2235        assert!(context_contains_batch(&own, &pending));
2236        // A batch with no calls can never match.
2237        let empty = PendingToolBatch {
2238            calls: vec![],
2239            ..pending
2240        };
2241        assert!(!context_contains_batch(&own, &empty));
2242    }
2243
2244    #[test]
2245    fn old_shape_tool_batch_json_still_parses() {
2246        // Archives written before the batch-journal fields existed carry
2247        // ToolBatch records without stage_index/iteration/response (and calls
2248        // without thought_signature); serde defaults fill them in.
2249        let json = br#"{"ToolBatch":{"calls":[{"id":"c1","name":"shell","arguments":"{}","result":"ok"}],"at":9}}"#;
2250        let mut buf = Vec::new();
2251        buf.extend_from_slice(&(json.len() as u64).to_be_bytes());
2252        buf.extend_from_slice(json);
2253        let record = read_record(&mut buf.as_slice()).unwrap().unwrap();
2254        assert_eq!(
2255            record,
2256            RunRecord::ToolBatch {
2257                calls: vec![call("c1", Some("ok"))],
2258                at: 9,
2259                stage_index: 0,
2260                iteration: 0,
2261                response: String::new(),
2262            }
2263        );
2264    }
2265
2266    // ── replay_points (context-window history) ──
2267
2268    /// Three context changes, so a windowing caller has something to page over.
2269    fn three_point_records() -> Vec<RunRecord> {
2270        let mut running = meta();
2271        running.status = RunStatus::Running;
2272        vec![
2273            header(),
2274            RunRecord::ContextCheckpoint {
2275                snapshot: snapshot("plan", vec![region("conv", vec![entry("first", 1)])]),
2276                at: 10,
2277            },
2278            RunRecord::ContextDiff {
2279                delta: ContextDelta {
2280                    stage_name: "plan".to_string(),
2281                    total_tokens: 2,
2282                    max_tokens: 10_000,
2283                    regions: vec![RegionDelta::Append {
2284                        name: "conv".to_string(),
2285                        entries: vec![entry("second", 1)],
2286                        current_tokens: 2,
2287                    }],
2288                },
2289                at: 20,
2290            },
2291            RunRecord::Progress {
2292                meta: Box::new(running),
2293                delta: ContextDelta {
2294                    stage_name: "code".to_string(),
2295                    total_tokens: 3,
2296                    max_tokens: 10_000,
2297                    regions: vec![RegionDelta::Append {
2298                        name: "conv".to_string(),
2299                        entries: vec![entry("third", 1)],
2300                        current_tokens: 3,
2301                    }],
2302                },
2303                at: 30,
2304            },
2305        ]
2306    }
2307
2308    #[test]
2309    fn visit_points_indexes_points_in_order_and_carries_the_running_window() {
2310        let records = three_point_records();
2311        let mut seen: Vec<(usize, i64, usize)> = Vec::new();
2312        visit_points(&records, &mut |point| {
2313            seen.push((
2314                point.index,
2315                point.at,
2316                point.context.regions[0].entries.len(),
2317            ));
2318            ControlFlow::Continue(())
2319        });
2320        // Index counts points, not records - the Header produces none.
2321        assert_eq!(seen, vec![(0, 10, 1), (1, 20, 2), (2, 30, 3)]);
2322    }
2323
2324    /// The reason this function exists: a caller wanting one window, or an
2325    /// answer to "does any point match", must be able to stop.
2326    #[test]
2327    fn visit_points_stops_at_the_first_break() {
2328        let records = three_point_records();
2329        let mut visits = 0;
2330        visit_points(&records, &mut |point| {
2331            visits += 1;
2332            if point.index == 1 {
2333                ControlFlow::Break(())
2334            } else {
2335                ControlFlow::Continue(())
2336            }
2337        });
2338        assert_eq!(
2339            visits, 2,
2340            "stopped at the breaking point, did not run the third"
2341        );
2342    }
2343
2344    #[test]
2345    fn visit_points_without_a_header_visits_nothing() {
2346        let mut visits = 0;
2347        {
2348            let mut count = |_: PointRef<'_>| {
2349                visits += 1;
2350                ControlFlow::Continue(())
2351            };
2352
2353            // A well-formed journal first, with the *same* visitor. Without
2354            // this the test would pass against a visitor that can never run at
2355            // all, which is exactly the reassurance it is not meant to give.
2356            visit_points(&three_point_records(), &mut count);
2357            // Neither of these starts with a Header, so neither is a replayable
2358            // journal and neither may produce a point.
2359            visit_points(&[], &mut count);
2360            visit_points(
2361                &[RunRecord::ContextCheckpoint {
2362                    snapshot: snapshot("plan", vec![]),
2363                    at: 1,
2364                }],
2365                &mut count,
2366            );
2367        }
2368        assert_eq!(visits, 3, "only the well-formed journal produced points");
2369    }
2370
2371    /// `replay_points` is now a thin collector over `visit_points`, so this
2372    /// pins the two together: if the reimplementation ever drifts, the borrowed
2373    /// walk and the materialized one stop agreeing here first.
2374    #[test]
2375    fn visit_points_and_replay_points_agree() {
2376        for records in [
2377            three_point_records(),
2378            vec![header()],
2379            vec![],
2380            vec![RunRecord::Message {
2381                message: MessageRecord {
2382                    role: "user".to_string(),
2383                    content: "x".to_string(),
2384                },
2385                at: 1,
2386            }],
2387        ] {
2388            let collected: Vec<RunPoint> = {
2389                let mut out = Vec::new();
2390                visit_points(&records, &mut |point| {
2391                    out.push(RunPoint {
2392                        meta: point.meta.clone(),
2393                        context: point.context.clone(),
2394                        at: point.at,
2395                    });
2396                    ControlFlow::Continue(())
2397                });
2398                out
2399            };
2400            assert_eq!(collected, replay_points(&records));
2401        }
2402    }
2403
2404    #[test]
2405    fn replay_points_requires_a_header() {
2406        assert!(replay_points(&[]).is_empty());
2407        assert!(
2408            replay_points(&[RunRecord::Message {
2409                message: MessageRecord {
2410                    role: "user".to_string(),
2411                    content: "x".to_string(),
2412                },
2413                at: 1,
2414            }])
2415            .is_empty()
2416        );
2417    }
2418
2419    #[test]
2420    fn replay_points_emits_a_snapshot_per_context_change() {
2421        // Header (no point) → checkpoint (point 1) → status (no point, but tracked)
2422        // → progress diff (point 2). Non-context records don't add points.
2423        let mut running = meta();
2424        running.status = RunStatus::Running;
2425        let records = vec![
2426            header(),
2427            RunRecord::Inference {
2428                stage: "plan".to_string(),
2429                iteration: 0,
2430                request: InferenceRequestRecord {
2431                    model: "m".to_string(),
2432                    system: vec![],
2433                    messages: vec![],
2434                    tool_names: vec![],
2435                    temperature: 0.7,
2436                    max_tokens: 10,
2437                },
2438                response: InferenceResponseRecord {
2439                    content: "ok".to_string(),
2440                    tool_calls: vec![],
2441                    prompt_tokens: 1,
2442                    completion_tokens: 1,
2443                    cached_tokens: 0,
2444                    cache_write_tokens: 0,
2445                },
2446                at: 1,
2447            },
2448            batch(0, vec![call("c1", None)]),
2449            RunRecord::ToolCallDone {
2450                iteration: 0,
2451                call_id: "c1".to_string(),
2452                result: "ran".to_string().into(),
2453                at: 1,
2454            },
2455            RunRecord::ContextCheckpoint {
2456                snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
2457                at: 2,
2458            },
2459            RunRecord::StatusChanged {
2460                status: RunStatus::Running,
2461                at: 3,
2462            },
2463            RunRecord::Progress {
2464                meta: Box::new(running),
2465                delta: ContextDelta {
2466                    stage_name: "implement".to_string(),
2467                    total_tokens: 3,
2468                    max_tokens: 10_000,
2469                    regions: vec![RegionDelta::Append {
2470                        name: "conv".to_string(),
2471                        entries: vec![entry("more", 2)],
2472                        current_tokens: 3,
2473                    }],
2474                },
2475                at: 4,
2476            },
2477        ];
2478        let points = replay_points(&records);
2479        assert_eq!(points.len(), 2, "one point per context change");
2480        // First point: the checkpoint window.
2481        assert_eq!(points[0].at, 2);
2482        assert_eq!(points[0].context.regions[0].entries.len(), 1);
2483        // Second point: the progress diff layered on, with the running status
2484        // carried from the StatusChanged + the progress meta.
2485        assert_eq!(points[1].at, 4);
2486        assert_eq!(points[1].context.regions[0].entries.len(), 2);
2487        assert_eq!(points[1].context.stage_name, "implement");
2488        assert_eq!(points[1].meta.status, RunStatus::Running);
2489    }
2490
2491    #[test]
2492    fn replay_points_handles_context_diff_and_a_later_header() {
2493        // A standalone ContextDiff is a point; a second Header refreshes meta
2494        // without adding a point.
2495        let mut relabeled = meta();
2496        relabeled.agent_name = "renamed".to_string();
2497        let records = vec![
2498            header(),
2499            RunRecord::ContextCheckpoint {
2500                snapshot: snapshot("plan", vec![region("conv", vec![entry("hi", 1)])]),
2501                at: 1,
2502            },
2503            RunRecord::Header {
2504                identity: identity(),
2505                meta: Box::new(relabeled),
2506            },
2507            RunRecord::ContextDiff {
2508                delta: ContextDelta {
2509                    stage_name: "plan".to_string(),
2510                    total_tokens: 3,
2511                    max_tokens: 10_000,
2512                    regions: vec![RegionDelta::Append {
2513                        name: "conv".to_string(),
2514                        entries: vec![entry("more", 2)],
2515                        current_tokens: 3,
2516                    }],
2517                },
2518                at: 2,
2519            },
2520        ];
2521        let points = replay_points(&records);
2522        assert_eq!(points.len(), 2); // checkpoint + diff (header adds no point)
2523        assert_eq!(points[1].context.regions[0].entries.len(), 2);
2524        // The later Header's meta is in effect at the diff point.
2525        assert_eq!(points[1].meta.agent_name, "renamed");
2526    }
2527
2528    #[test]
2529    fn replay_points_over_a_full_checkpoint() {
2530        // A `Checkpoint` (full meta+context) is also a point.
2531        let records = vec![
2532            header(),
2533            RunRecord::Checkpoint {
2534                meta: Box::new(meta()),
2535                context: snapshot("review", vec![region("conv", vec![entry("x", 4)])]),
2536                at: 9,
2537            },
2538        ];
2539        let points = replay_points(&records);
2540        assert_eq!(points.len(), 1);
2541        assert_eq!(points[0].context.stage_name, "review");
2542        assert_eq!(points[0].context.regions[0].entries[0].tokens, 4);
2543    }
2544
2545    /// Each kind has to survive the wire under its own name: the label is what
2546    /// a consumer groups a token chart by, so a rename that silently reordered
2547    /// the enum would re-attribute somebody's spend.
2548    #[test]
2549    fn every_inference_kind_has_a_distinct_label_and_serialized_name() {
2550        let all = [
2551            (InferenceKind::Stage, "stage"),
2552            (InferenceKind::Compaction, "compaction"),
2553            (InferenceKind::Title, "title"),
2554            (InferenceKind::Routing, "routing"),
2555        ];
2556        for (kind, label) in all {
2557            assert_eq!(kind.label(), label);
2558            assert_eq!(serde_json::to_value(kind).unwrap(), label);
2559        }
2560        let labels: std::collections::HashSet<_> = all.iter().map(|(k, _)| k.label()).collect();
2561        assert_eq!(labels.len(), all.len(), "labels must not collide");
2562    }
2563
2564    /// Only stage turns are work the agent asked for. The other three are
2565    /// machinery the runtime ran on its behalf, which is the split anything
2566    /// reporting "what did my agent actually do" needs.
2567    #[test]
2568    fn only_a_stage_turn_counts_as_stage_work() {
2569        assert!(InferenceKind::Stage.is_stage_work());
2570        for kind in [
2571            InferenceKind::Compaction,
2572            InferenceKind::Title,
2573            InferenceKind::Routing,
2574        ] {
2575            assert!(
2576                !kind.is_stage_work(),
2577                "{kind:?} is machinery, not stage work"
2578            );
2579        }
2580    }
2581
2582    /// A journal written before the field existed has to read back as stage
2583    /// work rather than failing to parse - every record in one is a stage turn,
2584    /// because nothing else was ever written.
2585    #[test]
2586    fn a_usage_record_without_a_kind_reads_back_as_stage_work() {
2587        let json = serde_json::json!({
2588            "InferenceUsage": {
2589                "stage": "plan",
2590                "iteration": 1,
2591                "provider": "anthropic",
2592                "model": "claude-sonnet-5",
2593                "prompt_tokens": 10,
2594                "completion_tokens": 2,
2595                "cached_tokens": 0,
2596                "cache_write_tokens": 0,
2597                "at": 5,
2598            }
2599        });
2600        // Compared whole rather than destructured: the point is that the
2601        // missing field defaults and every present one still lands, and a
2602        // destructure that pulled out `kind` alone would pass even if the rest
2603        // had been dropped.
2604        let record: RunRecord = serde_json::from_value(json).unwrap();
2605        assert_eq!(
2606            record,
2607            RunRecord::InferenceUsage {
2608                kind: InferenceKind::Stage,
2609                stage: "plan".to_string(),
2610                iteration: 1,
2611                provider: "anthropic".to_string(),
2612                model: "claude-sonnet-5".to_string(),
2613                prompt_tokens: 10,
2614                completion_tokens: 2,
2615                cached_tokens: 0,
2616                cache_write_tokens: 0,
2617                cost_usd: None,
2618                cost_reported_by_provider: None,
2619                at: 5,
2620            }
2621        );
2622    }
2623
2624    /// The point of the record. Folding keeps every call in order, so a
2625    /// consumer can ask what any single one cost - which the cumulative
2626    /// counters cannot answer, because two calls between two ticks arrive as
2627    /// their sum.
2628    #[test]
2629    fn folding_keeps_each_call_separate_instead_of_summing_them() {
2630        let usage = |kind, prompt, at| RunRecord::InferenceUsage {
2631            kind,
2632            stage: "plan".to_string(),
2633            iteration: 1,
2634            provider: "anthropic".to_string(),
2635            model: "claude-sonnet-5".to_string(),
2636            prompt_tokens: prompt,
2637            completion_tokens: 1,
2638            cached_tokens: 0,
2639            cache_write_tokens: 0,
2640            cost_usd: None,
2641            cost_reported_by_provider: None,
2642            at,
2643        };
2644        // The shape the issue reported: a compaction call and a stage call
2645        // landing between the same two progress ticks.
2646        let folded = fold(&[
2647            header(),
2648            usage(InferenceKind::Compaction, 7000, 1),
2649            usage(InferenceKind::Stage, 21_000, 2),
2650        ])
2651        .unwrap();
2652
2653        assert_eq!(folded.inference_count, 2);
2654        let seen: Vec<_> = folded
2655            .inference_usage
2656            .iter()
2657            .map(|u| (u.kind, u.prompt_tokens))
2658            .collect();
2659        assert_eq!(
2660            seen,
2661            vec![
2662                (InferenceKind::Compaction, 7000),
2663                (InferenceKind::Stage, 21_000)
2664            ]
2665        );
2666        // The whole reason this exists: their sum is 28k, and a reader of the
2667        // cumulative counter alone would see one 28k request and reasonably ask
2668        // whether a 32k window had been violated. Neither call came close.
2669        assert!(
2670            folded
2671                .inference_usage
2672                .iter()
2673                .all(|u| u.prompt_tokens < 32_000),
2674            "no single call exceeded the window, and the journal can now prove it"
2675        );
2676    }
2677
2678    /// The heavy variant and the light one both name one provider call, so a
2679    /// consumer counting calls should not have to know which the writer chose.
2680    #[test]
2681    fn both_inference_record_kinds_count_as_one_call_each() {
2682        let records = all_record_kinds();
2683        let folded = fold(&records).unwrap();
2684        let written = records
2685            .iter()
2686            .filter(|r| {
2687                matches!(
2688                    r,
2689                    RunRecord::Inference { .. } | RunRecord::InferenceUsage { .. }
2690                )
2691            })
2692            .count();
2693        assert_eq!(folded.inference_count, written);
2694        assert_eq!(folded.inference_usage.len(), 1);
2695    }
2696
2697    /// Replaying a journal has to land on exactly the state the run was in.
2698    ///
2699    /// Evictable regions are where a replay drifts: a folded `logs` region
2700    /// holding entries the live agent had already lost, or fewer than it
2701    /// held. This walks a region through the mutations a temporary
2702    /// region actually performs (append, evict-oldest, evict-and-append in one
2703    /// step, clear) and checks the digest -> delta -> apply chain reproduces
2704    /// every intermediate state exactly.
2705    #[test]
2706    fn probe_replay_matches_every_step() {
2707        fn region(name: &str, entries: &[(&str, usize)]) -> RegionSnapshot {
2708            RegionSnapshot {
2709                name: name.to_string(),
2710                kind: "temporary".to_string(),
2711                current_tokens: entries.iter().map(|(_, t)| *t).sum(),
2712                max_tokens: 1000,
2713                entries: entries
2714                    .iter()
2715                    .map(|(c, t)| RegionEntrySnapshot {
2716                        content: (*c).into(),
2717                        tokens: *t,
2718                        key: None,
2719                        kind: crate::region::EntryKind::Text,
2720                        metadata: None,
2721                        taint: crate::taint::TaintLevel::Public,
2722                        reasoning: None,
2723                    })
2724                    .collect(),
2725                description: None,
2726            }
2727        }
2728
2729        fn snap(entries: &[(&str, usize)]) -> ContextSnapshot {
2730            let r = region("logs", entries);
2731            ContextSnapshot {
2732                stage_name: "s".to_string(),
2733                total_tokens: r.current_tokens,
2734                max_tokens: 1000,
2735                regions: vec![r],
2736            }
2737        }
2738
2739        let steps: Vec<ContextSnapshot> = vec![
2740            snap(&[]),
2741            snap(&[("a", 10)]),
2742            snap(&[("a", 10), ("b", 20)]),
2743            snap(&[("a", 10), ("b", 20), ("c", 30)]),
2744            // Evict oldest.
2745            snap(&[("b", 20), ("c", 30)]),
2746            // Evict and append in one step - what a temporary region does when
2747            // an add pushes it over its bound.
2748            snap(&[("c", 30), ("d", 40)]),
2749            // Several evictions at once.
2750            snap(&[("d", 40)]),
2751            // Repeated identical content, the case a content hash cannot tell
2752            // apart by value alone.
2753            snap(&[("d", 40), ("x", 5)]),
2754            snap(&[("x", 5), ("x", 5)]),
2755            snap(&[("x", 5)]),
2756            snap(&[]),
2757        ];
2758
2759        // Fold exactly as the lane writes and the reader replays: retain a
2760        // digest, diff the next snapshot against it, apply to the running base.
2761        let mut base = steps[0].clone();
2762        let mut digest = digest_context(&steps[0]);
2763        for (i, next) in steps.iter().enumerate().skip(1) {
2764            let delta = diff_context_digest(&digest, next);
2765            apply_delta(&mut base, &delta);
2766            assert_eq!(
2767                base, *next,
2768                "step {i}: replay drifted from the live state\n  delta was {:?}",
2769                delta.regions
2770            );
2771            digest = digest_context(next);
2772        }
2773    }
2774
2775    /// One frame from a later build, written the way a later build would write
2776    /// it: correctly framed, with a variant name this one has never heard of.
2777    fn trailing_message() -> RunRecord {
2778        RunRecord::Message {
2779            message: MessageRecord {
2780                role: "user".to_string(),
2781                content: "after the unknown".to_string(),
2782            },
2783            at: 9,
2784        }
2785    }
2786
2787    /// Returns the archive and the size of the unknown frame's payload, so a
2788    /// caller can assert on the exact `Frame` it expects back.
2789    fn archive_with_an_unknown_record() -> (Vec<u8>, usize) {
2790        let mut buf = Vec::new();
2791        write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
2792        write_record(&mut buf, &header()).unwrap();
2793        let payload =
2794            serde_json::to_vec(&serde_json::json!({ "SomethingNew": { "whatever": 1 } })).unwrap();
2795        buf.extend_from_slice(&(payload.len() as u64).to_be_bytes());
2796        buf.extend_from_slice(&payload);
2797        write_record(&mut buf, &trailing_message()).unwrap();
2798        (buf, payload.len())
2799    }
2800
2801    /// A record's variant name, read off its serialized form. Avoids a match
2802    /// whose unreached arms would be uncovered, and pins the wire spelling.
2803    fn record_kind(record: &RunRecord) -> String {
2804        serde_json::to_value(record)
2805            .expect("a RunRecord always serializes")
2806            .as_object()
2807            .expect("externally tagged, so an object")
2808            .keys()
2809            .next()
2810            .expect("with exactly one key")
2811            .clone()
2812    }
2813
2814    /// The forward-compatibility guarantee, and the reason it is worth having:
2815    /// a reader that stopped at the first unknown record and returned the
2816    /// prefix would truncate the journal for every older reader, so a build
2817    /// predating `InferenceUsage` would read a 0.3.10 journal as "header, then
2818    /// nothing" - with no error.
2819    #[test]
2820    fn an_unknown_record_kind_is_stepped_over_not_treated_as_the_end() {
2821        let (buf, _) = archive_with_an_unknown_record();
2822        let (version, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
2823        assert_eq!(version, RUN_ARCHIVE_VERSION);
2824        let kinds: Vec<String> = records.iter().map(record_kind).collect();
2825        assert_eq!(
2826            kinds,
2827            vec!["Header".to_string(), "Message".to_string()],
2828            "the header, and the readable record after the gap"
2829        );
2830        assert_eq!(records[1], trailing_message(), "intact, not just present");
2831    }
2832
2833    /// The streaming reader skips the same way the buffering one does. It has
2834    /// its own loop, so "both apply the rule" is a claim that needs checking
2835    /// rather than assuming.
2836    #[test]
2837    fn the_streaming_reader_also_steps_over_an_unknown_record() {
2838        // A context record *after* the unknown frame: the only way to reach it
2839        // is to step over that frame, so the point count is the proof.
2840        let (mut buf, _) = archive_with_an_unknown_record();
2841        write_record(
2842            &mut buf,
2843            &RunRecord::ContextCheckpoint {
2844                snapshot: ContextSnapshot {
2845                    stage_name: "s".to_string(),
2846                    total_tokens: 1,
2847                    max_tokens: 10,
2848                    regions: vec![],
2849                },
2850                at: 11,
2851            },
2852        )
2853        .unwrap();
2854        let mut points = 0usize;
2855        visit_archive_points(&mut buf.as_slice(), &mut |_point| {
2856            points += 1;
2857            ControlFlow::Continue(())
2858        })
2859        .expect("a valid preamble");
2860        assert_eq!(points, 1, "the walk got past the unknown frame");
2861    }
2862
2863    /// The frame reader distinguishes "cannot parse this" from "cannot find
2864    /// the end of this". Only the second is fatal, and the difference is what
2865    /// makes stepping over the first safe.
2866    #[test]
2867    fn a_frame_reports_whether_its_payload_was_readable() {
2868        let (buf, unknown_bytes) = archive_with_an_unknown_record();
2869        let mut r = buf.as_slice();
2870        read_archive_start(&mut r).unwrap();
2871
2872        let mut frames = Vec::new();
2873        while let Some(frame) = read_frame(&mut r).expect("no torn frames here") {
2874            frames.push(frame);
2875        }
2876        assert_eq!(
2877            frames,
2878            vec![
2879                Frame::Record(Box::new(header())),
2880                // Stepped over, and it says how far.
2881                Frame::Unreadable {
2882                    bytes: unknown_bytes
2883                },
2884                Frame::Record(Box::new(trailing_message())),
2885            ],
2886            "one frame per record, with the unreadable one accounted for rather than ending the read"
2887        );
2888    }
2889
2890    /// A torn tail still ends the read. Skipping is for frames whose bytes are
2891    /// all present; a truncated one has no known length to step over.
2892    #[test]
2893    fn a_torn_frame_still_ends_the_read() {
2894        let mut buf = Vec::new();
2895        write_archive_start(&mut buf, RUN_ARCHIVE_VERSION).unwrap();
2896        write_record(&mut buf, &header()).unwrap();
2897        // A length prefix promising more than follows.
2898        buf.extend_from_slice(&999u64.to_be_bytes());
2899        buf.extend_from_slice(b"not enough");
2900
2901        let (_, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
2902        assert_eq!(records.len(), 1, "everything intact before the tear");
2903    }
2904
2905    /// An archive from a future framing generation is refused rather than
2906    /// misread. Reading it anyway would not fail cleanly - it would take
2907    /// whatever the length prefixes happened to say and produce nonsense.
2908    #[test]
2909    fn an_archive_from_a_newer_format_is_refused_with_both_versions_named() {
2910        let mut buf = Vec::new();
2911        write_archive_start(&mut buf, RUN_ARCHIVE_VERSION + 1).unwrap();
2912        write_record(&mut buf, &header()).unwrap();
2913
2914        let err = read_archive_lenient(&mut buf.as_slice()).unwrap_err();
2915        let message = err.to_string();
2916        assert!(
2917            message.contains(&(RUN_ARCHIVE_VERSION + 1).to_string()),
2918            "{message}"
2919        );
2920        assert!(message.contains("upgrade leviath"), "{message}");
2921    }
2922
2923    /// An older archive is read normally: framing has not changed under it, so
2924    /// the only difference is which record kinds it happens to contain.
2925    #[test]
2926    fn an_archive_from_an_older_format_still_reads() {
2927        let mut buf = Vec::new();
2928        write_archive_start(&mut buf, 0).unwrap();
2929        write_record(&mut buf, &header()).unwrap();
2930        let (version, records) = read_archive_lenient(&mut buf.as_slice()).unwrap();
2931        assert_eq!(version, 0);
2932        assert_eq!(records.len(), 1);
2933    }
2934}