use std::collections::{HashMap, HashSet};
use std::path::Path;
use chrono::{DateTime, Utc};
use super::{ExactUsage, RateOverrides, SOURCE_CLAUDE, resolve_rates, token_field};
pub fn sanitize_workdir(workdir: &str) -> String {
workdir
.chars()
.map(|c| if c.is_ascii_alphanumeric() { c } else { '-' })
.collect()
}
#[derive(Debug, Default)]
struct Totals {
input: u64,
output: u64,
cache_write: u64,
cache_read: u64,
cost_usd: f64,
unknown_model: bool,
records: u64,
by_model: HashMap<String, ModelAgg>,
}
#[derive(Debug, Default)]
struct ModelAgg {
tokens: u64,
cost_usd: f64,
unknown: bool,
}
fn apply_transcript_line(
line: &str,
since: DateTime<Utc>,
seen: &mut HashSet<String>,
totals: &mut Totals,
rates: &RateOverrides,
) {
let Ok(value) = serde_json::from_str::<serde_json::Value>(line) else {
return;
};
let Some(timestamp) = value
.get("timestamp")
.and_then(serde_json::Value::as_str)
.and_then(|ts| DateTime::parse_from_rfc3339(ts).ok())
else {
return;
};
if timestamp.with_timezone(&Utc) < since {
return;
}
let Some(message) = value.get("message") else {
return;
};
let Some(usage) = message.get("usage") else {
return;
};
let message_id = message.get("id").and_then(serde_json::Value::as_str);
let request_id = value.get("requestId").and_then(serde_json::Value::as_str);
if let (Some(mid), Some(rid)) = (message_id, request_id)
&& !seen.insert(format!("{mid}:{rid}"))
{
return;
}
let input = token_field(usage, "input_tokens");
let output = token_field(usage, "output_tokens");
let cache_read = token_field(usage, "cache_read_input_tokens");
let (five_min, one_hour) = usage.get("cache_creation").map_or_else(
|| (token_field(usage, "cache_creation_input_tokens"), 0),
|breakdown| {
(
token_field(breakdown, "ephemeral_5m_input_tokens"),
token_field(breakdown, "ephemeral_1h_input_tokens"),
)
},
);
totals.input += input;
totals.output += output;
totals.cache_read += cache_read;
totals.cache_write += five_min + one_hour;
totals.records += 1;
let tokens = input + output + cache_read + five_min + one_hour;
let model = message
.get("model")
.and_then(serde_json::Value::as_str)
.unwrap_or_default();
let agg = totals.by_model.entry(model.to_owned()).or_default();
agg.tokens += tokens;
#[allow(clippy::cast_precision_loss)]
match resolve_rates(model, rates) {
Some(rates) => {
let cost = (input as f64).mul_add(
rates.input,
(output as f64).mul_add(
rates.output,
(cache_read as f64).mul_add(
rates.cache_read,
(five_min as f64).mul_add(
rates.cache_write_5m,
(one_hour as f64) * rates.cache_write_1h,
),
),
),
) / 1_000_000.0;
totals.cost_usd += cost;
agg.cost_usd += cost;
}
None => {
if tokens > 0 {
totals.unknown_model = true;
agg.unknown = true;
}
}
}
}
pub fn read_usage(
claude_dir: &Path,
workdir: &str,
since: DateTime<Utc>,
rates: &RateOverrides,
) -> Option<ExactUsage> {
let project_dir = claude_dir.join("projects").join(sanitize_workdir(workdir));
read_usage_dir(&project_dir, since, rates).map(|d| d.usage)
}
pub(crate) struct DirUsage {
pub usage: ExactUsage,
pub cwd: Option<String>,
pub by_model: Vec<ModelUsage>,
}
pub(crate) struct ModelUsage {
pub model: String,
pub tokens: u64,
pub cost_usd: Option<f64>,
}
pub(crate) fn read_usage_dir(
project_dir: &Path,
since: DateTime<Utc>,
rates: &RateOverrides,
) -> Option<DirUsage> {
let entries = std::fs::read_dir(project_dir).ok()?;
let mut totals = Totals::default();
let mut seen = HashSet::new();
let mut cwd: Option<String> = None;
for entry in entries.flatten() {
let path = entry.path();
if path.extension().and_then(|ext| ext.to_str()) != Some("jsonl") {
continue;
}
if let Ok(file_meta) = entry.metadata()
&& let Ok(modified) = file_meta.modified()
{
let modified: DateTime<Utc> = modified.into();
if modified < since {
continue;
}
}
let Ok(content) = std::fs::read_to_string(&path) else {
continue;
};
for line in content.lines() {
apply_transcript_line(line, since, &mut seen, &mut totals, rates);
if cwd.is_none()
&& let Ok(value) = serde_json::from_str::<serde_json::Value>(line)
{
cwd = value
.get("cwd")
.and_then(serde_json::Value::as_str)
.map(str::to_owned);
}
}
}
if totals.records == 0 {
return None;
}
let by_model = totals
.by_model
.into_iter()
.map(|(model, agg)| ModelUsage {
model,
tokens: agg.tokens,
cost_usd: (!agg.unknown).then_some(agg.cost_usd),
})
.collect();
Some(DirUsage {
usage: ExactUsage {
source: SOURCE_CLAUDE,
input_tokens: totals.input,
output_tokens: totals.output,
cache_write_tokens: totals.cache_write,
cache_read_tokens: totals.cache_read,
cost_usd: (!totals.unknown_model).then_some(totals.cost_usd),
quota: None,
},
cwd,
by_model,
})
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::TimeDelta;
use std::fs;
fn transcript_line(timestamp: &str, message_id: &str, request_id: &str, model: &str) -> String {
format!(
r#"{{"timestamp":"{timestamp}","requestId":"{request_id}","type":"assistant","message":{{"id":"{message_id}","model":"{model}","usage":{{"input_tokens":1000,"output_tokens":500,"cache_read_input_tokens":2000,"cache_creation_input_tokens":300,"cache_creation":{{"ephemeral_5m_input_tokens":300,"ephemeral_1h_input_tokens":0}}}}}}}}"#
)
}
fn write_project_file(claude_dir: &Path, workdir: &str, name: &str, content: &str) {
let project_dir = claude_dir.join("projects").join(sanitize_workdir(workdir));
fs::create_dir_all(&project_dir).unwrap();
fs::write(project_dir.join(name), content).unwrap();
}
#[test]
fn test_sanitize_workdir_plain_path() {
assert_eq!(
sanitize_workdir("/Users/dario/Code/darioblanco/pulpo"),
"-Users-dario-Code-darioblanco-pulpo"
);
}
#[test]
fn test_sanitize_workdir_dots_become_dashes() {
assert_eq!(
sanitize_workdir("/Users/dario/.pulpo/worktrees/fix-1"),
"-Users-dario--pulpo-worktrees-fix-1"
);
}
#[test]
fn test_sanitize_workdir_underscores_and_spaces() {
assert_eq!(sanitize_workdir("/tmp/my_repo v2"), "-tmp-my-repo-v2");
}
#[test]
fn test_read_usage_missing_project_dir() {
let tmp = tempfile::tempdir().unwrap();
assert!(
read_usage(
tmp.path(),
"/tmp/repo",
Utc::now(),
&RateOverrides::default()
)
.is_none()
);
}
#[test]
fn test_read_usage_sums_records_and_computes_cost() {
let tmp = tempfile::tempdir().unwrap();
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
let content = format!(
"{}\n{}\n",
transcript_line(&ts, "msg_1", "req_1", "claude-fable-5"),
transcript_line(&ts, "msg_2", "req_2", "claude-fable-5"),
);
write_project_file(tmp.path(), "/tmp/repo", "abc.jsonl", &content);
let usage = read_usage(tmp.path(), "/tmp/repo", since, &RateOverrides::default()).unwrap();
assert_eq!(usage.source, SOURCE_CLAUDE);
assert_eq!(usage.input_tokens, 2000);
assert_eq!(usage.output_tokens, 1000);
assert_eq!(usage.cache_read_tokens, 4000);
assert_eq!(usage.cache_write_tokens, 600);
let expected = 2.0 * (10_000.0 + 25_000.0 + 2_000.0 + 3_750.0) / 1_000_000.0;
assert!((usage.cost_usd.unwrap() - expected).abs() < 1e-9);
}
#[test]
fn test_read_usage_sums_across_multiple_files_in_project_dir() {
let tmp = tempfile::tempdir().unwrap();
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
write_project_file(
tmp.path(),
"/tmp/repo",
"first.jsonl",
&transcript_line(&ts, "m1", "r1", "claude-opus-4-8"),
);
write_project_file(
tmp.path(),
"/tmp/repo",
"second.jsonl",
&transcript_line(&ts, "m2", "r2", "claude-opus-4-8"),
);
let usage = read_usage(tmp.path(), "/tmp/repo", since, &RateOverrides::default()).unwrap();
assert_eq!(usage.input_tokens, 2000);
assert_eq!(usage.output_tokens, 1000);
}
#[test]
fn test_read_usage_dedupes_repeated_message_and_request_id() {
let tmp = tempfile::tempdir().unwrap();
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
let line = transcript_line(&ts, "msg_1", "req_1", "claude-opus-4-8");
write_project_file(
tmp.path(),
"/tmp/repo",
"abc.jsonl",
&format!("{line}\n{line}\n"),
);
let usage = read_usage(tmp.path(), "/tmp/repo", since, &RateOverrides::default()).unwrap();
assert_eq!(usage.input_tokens, 1000);
assert_eq!(usage.output_tokens, 500);
}
#[test]
fn test_read_usage_skips_records_before_since() {
let tmp = tempfile::tempdir().unwrap();
let since = Utc::now() - TimeDelta::hours(1);
let old_ts = (Utc::now() - TimeDelta::hours(5)).to_rfc3339();
let new_ts = Utc::now().to_rfc3339();
let content = format!(
"{}\n{}\n",
transcript_line(&old_ts, "msg_old", "req_old", "claude-fable-5"),
transcript_line(&new_ts, "msg_new", "req_new", "claude-fable-5"),
);
write_project_file(tmp.path(), "/tmp/repo", "abc.jsonl", &content);
let usage = read_usage(tmp.path(), "/tmp/repo", since, &RateOverrides::default()).unwrap();
assert_eq!(usage.input_tokens, 1000);
}
#[test]
fn test_read_usage_all_records_too_old_returns_none() {
let tmp = tempfile::tempdir().unwrap();
let old_ts = (Utc::now() - TimeDelta::hours(5)).to_rfc3339();
let content = transcript_line(&old_ts, "msg_old", "req_old", "claude-fable-5");
write_project_file(tmp.path(), "/tmp/repo", "abc.jsonl", &content);
let since = Utc::now() - TimeDelta::hours(1);
assert!(read_usage(tmp.path(), "/tmp/repo", since, &RateOverrides::default()).is_none());
}
#[test]
fn test_read_usage_unknown_model_withholds_cost() {
let tmp = tempfile::tempdir().unwrap();
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
let content = format!(
"{}\n{}\n",
transcript_line(&ts, "msg_1", "req_1", "claude-fable-5"),
transcript_line(&ts, "msg_2", "req_2", "experimental-model"),
);
write_project_file(tmp.path(), "/tmp/repo", "abc.jsonl", &content);
let usage = read_usage(tmp.path(), "/tmp/repo", since, &RateOverrides::default()).unwrap();
assert_eq!(usage.input_tokens, 2000);
assert!(usage.cost_usd.is_none());
}
#[test]
fn test_read_usage_config_override_prices_unknown_model() {
let tmp = tempfile::tempdir().unwrap();
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
let content = transcript_line(&ts, "msg_1", "req_1", "brand-new-model");
write_project_file(tmp.path(), "/tmp/repo", "abc.jsonl", &content);
let bare = read_usage(tmp.path(), "/tmp/repo", since, &RateOverrides::default()).unwrap();
assert!(bare.cost_usd.is_none());
let overrides = RateOverrides::new([(
"brand-new-model".to_owned(),
crate::usage::ModelRates {
input: 2.0,
output: 8.0,
cache_read: 0.0,
cache_write_5m: 0.0,
cache_write_1h: 0.0,
},
)]);
let usage = read_usage(tmp.path(), "/tmp/repo", since, &overrides).unwrap();
let expected = (1000.0f64).mul_add(2.0, 500.0 * 8.0) / 1_000_000.0;
assert!((usage.cost_usd.unwrap() - expected).abs() < 1e-9);
}
#[test]
fn test_read_usage_config_override_reprices_known_model() {
let tmp = tempfile::tempdir().unwrap();
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
let content = transcript_line(&ts, "msg_1", "req_1", "claude-opus-4-8");
write_project_file(tmp.path(), "/tmp/repo", "abc.jsonl", &content);
let overrides = RateOverrides::new([(
"claude-opus-4-8".to_owned(),
crate::usage::ModelRates {
input: 99.0,
output: 0.0,
cache_read: 0.0,
cache_write_5m: 0.0,
cache_write_1h: 0.0,
},
)]);
let usage = read_usage(tmp.path(), "/tmp/repo", since, &overrides).unwrap();
let expected = 1000.0 * 99.0 / 1_000_000.0;
assert!((usage.cost_usd.unwrap() - expected).abs() < 1e-9);
}
#[test]
fn test_read_usage_prices_1h_cache_writes() {
let tmp = tempfile::tempdir().unwrap();
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
let line = format!(
r#"{{"timestamp":"{ts}","requestId":"req_1","message":{{"id":"msg_1","model":"claude-fable-5","usage":{{"input_tokens":0,"output_tokens":0,"cache_read_input_tokens":0,"cache_creation_input_tokens":1000000,"cache_creation":{{"ephemeral_5m_input_tokens":0,"ephemeral_1h_input_tokens":1000000}}}}}}}}"#
);
write_project_file(tmp.path(), "/tmp/repo", "abc.jsonl", &line);
let usage = read_usage(tmp.path(), "/tmp/repo", since, &RateOverrides::default()).unwrap();
assert_eq!(usage.cache_write_tokens, 1_000_000);
assert!((usage.cost_usd.unwrap() - 20.0).abs() < 1e-9);
}
#[test]
fn test_read_usage_flat_cache_creation_without_breakdown() {
let tmp = tempfile::tempdir().unwrap();
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
let line = format!(
r#"{{"timestamp":"{ts}","requestId":"req_1","message":{{"id":"msg_1","model":"claude-haiku-4-5","usage":{{"input_tokens":100,"output_tokens":50,"cache_creation_input_tokens":400}}}}}}"#
);
write_project_file(tmp.path(), "/tmp/repo", "abc.jsonl", &line);
let usage = read_usage(tmp.path(), "/tmp/repo", since, &RateOverrides::default()).unwrap();
assert_eq!(usage.cache_write_tokens, 400);
assert_eq!(usage.cache_read_tokens, 0);
}
#[test]
fn test_read_usage_skips_invalid_and_irrelevant_lines() {
let tmp = tempfile::tempdir().unwrap();
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
let content = format!(
"not json\n{{\"timestamp\":\"{ts}\",\"type\":\"user\"}}\n{{\"timestamp\":\"bad-ts\"}}\n{{\"no_timestamp\":true}}\n{{\"timestamp\":\"{ts}\",\"message\":{{\"id\":\"m\"}}}}\n{}\n",
transcript_line(&ts, "msg_1", "req_1", "claude-sonnet-4-6"),
);
write_project_file(tmp.path(), "/tmp/repo", "abc.jsonl", &content);
let usage = read_usage(tmp.path(), "/tmp/repo", since, &RateOverrides::default()).unwrap();
assert_eq!(usage.input_tokens, 1000);
}
#[test]
fn test_read_usage_counts_records_without_ids() {
let tmp = tempfile::tempdir().unwrap();
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
let line = format!(
r#"{{"timestamp":"{ts}","message":{{"model":"claude-opus-4-8","usage":{{"input_tokens":10,"output_tokens":5}}}}}}"#
);
write_project_file(
tmp.path(),
"/tmp/repo",
"abc.jsonl",
&format!("{line}\n{line}\n"),
);
let usage = read_usage(tmp.path(), "/tmp/repo", since, &RateOverrides::default()).unwrap();
assert_eq!(usage.input_tokens, 20);
}
#[test]
fn test_read_usage_ignores_non_jsonl_files() {
let tmp = tempfile::tempdir().unwrap();
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
write_project_file(
tmp.path(),
"/tmp/repo",
"notes.txt",
&transcript_line(&ts, "msg_1", "req_1", "claude-fable-5"),
);
assert!(read_usage(tmp.path(), "/tmp/repo", since, &RateOverrides::default()).is_none());
}
#[test]
fn test_read_usage_skips_files_untouched_since_spawn() {
let tmp = tempfile::tempdir().unwrap();
let ts = Utc::now().to_rfc3339();
write_project_file(
tmp.path(),
"/tmp/repo",
"old.jsonl",
&transcript_line(&ts, "msg_1", "req_1", "claude-fable-5"),
);
let file_path = tmp
.path()
.join("projects")
.join(sanitize_workdir("/tmp/repo"))
.join("old.jsonl");
let old_mtime = std::time::SystemTime::now() - std::time::Duration::from_secs(7200);
let file = fs::File::options().write(true).open(&file_path).unwrap();
file.set_modified(old_mtime).unwrap();
let since = Utc::now() - TimeDelta::hours(1);
assert!(read_usage(tmp.path(), "/tmp/repo", since, &RateOverrides::default()).is_none());
}
#[test]
fn test_read_usage_dir_breaks_down_by_model() {
let tmp = tempfile::tempdir().unwrap();
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
let content = format!(
"{}\n{}\n{}\n",
transcript_line(&ts, "m1", "r1", "claude-opus-4-8"),
transcript_line(&ts, "m2", "r2", "claude-opus-4-8"),
transcript_line(&ts, "m3", "r3", "claude-haiku-4-5"),
);
let project_dir = tmp
.path()
.join("projects")
.join(sanitize_workdir("/tmp/repo"));
fs::create_dir_all(&project_dir).unwrap();
fs::write(project_dir.join("abc.jsonl"), &content).unwrap();
let d = read_usage_dir(&project_dir, since, &RateOverrides::default()).unwrap();
assert_eq!(d.by_model.len(), 2);
let opus = d
.by_model
.iter()
.find(|m| m.model == "claude-opus-4-8")
.unwrap();
assert_eq!(opus.tokens, 7600);
assert!(opus.cost_usd.unwrap() > 0.0);
let haiku = d
.by_model
.iter()
.find(|m| m.model == "claude-haiku-4-5")
.unwrap();
assert_eq!(haiku.tokens, 3800);
}
#[test]
fn test_read_usage_dir_by_model_withholds_unknown_cost() {
let tmp = tempfile::tempdir().unwrap();
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
let content = transcript_line(&ts, "m1", "r1", "experimental-model");
let project_dir = tmp
.path()
.join("projects")
.join(sanitize_workdir("/tmp/repo"));
fs::create_dir_all(&project_dir).unwrap();
fs::write(project_dir.join("abc.jsonl"), &content).unwrap();
let d = read_usage_dir(&project_dir, since, &RateOverrides::default()).unwrap();
let m = d
.by_model
.iter()
.find(|m| m.model == "experimental-model")
.unwrap();
assert!(m.cost_usd.is_none());
assert_eq!(m.tokens, 3800);
}
#[test]
fn test_read_usage_sums_across_multiple_files() {
let tmp = tempfile::tempdir().unwrap();
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
write_project_file(
tmp.path(),
"/tmp/repo",
"a.jsonl",
&transcript_line(&ts, "msg_1", "req_1", "claude-fable-5"),
);
write_project_file(
tmp.path(),
"/tmp/repo",
"b.jsonl",
&transcript_line(&ts, "msg_2", "req_2", "claude-fable-5"),
);
let usage = read_usage(tmp.path(), "/tmp/repo", since, &RateOverrides::default()).unwrap();
assert_eq!(usage.input_tokens, 2000);
}
}