pub mod claude;
pub mod codex;
pub mod pi;
pub mod pool;
pub mod projection;
pub mod scan;
use std::path::{Path, PathBuf};
#[cfg(not(coverage))]
use std::sync::OnceLock;
use chrono::{DateTime, TimeDelta, Utc};
use pulpo_common::session::Session;
use crate::auth_info::agent_provider_for_command;
pub const SOURCE_CLAUDE: &str = "claude-jsonl";
pub const SOURCE_CODEX: &str = "codex-jsonl";
const SINCE_GRACE_SECS: i64 = 60;
pub(crate) struct ScanEntry {
pub cwd: String,
pub model: Option<String>,
pub tokens: u64,
pub cost_usd: Option<f64>,
}
pub(crate) fn token_field(usage: &serde_json::Value, key: &str) -> u64 {
usage
.get(key)
.and_then(serde_json::Value::as_u64)
.unwrap_or(0)
}
pub(crate) fn collect_jsonl_files(
root: std::path::PathBuf,
name_ok: impl Fn(&str) -> bool,
) -> Vec<std::path::PathBuf> {
let mut out = Vec::new();
let mut stack = vec![root];
while let Some(dir) = stack.pop() {
let Ok(entries) = std::fs::read_dir(&dir) else {
continue;
};
for entry in entries.flatten() {
let path = entry.path();
if path.is_dir() {
stack.push(path);
} else if path.extension().and_then(|e| e.to_str()) == Some("jsonl")
&& path
.file_name()
.and_then(|n| n.to_str())
.is_some_and(&name_ok)
{
out.push(path);
}
}
}
out
}
#[derive(Debug, Clone, PartialEq)]
pub struct ExactUsage {
pub source: &'static str,
pub input_tokens: u64,
pub output_tokens: u64,
pub cache_write_tokens: u64,
pub cache_read_tokens: u64,
pub cost_usd: Option<f64>,
pub quota: Option<QuotaSnapshot>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct QuotaWindow {
pub used_percent: f64,
pub window_minutes: Option<u64>,
pub resets_at: Option<i64>,
}
#[derive(Debug, Clone, Default, PartialEq)]
pub struct QuotaSnapshot {
pub primary: Option<QuotaWindow>,
pub secondary: Option<QuotaWindow>,
pub plan: Option<String>,
}
#[derive(Debug, Clone, Copy)]
pub struct ModelRates {
pub input: f64,
pub output: f64,
pub cache_read: f64,
pub cache_write_5m: f64,
pub cache_write_1h: f64,
}
const FABLE_RATES: ModelRates = ModelRates {
input: 10.0,
output: 50.0,
cache_read: 1.0,
cache_write_5m: 12.5,
cache_write_1h: 20.0,
};
const OPUS_RATES: ModelRates = ModelRates {
input: 5.0,
output: 25.0,
cache_read: 0.5,
cache_write_5m: 6.25,
cache_write_1h: 10.0,
};
const SONNET_RATES: ModelRates = ModelRates {
input: 3.0,
output: 15.0,
cache_read: 0.3,
cache_write_5m: 3.75,
cache_write_1h: 6.0,
};
const HAIKU_RATES: ModelRates = ModelRates {
input: 1.0,
output: 5.0,
cache_read: 0.1,
cache_write_5m: 1.25,
cache_write_1h: 2.0,
};
pub fn rates_for_model(model: &str) -> Option<ModelRates> {
let lower = model.to_lowercase();
if lower.contains("fable") || lower.contains("mythos") {
Some(FABLE_RATES)
} else if lower.contains("opus") {
Some(OPUS_RATES)
} else if lower.contains("sonnet") {
Some(SONNET_RATES)
} else if lower.contains("haiku") {
Some(HAIKU_RATES)
} else {
None
}
}
#[derive(Debug, Clone, Default)]
pub struct RateOverrides {
entries: Vec<(String, ModelRates)>,
}
impl RateOverrides {
pub fn new(pairs: impl IntoIterator<Item = (String, ModelRates)>) -> Self {
let mut entries: Vec<(String, ModelRates)> = pairs
.into_iter()
.map(|(k, v)| (k.to_lowercase(), v))
.filter(|(k, _)| !k.is_empty())
.collect();
entries.sort_by(|a, b| b.0.len().cmp(&a.0.len()).then_with(|| a.0.cmp(&b.0)));
Self { entries }
}
pub const fn is_empty(&self) -> bool {
self.entries.is_empty()
}
fn lookup(&self, model_lower: &str) -> Option<ModelRates> {
self.entries
.iter()
.find(|(key, _)| model_lower.contains(key.as_str()))
.map(|(_, rates)| *rates)
}
}
pub fn resolve_rates(model: &str, overrides: &RateOverrides) -> Option<ModelRates> {
overrides
.lookup(&model.to_lowercase())
.or_else(|| rates_for_model(model))
}
#[cfg(not(coverage))]
static RATE_OVERRIDES: OnceLock<RateOverrides> = OnceLock::new();
#[cfg(not(coverage))]
pub fn set_rate_overrides(overrides: RateOverrides) {
let _ = RATE_OVERRIDES.set(overrides);
}
#[cfg(coverage)]
pub fn set_rate_overrides(_overrides: RateOverrides) {}
#[cfg(not(coverage))]
fn active_rate_overrides() -> &'static RateOverrides {
RATE_OVERRIDES.get_or_init(RateOverrides::default)
}
pub fn read_exact_usage(
command: &str,
workdir: &str,
since: DateTime<Utc>,
now: DateTime<Utc>,
claude_dir: &std::path::Path,
codex_dir: &std::path::Path,
rates: &RateOverrides,
) -> Option<ExactUsage> {
let since = since - TimeDelta::seconds(SINCE_GRACE_SECS);
match agent_provider_for_command(command)? {
"claude.ai" => claude::read_usage(claude_dir, workdir, since, rates),
"openai" => codex::read_usage(codex_dir, workdir, since, now),
_ => None,
}
}
#[allow(clippy::too_many_arguments)]
pub fn read_exact_usage_with_harness(
command: &str,
workdir: &str,
since: DateTime<Utc>,
now: DateTime<Utc>,
claude_dir: &std::path::Path,
codex_dir: &std::path::Path,
rates: &RateOverrides,
harness: Option<&str>,
harness_session_id: Option<&str>,
) -> Option<ExactUsage> {
if harness == Some("claude")
&& let Some(session_id) = harness_session_id
{
let exact_path = claude_dir
.join("projects")
.join(claude::sanitize_workdir(workdir))
.join(format!("{session_id}.jsonl"));
if exact_path.is_file() {
let since_with_grace = since - TimeDelta::seconds(SINCE_GRACE_SECS);
return claude::read_usage_file(
claude_dir,
workdir,
session_id,
since_with_grace,
rates,
);
}
}
read_exact_usage(command, workdir, since, now, claude_dir, codex_dir, rates)
}
pub fn effective_usage_dir(session: &Session) -> &str {
session.worktree_path.as_deref().unwrap_or(&session.workdir)
}
fn codex_dir_candidates(session: &Session, home: &Path, data_dir: &Path) -> Vec<PathBuf> {
let default_dir = home.join(".codex");
if session.harness.as_deref() == Some("codex") {
let per_session = data_dir
.join("harness")
.join(session.id.to_string())
.join("codex-home");
vec![per_session, default_dir]
} else {
vec![default_dir]
}
}
#[cfg(not(coverage))]
pub fn read_exact_usage_for_session(session: &Session, data_dir: &Path) -> Option<ExactUsage> {
let home = dirs::home_dir()?;
for codex_dir in codex_dir_candidates(session, &home, data_dir) {
if let Some(usage) = read_exact_usage_with_harness(
&session.command,
effective_usage_dir(session),
session.created_at,
Utc::now(),
&home.join(".claude"),
&codex_dir,
active_rate_overrides(),
session.harness.as_deref(),
session.harness_session_id.as_deref(),
) {
return Some(usage);
}
}
None
}
#[cfg(coverage)]
pub fn read_exact_usage_for_session(_session: &Session, _data_dir: &Path) -> Option<ExactUsage> {
None
}
pub(crate) fn codex_harness_home_dirs(data_dir: &Path) -> Vec<PathBuf> {
let Ok(entries) = std::fs::read_dir(data_dir.join("harness")) else {
return Vec::new();
};
entries
.flatten()
.map(|entry| entry.path().join("codex-home"))
.filter(|path| path.is_dir())
.collect()
}
#[cfg(not(coverage))]
pub fn scan_local_usage(
node_name: &str,
by_worktree: bool,
since_days: Option<u32>,
data_dir: &Path,
) -> Option<pulpo_common::api::UsageScanResponse> {
let home = dirs::home_dir()?;
let claude_dir = home.join(".claude");
let codex_dir = home.join(".codex");
let pi_dir = home.join(".pi");
let harness_codex_dirs = codex_harness_home_dirs(data_dir);
let codex_dirs: Vec<&Path> = std::iter::once(codex_dir.as_path())
.chain(harness_codex_dirs.iter().map(PathBuf::as_path))
.collect();
let dirs = scan::ScanDirs {
claude: &claude_dir,
codex: &codex_dirs,
pi: &pi_dir,
};
let rates = active_rate_overrides();
let now = Utc::now();
let resp = if by_worktree {
scan::scan_usage(&dirs, rates, node_name, now, since_days, |cwd: &str| {
cwd.to_owned()
})
} else {
scan::scan_usage(
&dirs,
rates,
node_name,
now,
since_days,
scan::canonical_repo,
)
};
Some(resp)
}
#[cfg(coverage)]
pub fn scan_local_usage(
_node_name: &str,
_by_worktree: bool,
_since_days: Option<u32>,
_data_dir: &Path,
) -> Option<pulpo_common::api::UsageScanResponse> {
None
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::Datelike;
#[test]
fn test_codex_harness_home_dirs_missing_harness_dir_is_empty() {
let tmp = tempfile::tempdir().unwrap();
assert!(codex_harness_home_dirs(tmp.path()).is_empty());
}
#[test]
fn test_codex_harness_home_dirs_collects_each_session_codex_home() {
let tmp = tempfile::tempdir().unwrap();
let harness_dir = tmp.path().join("harness");
let session_a = harness_dir.join("session-a").join("codex-home");
let session_b = harness_dir.join("session-b").join("codex-home");
std::fs::create_dir_all(&session_a).unwrap();
std::fs::create_dir_all(&session_b).unwrap();
std::fs::create_dir_all(harness_dir.join("session-c")).unwrap();
let mut dirs = codex_harness_home_dirs(tmp.path());
dirs.sort();
assert_eq!(dirs, vec![session_a, session_b]);
}
#[test]
#[allow(clippy::float_cmp)]
fn test_rates_for_model_families() {
assert_eq!(rates_for_model("claude-fable-5").unwrap().input, 10.0);
assert_eq!(rates_for_model("claude-mythos-5").unwrap().output, 50.0);
assert_eq!(rates_for_model("claude-opus-4-8").unwrap().input, 5.0);
assert_eq!(rates_for_model("claude-sonnet-4-6").unwrap().output, 15.0);
assert_eq!(rates_for_model("claude-haiku-4-5").unwrap().input, 1.0);
}
#[test]
fn test_rates_for_model_unknown() {
assert!(rates_for_model("gpt-5").is_none());
assert!(rates_for_model("").is_none());
}
#[test]
#[allow(clippy::float_cmp)]
fn test_resolve_rates_falls_back_to_builtin() {
let empty = RateOverrides::default();
assert!(empty.is_empty());
assert_eq!(resolve_rates("claude-opus-4-8", &empty).unwrap().input, 5.0);
assert!(resolve_rates("gpt-5", &empty).is_none());
assert!(resolve_rates("", &empty).is_none());
}
#[test]
#[allow(clippy::float_cmp)]
fn test_resolve_rates_override_beats_builtin_and_most_specific_wins() {
let rate = |input: f64| ModelRates {
input,
output: 0.0,
cache_read: 0.0,
cache_write_5m: 0.0,
cache_write_1h: 0.0,
};
let overrides = RateOverrides::new([
("opus".to_owned(), rate(1.0)),
("claude-opus-4-8".to_owned(), rate(7.0)),
]);
assert_eq!(
resolve_rates("claude-opus-4-8", &overrides).unwrap().input,
7.0
);
assert_eq!(
resolve_rates("claude-opus-4-9", &overrides).unwrap().input,
1.0
);
assert!(resolve_rates("gpt-5", &overrides).is_none());
}
#[test]
fn test_rate_overrides_new_drops_empty_keys_and_lowercases() {
let overrides = RateOverrides::new([
(String::new(), HAIKU_RATES),
("GPT-6".to_owned(), HAIKU_RATES),
]);
assert!(!overrides.is_empty());
assert!(resolve_rates("gpt-6-turbo", &overrides).is_some());
}
#[test]
fn test_set_rate_overrides_is_callable() {
set_rate_overrides(RateOverrides::new([("zzz-smoke".to_owned(), HAIKU_RATES)]));
}
#[test]
fn test_effective_usage_dir_prefers_worktree() {
let session = Session {
workdir: "/repo".into(),
worktree_path: Some("/home/u/.pulpo/worktrees/fix".into()),
..Default::default()
};
assert_eq!(
effective_usage_dir(&session),
"/home/u/.pulpo/worktrees/fix"
);
let plain = Session {
workdir: "/repo".into(),
worktree_path: None,
..Default::default()
};
assert_eq!(effective_usage_dir(&plain), "/repo");
}
#[test]
fn test_read_exact_usage_unknown_agent() {
let tmp = tempfile::tempdir().unwrap();
let result = read_exact_usage(
"cargo test",
"/tmp/repo",
Utc::now(),
Utc::now(),
tmp.path(),
tmp.path(),
&RateOverrides::default(),
);
assert!(result.is_none());
}
#[test]
fn test_read_exact_usage_gemini_has_no_reader() {
let tmp = tempfile::tempdir().unwrap();
let result = read_exact_usage(
"gemini chat",
"/tmp/repo",
Utc::now(),
Utc::now(),
tmp.path(),
tmp.path(),
&RateOverrides::default(),
);
assert!(result.is_none());
}
#[test]
fn test_read_exact_usage_claude_without_files() {
let tmp = tempfile::tempdir().unwrap();
let result = read_exact_usage(
"claude -p 'fix'",
"/tmp/repo",
Utc::now(),
Utc::now(),
tmp.path(),
tmp.path(),
&RateOverrides::default(),
);
assert!(result.is_none());
}
#[test]
fn test_read_exact_usage_codex_without_files() {
let tmp = tempfile::tempdir().unwrap();
let result = read_exact_usage(
"codex exec 'fix'",
"/tmp/repo",
Utc::now(),
Utc::now(),
tmp.path(),
tmp.path(),
&RateOverrides::default(),
);
assert!(result.is_none());
}
#[test]
fn test_read_exact_usage_with_harness_prefers_exact_file() {
let tmp = tempfile::tempdir().unwrap();
let workdir = "/tmp/repo";
let session_id = "44444444-4444-4444-4444-444444444444";
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
let project_dir = tmp
.path()
.join("projects")
.join(claude::sanitize_workdir(workdir));
std::fs::create_dir_all(&project_dir).unwrap();
std::fs::write(
project_dir.join(format!("{session_id}.jsonl")),
format!(
r#"{{"timestamp":"{ts}","requestId":"r1","type":"assistant","message":{{"id":"m1","model":"claude-opus-4-8","usage":{{"input_tokens":1000,"output_tokens":500}}}}}}"#
),
)
.unwrap();
std::fs::write(
project_dir.join("other-session.jsonl"),
format!(
r#"{{"timestamp":"{ts}","requestId":"r2","type":"assistant","message":{{"id":"m2","model":"claude-opus-4-8","usage":{{"input_tokens":9000,"output_tokens":9000}}}}}}"#
),
)
.unwrap();
let usage = read_exact_usage_with_harness(
"claude -p 'fix'",
workdir,
since,
Utc::now(),
tmp.path(),
tmp.path(),
&RateOverrides::default(),
Some("claude"),
Some(session_id),
)
.unwrap();
assert_eq!(usage.input_tokens, 1000);
assert_eq!(usage.output_tokens, 500);
}
#[test]
fn test_read_exact_usage_with_harness_falls_back_when_file_missing() {
let tmp = tempfile::tempdir().unwrap();
let workdir = "/tmp/repo";
let since = Utc::now() - TimeDelta::hours(1);
let ts = Utc::now().to_rfc3339();
let project_dir = tmp
.path()
.join("projects")
.join(claude::sanitize_workdir(workdir));
std::fs::create_dir_all(&project_dir).unwrap();
std::fs::write(
project_dir.join("some-other-name.jsonl"),
format!(
r#"{{"timestamp":"{ts}","requestId":"r1","type":"assistant","message":{{"id":"m1","model":"claude-opus-4-8","usage":{{"input_tokens":42,"output_tokens":7}}}}}}"#
),
)
.unwrap();
let usage = read_exact_usage_with_harness(
"claude -p 'fix'",
workdir,
since,
Utc::now(),
tmp.path(),
tmp.path(),
&RateOverrides::default(),
Some("claude"),
Some("55555555-5555-5555-5555-555555555555"),
)
.unwrap();
assert_eq!(usage.input_tokens, 42);
assert_eq!(usage.output_tokens, 7);
}
#[test]
fn test_read_exact_usage_with_harness_no_harness_id_uses_heuristic() {
let tmp = tempfile::tempdir().unwrap();
let result = read_exact_usage_with_harness(
"claude -p 'fix'",
"/tmp/repo",
Utc::now(),
Utc::now(),
tmp.path(),
tmp.path(),
&RateOverrides::default(),
None,
None,
);
assert!(result.is_none());
}
#[test]
fn test_read_exact_usage_for_session_no_agent_command() {
let session = Session {
command: "cargo build".into(),
..Default::default()
};
assert!(
read_exact_usage_for_session(&session, Path::new("/nonexistent-data-dir")).is_none()
);
}
#[test]
fn test_codex_dir_candidates_prefers_per_session_codex_home_for_codex_harness() {
let session = Session {
id: "33333333-3333-3333-3333-333333333333".parse().unwrap(),
harness: Some("codex".into()),
..Default::default()
};
let home = Path::new("/home/user");
let data_dir = Path::new("/data");
let candidates = codex_dir_candidates(&session, home, data_dir);
assert_eq!(
candidates,
vec![
PathBuf::from("/data/harness/33333333-3333-3333-3333-333333333333/codex-home"),
PathBuf::from("/home/user/.codex"),
]
);
}
#[test]
fn test_codex_dir_candidates_non_codex_harness_only_default() {
let session = Session {
harness: Some("claude".into()),
..Default::default()
};
let home = Path::new("/home/user");
let data_dir = Path::new("/data");
assert_eq!(
codex_dir_candidates(&session, home, data_dir),
vec![PathBuf::from("/home/user/.codex")]
);
}
#[test]
fn test_codex_dir_candidates_no_harness_only_default() {
let session = Session::default();
let home = Path::new("/home/user");
let data_dir = Path::new("/data");
assert_eq!(
codex_dir_candidates(&session, home, data_dir),
vec![PathBuf::from("/home/user/.codex")]
);
}
#[test]
fn test_read_exact_usage_with_harness_finds_rollout_under_per_session_codex_home() {
let tmp = tempfile::tempdir().unwrap();
let data_dir = tmp.path();
let session_id = "44444444-4444-4444-4444-444444444444";
let workdir = "/tmp/codex-usage-repo";
let session = Session {
id: session_id.parse().unwrap(),
harness: Some("codex".into()),
..Default::default()
};
let codex_home = codex_dir_candidates(&session, Path::new("/home/user"), data_dir)
.into_iter()
.next()
.unwrap();
let now = Utc::now();
let day_dir = codex_home
.join("sessions")
.join(format!("{:04}", now.year()))
.join(format!("{:02}", now.month()))
.join(format!("{:02}", now.day()));
std::fs::create_dir_all(&day_dir).unwrap();
let rollout = format!(
r#"{{"timestamp":"{ts}","type":"session_meta","payload":{{"id":"abc","timestamp":"{ts}","cwd":"{workdir}","originator":"codex_cli_rs"}}}}
{{"timestamp":"{ts}","type":"event_msg","payload":{{"type":"token_count","info":{{"total_token_usage":{{"input_tokens":100,"cached_input_tokens":0,"output_tokens":10,"total_tokens":110}}}}}}}}
"#,
ts = now.to_rfc3339()
);
std::fs::write(day_dir.join("rollout-a.jsonl"), rollout).unwrap();
let usage = read_exact_usage_with_harness(
"codex -p hi",
workdir,
now - TimeDelta::hours(1),
now,
Path::new("/nonexistent-claude-dir"),
&codex_home,
&RateOverrides::default(),
Some("codex"),
None,
)
.unwrap();
assert_eq!(usage.source, SOURCE_CODEX);
assert_eq!(usage.input_tokens, 100);
assert_eq!(usage.output_tokens, 10);
}
}