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;
8
9use crate::model::{Activity, CostBreakdown, Harness, SpanKind, TokenUsage, ToolSpan};
10use std::collections::VecDeque;
11use std::path::{Path, PathBuf};
12use std::time::SystemTime;
13
14/// Evidence that the parser still understands the file it is reading.
15///
16/// Every field is read with a fallback to zero, which is the right behaviour
17/// for a genuinely absent field and the wrong behaviour for a renamed one: a
18/// harness that renames `usage` next week would show a user 0 tokens and $0.00
19/// with no error at all. So count the usage records seen and how many of them
20/// yielded nothing. Records present and all of them empty is not a quiet
21/// session, it is a parser that has fallen behind the format.
22#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
23pub struct ParseHealth {
24    /// Model responses seen. Each of these should account for some tokens.
25    pub billable_messages: u64,
26    /// Usage records found on them. Zero of these, with messages present, means
27    /// the record itself moved or was renamed.
28    pub usage_records: u64,
29    /// Records found but yielding nothing, which is what a renamed field inside
30    /// an intact record looks like.
31    pub empty_usage_records: u64,
32}
33
34impl ParseHealth {
35    /// Enough responses to accuse the parser rather than the session. A couple
36    /// of odd messages must not raise the alarm.
37    const MIN_EVIDENCE: u64 = 3;
38
39    /// The session did work that must have cost tokens, and we read none.
40    ///
41    /// Covers both ways a format change reaches us: the usage record moving or
42    /// being renamed, so we never find one, and the fields inside it being
43    /// renamed, so we find records that read as empty. Neither raises an error
44    /// on its own, because every field falls back to zero.
45    pub fn fields_unrecognised(&self) -> bool {
46        self.billable_messages >= Self::MIN_EVIDENCE && self.usage_records == self.empty_usage_records
47    }
48}
49
50#[derive(Debug, Clone, Default)]
51pub struct SessionSummary {
52    pub harness: Option<Harness>,
53    pub session_id: Option<String>,
54    pub cwd: Option<PathBuf>,
55    pub model: Option<String>,
56    pub harness_version: Option<String>,
57    pub usage: TokenUsage,
58    pub cost_usd: f64,
59    /// `cost_usd` by kind of token. See `Agent::cost_breakdown`.
60    pub cost_breakdown: CostBreakdown,
61    pub unpriced_tokens: u64,
62    pub turns: u64,
63    pub subagent_turns: u64,
64    pub tool_calls: u64,
65    /// See `Agent::web_searches`.
66    pub web_searches: u64,
67    pub spans: SpanLog,
68    pub health: ParseHealth,
69    pub activity: Activity,
70    pub started_at: Option<SystemTime>,
71    pub last_activity: Option<SystemTime>,
72}
73
74/// Spans kept per session by the live tracker. A screenful of waterfall is
75/// a few dozen rows; the rest is history nobody scrolls to in a live view, and
76/// every span costs a clone on each refresh. An export wants the whole session
77/// and uses `SpanLog::unbounded` in a separate pass; see `SpanRetention`.
78///
79/// Sized for roughly a hundred tool calls: each model response also adds an
80/// inference span and each human prompt a turn span.
81pub const MAX_SPANS: usize = 256;
82
83/// A bounded, in-order log of tool spans, built by pairing a harness's
84/// "call started" and "call finished" records by call id.
85///
86/// Records arrive interleaved and out of order (agents run tools in parallel),
87/// so a span is closed by searching back for the still-open span with that id
88/// rather than assuming the most recent one.
89#[derive(Debug, Clone)]
90pub struct SpanLog {
91    spans: VecDeque<ToolSpan>,
92    cap: usize,
93}
94
95impl Default for SpanLog {
96    fn default() -> Self {
97        SpanLog { spans: VecDeque::new(), cap: MAX_SPANS }
98    }
99}
100
101impl SpanLog {
102    /// A log that keeps every span. For a one-shot pass over a whole
103    /// transcript, never for the live tracker, where the memory and the clone
104    /// per refresh would grow with the session.
105    pub fn unbounded() -> Self {
106        SpanLog { spans: VecDeque::new(), cap: usize::MAX }
107    }
108
109    /// Record the start of a tool call. Ignored when that id is already open,
110    /// so a transcript line replayed by the harness does not double-count.
111    pub fn open(&mut self, id: String, name: String, at: SystemTime, sidechain: bool) {
112        self.open_kind(id, name, at, sidechain, SpanKind::Tool);
113    }
114
115    /// `open`, for any kind of span.
116    pub fn open_kind(&mut self, id: String, name: String, at: SystemTime, sidechain: bool, kind: SpanKind) {
117        if id.is_empty() || self.spans.iter().any(|s| s.is_open() && s.id == id) {
118            return;
119        }
120        if self.spans.len() >= self.cap {
121            self.spans.pop_front();
122        }
123        self.spans.push_back(ToolSpan { id, name, started_at: at, duration_ms: None, sidechain, error: false, kind });
124    }
125
126    /// Move the end of the newest span with this id to `at`, open or not. An
127    /// inference span grows as the response streams in, one content block
128    /// per line, and its end is wherever the last block landed.
129    pub fn end_at(&mut self, id: &str, at: SystemTime) {
130        let Some(s) = self.spans.iter_mut().rev().find(|s| s.id == id) else { return };
131        s.duration_ms = Some(at.duration_since(s.started_at).map(|d| d.as_millis() as u64).unwrap_or(0));
132    }
133
134    /// The newest open span of this kind, if any.
135    pub fn open_of_kind(&self, kind: SpanKind) -> Option<&ToolSpan> {
136        self.spans.iter().rev().find(|s| s.is_open() && s.kind == kind)
137    }
138
139    /// Remove the newest span with this id if it is still open. For a span
140    /// that turned out not to be one: an inference that never produced a
141    /// reply because the user interrupted or submitted again.
142    pub fn discard_open(&mut self, id: &str) {
143        if let Some(i) = self.spans.iter().rposition(|s| s.is_open() && s.id == id) {
144            self.spans.remove(i);
145        }
146    }
147
148    /// Close the open call with this id. A result whose call scrolled out of
149    /// the window, or that we never saw start, is dropped.
150    pub fn close(&mut self, id: &str, at: SystemTime, error: bool) {
151        let Some(s) = self.spans.iter_mut().rev().find(|s| s.is_open() && s.id == id) else { return };
152        s.duration_ms = Some(at.duration_since(s.started_at).map(|d| d.as_millis() as u64).unwrap_or(0));
153        s.error = error;
154    }
155
156    pub fn len(&self) -> usize {
157        self.spans.len()
158    }
159
160    pub fn is_empty(&self) -> bool {
161        self.spans.is_empty()
162    }
163
164    /// Oldest first. Double-ended so callers can take the newest spans without
165    /// collecting the whole log first, which the UI and the golden tests both do.
166    pub fn iter(&self) -> impl DoubleEndedIterator<Item = &ToolSpan> + ExactSizeIterator {
167        self.spans.iter()
168    }
169
170    pub fn to_vec(&self) -> Vec<ToolSpan> {
171        self.spans.iter().cloned().collect()
172    }
173
174    /// One log from several, ordered by start time, keeping the newest `cap`.
175    /// A parent's spans and its subagents' spans interleave in wall-clock
176    /// order, which is what a waterfall wants.
177    pub fn merged<'a>(logs: impl IntoIterator<Item = &'a SpanLog>, cap: usize) -> SpanLog {
178        let mut spans: Vec<ToolSpan> = logs.into_iter().flat_map(|l| l.spans.iter().cloned()).collect();
179        spans.sort_by_key(|s| s.started_at);
180        if spans.len() > cap {
181            spans.drain(..spans.len() - cap);
182        }
183        SpanLog { spans: spans.into(), cap }
184    }
185
186    pub fn cap(&self) -> usize {
187        self.cap
188    }
189}
190
191/// How many of a session's tool spans a tracker keeps.
192#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
193pub enum SpanRetention {
194    /// The newest `MAX_SPANS`, enough for the live waterfall. The default.
195    #[default]
196    Recent,
197    /// Every span in the transcript, for a trace export. Memory grows with
198    /// the session, so this is for a single pass, not a tracker kept across
199    /// refreshes.
200    All,
201}
202
203impl SpanRetention {
204    pub(crate) fn log(self) -> SpanLog {
205        match self {
206            SpanRetention::Recent => SpanLog::default(),
207            SpanRetention::All => SpanLog::unbounded(),
208        }
209    }
210}
211
212pub trait SessionTracker {
213    /// Ingest whatever was appended since the last call. Returns true when
214    /// there is still unread data (the byte budget was exhausted).
215    fn refresh(&mut self) -> anyhow::Result<bool>;
216    fn summary(&self) -> &SessionSummary;
217    fn path(&self) -> &Path;
218
219    /// Ingest the whole file, however many refreshes that takes. For a
220    /// one-shot read such as an export; the live collector spreads a large
221    /// transcript over several ticks instead.
222    fn refresh_all(&mut self) -> anyhow::Result<()> {
223        while self.refresh()? {}
224        Ok(())
225    }
226}
227
228/// Which harness wrote a transcript, judged from its first few lines. Codex
229/// opens every rollout with a `session_meta` record; Claude Code lines carry
230/// `sessionId`. Anything else is not a transcript agent-top reads.
231pub fn detect(path: &Path) -> Option<Harness> {
232    use std::io::{BufRead, BufReader};
233    let f = std::fs::File::open(path).ok()?;
234    for line in BufReader::new(f).lines().map_while(Result::ok).take(5) {
235        let Ok(v) = serde_json::from_str::<serde_json::Value>(&line) else { continue };
236        if v.get("type").and_then(serde_json::Value::as_str) == Some("session_meta") {
237            return Some(Harness::Codex);
238        }
239        if v.get("sessionId").is_some() || v.get("parentUuid").is_some() {
240            return Some(Harness::Claude);
241        }
242    }
243    None
244}
245
246/// A tracker for a transcript whose harness is already known.
247pub fn open_transcript(path: &Path, harness: Harness, spans: SpanRetention) -> Box<dyn SessionTracker> {
248    match harness {
249        Harness::Codex => Box::new(codex::CodexTranscript::new(path).with_spans(spans)),
250        _ => Box::new(claude::ClaudeTranscript::new(path).with_spans(spans)),
251    }
252}
253
254/// Bytes ingested per tracker per refresh. Keeps a cold start on a 100 MB
255/// transcript from freezing the first frame; the rest streams in on later ticks.
256pub const REFRESH_BUDGET_BYTES: usize = 8 * 1024 * 1024;
257
258/// Parse an RFC 3339 timestamp like `2026-09-03T07:15:34.123Z` into SystemTime
259/// without pulling in a date crate. Only the UTC `Z` form is handled, which is
260/// what both Claude Code and Codex write.
261pub fn parse_rfc3339_utc(s: &str) -> Option<SystemTime> {
262    let s = s.strip_suffix('Z')?;
263    let (date, time) = s.split_once('T')?;
264    let mut d = date.split('-');
265    let (y, mo, da) = (d.next()?.parse::<i64>().ok()?, d.next()?.parse::<u32>().ok()?, d.next()?.parse::<u32>().ok()?);
266    let mut t = time.split(':');
267    let (h, mi) = (t.next()?.parse::<u64>().ok()?, t.next()?.parse::<u64>().ok()?);
268    let sec_str = t.next()?;
269    let (sec, frac) = match sec_str.split_once('.') {
270        Some((s, f)) => (s.parse::<u64>().ok()?, f),
271        None => (sec_str.parse::<u64>().ok()?, ""),
272    };
273    let nanos: u32 = if frac.is_empty() {
274        0
275    } else {
276        let mut f = frac.to_string();
277        f.truncate(9);
278        while f.len() < 9 {
279            f.push('0');
280        }
281        f.parse().ok()?
282    };
283    let days = days_from_civil(y, mo, da);
284    let secs = days * 86_400 + (h * 3600 + mi * 60 + sec) as i64;
285    if secs < 0 {
286        return None;
287    }
288    Some(SystemTime::UNIX_EPOCH + std::time::Duration::new(secs as u64, nanos))
289}
290
291// Howard Hinnant's days-from-civil algorithm.
292fn days_from_civil(y: i64, m: u32, d: u32) -> i64 {
293    let y = if m <= 2 { y - 1 } else { y };
294    let era = if y >= 0 { y } else { y - 399 } / 400;
295    let yoe = y - era * 400;
296    let mp = (m as i64 + 9) % 12;
297    let doy = (153 * mp + 2) / 5 + d as i64 - 1;
298    let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
299    era * 146_097 + doe - 719_468
300}
301
302#[cfg(test)]
303mod tests {
304    use super::*;
305    use std::time::Duration;
306
307    fn at(secs: u64) -> SystemTime {
308        SystemTime::UNIX_EPOCH + Duration::from_secs(secs)
309    }
310
311    #[test]
312    fn pairs_spans_by_id_out_of_order() {
313        let mut log = SpanLog::default();
314        log.open("a".into(), "Bash".into(), at(10), false);
315        log.open("b".into(), "Read".into(), at(11), true);
316        // Replay of the same start line must not open a second span.
317        log.open("a".into(), "Bash".into(), at(10), false);
318        // Results come back in the other order.
319        log.close("b", at(12), false);
320        log.close("a", at(14), true);
321        // A result with no matching call is ignored.
322        log.close("zzz", at(15), false);
323        let v = log.to_vec();
324        assert_eq!(v.len(), 2);
325        assert_eq!(v[0].name, "Bash");
326        assert_eq!(v[0].duration_ms, Some(4_000));
327        assert!(v[0].error);
328        assert_eq!(v[1].duration_ms, Some(1_000));
329        assert!(v[1].sidechain);
330        assert!(!v[1].error);
331    }
332
333    #[test]
334    fn keeps_the_newest_spans_and_reports_open_ones() {
335        let mut log = SpanLog::default();
336        for i in 0..(MAX_SPANS + 10) {
337            log.open(format!("id{i}"), "T".into(), at(i as u64), false);
338            log.close(&format!("id{i}"), at(i as u64), false);
339        }
340        assert_eq!(log.len(), MAX_SPANS);
341        assert_eq!(log.iter().next().unwrap().id, "id10");
342        log.open("live".into(), "Bash".into(), at(500), false);
343        let last = log.to_vec().pop().unwrap();
344        assert!(last.is_open());
345        assert_eq!(last.elapsed_ms(at(503)), 3_000);
346    }
347
348    #[test]
349    fn end_at_moves_the_end_of_any_kind_of_span() {
350        let mut log = SpanLog::default();
351        log.open_kind("inference:1".into(), "inference".into(), at(10), false, SpanKind::Inference);
352        assert!(log.open_of_kind(SpanKind::Inference).is_some());
353        assert!(log.open_of_kind(SpanKind::Turn).is_none());
354        // The response streams in over three lines; the span ends at the last one.
355        log.end_at("inference:1", at(11));
356        log.end_at("inference:1", at(13));
357        log.end_at("nope", at(99));
358        let v = log.to_vec();
359        assert_eq!(v[0].duration_ms, Some(3_000));
360        assert_eq!(v[0].kind, SpanKind::Inference);
361        assert!(log.open_of_kind(SpanKind::Inference).is_none());
362        // Discarding only removes open spans; the ended one stays.
363        log.discard_open("inference:1");
364        assert_eq!(log.len(), 1);
365        log.open_kind("inference:2".into(), "inference".into(), at(20), false, SpanKind::Inference);
366        log.discard_open("inference:2");
367        assert_eq!(log.len(), 1);
368    }
369
370    #[test]
371    fn unbounded_log_keeps_everything() {
372        let mut log = SpanLog::unbounded();
373        for i in 0..(MAX_SPANS * 3) {
374            log.open(format!("id{i}"), "T".into(), at(i as u64), false);
375            log.close(&format!("id{i}"), at(i as u64 + 1), false);
376        }
377        assert_eq!(log.len(), MAX_SPANS * 3);
378        assert_eq!(log.iter().next().unwrap().id, "id0");
379        assert_eq!(SpanRetention::default(), SpanRetention::Recent);
380    }
381
382    #[test]
383    fn detects_the_harness_from_the_first_lines() {
384        let dir = std::env::temp_dir().join(format!("agent-top-detect-{}", std::process::id()));
385        std::fs::create_dir_all(&dir).unwrap();
386        let codex = dir.join("rollout.jsonl");
387        std::fs::write(&codex, "{\"type\":\"session_meta\",\"payload\":{\"id\":\"x\"}}\n").unwrap();
388        let claude = dir.join("s.jsonl");
389        // A summary line first, as Claude Code writes on resume, then a real one.
390        std::fs::write(&claude, "{\"type\":\"summary\",\"leafUuid\":\"u\"}\n{\"type\":\"user\",\"sessionId\":\"abc\"}\n").unwrap();
391        let other = dir.join("other.jsonl");
392        std::fs::write(&other, "{\"hello\":1}\nnot json\n").unwrap();
393        assert_eq!(detect(&codex), Some(Harness::Codex));
394        assert_eq!(detect(&claude), Some(Harness::Claude));
395        assert_eq!(detect(&other), None);
396        assert_eq!(detect(&dir.join("missing.jsonl")), None);
397        let _ = std::fs::remove_dir_all(&dir);
398    }
399
400    #[test]
401    fn accuses_the_parser_only_with_enough_evidence() {
402        // Healthy: records found on the messages, tokens read from them.
403        let h = ParseHealth { billable_messages: 40, usage_records: 40, empty_usage_records: 0 };
404        assert!(!h.fields_unrecognised());
405        // One odd message among many is a message, not a format change.
406        let h = ParseHealth { billable_messages: 40, usage_records: 40, empty_usage_records: 39 };
407        assert!(!h.fields_unrecognised());
408        // Fields inside the record renamed: records found, all of them empty.
409        let h = ParseHealth { billable_messages: 40, usage_records: 40, empty_usage_records: 40 };
410        assert!(h.fields_unrecognised());
411        // The record itself renamed or moved: messages, but no records at all.
412        // This is the case a naive check misses, because there is nothing to count.
413        let h = ParseHealth { billable_messages: 40, usage_records: 0, empty_usage_records: 0 };
414        assert!(h.fields_unrecognised());
415        // Too early to tell: a session that has barely started.
416        let h = ParseHealth { billable_messages: 2, usage_records: 0, empty_usage_records: 0 };
417        assert!(!h.fields_unrecognised());
418        // Nothing parsed at all is silence, not evidence.
419        assert!(!ParseHealth::default().fields_unrecognised());
420    }
421
422    #[test]
423    fn parses_timestamps() {
424        let t = parse_rfc3339_utc("1970-01-02T00:00:00.000Z").unwrap();
425        assert_eq!(t, SystemTime::UNIX_EPOCH + Duration::from_secs(86_400));
426        let t = parse_rfc3339_utc("2026-09-03T07:15:34.5Z").unwrap();
427        let secs = t.duration_since(SystemTime::UNIX_EPOCH).unwrap();
428        assert_eq!(secs.as_secs(), 1_788_419_734);
429        assert_eq!(secs.subsec_millis(), 500);
430        assert!(parse_rfc3339_utc("nope").is_none());
431    }
432}