1pub 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#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
27pub struct ParseHealth {
28 pub billable_messages: u64,
30 pub usage_records: u64,
33 pub empty_usage_records: u64,
36}
37
38impl ParseHealth {
39 const MIN_EVIDENCE: u64 = 3;
42
43 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 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 pub cost_breakdown: CostBreakdown,
67 pub price_source: Option<crate::model::PriceSource>,
69 pub unpriced_tokens: u64,
70 pub turns: u64,
71 pub subagent_turns: u64,
72 pub folds_child_usage: bool,
76 pub tool_calls: u64,
77 pub tool_calls_lower_bound: bool,
79 pub web_searches: u64,
81 pub spans: SpanLog,
82 pub mcp: BTreeMap<String, McpUsage>,
84 pub context: ContextLedger,
87 pub health: ParseHealth,
88 pub activity: Activity,
89 pub started_at: Option<SystemTime>,
90 pub last_activity: Option<SystemTime>,
91 pub rate_limit: Option<crate::model::RateLimit>,
93 pub harness_cost: Option<crate::model::HarnessCost>,
95}
96
97#[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#[derive(Debug, Clone, Default)]
140pub struct ContextLedger {
141 shares: BTreeMap<ContextKey, ContextShare>,
142 live: BTreeMap<ContextKey, u64>,
144 pending: Vec<(String, ContextWeights)>,
147 prev: Option<(u64, u64)>,
149}
150
151type ContextKey = (ContextOrigin, String);
152
153pub(super) type ContextWeights = Vec<(ContextOrigin, String, u64)>;
155
156#[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 pub const OTHER: &str = "other";
167
168 fn other() -> ContextKey {
169 (ContextOrigin::Other, Self::OTHER.to_string())
170 }
171
172 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 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 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 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 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 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 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 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 pub fn compacted(&mut self) {
286 self.live.clear();
287 self.prev = None;
288 }
289
290 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 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
323pub 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
336pub const MAX_SPANS: usize = 256;
344
345#[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 pub fn unbounded() -> Self {
368 SpanLog { spans: VecDeque::new(), cap: usize::MAX }
369 }
370
371 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 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 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 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 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 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 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 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
455pub enum SpanRetention {
456 #[default]
458 Recent,
459 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 fn refresh(&mut self) -> anyhow::Result<bool>;
478 fn summary(&self) -> &SessionSummary;
479 fn path(&self) -> &Path;
480
481 fn refresh_all(&mut self) -> anyhow::Result<()> {
485 while self.refresh()? {}
486 Ok(())
487 }
488}
489
490#[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 pub status: Option<String>,
502}
503
504pub struct AttributeContext<'a> {
507 pub cwd: Option<&'a Path>,
508 pub proc_start: SystemTime,
509 pub now: SystemTime,
510 pub attached: &'a HashSet<PathBuf>,
513 pub activity_timeout: Duration,
516}
517
518pub trait HarnessAdapter {
525 fn harness(&self) -> Harness;
526
527 fn rescan(&mut self, since: SystemTime);
530
531 fn prepare(&mut self, _roots: &[&ProcNode], _by_pid: &std::collections::HashMap<u32, &RawProc>) {}
536
537 fn hints(&self, _pid: u32) -> Option<RegistryHints> {
539 None
540 }
541
542 fn attribute(&self, root: &ProcNode, raw: Option<&RawProc>, ctx: &AttributeContext) -> (Vec<PathBuf>, Attribution);
546
547 fn unowned(&self, attached: &HashSet<PathBuf>) -> Vec<PathBuf>;
549
550 fn open(&self, path: &Path, spans: SpanRetention) -> Box<dyn SessionTracker>;
552
553 fn detect(&self, path: &Path) -> bool;
555
556 fn transcripts(&self) -> Vec<(String, PathBuf)>;
559
560 fn stamp(&self, path: &Path) -> Option<SourceStamp> {
565 SourceStamp::of_file(path)
566 }
567}
568
569#[derive(Debug, Clone, Copy, PartialEq, Eq)]
572pub struct SourceStamp {
573 pub size: u64,
574 pub mtime_ms: i64,
575}
576
577impl SourceStamp {
578 pub fn of_file(path: &Path) -> Option<SourceStamp> {
579 let md = std::fs::metadata(path).ok().filter(|m| m.is_file())?;
580 let mtime_ms = md.modified().ok()?.duration_since(SystemTime::UNIX_EPOCH).ok()?.as_millis() as i64;
581 Some(SourceStamp { size: md.len(), mtime_ms })
582 }
583
584 pub fn with(self, other: SourceStamp) -> SourceStamp {
586 SourceStamp { size: self.size + other.size, mtime_ms: self.mtime_ms.max(other.mtime_ms) }
587 }
588}
589
590pub fn adapters() -> Vec<Box<dyn HarnessAdapter>> {
594 vec![
595 Box::new(codex::CodexAdapter::default()),
596 Box::new(gemini::GeminiAdapter::default()),
597 Box::new(opencode::OpenCodeAdapter::default()),
598 Box::new(kodelet::KodeletAdapter::default()),
599 Box::new(claude::ClaudeAdapter::default()),
600 ]
601}
602
603pub fn adapter_for(harness: Harness) -> Option<Box<dyn HarnessAdapter>> {
605 adapters().into_iter().find(|a| a.harness() == harness)
606}
607
608pub fn detect(path: &Path) -> Option<Harness> {
611 adapters().iter().find(|a| a.detect(path)).map(|a| a.harness())
612}
613
614pub fn open_transcript(path: &Path, harness: Harness, spans: SpanRetention) -> Option<Box<dyn SessionTracker>> {
617 adapter_for(harness).map(|a| a.open(path, spans))
618}
619
620pub(crate) fn head_lines(path: &Path) -> Vec<serde_json::Value> {
622 use std::io::{BufRead, BufReader};
623 let Ok(f) = std::fs::File::open(path) else { return Vec::new() };
624 BufReader::new(f).lines().map_while(Result::ok).take(5).filter_map(|l| serde_json::from_str(&l).ok()).collect()
625}
626
627pub const REFRESH_BUDGET_BYTES: usize = 8 * 1024 * 1024;
630
631pub fn parse_rfc3339_utc(s: &str) -> Option<SystemTime> {
637 let s = s.trim();
638 let (s, offset_secs) = match s.strip_suffix(['Z', 'z']) {
640 Some(rest) => (rest, 0i64),
641 None => {
642 let t_pos = s.find('T')?;
643 let sign_pos = s[t_pos..].rfind(['+', '-'])? + t_pos;
644 let (rest, zone) = s.split_at(sign_pos);
645 let sign = if zone.starts_with('-') { -1 } else { 1 };
646 let digits: String = zone[1..].chars().filter(|c| c.is_ascii_digit()).collect();
647 if digits.len() != 4 {
648 return None;
649 }
650 let oh = digits[..2].parse::<i64>().ok()?;
651 let om = digits[2..].parse::<i64>().ok()?;
652 (rest, sign * (oh * 3600 + om * 60))
653 }
654 };
655 let (date, time) = s.split_once('T')?;
656 let mut d = date.split('-');
657 let (y, mo, da) = (d.next()?.parse::<i64>().ok()?, d.next()?.parse::<u32>().ok()?, d.next()?.parse::<u32>().ok()?);
658 let mut t = time.split(':');
659 let (h, mi) = (t.next()?.parse::<u64>().ok()?, t.next()?.parse::<u64>().ok()?);
660 let sec_str = t.next()?;
661 let (sec, frac) = match sec_str.split_once('.') {
662 Some((s, f)) => (s.parse::<u64>().ok()?, f),
663 None => (sec_str.parse::<u64>().ok()?, ""),
664 };
665 let nanos: u32 = if frac.is_empty() {
666 0
667 } else {
668 let mut f = frac.to_string();
669 f.truncate(9);
670 while f.len() < 9 {
671 f.push('0');
672 }
673 f.parse().ok()?
674 };
675 let days = days_from_civil(y, mo, da);
676 let secs = days * 86_400 + (h * 3600 + mi * 60 + sec) as i64 - offset_secs;
678 if secs < 0 {
679 return None;
680 }
681 Some(SystemTime::UNIX_EPOCH + std::time::Duration::new(secs as u64, nanos))
682}
683
684fn days_from_civil(y: i64, m: u32, d: u32) -> i64 {
686 let y = if m <= 2 { y - 1 } else { y };
687 let era = if y >= 0 { y } else { y - 399 } / 400;
688 let yoe = y - era * 400;
689 let mp = (m as i64 + 9) % 12;
690 let doy = (153 * mp + 2) / 5 + d as i64 - 1;
691 let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy;
692 era * 146_097 + doe - 719_468
693}
694
695#[cfg(test)]
696mod tests {
697 use super::*;
698 use std::time::Duration;
699
700 fn at(secs: u64) -> SystemTime {
701 SystemTime::UNIX_EPOCH + Duration::from_secs(secs)
702 }
703
704 #[test]
705 fn pairs_spans_by_id_out_of_order() {
706 let mut log = SpanLog::default();
707 log.open("a".into(), "Bash".into(), at(10), false);
708 log.open("b".into(), "Read".into(), at(11), true);
709 log.open("a".into(), "Bash".into(), at(10), false);
711 log.close("b", at(12), false);
713 log.close("a", at(14), true);
714 log.close("zzz", at(15), false);
716 let v = log.to_vec();
717 assert_eq!(v.len(), 2);
718 assert_eq!(v[0].name, "Bash");
719 assert_eq!(v[0].duration_ms, Some(4_000));
720 assert!(v[0].error);
721 assert_eq!(v[1].duration_ms, Some(1_000));
722 assert!(v[1].sidechain);
723 assert!(!v[1].error);
724 }
725
726 #[test]
727 fn keeps_the_newest_spans_and_reports_open_ones() {
728 let mut log = SpanLog::default();
729 for i in 0..(MAX_SPANS + 10) {
730 log.open(format!("id{i}"), "T".into(), at(i as u64), false);
731 log.close(&format!("id{i}"), at(i as u64), false);
732 }
733 assert_eq!(log.len(), MAX_SPANS);
734 assert_eq!(log.iter().next().unwrap().id, "id10");
735 log.open("live".into(), "Bash".into(), at(500), false);
736 let last = log.to_vec().pop().unwrap();
737 assert!(last.is_open());
738 assert_eq!(last.elapsed_ms(at(503)), 3_000);
739 }
740
741 #[test]
742 fn end_at_moves_the_end_of_any_kind_of_span() {
743 let mut log = SpanLog::default();
744 log.open_kind("inference:1".into(), "inference".into(), at(10), false, SpanKind::Inference);
745 assert!(log.open_of_kind(SpanKind::Inference).is_some());
746 assert!(log.open_of_kind(SpanKind::Turn).is_none());
747 log.end_at("inference:1", at(11));
749 log.end_at("inference:1", at(13));
750 log.end_at("nope", at(99));
751 let v = log.to_vec();
752 assert_eq!(v[0].duration_ms, Some(3_000));
753 assert_eq!(v[0].kind, SpanKind::Inference);
754 assert!(log.open_of_kind(SpanKind::Inference).is_none());
755 log.discard_open("inference:1");
757 assert_eq!(log.len(), 1);
758 log.open_kind("inference:2".into(), "inference".into(), at(20), false, SpanKind::Inference);
759 log.discard_open("inference:2");
760 assert_eq!(log.len(), 1);
761 }
762
763 #[test]
764 fn unbounded_log_keeps_everything() {
765 let mut log = SpanLog::unbounded();
766 for i in 0..(MAX_SPANS * 3) {
767 log.open(format!("id{i}"), "T".into(), at(i as u64), false);
768 log.close(&format!("id{i}"), at(i as u64 + 1), false);
769 }
770 assert_eq!(log.len(), MAX_SPANS * 3);
771 assert_eq!(log.iter().next().unwrap().id, "id0");
772 assert_eq!(SpanRetention::default(), SpanRetention::Recent);
773 }
774
775 #[test]
776 fn detects_the_harness_from_the_first_lines() {
777 let dir = std::env::temp_dir().join(format!("agent-top-detect-{}", std::process::id()));
778 std::fs::create_dir_all(&dir).unwrap();
779 let codex = dir.join("rollout.jsonl");
780 std::fs::write(&codex, "{\"type\":\"session_meta\",\"payload\":{\"id\":\"x\"}}\n").unwrap();
781 let claude = dir.join("s.jsonl");
782 std::fs::write(&claude, "{\"type\":\"summary\",\"leafUuid\":\"u\"}\n{\"type\":\"user\",\"sessionId\":\"abc\"}\n").unwrap();
784 let other = dir.join("other.jsonl");
785 std::fs::write(&other, "{\"hello\":1}\nnot json\n").unwrap();
786 assert_eq!(detect(&codex), Some(Harness::Codex));
787 assert_eq!(detect(&claude), Some(Harness::Claude));
788 assert_eq!(detect(&other), None);
789 assert_eq!(detect(&dir.join("missing.jsonl")), None);
790 let _ = std::fs::remove_dir_all(&dir);
791 }
792
793 #[test]
794 fn names_the_server_behind_an_mcp_tool() {
795 assert_eq!(mcp_server_of("mcp__filesystem__read_file"), Some("filesystem"));
796 assert_eq!(mcp_server_of("mcp__chrome-devtools__take_screenshot"), Some("chrome-devtools"));
797 assert_eq!(mcp_server_of("mcp__claude_ai_Gmail__authenticate"), Some("claude_ai_Gmail"));
798 assert_eq!(mcp_server_of("mcp__odd"), Some("odd"));
799 assert_eq!(mcp_server_of("mcp____x"), None);
800 assert_eq!(mcp_server_of("Bash"), None);
801 }
802
803 #[test]
804 fn accuses_the_parser_only_with_enough_evidence() {
805 let h = ParseHealth { billable_messages: 40, usage_records: 40, empty_usage_records: 0 };
807 assert!(!h.fields_unrecognised());
808 let h = ParseHealth { billable_messages: 40, usage_records: 40, empty_usage_records: 39 };
810 assert!(!h.fields_unrecognised());
811 let h = ParseHealth { billable_messages: 40, usage_records: 40, empty_usage_records: 40 };
813 assert!(h.fields_unrecognised());
814 let h = ParseHealth { billable_messages: 40, usage_records: 0, empty_usage_records: 0 };
817 assert!(h.fields_unrecognised());
818 let h = ParseHealth { billable_messages: 2, usage_records: 0, empty_usage_records: 0 };
820 assert!(!h.fields_unrecognised());
821 assert!(!ParseHealth::default().fields_unrecognised());
823 }
824
825 fn usage(prompt: u64, output: u64) -> TokenUsage {
826 TokenUsage { cache_read: prompt, output, ..Default::default() }
827 }
828
829 fn cost(prompt: u64) -> CostBreakdown {
832 CostBreakdown { cache_read: prompt as f64 / 1e6, ..Default::default() }
833 }
834
835 fn share<'a>(v: &'a [ContextSource], name: &str) -> &'a ContextSource {
836 v.iter().find(|s| s.name == name).unwrap_or_else(|| panic!("no source {name}"))
837 }
838
839 #[test]
840 fn context_ledger_files_prompt_growth_under_the_results_that_caused_it() {
841 let mut l = ContextLedger::default();
842 l.response(&usage(1_000, 100), &cost(1_000));
844 l.result("a", ContextOrigin::Tool, "Read");
847 l.result("b", ContextOrigin::Mcp, "fs");
848 l.response(&usage(3_300, 50), &cost(3_300));
849 let v = l.sources();
850 assert_eq!(share(&v, "Read").tokens, 1_100);
851 assert_eq!(share(&v, "fs").tokens, 1_100);
852 assert_eq!(share(&v, "fs").origin, ContextOrigin::Mcp);
853 assert_eq!(share(&v, "fs").calls, 1);
854 let other = share(&v, ContextLedger::OTHER);
855 assert_eq!((other.tokens, other.calls), (1_100, 0));
856 assert!((other.cost_usd - 2_100e-6).abs() < 1e-12, "{}", other.cost_usd);
858 assert!((share(&v, "Read").cost_usd - 1_100e-6).abs() < 1e-12);
859 let total: f64 = v.iter().map(|s| s.cost_usd).sum();
861 assert!((total - 4_300e-6).abs() < 1e-12, "{total}");
862 assert_eq!(v[0].tokens, 1_100, "largest first");
863 }
864
865 #[test]
866 fn context_ledger_weights_preserve_equal_shares_between_result_records() {
867 let mut l = ContextLedger::default();
868 l.response(&usage(1_000, 100), &cost(1_000));
869 l.result_weighted("wrapper", vec![(ContextOrigin::Tool, "exec_command".into(), 1), (ContextOrigin::Tool, "apply_patch".into(), 3)]);
870 l.result("direct", ContextOrigin::Mcp, "fs");
871 l.response(&usage(3_100, 0), &cost(3_100));
872 let v = l.sources();
873 assert_eq!(share(&v, "exec_command").tokens, 250);
874 assert_eq!(share(&v, "apply_patch").tokens, 750);
875 assert_eq!(share(&v, "fs").tokens, 1_000);
876 assert_eq!(v.iter().map(|s| s.calls).sum::<u64>(), 3);
877 assert_eq!(v.iter().map(|s| s.tokens).sum::<u64>(), 3_100);
878 assert!((v.iter().map(|s| s.cost_usd).sum::<f64>() - 4_100e-6).abs() < 1e-12);
879 }
880
881 #[test]
882 fn context_ledger_weighted_rounding_handles_zero_and_equal_weights() {
883 for (weights, tokens, expected) in [
884 ([0, 0, 0], 2, [1, 1, 0]),
885 ([1, 1, 1], 2, [1, 1, 0]), ([0, 1, 3], 1, [0, 0, 1]),
887 ([0, 1, 3], 2, [0, 1, 1]),
888 ([0, 1, 3], 0, [0, 0, 0]),
889 ] {
890 let mut l = ContextLedger::default();
891 l.response(&usage(100, 0), &cost(100));
892 l.result_weighted("wrapper", weights.into_iter().enumerate().map(|(i, w)| (ContextOrigin::Tool, i.to_string(), w)).collect());
893 l.response(&usage(100 + tokens, 0), &cost(100 + tokens));
894 let v = l.sources();
895 for (i, tokens) in expected.into_iter().enumerate() {
896 let source = share(&v, &i.to_string());
897 assert_eq!((source.calls, source.tokens), (1, tokens), "weights: {weights:?}");
898 }
899 assert_eq!(v.iter().map(|s| s.tokens).sum::<u64>(), 100 + tokens);
900 assert!((v.iter().map(|s| s.cost_usd).sum::<f64>() - (200 + tokens) as f64 / 1e6).abs() < 1e-12);
901 }
902 }
903
904 #[test]
905 fn context_ledger_weighted_empty_groups_do_not_lose_tokens() {
906 let mut l = ContextLedger::default();
907 l.result_weighted("empty", Vec::new());
908 l.response(&usage(100, 0), &cost(100));
909 assert_eq!(share(&l.sources(), ContextLedger::OTHER).tokens, 100);
910 l.result_weighted("empty", Vec::new());
911 l.result("direct", ContextOrigin::Tool, "Read");
912 l.response(&usage(200, 0), &cost(200));
913 let v = l.sources();
914 assert_eq!((share(&v, "Read").calls, share(&v, "Read").tokens), (1, 100));
915 assert_eq!(v.iter().map(|s| s.tokens).sum::<u64>(), 200);
916 }
917
918 #[test]
919 fn context_ledger_weights_aggregate_repeated_names_and_survive_retagging() {
920 let mut l = ContextLedger::default();
921 l.response(&usage(100, 0), &cost(100));
922 let sources = vec![(ContextOrigin::Tool, "Read".into(), 1), (ContextOrigin::Tool, "Read".into(), 3)];
923 l.result_weighted("first", sources.clone());
924 l.response(&usage(140, 0), &cost(140));
925 assert_eq!((share(&l.sources(), "Read").calls, share(&l.sources(), "Read").tokens), (2, 40));
926 l.result_weighted("second", sources);
927 l.retag("second", ContextOrigin::Mcp, "fs");
928 assert_eq!(l.pending[0].1, vec![(ContextOrigin::Mcp, "fs".into(), 1), (ContextOrigin::Mcp, "fs".into(), 3)]);
929 l.response(&usage(180, 0), &cost(180));
930 let v = l.sources();
931 assert_eq!((share(&v, "fs").calls, share(&v, "fs").tokens), (2, 40));
932 assert_eq!(share(&v, "fs").origin, ContextOrigin::Mcp);
933 assert_eq!((share(&v, "Read").calls, share(&v, "Read").tokens), (2, 40));
934 assert_eq!(v.iter().map(|s| s.tokens).sum::<u64>(), 180);
935 }
936
937 #[test]
938 fn context_ledger_weights_use_wide_products_and_sums() {
939 let mut l = ContextLedger::default();
940 l.result_weighted("wrapper", vec![(ContextOrigin::Tool, "a".into(), u64::MAX), (ContextOrigin::Tool, "b".into(), u64::MAX - 1)]);
941 l.response(&usage(u64::MAX, 0), &cost(u64::MAX));
942 let v = l.sources();
943 assert_eq!(share(&v, "a").tokens, u64::MAX / 2 + 1);
944 assert_eq!(share(&v, "b").tokens, u64::MAX / 2);
945 assert_eq!(v.iter().map(|s| s.tokens).sum::<u64>(), u64::MAX);
946 assert_eq!(v.iter().map(|s| s.calls).sum::<u64>(), 2);
947 }
948
949 #[test]
950 fn context_ledger_weighted_costs_conserve_totals_across_rereads_and_compaction() {
951 let mut l = ContextLedger::default();
952 l.response(&usage(1_000, 100), &cost(1_000));
953 l.result_weighted("wrapper", vec![(ContextOrigin::Tool, "exec_command".into(), 1), (ContextOrigin::Tool, "apply_patch".into(), 3)]);
954 l.response(&usage(1_500, 50), &cost(1_500));
955 let v = l.sources();
956 assert_eq!(share(&v, "exec_command").tokens, 100);
957 assert_eq!(share(&v, "apply_patch").tokens, 300);
958 l.response(&usage(1_550, 10), &cost(1_550));
959 let v = l.sources();
960 assert!((share(&v, "exec_command").cost_usd - 200e-6).abs() < 1e-12);
961 assert!((share(&v, "apply_patch").cost_usd - 600e-6).abs() < 1e-12);
962 assert!((v.iter().map(|s| s.cost_usd).sum::<f64>() - 4_050e-6).abs() < 1e-12);
963
964 l.response(&usage(500, 10), &cost(500)); l.response(&usage(510, 0), &cost(510));
966 l.compacted();
967 l.response(&usage(700, 0), &cost(700));
968 let v = l.sources();
969 assert!((share(&v, "exec_command").cost_usd - 200e-6).abs() < 1e-12);
970 assert!((share(&v, "apply_patch").cost_usd - 600e-6).abs() < 1e-12);
971 assert_eq!(v.iter().map(|s| s.calls).sum::<u64>(), 2);
972 assert_eq!(v.iter().map(|s| s.tokens).sum::<u64>(), 2_760);
973 assert!((v.iter().map(|s| s.cost_usd).sum::<f64>() - 5_760e-6).abs() < 1e-12);
974 }
975
976 #[test]
977 fn context_ledger_takes_a_shrink_off_other_and_a_halving_as_compaction() {
978 let mut l = ContextLedger::default();
979 l.response(&usage(10_000, 2_000), &cost(10_000));
980 l.result("a", ContextOrigin::Tool, "Bash");
981 l.response(&usage(12_500, 3_000), &cost(12_500)); l.response(&usage(11_500, 10), &cost(11_500));
985 let v = l.sources();
986 assert_eq!(share(&v, "Bash").tokens, 500);
987 assert_eq!(share(&v, ContextLedger::OTHER).tokens, 12_000);
988 assert!((share(&v, "Bash").cost_usd - 1_000e-6).abs() < 1e-12);
990 l.response(&usage(3_000, 10), &cost(3_000));
993 let v = l.sources();
994 assert!((share(&v, "Bash").cost_usd - 1_000e-6).abs() < 1e-12, "not charged after compaction");
995 assert_eq!(share(&v, ContextLedger::OTHER).tokens, 15_000);
996 l.compacted();
998 l.response(&usage(4_000, 10), &cost(4_000));
999 let v = l.sources();
1000 assert_eq!(share(&v, ContextLedger::OTHER).tokens, 19_000);
1001 assert!((share(&v, "Bash").cost_usd - 1_000e-6).abs() < 1e-12);
1002 }
1003
1004 #[test]
1005 fn context_ledger_retags_pending_results_and_merges() {
1006 let mut l = ContextLedger::default();
1007 l.response(&usage(100, 0), &cost(100));
1008 l.result("c1", ContextOrigin::Tool, "fetch");
1009 l.retag("c1", ContextOrigin::Mcp, "apps");
1010 l.retag("zzz", ContextOrigin::Mcp, "nope");
1011 l.response(&TokenUsage::default(), &CostBreakdown::default());
1013 l.response(&usage(300, 0), &cost(300));
1014 let v = l.sources();
1015 assert_eq!(v.len(), 2);
1016 assert_eq!((share(&v, "apps").origin, share(&v, "apps").tokens), (ContextOrigin::Mcp, 200));
1017 assert!(v.iter().all(|s| s.name != "fetch"));
1018
1019 let mut sub = ContextLedger::default();
1020 sub.response(&usage(50, 0), &cost(50));
1021 sub.result("x", ContextOrigin::Mcp, "apps");
1022 sub.response(&usage(70, 0), &cost(70));
1023 l.merge(&sub);
1024 let v = l.sources();
1025 assert_eq!(share(&v, "apps").tokens, 220);
1026 assert_eq!(share(&v, "apps").calls, 2);
1027 assert_eq!(share(&v, ContextLedger::OTHER).tokens, 150);
1028 assert!(!l.is_empty() && ContextLedger::default().is_empty());
1029 }
1030
1031 #[test]
1032 fn parses_timestamps() {
1033 let t = parse_rfc3339_utc("1970-01-02T00:00:00.000Z").unwrap();
1034 assert_eq!(t, SystemTime::UNIX_EPOCH + Duration::from_secs(86_400));
1035 let t = parse_rfc3339_utc("2026-09-03T07:15:34.5Z").unwrap();
1036 let secs = t.duration_since(SystemTime::UNIX_EPOCH).unwrap();
1037 assert_eq!(secs.as_secs(), 1_788_419_734);
1038 assert_eq!(secs.subsec_millis(), 500);
1039 assert!(parse_rfc3339_utc("nope").is_none());
1040 }
1041}
1042
1043#[cfg(test)]
1044mod rfc3339_tests {
1045 use super::parse_rfc3339_utc;
1046 use std::time::{Duration, UNIX_EPOCH};
1047
1048 fn secs(s: &str) -> u64 {
1049 parse_rfc3339_utc(s).unwrap().duration_since(UNIX_EPOCH).unwrap().as_secs()
1050 }
1051
1052 #[test]
1055 fn offsets_are_folded_into_utc() {
1056 let z = secs("2026-09-03T07:15:34Z");
1057 assert_eq!(secs("2026-09-03T08:15:34+01:00"), z);
1058 assert_eq!(secs("2026-09-03T00:15:34-07:00"), z);
1059 assert_eq!(secs("2026-09-03T08:15:34+0100"), z, "no colon");
1060 assert_eq!(secs("2026-09-03T07:15:34+00:00"), z);
1061 assert_eq!(secs("2026-09-03T12:45:34+05:30"), z, "half-hour zone");
1062 let ms = parse_rfc3339_utc("2026-09-03T08:15:34.250+01:00").unwrap();
1063 assert_eq!(ms, UNIX_EPOCH + Duration::new(z, 250_000_000));
1064 assert_eq!(secs("2026-09-03T07:15:34.5Z"), z);
1066 assert!(parse_rfc3339_utc("2026-09-03T07:15:34").is_none(), "no zone at all");
1068 assert!(parse_rfc3339_utc("2026-09-03T07:15:34+1").is_none());
1069 assert!(parse_rfc3339_utc("garbage").is_none());
1070 }
1071}