use std::collections::HashMap;
use std::path::Path;
use chrono::{DateTime, TimeDelta, Utc};
use pulpo_common::api::{ScanRollup, UsageScanResponse};
use super::{ExactUsage, RateOverrides, ScanEntry, claude, codex, pi};
const fn exact_total_tokens(u: &ExactUsage) -> u64 {
u.input_tokens + u.output_tokens + u.cache_write_tokens + u.cache_read_tokens
}
fn accumulate(
map: &mut HashMap<String, (u64, Option<f64>)>,
label: String,
tokens: u64,
cost: Option<f64>,
) {
let e = map.entry(label).or_insert((0, None));
e.0 += tokens;
if let Some(c) = cost {
e.1 = Some(e.1.unwrap_or(0.0) + c);
}
}
fn into_rollups(map: HashMap<String, (u64, Option<f64>)>) -> Vec<ScanRollup> {
let mut rows: Vec<ScanRollup> = map
.into_iter()
.filter(|(_, (tokens, _))| *tokens > 0)
.map(|(label, (total_tokens, total_cost_usd))| ScanRollup {
label,
total_tokens,
total_cost_usd,
})
.collect();
sort_rollups(&mut rows);
rows
}
pub fn scan_usage(
dirs: &ScanDirs<'_>,
rates: &RateOverrides,
node_name: &str,
now: DateTime<Utc>,
window_days: Option<u32>,
resolve_repo: impl Fn(&str) -> String,
) -> UsageScanResponse {
let since = window_days.map_or_else(
|| DateTime::<Utc>::from_timestamp(0, 0).unwrap_or(now),
|d| now - TimeDelta::days(i64::from(d)),
);
let mut repo_cache: HashMap<String, String> = HashMap::new();
let mut resolve = |cwd: String| -> String {
if let Some(label) = repo_cache.get(&cwd) {
return label.clone();
}
let label = resolve_repo(&cwd);
repo_cache.insert(cwd, label.clone());
label
};
let mut acc = ScanAccumulator::default();
fold_claude(dirs.claude, since, rates, &mut resolve, &mut acc);
fold_entries(
"codex",
codex::scan_rollouts(dirs.codex, since),
&mut resolve,
&mut acc,
);
fold_entries(
"pi",
pi::scan_sessions(dirs.pi, since),
&mut resolve,
&mut acc,
);
let total_tokens = acc.agents.values().map(|(tokens, _)| *tokens).sum();
let total_cost_usd = acc
.agents
.values()
.filter_map(|(_, cost)| *cost)
.reduce(|a, b| a + b);
UsageScanResponse {
node_name: node_name.to_owned(),
generated_at: now.to_rfc3339(),
window_days,
total_tokens,
total_cost_usd,
by_agent: into_rollups(acc.agents),
by_model: into_rollups(acc.models),
by_repo: into_rollups(acc.repos),
}
}
#[derive(Clone, Copy)]
pub struct ScanDirs<'a> {
pub claude: &'a Path,
pub codex: &'a Path,
pub pi: &'a Path,
}
#[derive(Default)]
struct ScanAccumulator {
agents: HashMap<String, (u64, Option<f64>)>,
repos: HashMap<String, (u64, Option<f64>)>,
models: HashMap<String, (u64, Option<f64>)>,
}
fn fold_claude(
claude_dir: &Path,
since: DateTime<Utc>,
rates: &RateOverrides,
resolve: &mut impl FnMut(String) -> String,
acc: &mut ScanAccumulator,
) {
if let Ok(entries) = std::fs::read_dir(claude_dir.join("projects")) {
for entry in entries.flatten() {
let dir = entry.path();
if !dir.is_dir() {
continue;
}
let Some(d) = claude::read_usage_dir(&dir, since, rates) else {
continue;
};
let tokens = exact_total_tokens(&d.usage);
let cost = d.usage.cost_usd;
accumulate(&mut acc.agents, "claude".into(), tokens, cost);
let raw = d.cwd.unwrap_or_else(|| {
dir.file_name()
.and_then(|n| n.to_str())
.unwrap_or("unknown")
.to_owned()
});
let repo = resolve(raw);
accumulate(&mut acc.repos, repo, tokens, cost);
for m in d.by_model {
accumulate(&mut acc.models, m.model, m.tokens, m.cost_usd);
}
}
}
}
fn fold_entries(
agent: &str,
entries: Vec<ScanEntry>,
resolve: &mut impl FnMut(String) -> String,
acc: &mut ScanAccumulator,
) {
for entry in entries {
accumulate(
&mut acc.agents,
agent.to_owned(),
entry.tokens,
entry.cost_usd,
);
let repo = resolve(entry.cwd);
accumulate(&mut acc.repos, repo, entry.tokens, entry.cost_usd);
let model = entry.model.unwrap_or_else(|| agent.to_owned());
accumulate(&mut acc.models, model, entry.tokens, entry.cost_usd);
}
}
fn sort_rollups(rollups: &mut [ScanRollup]) {
rollups.sort_by(|a, b| {
b.total_cost_usd
.unwrap_or(-1.0)
.partial_cmp(&a.total_cost_usd.unwrap_or(-1.0))
.unwrap_or(std::cmp::Ordering::Equal)
.then(b.total_tokens.cmp(&a.total_tokens))
.then(a.label.cmp(&b.label))
});
}
#[cfg(not(coverage))]
pub(crate) fn canonical_repo(cwd: &str) -> String {
use std::process::Command;
let output = Command::new("git")
.args([
"-C",
cwd,
"rev-parse",
"--path-format=absolute",
"--git-common-dir",
])
.output();
if let Ok(output) = output
&& output.status.success()
{
let common = String::from_utf8_lossy(&output.stdout);
if let Some(root) = Path::new(common.trim()).parent().and_then(Path::to_str)
&& !root.is_empty()
{
return root.to_owned();
}
}
cwd.to_owned()
}
#[cfg(test)]
mod tests {
use super::*;
use std::fs;
fn claude_record(cwd: Option<&str>, model: &str, input: u64, output: u64) -> String {
let ts = Utc::now().to_rfc3339();
let cwd_field = cwd.map(|c| format!(r#""cwd":"{c}","#)).unwrap_or_default();
format!(
r#"{{"timestamp":"{ts}",{cwd_field}"requestId":"r1","type":"assistant","message":{{"id":"m1","model":"{model}","usage":{{"input_tokens":{input},"output_tokens":{output},"cache_read_input_tokens":0,"cache_creation_input_tokens":0}}}}}}"#
)
}
fn write_claude_project(claude_dir: &Path, dir_name: &str, content: &str) {
let d = claude_dir.join("projects").join(dir_name);
fs::create_dir_all(&d).unwrap();
fs::write(d.join("session.jsonl"), content).unwrap();
}
fn codex_rollout(cwd: &str, input: u64, cached: u64, output: u64) -> String {
let meta = format!(
r#"{{"timestamp":"2026-06-12T10:00:00Z","type":"session_meta","payload":{{"id":"abc","timestamp":"2026-06-12T10:00:00Z","cwd":"{cwd}","originator":"codex_cli_rs"}}}}"#
);
let tc = format!(
r#"{{"timestamp":"2026-06-12T10:00:00Z","type":"event_msg","payload":{{"type":"token_count","info":{{"total_token_usage":{{"input_tokens":{input},"cached_input_tokens":{cached},"output_tokens":{output},"total_tokens":{}}}}}}}}}"#,
input + output
);
format!("{meta}\n{tc}\n")
}
fn pi_session(cwd: &str, model: &str, input: u64, output: u64, cost: f64) -> String {
let header = format!(
r#"{{"type":"session","version":3,"id":"0197-abc","timestamp":"2026-06-12T10:00:00.000Z","cwd":"{cwd}"}}"#
);
let msg = format!(
r#"{{"type":"message","id":"bbbb2222","parentId":null,"message":{{"role":"assistant","content":[],"model":"{model}","usage":{{"input":{input},"output":{output},"cacheRead":0,"cacheWrite":0,"totalTokens":{},"cost":{{"input":0.0,"output":0.0,"cacheRead":0.0,"cacheWrite":0.0,"total":{cost}}}}},"stopReason":"stop","timestamp":{}}}}}"#,
input + output,
Utc::now().timestamp_millis(),
);
format!("{header}\n{msg}\n")
}
fn write_pi_session(pi_dir: &Path, dir_name: &str, content: &str) {
let d = pi_dir.join("agent").join("sessions").join(dir_name);
fs::create_dir_all(&d).unwrap();
fs::write(d.join("0197-abc.jsonl"), content).unwrap();
}
fn write_codex_rollout(codex_dir: &Path, name: &str, content: &str) {
let d = codex_dir
.join("sessions")
.join("2026")
.join("06")
.join("12");
fs::create_dir_all(&d).unwrap();
fs::write(d.join(name), content).unwrap();
}
#[test]
fn test_scan_empty_dirs() {
let claude = tempfile::tempdir().unwrap();
let codex = tempfile::tempdir().unwrap();
let pi_dir = tempfile::tempdir().unwrap();
let r = scan_usage(
&ScanDirs {
claude: claude.path(),
codex: codex.path(),
pi: pi_dir.path(),
},
&RateOverrides::default(),
"n",
Utc::now(),
None,
|s: &str| s.to_owned(),
);
assert_eq!(r.total_tokens, 0);
assert!(r.by_agent.is_empty());
assert!(r.by_repo.is_empty());
assert_eq!(r.total_cost_usd, None);
}
#[test]
fn test_scan_merges_agents_by_repo() {
let claude = tempfile::tempdir().unwrap();
let codex = tempfile::tempdir().unwrap();
let pi_dir = tempfile::tempdir().unwrap();
write_claude_project(
claude.path(),
"proj-api",
&claude_record(Some("/repos/api"), "claude-opus-4-8", 1000, 500),
);
write_claude_project(
claude.path(),
"proj-web",
&claude_record(Some("/repos/web"), "claude-opus-4-8", 200, 100),
);
write_codex_rollout(
codex.path(),
"rollout-2026-06-12-a.jsonl",
&codex_rollout("/repos/api", 800, 0, 200),
);
let r = scan_usage(
&ScanDirs {
claude: claude.path(),
codex: codex.path(),
pi: pi_dir.path(),
},
&RateOverrides::default(),
"node-x",
Utc::now(),
None,
|s: &str| s.to_owned(),
);
assert_eq!(r.by_agent.len(), 2);
assert_eq!(r.total_tokens, 2800);
let expected = 21_000.0 / 1_000_000.0;
assert!((r.total_cost_usd.unwrap() - expected).abs() < 1e-9);
assert_eq!(r.by_repo[0].label, "/repos/api");
assert_eq!(r.by_repo[0].total_tokens, 2500);
let web = r.by_repo.iter().find(|x| x.label == "/repos/web").unwrap();
assert_eq!(web.total_tokens, 300);
assert_eq!(r.by_model.len(), 2);
let opus = r
.by_model
.iter()
.find(|m| m.label == "claude-opus-4-8")
.unwrap();
assert_eq!(opus.total_tokens, 1800);
assert!(opus.total_cost_usd.unwrap() > 0.0);
let codex = r.by_model.iter().find(|m| m.label == "codex").unwrap();
assert_eq!(codex.total_tokens, 1000);
assert!(codex.total_cost_usd.is_none());
assert_eq!(r.by_model[0].label, "claude-opus-4-8");
}
#[test]
fn test_scan_includes_pi_with_exact_cost_and_merges_repo() {
let claude = tempfile::tempdir().unwrap();
let codex = tempfile::tempdir().unwrap();
let pi_dir = tempfile::tempdir().unwrap();
write_claude_project(
claude.path(),
"proj-api",
&claude_record(Some("/repos/api"), "claude-opus-4-8", 1000, 500),
);
write_pi_session(
pi_dir.path(),
"--repos-api--",
&pi_session("/repos/api", "gemini-3-pro", 800, 200, 0.05),
);
let r = scan_usage(
&ScanDirs {
claude: claude.path(),
codex: codex.path(),
pi: pi_dir.path(),
},
&RateOverrides::default(),
"n",
Utc::now(),
None,
|s: &str| s.to_owned(),
);
let pi_row = r.by_agent.iter().find(|a| a.label == "pi").unwrap();
assert_eq!(pi_row.total_tokens, 1000);
assert!((pi_row.total_cost_usd.unwrap() - 0.05).abs() < 1e-12);
let claude_cost = 17_500.0 / 1_000_000.0;
assert!((r.total_cost_usd.unwrap() - (claude_cost + 0.05)).abs() < 1e-9);
assert_eq!(r.total_tokens, 2500);
assert_eq!(r.by_repo.len(), 1);
assert_eq!(r.by_repo[0].label, "/repos/api");
assert_eq!(r.by_repo[0].total_tokens, 2500);
assert!((r.by_repo[0].total_cost_usd.unwrap() - (claude_cost + 0.05)).abs() < 1e-9);
let gem = r
.by_model
.iter()
.find(|m| m.label == "gemini-3-pro")
.unwrap();
assert_eq!(gem.total_tokens, 1000);
assert!((gem.total_cost_usd.unwrap() - 0.05).abs() < 1e-12);
}
#[test]
fn test_scan_window_days_filters_old_records_and_sets_field() {
let claude = tempfile::tempdir().unwrap();
let codex = tempfile::tempdir().unwrap();
let pi_dir = tempfile::tempdir().unwrap();
let recent = Utc::now().to_rfc3339();
let old = (Utc::now() - TimeDelta::days(10)).to_rfc3339();
let line = |ts: &str, id: &str| {
format!(
r#"{{"timestamp":"{ts}","cwd":"/repos/api","requestId":"{id}","type":"assistant","message":{{"id":"{id}","model":"claude-opus-4-8","usage":{{"input_tokens":100,"output_tokens":0,"cache_read_input_tokens":0,"cache_creation_input_tokens":0}}}}}}"#
)
};
write_claude_project(
claude.path(),
"proj-api",
&format!("{}\n{}\n", line(&recent, "a"), line(&old, "b")),
);
let r = scan_usage(
&ScanDirs {
claude: claude.path(),
codex: codex.path(),
pi: pi_dir.path(),
},
&RateOverrides::default(),
"n",
Utc::now(),
Some(3),
|s: &str| s.to_owned(),
);
assert_eq!(r.window_days, Some(3));
assert_eq!(r.total_tokens, 100);
let all = scan_usage(
&ScanDirs {
claude: claude.path(),
codex: codex.path(),
pi: pi_dir.path(),
},
&RateOverrides::default(),
"n",
Utc::now(),
None,
|s: &str| s.to_owned(),
);
assert_eq!(all.window_days, None);
assert_eq!(all.total_tokens, 200);
}
#[test]
fn test_scan_skips_non_dir_and_empty_project_entries() {
let claude = tempfile::tempdir().unwrap();
let codex = tempfile::tempdir().unwrap();
let pi_dir = tempfile::tempdir().unwrap();
let projects = claude.path().join("projects");
fs::create_dir_all(&projects).unwrap();
fs::write(projects.join("stray.txt"), "x").unwrap();
fs::create_dir_all(projects.join("empty-proj")).unwrap();
let r = scan_usage(
&ScanDirs {
claude: claude.path(),
codex: codex.path(),
pi: pi_dir.path(),
},
&RateOverrides::default(),
"n",
Utc::now(),
None,
|s: &str| s.to_owned(),
);
assert_eq!(r.total_tokens, 0);
assert!(r.by_agent.is_empty());
assert!(r.by_model.is_empty());
}
#[test]
fn test_scan_drops_zero_token_model_rows() {
let claude = tempfile::tempdir().unwrap();
let codex = tempfile::tempdir().unwrap();
let pi_dir = tempfile::tempdir().unwrap();
let ts = Utc::now().to_rfc3339();
let line = |model: &str, id: &str, input: u64| {
format!(
r#"{{"timestamp":"{ts}","cwd":"/repos/api","requestId":"{id}","type":"assistant","message":{{"id":"{id}","model":"{model}","usage":{{"input_tokens":{input},"output_tokens":0,"cache_read_input_tokens":0,"cache_creation_input_tokens":0}}}}}}"#
)
};
write_claude_project(
claude.path(),
"proj-api",
&format!(
"{}\n{}\n",
line("claude-opus-4-8", "a", 100),
line("<synthetic>", "b", 0)
),
);
let r = scan_usage(
&ScanDirs {
claude: claude.path(),
codex: codex.path(),
pi: pi_dir.path(),
},
&RateOverrides::default(),
"n",
Utc::now(),
None,
|s: &str| s.to_owned(),
);
assert_eq!(r.by_model.len(), 1);
assert_eq!(r.by_model[0].label, "claude-opus-4-8");
}
#[test]
fn test_scan_falls_back_to_dir_name_without_cwd() {
let claude = tempfile::tempdir().unwrap();
let codex = tempfile::tempdir().unwrap();
let pi_dir = tempfile::tempdir().unwrap();
write_claude_project(
claude.path(),
"-Users-x-repo",
&claude_record(None, "claude-opus-4-8", 10, 5),
);
let r = scan_usage(
&ScanDirs {
claude: claude.path(),
codex: codex.path(),
pi: pi_dir.path(),
},
&RateOverrides::default(),
"n",
Utc::now(),
None,
|s: &str| s.to_owned(),
);
assert_eq!(r.by_repo.len(), 1);
assert_eq!(r.by_repo[0].label, "-Users-x-repo");
}
#[test]
fn test_scan_collapses_worktrees_via_resolver() {
let claude = tempfile::tempdir().unwrap();
let codex = tempfile::tempdir().unwrap();
let pi_dir = tempfile::tempdir().unwrap();
write_claude_project(
claude.path(),
"proj-main",
&claude_record(Some("/repos/api"), "claude-opus-4-8", 1000, 0),
);
write_claude_project(
claude.path(),
"proj-wt",
&claude_record(Some("/repos/api-worktrees/feat"), "claude-opus-4-8", 500, 0),
);
write_codex_rollout(
codex.path(),
"rollout-2026-06-12-a.jsonl",
&codex_rollout("/repos/api/src", 200, 0, 0),
);
let resolve = |cwd: &str| {
if cwd.starts_with("/repos/api") {
"/repos/api".to_owned()
} else {
cwd.to_owned()
}
};
let r = scan_usage(
&ScanDirs {
claude: claude.path(),
codex: codex.path(),
pi: pi_dir.path(),
},
&RateOverrides::default(),
"n",
Utc::now(),
None,
resolve,
);
assert_eq!(r.by_repo.len(), 1);
assert_eq!(r.by_repo[0].label, "/repos/api");
assert_eq!(r.by_repo[0].total_tokens, 1700);
}
#[cfg(not(coverage))]
#[test]
fn test_canonical_repo_collapses_worktree_and_subdir() {
use std::process::Command;
let tmp = tempfile::tempdir().unwrap();
let repo = tmp.path().join("origin");
fs::create_dir_all(&repo).unwrap();
let git = |args: &[&str], dir: &Path| {
let ok = Command::new("git")
.args(args)
.current_dir(dir)
.output()
.unwrap()
.status
.success();
assert!(ok, "git {args:?} failed");
};
git(&["init", "-q"], &repo);
git(&["config", "user.email", "t@t"], &repo);
git(&["config", "user.name", "t"], &repo);
fs::write(repo.join("f"), "x").unwrap();
git(&["add", "-A"], &repo);
git(&["commit", "-qm", "init"], &repo);
let sub = repo.join("src");
fs::create_dir_all(&sub).unwrap();
let wt = tmp.path().join("wt");
git(&["worktree", "add", "-q", wt.to_str().unwrap()], &repo);
let root = canonical_repo(repo.to_str().unwrap());
assert_eq!(canonical_repo(sub.to_str().unwrap()), root);
assert_eq!(canonical_repo(wt.to_str().unwrap()), root);
let plain = tmp.path().join("plain");
fs::create_dir_all(&plain).unwrap();
assert_eq!(
canonical_repo(plain.to_str().unwrap()),
plain.to_str().unwrap()
);
}
}