Skip to main content

agent_top_core/harness/
mod.rs

1//! Per-harness transcript readers.
2//!
3//! Each harness writes a different append-only log. A `SessionTracker` turns
4//! one of those logs into the harness-neutral `SessionSummary` incrementally.
5
6pub mod claude;
7pub mod codex;
8pub mod gemini;
9pub mod kodelet;
10pub mod opencode;
11
12use crate::model::{Activity, Attribution, ContextOrigin, ContextSource, CostBreakdown, Harness, ProcNode, SpanKind, TokenUsage, ToolSpan};
13use crate::process::RawProc;
14use std::collections::{BTreeMap, HashSet, VecDeque};
15use std::path::{Path, PathBuf};
16use std::time::{Duration, SystemTime};
17
18/// Evidence that the parser still understands the file it is reading.
19///
20/// Every field is read with a fallback to zero, which is the right behaviour
21/// for a genuinely absent field and the wrong behaviour for a renamed one: a
22/// harness that renames `usage` next week would show a user 0 tokens and $0.00
23/// with no error at all. So count the usage records seen and how many of them
24/// yielded nothing. Records present and all of them empty is not a quiet
25/// session, it is a parser that has fallen behind the format.
26#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
27pub struct ParseHealth {
28    /// Model responses seen. Each of these should account for some tokens.
29    pub billable_messages: u64,
30    /// Usage records found on them. Zero of these, with messages present, means
31    /// the record itself moved or was renamed.
32    pub usage_records: u64,
33    /// Records found but yielding nothing, which is what a renamed field inside
34    /// an intact record looks like.
35    pub empty_usage_records: u64,
36}
37
38impl ParseHealth {
39    /// Enough responses to accuse the parser rather than the session. A couple
40    /// of odd messages must not raise the alarm.
41    const MIN_EVIDENCE: u64 = 3;
42
43    /// The session did work that must have cost tokens, and we read none.
44    ///
45    /// Covers both ways a format change reaches us: the usage record moving or
46    /// being renamed, so we never find one, and the fields inside it being
47    /// renamed, so we find records that read as empty. Neither raises an error
48    /// on its own, because every field falls back to zero.
49    pub fn fields_unrecognised(&self) -> bool {
50        self.billable_messages >= Self::MIN_EVIDENCE && self.usage_records == self.empty_usage_records
51    }
52}
53
54#[derive(Debug, Clone, Default)]
55pub struct SessionSummary {
56    pub harness: Option<Harness>,
57    pub session_id: Option<String>,
58    /// Logical subagent metadata from the transcript, never from process ancestry.
59    pub subagent: Option<crate::model::SubagentInfo>,
60    pub cwd: Option<PathBuf>,
61    pub model: Option<String>,
62    pub harness_version: Option<String>,
63    pub usage: TokenUsage,
64    pub cost_usd: f64,
65    /// `cost_usd` by kind of token. See `Agent::cost_breakdown`.
66    pub cost_breakdown: CostBreakdown,
67    /// Set when the harness records its own costs instead of using our table.
68    pub price_source: Option<crate::model::PriceSource>,
69    pub unpriced_tokens: u64,
70    pub turns: u64,
71    pub subagent_turns: u64,
72    /// Child sessions are counted in this row's own totals, so the subagent
73    /// share of them is worth breaking out. A harness that gives each child a
74    /// row of its own leaves this false: nothing was folded in to separate.
75    pub folds_child_usage: bool,
76    pub tool_calls: u64,
77    /// Earlier calls may have been removed from the harness's retained history.
78    pub tool_calls_lower_bound: bool,
79    /// See `Agent::web_searches`.
80    pub web_searches: u64,
81    pub spans: SpanLog,
82    /// Calls to each MCP server, by the server's name.
83    pub mcp: BTreeMap<String, McpUsage>,
84    /// What each tool's results added to the prompt, and what carrying it
85    /// has cost. See `ContextLedger`.
86    pub context: ContextLedger,
87    pub health: ParseHealth,
88    pub activity: Activity,
89    pub started_at: Option<SystemTime>,
90    pub last_activity: Option<SystemTime>,
91    /// How close the session is to its rate limit, when the harness writes it.
92    pub rate_limit: Option<crate::model::RateLimit>,
93    /// The harness's own cost figure, when it writes one. See `HarnessCost`.
94    pub harness_cost: Option<crate::model::HarnessCost>,
95}
96
97/// What a transcript says about one MCP server: how often it was called,
98/// how often that failed, and when it was last called.
99#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
100pub struct McpUsage {
101    pub calls: u64,
102    pub errors: u64,
103    pub last_call: Option<SystemTime>,
104}
105
106impl McpUsage {
107    pub fn add(&mut self, o: &McpUsage) {
108        self.calls += o.calls;
109        self.errors += o.errors;
110        self.last_call = self.last_call.max(o.last_call);
111    }
112}
113
114/// What each source has added to a session's context, and what carrying it
115/// has cost. Built incrementally by an adapter from two facts it already
116/// has: the tool results it sees submitted, and the usage on each response.
117///
118/// The arithmetic. A response's prompt is the previous response's prompt
119/// plus everything appended since: the previous reply, and the tool results
120/// that answered it. So `prompt_n - prompt_n-1` is the new material, the
121/// previous reply's `output` is the part of it the model wrote itself, and
122/// the rest is the tool results submitted in between. Those tokens go to
123/// the tools that produced them, split evenly across result records. Each
124/// record may split its share by nested tools' byte weights, a heuristic. The
125/// first response's whole prompt, the replies, and any growth with no
126/// result to explain it (the user's own messages) are `Other`.
127///
128/// Every response then re-reads the whole context, so each source's live
129/// tokens are charged at that response's prompt rate: its prompt-side cost
130/// over its prompt tokens, which is mostly the cache-read price with some
131/// fresh input mixed in. The first read is charged the same way, so the
132/// sources' costs sum to the session's prompt-side cost.
133///
134/// A compaction replaces the context. The live set is cleared when the
135/// harness says one happened (`compacted`), and, for a harness that does
136/// not say, when the prompt halves, which nothing else does. Thinking
137/// blocks a harness drops between turns shrink the prompt by less than
138/// that; the shrink is taken off `Other`, whose replies they were.
139#[derive(Debug, Clone, Default)]
140pub struct ContextLedger {
141    shares: BTreeMap<ContextKey, ContextShare>,
142    /// Tokens each source has in the context now; cleared at a compaction.
143    live: BTreeMap<ContextKey, u64>,
144    /// Results submitted since the last response, by call id, awaiting the
145    /// response that will say how big they were.
146    pending: Vec<(String, ContextWeights)>,
147    /// The last response's prompt and output, for the next delta.
148    prev: Option<(u64, u64)>,
149}
150
151type ContextKey = (ContextOrigin, String);
152
153/// Per-call byte weights within one result record, not across records.
154pub(super) type ContextWeights = Vec<(ContextOrigin, String, u64)>;
155
156/// One source's running totals.
157#[derive(Debug, Clone, Copy, Default, PartialEq)]
158pub struct ContextShare {
159    pub calls: u64,
160    pub tokens: u64,
161    pub cost_usd: f64,
162}
163
164impl ContextLedger {
165    /// The name every non-tool share is filed under.
166    pub const OTHER: &str = "other";
167
168    fn other() -> ContextKey {
169        (ContextOrigin::Other, Self::OTHER.to_string())
170    }
171
172    /// A tool result was submitted to the model. `id` is the harness's call
173    /// id, so a harness that names the MCP server only after the result is
174    /// written can `retag` it before the response arrives.
175    pub fn result(&mut self, id: &str, origin: ContextOrigin, name: &str) {
176        self.result_weighted(id, vec![(origin, name.to_string(), 1)]);
177    }
178
179    /// Split one result's share by weight; all-zero weights split equally.
180    /// Callers resolve missing weights. Empty groups are ignored.
181    pub(super) fn result_weighted(&mut self, id: &str, sources: ContextWeights) {
182        if !sources.is_empty() {
183            self.pending.push((id.to_string(), sources));
184        }
185    }
186
187    /// Re-file a pending result under another source. A no-op once the
188    /// response that sized it has been seen.
189    pub fn retag(&mut self, id: &str, origin: ContextOrigin, name: &str) {
190        if let Some((_, sources)) = self.pending.iter_mut().find(|(i, _)| i == id) {
191            for (source_origin, source_name, _) in sources {
192                *source_origin = origin;
193                *source_name = name.to_string();
194            }
195        }
196    }
197
198    /// A response came back. Sizes and files whatever was submitted since
199    /// the last one, then charges everything in the context for this read.
200    /// A record with no prompt is not an API response and is ignored.
201    pub fn response(&mut self, usage: &TokenUsage, cost: &CostBreakdown) {
202        let prompt = usage.prompt();
203        if prompt == 0 {
204            return;
205        }
206        let pending = std::mem::take(&mut self.pending);
207        match self.prev {
208            Some((p, _)) if prompt < p / 2 => {
209                // Compacted: what is in the context now is all new.
210                self.live.clear();
211                self.file(prompt, 0, pending);
212            }
213            Some((p, o)) if prompt >= p => {
214                let growth = prompt - p;
215                let reply = o.min(growth);
216                self.file(growth - reply, reply, pending);
217            }
218            Some(_) => {
219                // Shrunk, but not compacted: thinking blocks dropped between
220                // turns. They were part of a reply, so the shrink is Other's.
221                let e = self.live.entry(Self::other()).or_default();
222                *e = e.saturating_sub(self.prev.map(|(p, _)| p - prompt).unwrap_or(0));
223                self.file(0, 0, pending);
224            }
225            None => self.file(prompt, 0, pending),
226        }
227        let prompt_cost = cost.input + cost.cache_read + cost.cache_write_5m + cost.cache_write_1h;
228        if prompt_cost > 0.0 {
229            let rate = prompt_cost / prompt as f64;
230            for (k, t) in &self.live {
231                self.shares.entry(k.clone()).or_default().cost_usd += *t as f64 * rate;
232            }
233        }
234        self.prev = Some((prompt, usage.output));
235    }
236
237    /// Split evenly across records, then by weight within each record.
238    /// No pending result puts everything under `Other`.
239    fn file(&mut self, results: u64, other: u64, pending: Vec<(String, ContextWeights)>) {
240        if pending.is_empty() {
241            self.add(Self::other(), results + other, 0);
242            return;
243        }
244        self.add(Self::other(), other, 0);
245        let n = pending.len() as u64;
246        let (each, mut rem) = (results / n, results % n);
247        for (_, sources) in pending {
248            let t = each + u64::from(rem > 0);
249            rem = rem.saturating_sub(1);
250            let total: u128 = sources.iter().map(|(_, _, weight)| u128::from(*weight)).sum();
251            let denominator = if total == 0 { sources.len() as u128 } else { total };
252            let mut remaining = t;
253            let mut parts: Vec<_> = sources
254                .into_iter()
255                .map(|(origin, name, weight)| {
256                    let weight = if total == 0 { 1 } else { u128::from(weight) };
257                    let numerator = u128::from(t) * weight;
258                    let tokens = (numerator / denominator) as u64;
259                    remaining -= tokens;
260                    ((origin, name), tokens, numerator % denominator)
261                })
262                .collect();
263            // Largest remainders first; stable ties follow contribution order.
264            parts.sort_by_key(|part| std::cmp::Reverse(part.2));
265            for (key, tokens, _) in parts {
266                let tokens = tokens + u64::from(remaining > 0);
267                remaining = remaining.saturating_sub(1);
268                self.add(key, tokens, 1);
269            }
270        }
271    }
272
273    fn add(&mut self, key: ContextKey, tokens: u64, calls: u64) {
274        if tokens == 0 && calls == 0 {
275            return;
276        }
277        let sh = self.shares.entry(key.clone()).or_default();
278        sh.tokens += tokens;
279        sh.calls += calls;
280        *self.live.entry(key).or_default() += tokens;
281    }
282
283    /// The harness replaced the context with a summary. Whatever the next
284    /// response carries is new.
285    pub fn compacted(&mut self) {
286        self.live.clear();
287        self.prev = None;
288    }
289
290    /// Fold another ledger's totals in, as a parent folds its subagents.
291    /// The running state is not merged: a fold is read, never fed.
292    pub fn merge(&mut self, other: &ContextLedger) {
293        for (k, sh) in &other.shares {
294            let e = self.shares.entry(k.clone()).or_default();
295            e.calls += sh.calls;
296            e.tokens += sh.tokens;
297            e.cost_usd += sh.cost_usd;
298        }
299    }
300
301    pub fn is_empty(&self) -> bool {
302        self.shares.is_empty()
303    }
304
305    /// Every source, largest first.
306    pub fn sources(&self) -> Vec<ContextSource> {
307        let mut v: Vec<ContextSource> = self
308            .shares
309            .iter()
310            .map(|((origin, name), sh)| ContextSource {
311                name: name.clone(),
312                origin: *origin,
313                calls: sh.calls,
314                tokens: sh.tokens,
315                cost_usd: sh.cost_usd,
316            })
317            .collect();
318        v.sort_by(|a, b| b.tokens.cmp(&a.tokens).then_with(|| a.name.cmp(&b.name)));
319        v
320    }
321}
322
323/// The server behind an MCP tool name. Claude Code names them
324/// `mcp__<server>__<tool>`; the server part may itself contain underscores,
325/// so the split is on the first double underscore after the prefix and the
326/// tool part is whatever follows the last one.
327pub fn mcp_server_of(tool_name: &str) -> Option<&str> {
328    let rest = tool_name.strip_prefix("mcp__")?;
329    let server = match rest.rfind("__") {
330        Some(i) => &rest[..i],
331        None => rest,
332    };
333    if server.is_empty() { None } else { Some(server) }
334}
335
336/// Spans kept per session by the live tracker. A screenful of waterfall is
337/// a few dozen rows; the rest is history nobody scrolls to in a live view, and
338/// every span costs a clone on each refresh. An export wants the whole session
339/// and uses `SpanLog::unbounded` in a separate pass; see `SpanRetention`.
340///
341/// Sized for roughly a hundred tool calls: each model response also adds an
342/// inference span and each human prompt a turn span.
343pub const MAX_SPANS: usize = 256;
344
345/// A bounded, in-order log of tool spans, built by pairing a harness's
346/// "call started" and "call finished" records by call id.
347///
348/// Records arrive interleaved and out of order (agents run tools in parallel),
349/// so a span is closed by searching back for the still-open span with that id
350/// rather than assuming the most recent one.
351#[derive(Debug, Clone)]
352pub struct SpanLog {
353    spans: VecDeque<ToolSpan>,
354    cap: usize,
355}
356
357impl Default for SpanLog {
358    fn default() -> Self {
359        SpanLog { spans: VecDeque::new(), cap: MAX_SPANS }
360    }
361}
362
363impl SpanLog {
364    /// A log that keeps every span. For a one-shot pass over a whole
365    /// transcript, never for the live tracker, where the memory and the clone
366    /// per refresh would grow with the session.
367    pub fn unbounded() -> Self {
368        SpanLog { spans: VecDeque::new(), cap: usize::MAX }
369    }
370
371    /// Record the start of a tool call. Ignored when that id is already open,
372    /// so a transcript line replayed by the harness does not double-count.
373    pub fn open(&mut self, id: String, name: String, at: SystemTime, sidechain: bool) {
374        self.open_kind(id, name, at, sidechain, SpanKind::Tool);
375    }
376
377    /// `open`, for any kind of span.
378    pub fn open_kind(&mut self, id: String, name: String, at: SystemTime, sidechain: bool, kind: SpanKind) {
379        if id.is_empty() || self.spans.iter().any(|s| s.is_open() && s.id == id) {
380            return;
381        }
382        if self.spans.len() >= self.cap {
383            self.spans.pop_front();
384        }
385        self.spans.push_back(ToolSpan { id, name, started_at: at, duration_ms: None, sidechain, error: false, kind });
386    }
387
388    /// Move the end of the newest span with this id to `at`, open or not. An
389    /// inference span grows as the response streams in, one content block
390    /// per line, and its end is wherever the last block landed.
391    pub fn end_at(&mut self, id: &str, at: SystemTime) {
392        let Some(s) = self.spans.iter_mut().rev().find(|s| s.id == id) else { return };
393        s.duration_ms = Some(at.duration_since(s.started_at).map(|d| d.as_millis() as u64).unwrap_or(0));
394    }
395
396    /// The newest open span of this kind, if any.
397    pub fn open_of_kind(&self, kind: SpanKind) -> Option<&ToolSpan> {
398        self.spans.iter().rev().find(|s| s.is_open() && s.kind == kind)
399    }
400
401    /// Remove the newest span with this id if it is still open. For a span
402    /// that turned out not to be one: an inference that never produced a
403    /// reply because the user interrupted or submitted again.
404    pub fn discard_open(&mut self, id: &str) {
405        if let Some(i) = self.spans.iter().rposition(|s| s.is_open() && s.id == id) {
406            self.spans.remove(i);
407        }
408    }
409
410    /// Close the open call with this id. A result whose call scrolled out of
411    /// the window, or that we never saw start, is dropped.
412    pub fn close(&mut self, id: &str, at: SystemTime, error: bool) {
413        let Some(s) = self.spans.iter_mut().rev().find(|s| s.is_open() && s.id == id) else { return };
414        s.duration_ms = Some(at.duration_since(s.started_at).map(|d| d.as_millis() as u64).unwrap_or(0));
415        s.error = error;
416    }
417
418    pub fn len(&self) -> usize {
419        self.spans.len()
420    }
421
422    pub fn is_empty(&self) -> bool {
423        self.spans.is_empty()
424    }
425
426    /// Oldest first. Double-ended so callers can take the newest spans without
427    /// collecting the whole log first, which the UI and the golden tests both do.
428    pub fn iter(&self) -> impl DoubleEndedIterator<Item = &ToolSpan> + ExactSizeIterator {
429        self.spans.iter()
430    }
431
432    pub fn to_vec(&self) -> Vec<ToolSpan> {
433        self.spans.iter().cloned().collect()
434    }
435
436    /// One log from several, ordered by start time, keeping the newest `cap`.
437    /// A parent's spans and its subagents' spans interleave in wall-clock
438    /// order, which is what a waterfall wants.
439    pub fn merged<'a>(logs: impl IntoIterator<Item = &'a SpanLog>, cap: usize) -> SpanLog {
440        let mut spans: Vec<ToolSpan> = logs.into_iter().flat_map(|l| l.spans.iter().cloned()).collect();
441        spans.sort_by_key(|s| s.started_at);
442        if spans.len() > cap {
443            spans.drain(..spans.len() - cap);
444        }
445        SpanLog { spans: spans.into(), cap }
446    }
447
448    pub fn cap(&self) -> usize {
449        self.cap
450    }
451}
452
453/// How many of a session's tool spans a tracker keeps.
454#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
455pub enum SpanRetention {
456    /// The newest `MAX_SPANS`, enough for the live waterfall. The default.
457    #[default]
458    Recent,
459    /// Every span in the transcript, for a trace export. Memory grows with
460    /// the session, so this is for a single pass, not a tracker kept across
461    /// refreshes.
462    All,
463}
464
465impl SpanRetention {
466    pub(crate) fn log(self) -> SpanLog {
467        match self {
468            SpanRetention::Recent => SpanLog::default(),
469            SpanRetention::All => SpanLog::unbounded(),
470        }
471    }
472}
473
474pub trait SessionTracker {
475    /// Ingest whatever was appended since the last call. Returns true when
476    /// there is still unread data (the byte budget was exhausted).
477    fn refresh(&mut self) -> anyhow::Result<bool>;
478    fn summary(&self) -> &SessionSummary;
479    fn path(&self) -> &Path;
480
481    /// Ingest the whole file, however many refreshes that takes. For a
482    /// one-shot read such as an export; the live collector spreads a large
483    /// transcript over several ticks instead.
484    fn refresh_all(&mut self) -> anyhow::Result<()> {
485        while self.refresh()? {}
486        Ok(())
487    }
488}
489
490/// What a harness's own registry says about one of its processes, when it
491/// keeps one (Claude Code's `~/.claude/sessions/<pid>.json`). Every field is
492/// optional; a harness with no registry returns none of this.
493#[derive(Debug, Clone, Default, PartialEq, Eq)]
494pub struct RegistryHints {
495    pub name: Option<String>,
496    pub session_id: Option<String>,
497    pub cwd: Option<PathBuf>,
498    pub version: Option<String>,
499    /// The harness's own word for its state (`busy`, `idle`, ...), which beats
500    /// any transcript heuristic.
501    pub status: Option<String>,
502}
503
504/// What the collector knows about a process when it asks an adapter which
505/// transcript is the process's.
506pub struct AttributeContext<'a> {
507    pub cwd: Option<&'a Path>,
508    pub proc_start: SystemTime,
509    pub now: SystemTime,
510    /// Transcripts already given to another process this pass. An adapter
511    /// must not hand one out twice.
512    pub attached: &'a HashSet<PathBuf>,
513    /// A transcript idle for longer than this is a finished conversation, not
514    /// a thread of a process that cannot otherwise be matched.
515    pub activity_timeout: Duration,
516}
517
518/// One harness, as the collector sees it: where its transcripts are, which
519/// belongs to which process, and how to read one. The collector holds a list
520/// of these and never names a harness itself, so adding a harness is one
521/// module and one line in `adapters()`. Process recognition stays in
522/// `process::classify_agent`, which also knows the harnesses that have no
523/// transcript adapter yet.
524pub trait HarnessAdapter {
525    fn harness(&self) -> Harness;
526
527    /// Re-list the transcripts written since `since`. Called every
528    /// `fs_scan_interval`, not every tick.
529    fn rescan(&mut self, since: SystemTime);
530
531    /// Called once per pass with this harness's root processes, before any
532    /// of them is attributed. For work that must see every process at once:
533    /// Codex reads which rollouts each process holds open here, so that no
534    /// process's fallback can claim a thread another is demonstrably writing.
535    fn prepare(&mut self, _roots: &[&ProcNode], _by_pid: &std::collections::HashMap<u32, &RawProc>) {}
536
537    /// The harness's own registry entry for a process, if it keeps one.
538    fn hints(&self, _pid: u32) -> Option<RegistryHints> {
539        None
540    }
541
542    /// The transcripts this process is writing, newest activity first, and
543    /// how sure the adapter is. One per conversation: a Codex app-server hosts
544    /// many, a CLI runs one, a process with none gets an empty list.
545    fn attribute(&self, root: &ProcNode, raw: Option<&RawProc>, ctx: &AttributeContext) -> (Vec<PathBuf>, Attribution);
546
547    /// Recently written transcripts no process owns: the stopped list.
548    fn unowned(&self, attached: &HashSet<PathBuf>) -> Vec<PathBuf>;
549
550    /// A tracker for one transcript.
551    fn open(&self, path: &Path, spans: SpanRetention) -> Box<dyn SessionTracker>;
552
553    /// Whether this harness wrote the file, judged from its first few lines.
554    fn detect(&self, path: &Path) -> bool;
555
556    /// Every transcript on disk, however old, with the id a user would type
557    /// to name it. For `agent-top trace --session <id>`.
558    fn transcripts(&self) -> Vec<(String, PathBuf)>;
559}
560
561/// Every harness that has a transcript adapter, in the order they are asked.
562/// The order matters to `detect` alone: Gemini's metadata line carries a
563/// `sessionId` like Claude Code's lines do, so it is asked first.
564pub fn adapters() -> Vec<Box<dyn HarnessAdapter>> {
565    vec![
566        Box::new(codex::CodexAdapter::default()),
567        Box::new(gemini::GeminiAdapter::default()),
568        Box::new(opencode::OpenCodeAdapter::default()),
569        Box::new(kodelet::KodeletAdapter::default()),
570        Box::new(claude::ClaudeAdapter::default()),
571    ]
572}
573
574/// The adapter for one harness, or none when it has only a process table entry.
575pub fn adapter_for(harness: Harness) -> Option<Box<dyn HarnessAdapter>> {
576    adapters().into_iter().find(|a| a.harness() == harness)
577}
578
579/// Which harness wrote a transcript, judged from its first few lines. Anything
580/// no adapter recognises is not a transcript agent-top reads.
581pub fn detect(path: &Path) -> Option<Harness> {
582    adapters().iter().find(|a| a.detect(path)).map(|a| a.harness())
583}
584
585/// A tracker for a transcript whose harness is already known, or none when
586/// that harness has no transcript adapter.
587pub fn open_transcript(path: &Path, harness: Harness, spans: SpanRetention) -> Option<Box<dyn SessionTracker>> {
588    adapter_for(harness).map(|a| a.open(path, spans))
589}
590
591/// The first few lines of a file, parsed, for `HarnessAdapter::detect`.
592pub(crate) fn head_lines(path: &Path) -> Vec<serde_json::Value> {
593    use std::io::{BufRead, BufReader};
594    let Ok(f) = std::fs::File::open(path) else { return Vec::new() };
595    BufReader::new(f).lines().map_while(Result::ok).take(5).filter_map(|l| serde_json::from_str(&l).ok()).collect()
596}
597
598/// Bytes ingested per tracker per refresh. Keeps a cold start on a 100 MB
599/// transcript from freezing the first frame; the rest streams in on later ticks.
600pub const REFRESH_BUDGET_BYTES: usize = 8 * 1024 * 1024;
601
602/// Parse an RFC 3339 timestamp like `2026-09-03T07:15:34.123Z` into SystemTime
603/// without pulling in a date crate. The UTC `Z` form is what Claude Code,
604/// Codex and Gemini write; a numeric offset (`+01:00`, `-0700`) is accepted
605/// too, since a harness that writes local time would otherwise lose its
606/// last-activity time and with it the idle clock and the mtime fallbacks.
607pub fn parse_rfc3339_utc(s: &str) -> Option<SystemTime> {
608    let s = s.trim();
609    // Split off the zone: `Z`, or a signed offset after the time.
610    let (s, offset_secs) = match s.strip_suffix(['Z', 'z']) {
611        Some(rest) => (rest, 0i64),
612        None => {
613            let t_pos = s.find('T')?;
614            let sign_pos = s[t_pos..].rfind(['+', '-'])? + t_pos;
615            let (rest, zone) = s.split_at(sign_pos);
616            let sign = if zone.starts_with('-') { -1 } else { 1 };
617            let digits: String = zone[1..].chars().filter(|c| c.is_ascii_digit()).collect();
618            if digits.len() != 4 {
619                return None;
620            }
621            let oh = digits[..2].parse::<i64>().ok()?;
622            let om = digits[2..].parse::<i64>().ok()?;
623            (rest, sign * (oh * 3600 + om * 60))
624        }
625    };
626    let (date, time) = s.split_once('T')?;
627    let mut d = date.split('-');
628    let (y, mo, da) = (d.next()?.parse::<i64>().ok()?, d.next()?.parse::<u32>().ok()?, d.next()?.parse::<u32>().ok()?);
629    let mut t = time.split(':');
630    let (h, mi) = (t.next()?.parse::<u64>().ok()?, t.next()?.parse::<u64>().ok()?);
631    let sec_str = t.next()?;
632    let (sec, frac) = match sec_str.split_once('.') {
633        Some((s, f)) => (s.parse::<u64>().ok()?, f),
634        None => (sec_str.parse::<u64>().ok()?, ""),
635    };
636    let nanos: u32 = if frac.is_empty() {
637        0
638    } else {
639        let mut f = frac.to_string();
640        f.truncate(9);
641        while f.len() < 9 {
642            f.push('0');
643        }
644        f.parse().ok()?
645    };
646    let days = days_from_civil(y, mo, da);
647    // Local time minus its offset is UTC.
648    let secs = days * 86_400 + (h * 3600 + mi * 60 + sec) as i64 - offset_secs;
649    if secs < 0 {
650        return None;
651    }
652    Some(SystemTime::UNIX_EPOCH + std::time::Duration::new(secs as u64, nanos))
653}
654
655// Howard Hinnant's days-from-civil algorithm.
656fn days_from_civil(y: i64, m: u32, d: u32) -> i64 {
657    let y = if m <= 2 { y - 1 } else { y };
658    let era = if y >= 0 { y } else { y - 399 } / 400;
659    let yoe = y - era * 400;
660    let mp = (m as i64 + 9) % 12;
661    let doy = (153 * mp + 2) / 5 + d as i64 - 1;
662    let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
663    era * 146_097 + doe - 719_468
664}
665
666#[cfg(test)]
667mod tests {
668    use super::*;
669    use std::time::Duration;
670
671    fn at(secs: u64) -> SystemTime {
672        SystemTime::UNIX_EPOCH + Duration::from_secs(secs)
673    }
674
675    #[test]
676    fn pairs_spans_by_id_out_of_order() {
677        let mut log = SpanLog::default();
678        log.open("a".into(), "Bash".into(), at(10), false);
679        log.open("b".into(), "Read".into(), at(11), true);
680        // Replay of the same start line must not open a second span.
681        log.open("a".into(), "Bash".into(), at(10), false);
682        // Results come back in the other order.
683        log.close("b", at(12), false);
684        log.close("a", at(14), true);
685        // A result with no matching call is ignored.
686        log.close("zzz", at(15), false);
687        let v = log.to_vec();
688        assert_eq!(v.len(), 2);
689        assert_eq!(v[0].name, "Bash");
690        assert_eq!(v[0].duration_ms, Some(4_000));
691        assert!(v[0].error);
692        assert_eq!(v[1].duration_ms, Some(1_000));
693        assert!(v[1].sidechain);
694        assert!(!v[1].error);
695    }
696
697    #[test]
698    fn keeps_the_newest_spans_and_reports_open_ones() {
699        let mut log = SpanLog::default();
700        for i in 0..(MAX_SPANS + 10) {
701            log.open(format!("id{i}"), "T".into(), at(i as u64), false);
702            log.close(&format!("id{i}"), at(i as u64), false);
703        }
704        assert_eq!(log.len(), MAX_SPANS);
705        assert_eq!(log.iter().next().unwrap().id, "id10");
706        log.open("live".into(), "Bash".into(), at(500), false);
707        let last = log.to_vec().pop().unwrap();
708        assert!(last.is_open());
709        assert_eq!(last.elapsed_ms(at(503)), 3_000);
710    }
711
712    #[test]
713    fn end_at_moves_the_end_of_any_kind_of_span() {
714        let mut log = SpanLog::default();
715        log.open_kind("inference:1".into(), "inference".into(), at(10), false, SpanKind::Inference);
716        assert!(log.open_of_kind(SpanKind::Inference).is_some());
717        assert!(log.open_of_kind(SpanKind::Turn).is_none());
718        // The response streams in over three lines; the span ends at the last one.
719        log.end_at("inference:1", at(11));
720        log.end_at("inference:1", at(13));
721        log.end_at("nope", at(99));
722        let v = log.to_vec();
723        assert_eq!(v[0].duration_ms, Some(3_000));
724        assert_eq!(v[0].kind, SpanKind::Inference);
725        assert!(log.open_of_kind(SpanKind::Inference).is_none());
726        // Discarding only removes open spans; the ended one stays.
727        log.discard_open("inference:1");
728        assert_eq!(log.len(), 1);
729        log.open_kind("inference:2".into(), "inference".into(), at(20), false, SpanKind::Inference);
730        log.discard_open("inference:2");
731        assert_eq!(log.len(), 1);
732    }
733
734    #[test]
735    fn unbounded_log_keeps_everything() {
736        let mut log = SpanLog::unbounded();
737        for i in 0..(MAX_SPANS * 3) {
738            log.open(format!("id{i}"), "T".into(), at(i as u64), false);
739            log.close(&format!("id{i}"), at(i as u64 + 1), false);
740        }
741        assert_eq!(log.len(), MAX_SPANS * 3);
742        assert_eq!(log.iter().next().unwrap().id, "id0");
743        assert_eq!(SpanRetention::default(), SpanRetention::Recent);
744    }
745
746    #[test]
747    fn detects_the_harness_from_the_first_lines() {
748        let dir = std::env::temp_dir().join(format!("agent-top-detect-{}", std::process::id()));
749        std::fs::create_dir_all(&dir).unwrap();
750        let codex = dir.join("rollout.jsonl");
751        std::fs::write(&codex, "{\"type\":\"session_meta\",\"payload\":{\"id\":\"x\"}}\n").unwrap();
752        let claude = dir.join("s.jsonl");
753        // A summary line first, as Claude Code writes on resume, then a real one.
754        std::fs::write(&claude, "{\"type\":\"summary\",\"leafUuid\":\"u\"}\n{\"type\":\"user\",\"sessionId\":\"abc\"}\n").unwrap();
755        let other = dir.join("other.jsonl");
756        std::fs::write(&other, "{\"hello\":1}\nnot json\n").unwrap();
757        assert_eq!(detect(&codex), Some(Harness::Codex));
758        assert_eq!(detect(&claude), Some(Harness::Claude));
759        assert_eq!(detect(&other), None);
760        assert_eq!(detect(&dir.join("missing.jsonl")), None);
761        let _ = std::fs::remove_dir_all(&dir);
762    }
763
764    #[test]
765    fn names_the_server_behind_an_mcp_tool() {
766        assert_eq!(mcp_server_of("mcp__filesystem__read_file"), Some("filesystem"));
767        assert_eq!(mcp_server_of("mcp__chrome-devtools__take_screenshot"), Some("chrome-devtools"));
768        assert_eq!(mcp_server_of("mcp__claude_ai_Gmail__authenticate"), Some("claude_ai_Gmail"));
769        assert_eq!(mcp_server_of("mcp__odd"), Some("odd"));
770        assert_eq!(mcp_server_of("mcp____x"), None);
771        assert_eq!(mcp_server_of("Bash"), None);
772    }
773
774    #[test]
775    fn accuses_the_parser_only_with_enough_evidence() {
776        // Healthy: records found on the messages, tokens read from them.
777        let h = ParseHealth { billable_messages: 40, usage_records: 40, empty_usage_records: 0 };
778        assert!(!h.fields_unrecognised());
779        // One odd message among many is a message, not a format change.
780        let h = ParseHealth { billable_messages: 40, usage_records: 40, empty_usage_records: 39 };
781        assert!(!h.fields_unrecognised());
782        // Fields inside the record renamed: records found, all of them empty.
783        let h = ParseHealth { billable_messages: 40, usage_records: 40, empty_usage_records: 40 };
784        assert!(h.fields_unrecognised());
785        // The record itself renamed or moved: messages, but no records at all.
786        // This is the case a naive check misses, because there is nothing to count.
787        let h = ParseHealth { billable_messages: 40, usage_records: 0, empty_usage_records: 0 };
788        assert!(h.fields_unrecognised());
789        // Too early to tell: a session that has barely started.
790        let h = ParseHealth { billable_messages: 2, usage_records: 0, empty_usage_records: 0 };
791        assert!(!h.fields_unrecognised());
792        // Nothing parsed at all is silence, not evidence.
793        assert!(!ParseHealth::default().fields_unrecognised());
794    }
795
796    fn usage(prompt: u64, output: u64) -> TokenUsage {
797        TokenUsage { cache_read: prompt, output, ..Default::default() }
798    }
799
800    /// $1 per million prompt tokens, so a source's cost is its live tokens
801    /// summed over the responses that read them, in micro-dollars.
802    fn cost(prompt: u64) -> CostBreakdown {
803        CostBreakdown { cache_read: prompt as f64 / 1e6, ..Default::default() }
804    }
805
806    fn share<'a>(v: &'a [ContextSource], name: &str) -> &'a ContextSource {
807        v.iter().find(|s| s.name == name).unwrap_or_else(|| panic!("no source {name}"))
808    }
809
810    #[test]
811    fn context_ledger_files_prompt_growth_under_the_results_that_caused_it() {
812        let mut l = ContextLedger::default();
813        // First response: the whole prompt is the system prompt and the ask.
814        l.response(&usage(1_000, 100), &cost(1_000));
815        // Two tool results, then a response 2_300 bigger: 100 of that is the
816        // reply, 2_200 the results, split evenly.
817        l.result("a", ContextOrigin::Tool, "Read");
818        l.result("b", ContextOrigin::Mcp, "fs");
819        l.response(&usage(3_300, 50), &cost(3_300));
820        let v = l.sources();
821        assert_eq!(share(&v, "Read").tokens, 1_100);
822        assert_eq!(share(&v, "fs").tokens, 1_100);
823        assert_eq!(share(&v, "fs").origin, ContextOrigin::Mcp);
824        assert_eq!(share(&v, "fs").calls, 1);
825        let other = share(&v, ContextLedger::OTHER);
826        assert_eq!((other.tokens, other.calls), (1_100, 0));
827        // Costs: Other was read twice (1_000 then 1_100), the results once.
828        assert!((other.cost_usd - 2_100e-6).abs() < 1e-12, "{}", other.cost_usd);
829        assert!((share(&v, "Read").cost_usd - 1_100e-6).abs() < 1e-12);
830        // The sources sum to the prompt-side cost.
831        let total: f64 = v.iter().map(|s| s.cost_usd).sum();
832        assert!((total - 4_300e-6).abs() < 1e-12, "{total}");
833        assert_eq!(v[0].tokens, 1_100, "largest first");
834    }
835
836    #[test]
837    fn context_ledger_weights_preserve_equal_shares_between_result_records() {
838        let mut l = ContextLedger::default();
839        l.response(&usage(1_000, 100), &cost(1_000));
840        l.result_weighted("wrapper", vec![(ContextOrigin::Tool, "exec_command".into(), 1), (ContextOrigin::Tool, "apply_patch".into(), 3)]);
841        l.result("direct", ContextOrigin::Mcp, "fs");
842        l.response(&usage(3_100, 0), &cost(3_100));
843        let v = l.sources();
844        assert_eq!(share(&v, "exec_command").tokens, 250);
845        assert_eq!(share(&v, "apply_patch").tokens, 750);
846        assert_eq!(share(&v, "fs").tokens, 1_000);
847        assert_eq!(v.iter().map(|s| s.calls).sum::<u64>(), 3);
848        assert_eq!(v.iter().map(|s| s.tokens).sum::<u64>(), 3_100);
849        assert!((v.iter().map(|s| s.cost_usd).sum::<f64>() - 4_100e-6).abs() < 1e-12);
850    }
851
852    #[test]
853    fn context_ledger_weighted_rounding_handles_zero_and_equal_weights() {
854        for (weights, tokens, expected) in [
855            ([0, 0, 0], 2, [1, 1, 0]),
856            ([1, 1, 1], 2, [1, 1, 0]), // Caller-resolved missing weights.
857            ([0, 1, 3], 1, [0, 0, 1]),
858            ([0, 1, 3], 2, [0, 1, 1]),
859            ([0, 1, 3], 0, [0, 0, 0]),
860        ] {
861            let mut l = ContextLedger::default();
862            l.response(&usage(100, 0), &cost(100));
863            l.result_weighted("wrapper", weights.into_iter().enumerate().map(|(i, w)| (ContextOrigin::Tool, i.to_string(), w)).collect());
864            l.response(&usage(100 + tokens, 0), &cost(100 + tokens));
865            let v = l.sources();
866            for (i, tokens) in expected.into_iter().enumerate() {
867                let source = share(&v, &i.to_string());
868                assert_eq!((source.calls, source.tokens), (1, tokens), "weights: {weights:?}");
869            }
870            assert_eq!(v.iter().map(|s| s.tokens).sum::<u64>(), 100 + tokens);
871            assert!((v.iter().map(|s| s.cost_usd).sum::<f64>() - (200 + tokens) as f64 / 1e6).abs() < 1e-12);
872        }
873    }
874
875    #[test]
876    fn context_ledger_weighted_empty_groups_do_not_lose_tokens() {
877        let mut l = ContextLedger::default();
878        l.result_weighted("empty", Vec::new());
879        l.response(&usage(100, 0), &cost(100));
880        assert_eq!(share(&l.sources(), ContextLedger::OTHER).tokens, 100);
881        l.result_weighted("empty", Vec::new());
882        l.result("direct", ContextOrigin::Tool, "Read");
883        l.response(&usage(200, 0), &cost(200));
884        let v = l.sources();
885        assert_eq!((share(&v, "Read").calls, share(&v, "Read").tokens), (1, 100));
886        assert_eq!(v.iter().map(|s| s.tokens).sum::<u64>(), 200);
887    }
888
889    #[test]
890    fn context_ledger_weights_aggregate_repeated_names_and_survive_retagging() {
891        let mut l = ContextLedger::default();
892        l.response(&usage(100, 0), &cost(100));
893        let sources = vec![(ContextOrigin::Tool, "Read".into(), 1), (ContextOrigin::Tool, "Read".into(), 3)];
894        l.result_weighted("first", sources.clone());
895        l.response(&usage(140, 0), &cost(140));
896        assert_eq!((share(&l.sources(), "Read").calls, share(&l.sources(), "Read").tokens), (2, 40));
897        l.result_weighted("second", sources);
898        l.retag("second", ContextOrigin::Mcp, "fs");
899        assert_eq!(l.pending[0].1, vec![(ContextOrigin::Mcp, "fs".into(), 1), (ContextOrigin::Mcp, "fs".into(), 3)]);
900        l.response(&usage(180, 0), &cost(180));
901        let v = l.sources();
902        assert_eq!((share(&v, "fs").calls, share(&v, "fs").tokens), (2, 40));
903        assert_eq!(share(&v, "fs").origin, ContextOrigin::Mcp);
904        assert_eq!((share(&v, "Read").calls, share(&v, "Read").tokens), (2, 40));
905        assert_eq!(v.iter().map(|s| s.tokens).sum::<u64>(), 180);
906    }
907
908    #[test]
909    fn context_ledger_weights_use_wide_products_and_sums() {
910        let mut l = ContextLedger::default();
911        l.result_weighted("wrapper", vec![(ContextOrigin::Tool, "a".into(), u64::MAX), (ContextOrigin::Tool, "b".into(), u64::MAX - 1)]);
912        l.response(&usage(u64::MAX, 0), &cost(u64::MAX));
913        let v = l.sources();
914        assert_eq!(share(&v, "a").tokens, u64::MAX / 2 + 1);
915        assert_eq!(share(&v, "b").tokens, u64::MAX / 2);
916        assert_eq!(v.iter().map(|s| s.tokens).sum::<u64>(), u64::MAX);
917        assert_eq!(v.iter().map(|s| s.calls).sum::<u64>(), 2);
918    }
919
920    #[test]
921    fn context_ledger_weighted_costs_conserve_totals_across_rereads_and_compaction() {
922        let mut l = ContextLedger::default();
923        l.response(&usage(1_000, 100), &cost(1_000));
924        l.result_weighted("wrapper", vec![(ContextOrigin::Tool, "exec_command".into(), 1), (ContextOrigin::Tool, "apply_patch".into(), 3)]);
925        l.response(&usage(1_500, 50), &cost(1_500));
926        let v = l.sources();
927        assert_eq!(share(&v, "exec_command").tokens, 100);
928        assert_eq!(share(&v, "apply_patch").tokens, 300);
929        l.response(&usage(1_550, 10), &cost(1_550));
930        let v = l.sources();
931        assert!((share(&v, "exec_command").cost_usd - 200e-6).abs() < 1e-12);
932        assert!((share(&v, "apply_patch").cost_usd - 600e-6).abs() < 1e-12);
933        assert!((v.iter().map(|s| s.cost_usd).sum::<f64>() - 4_050e-6).abs() < 1e-12);
934
935        l.response(&usage(500, 10), &cost(500)); // Implicit compaction.
936        l.response(&usage(510, 0), &cost(510));
937        l.compacted();
938        l.response(&usage(700, 0), &cost(700));
939        let v = l.sources();
940        assert!((share(&v, "exec_command").cost_usd - 200e-6).abs() < 1e-12);
941        assert!((share(&v, "apply_patch").cost_usd - 600e-6).abs() < 1e-12);
942        assert_eq!(v.iter().map(|s| s.calls).sum::<u64>(), 2);
943        assert_eq!(v.iter().map(|s| s.tokens).sum::<u64>(), 2_760);
944        assert!((v.iter().map(|s| s.cost_usd).sum::<f64>() - 5_760e-6).abs() < 1e-12);
945    }
946
947    #[test]
948    fn context_ledger_takes_a_shrink_off_other_and_a_halving_as_compaction() {
949        let mut l = ContextLedger::default();
950        l.response(&usage(10_000, 2_000), &cost(10_000));
951        l.result("a", ContextOrigin::Tool, "Bash");
952        l.response(&usage(12_500, 3_000), &cost(12_500)); // Bash gets 500
953        // Thinking dropped at the new turn: 1_000 smaller. Bash keeps its 500;
954        // Other's live share takes the shrink and no source grows.
955        l.response(&usage(11_500, 10), &cost(11_500));
956        let v = l.sources();
957        assert_eq!(share(&v, "Bash").tokens, 500);
958        assert_eq!(share(&v, ContextLedger::OTHER).tokens, 12_000);
959        // Bash was read by the two responses since it arrived: 500 * 2.
960        assert!((share(&v, "Bash").cost_usd - 1_000e-6).abs() < 1e-12);
961        // Compaction halves the prompt: the old context is gone, the summary
962        // is Other, and Bash is no longer charged.
963        l.response(&usage(3_000, 10), &cost(3_000));
964        let v = l.sources();
965        assert!((share(&v, "Bash").cost_usd - 1_000e-6).abs() < 1e-12, "not charged after compaction");
966        assert_eq!(share(&v, ContextLedger::OTHER).tokens, 15_000);
967        // A harness that says so resets the same way, mid-growth.
968        l.compacted();
969        l.response(&usage(4_000, 10), &cost(4_000));
970        let v = l.sources();
971        assert_eq!(share(&v, ContextLedger::OTHER).tokens, 19_000);
972        assert!((share(&v, "Bash").cost_usd - 1_000e-6).abs() < 1e-12);
973    }
974
975    #[test]
976    fn context_ledger_retags_pending_results_and_merges() {
977        let mut l = ContextLedger::default();
978        l.response(&usage(100, 0), &cost(100));
979        l.result("c1", ContextOrigin::Tool, "fetch");
980        l.retag("c1", ContextOrigin::Mcp, "apps");
981        l.retag("zzz", ContextOrigin::Mcp, "nope");
982        // No usage record: not a response, nothing filed.
983        l.response(&TokenUsage::default(), &CostBreakdown::default());
984        l.response(&usage(300, 0), &cost(300));
985        let v = l.sources();
986        assert_eq!(v.len(), 2);
987        assert_eq!((share(&v, "apps").origin, share(&v, "apps").tokens), (ContextOrigin::Mcp, 200));
988        assert!(v.iter().all(|s| s.name != "fetch"));
989
990        let mut sub = ContextLedger::default();
991        sub.response(&usage(50, 0), &cost(50));
992        sub.result("x", ContextOrigin::Mcp, "apps");
993        sub.response(&usage(70, 0), &cost(70));
994        l.merge(&sub);
995        let v = l.sources();
996        assert_eq!(share(&v, "apps").tokens, 220);
997        assert_eq!(share(&v, "apps").calls, 2);
998        assert_eq!(share(&v, ContextLedger::OTHER).tokens, 150);
999        assert!(!l.is_empty() && ContextLedger::default().is_empty());
1000    }
1001
1002    #[test]
1003    fn parses_timestamps() {
1004        let t = parse_rfc3339_utc("1970-01-02T00:00:00.000Z").unwrap();
1005        assert_eq!(t, SystemTime::UNIX_EPOCH + Duration::from_secs(86_400));
1006        let t = parse_rfc3339_utc("2026-09-03T07:15:34.5Z").unwrap();
1007        let secs = t.duration_since(SystemTime::UNIX_EPOCH).unwrap();
1008        assert_eq!(secs.as_secs(), 1_788_419_734);
1009        assert_eq!(secs.subsec_millis(), 500);
1010        assert!(parse_rfc3339_utc("nope").is_none());
1011    }
1012}
1013
1014#[cfg(test)]
1015mod rfc3339_tests {
1016    use super::parse_rfc3339_utc;
1017    use std::time::{Duration, UNIX_EPOCH};
1018
1019    fn secs(s: &str) -> u64 {
1020        parse_rfc3339_utc(s).unwrap().duration_since(UNIX_EPOCH).unwrap().as_secs()
1021    }
1022
1023    /// Debt #5: the same instant written with an offset parses to the same
1024    /// time as its `Z` form, and the fraction and the sign survive.
1025    #[test]
1026    fn offsets_are_folded_into_utc() {
1027        let z = secs("2026-09-03T07:15:34Z");
1028        assert_eq!(secs("2026-09-03T08:15:34+01:00"), z);
1029        assert_eq!(secs("2026-09-03T00:15:34-07:00"), z);
1030        assert_eq!(secs("2026-09-03T08:15:34+0100"), z, "no colon");
1031        assert_eq!(secs("2026-09-03T07:15:34+00:00"), z);
1032        assert_eq!(secs("2026-09-03T12:45:34+05:30"), z, "half-hour zone");
1033        let ms = parse_rfc3339_utc("2026-09-03T08:15:34.250+01:00").unwrap();
1034        assert_eq!(ms, UNIX_EPOCH + Duration::new(z, 250_000_000));
1035        // A date's own hyphens are not mistaken for a zone sign.
1036        assert_eq!(secs("2026-09-03T07:15:34.5Z"), z);
1037        // Not timestamps.
1038        assert!(parse_rfc3339_utc("2026-09-03T07:15:34").is_none(), "no zone at all");
1039        assert!(parse_rfc3339_utc("2026-09-03T07:15:34+1").is_none());
1040        assert!(parse_rfc3339_utc("garbage").is_none());
1041    }
1042}