use std::collections::{BTreeMap, BTreeSet, HashMap};
use std::io::{Read, Seek, SeekFrom};
use std::path::Path;
use serde_json::{json, Map, Value};
use supercode_interchange::sidecar::{ms_to_rfc3339, rfc3339_to_ms};
use crate::pricing::{built_in, RecordedTokens};
use crate::HarnessHomes;
#[derive(Debug, Clone)]
pub struct UseRecord {
pub at: Option<String>,
pub at_ms: Option<i64>,
pub model: String,
pub tokens: RecordedTokens,
pub fast: bool,
pub us_only: bool,
}
#[derive(Debug, Clone, Default)]
pub struct Spend {
pub tokens: RecordedTokens,
pub cost_usd: f64,
pub priced: bool,
pub unpriced: bool,
}
impl Spend {
pub fn record(&mut self, record: &UseRecord) {
self.tokens.add(&record.tokens);
match built_in(&record.model) {
Some(price) => {
self.cost_usd +=
price.recorded_cost_usd(&record.tokens, record.fast, record.us_only);
self.priced = true;
}
None => self.unpriced = true,
}
}
pub fn add(&mut self, other: &Spend) {
self.tokens.add(&other.tokens);
self.cost_usd += other.cost_usd;
self.priced |= other.priced;
self.unpriced |= other.unpriced;
}
pub fn fields(&self) -> Map<String, Value> {
let t = &self.tokens;
let mut fields = Map::new();
fields.insert("input_tokens".into(), json!(t.input));
fields.insert("output_tokens".into(), json!(t.output));
fields.insert("cache_read_tokens".into(), json!(t.cache_read));
fields.insert(
"cache_write_tokens".into(),
json!(t.cache_write_5m + t.cache_write_1h),
);
fields.insert("cache_write_5m_tokens".into(), json!(t.cache_write_5m));
fields.insert("cache_write_1h_tokens".into(), json!(t.cache_write_1h));
if self.priced {
fields.insert(
"cost".into(),
json!({"amount": (self.cost_usd * 1e6).round() / 1e6, "currency": "USD", "source": "rate_table"}),
);
}
if self.unpriced {
fields.insert("cost_unrecorded".into(), json!(true));
}
fields
}
}
pub fn claude_use(record: &Value) -> Option<(String, UseRecord)> {
let message = record.get("message")?;
let id = message.get("id")?.as_str()?;
let usage = message.get("usage")?;
let model = message
.get("model")
.and_then(Value::as_str)
.unwrap_or("unknown");
if model == "<synthetic>" {
return None;
}
let n = |value: Option<&Value>| value.and_then(Value::as_u64).unwrap_or(0);
let written = n(usage.get("cache_creation_input_tokens"));
let (write_5m, write_1h) = match usage.get("cache_creation") {
Some(split) => {
let write_5m = n(split.get("ephemeral_5m_input_tokens"));
let write_1h = n(split.get("ephemeral_1h_input_tokens"));
(
write_5m + written.saturating_sub(write_5m + write_1h),
write_1h,
)
}
None => (written, 0),
};
let at = record
.get("timestamp")
.and_then(Value::as_str)
.map(str::to_string);
Some((
id.to_string(),
UseRecord {
at_ms: at.as_deref().and_then(rfc3339_to_ms),
at,
model: model.to_string(),
tokens: RecordedTokens {
input: n(usage.get("input_tokens")),
output: n(usage.get("output_tokens")),
cache_read: n(usage.get("cache_read_input_tokens")),
cache_write_5m: write_5m,
cache_write_1h: write_1h,
},
fast: usage.get("speed").and_then(Value::as_str) == Some("fast"),
us_only: usage.get("inference_geo").and_then(Value::as_str) == Some("us"),
},
))
}
#[derive(Debug, Clone)]
pub struct CodexUse {
model: String,
previous: [u64; 3],
}
impl CodexUse {
pub fn new(model: Option<String>) -> Self {
CodexUse {
model: model.unwrap_or_else(|| "unknown".into()),
previous: [0; 3],
}
}
pub fn next(&mut self, record: &Value) -> Option<UseRecord> {
let payload = record.get("payload").unwrap_or(&Value::Null);
if record.get("type").and_then(Value::as_str) == Some("turn_context") {
if let Some(model) = payload.get("model").and_then(Value::as_str) {
self.model = model.to_string();
}
return None;
}
if payload.get("type").and_then(Value::as_str) != Some("token_count") {
return None;
}
let total = payload.get("info")?.get("total_token_usage")?;
let n = |value: Option<&Value>| value.and_then(Value::as_u64).unwrap_or(0);
let now = [
n(total.get("input_tokens")),
n(total.get("output_tokens")),
n(total.get("cached_input_tokens")),
];
if now == self.previous {
return None;
}
let [input, output, cached] = [0, 1, 2].map(|i| now[i].saturating_sub(self.previous[i]));
self.previous = now;
let at = record
.get("timestamp")
.and_then(Value::as_str)
.map(str::to_string);
Some(UseRecord {
at_ms: at.as_deref().and_then(rfc3339_to_ms),
at,
model: self.model.clone(),
tokens: RecordedTokens {
input: input.saturating_sub(cached),
output,
cache_read: cached,
..RecordedTokens::default()
},
fast: false,
us_only: false,
})
}
}
pub fn session_spend(
source: &crate::SessionSource,
model: Option<String>,
raw: &[String],
since_ms: Option<i64>,
) -> Vec<Value> {
let mut days: BTreeMap<(String, String), Spend> = BTreeMap::new();
let mut count = |record: &UseRecord| {
if since_ms.is_some_and(|since| record.at_ms.is_none_or(|at| at < since)) {
return;
}
let day = record
.at
.as_deref()
.and_then(|at| at.get(..10))
.unwrap_or("unknown")
.to_string();
days.entry((day, record.model.clone()))
.or_default()
.record(record);
};
let records = raw
.iter()
.filter_map(|line| serde_json::from_str::<Value>(line).ok());
match source {
crate::SessionSource::ClaudeCode => {
let responses: BTreeMap<String, UseRecord> =
records.filter_map(|record| claude_use(&record)).collect();
responses.values().for_each(&mut count);
}
crate::SessionSource::Codex => {
let mut codex = CodexUse::new(model);
for record in records {
if let Some(used) = codex.next(&record) {
count(&used);
}
}
}
_ => {}
}
days.into_iter()
.map(|((day, model), spend)| {
let mut row = Map::new();
row.insert("at".into(), json!(format!("{day}T00:00:00Z")));
row.insert("model".into(), json!(model));
row.extend(spend.fields());
Value::Object(row)
})
.collect()
}
pub fn parse_since(text: &str, now_ms: i64) -> Result<i64, String> {
let text = text.trim();
if let Some(at) = rfc3339_to_ms(text) {
return Ok(at);
}
let (number, unit) = text.split_at(text.find(|c: char| !c.is_ascii_digit()).unwrap_or(text.len()));
let amount: i64 = number
.parse()
.map_err(|_| format!("`{text}` is neither a duration (10m, 1h, 2d) nor an RFC 3339 moment"))?;
let unit_ms = match unit {
"s" => 1_000,
"m" => 60_000,
"h" => 3_600_000,
"d" => 86_400_000,
_ => {
return Err(format!(
"`{text}`: a duration's unit is s, m, h or d (10m, 1h, 2d)"
))
}
};
Ok(now_ms - amount * unit_ms)
}
#[derive(Debug, Default)]
struct SessionSpend {
by_model: BTreeMap<String, Spend>,
cwd: Option<String>,
subagents: BTreeSet<String>,
last_ms: i64,
}
impl SessionSpend {
fn record(&mut self, record: &UseRecord) {
self.by_model
.entry(record.model.clone())
.or_default()
.record(record);
self.last_ms = self.last_ms.max(record.at_ms.unwrap_or_default());
}
fn total(&self) -> Spend {
let mut total = Spend::default();
for spend in self.by_model.values() {
total.add(spend);
}
total
}
}
pub fn machine_spend(
homes: &HarnessHomes,
since_ms: i64,
now_ms: i64,
harnesses: &[String],
sort: &str,
limit: Option<usize>,
) -> Value {
let wants = |harness: &str| harnesses.is_empty() || harnesses.iter().any(|h| h == harness);
let mut sessions: BTreeMap<(String, String), SessionSpend> = BTreeMap::new();
if wants(crate::HarnessId::CLAUDE_CODE) {
claude_spend(&homes.claude_code, since_ms, &mut sessions);
}
if wants(crate::HarnessId::CODEX) {
codex_spend(&homes.codex, since_ms, &mut sessions);
}
let names: HashMap<String, String> =
crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes))
.into_iter()
.map(|session| (session.session_id, session.name))
.collect();
let machine = crate::mailbox::local_machine_name();
let mut rows: Vec<(Spend, Value)> = sessions
.into_iter()
.map(|((harness, session_id), spend)| {
let total = spend.total();
let name = match harness.as_str() {
crate::HarnessId::CODEX => Some(crate::mail_route::codex_name(&session_id)),
_ => names.get(&session_id).cloned(),
};
let mut row = Map::new();
row.insert("harness".into(), json!(harness));
row.insert("session_id".into(), json!(session_id));
row.insert(
"address".into(),
json!(crate::mailbox::MailAddress::new(&machine, &harness, &session_id)
.ok()
.map(|address| address.to_string())),
);
row.insert("name".into(), json!(name));
row.insert("cwd".into(), json!(spend.cwd));
row.insert("subagents".into(), json!(spend.subagents.len()));
row.insert(
"last_at".into(),
json!((spend.last_ms > 0).then(|| ms_to_rfc3339(spend.last_ms))),
);
row.extend(total.fields());
row.insert(
"models".into(),
Value::Array(
spend
.by_model
.iter()
.map(|(model, spend)| {
let mut entry = Map::new();
entry.insert("model".into(), json!(model));
entry.extend(spend.fields());
Value::Object(entry)
})
.collect(),
),
);
(total, Value::Object(row))
})
.collect();
let by_tokens = sort == "tokens";
rows.sort_by(|(a, _), (b, _)| {
if by_tokens || a.priced != b.priced || !a.priced {
(b.priced && !by_tokens)
.cmp(&(a.priced && !by_tokens))
.then(b.tokens.total().cmp(&a.tokens.total()))
} else {
b.cost_usd.total_cmp(&a.cost_usd)
}
});
let mut total = Spend::default();
for (spend, _) in &rows {
total.add(spend);
}
let count = rows.len();
let listed: Vec<Value> = rows
.into_iter()
.take(limit.unwrap_or(usize::MAX))
.map(|(_, row)| row)
.collect();
let mut answer = Map::new();
answer.insert("since".into(), json!(ms_to_rfc3339(since_ms)));
answer.insert("until".into(), json!(ms_to_rfc3339(now_ms)));
answer.insert("session_count".into(), json!(count));
answer.insert("sessions".into(), Value::Array(listed));
answer.insert("total".into(), Value::Object(total.fields()));
Value::Object(answer)
}
const OUT_OF_ORDER_MS: i64 = 10 * 60 * 1000;
fn modified_since(path: &Path, since_ms: i64) -> bool {
std::fs::metadata(path)
.and_then(|meta| meta.modified())
.ok()
.and_then(|at| at.duration_since(std::time::UNIX_EPOCH).ok())
.is_some_and(|at| at.as_millis() as i64 >= since_ms)
}
fn record_ms(record: &Value) -> Option<i64> {
record
.get("timestamp")
.and_then(Value::as_str)
.and_then(rfc3339_to_ms)
}
fn claude_spend(
projects: &Path,
since_ms: i64,
sessions: &mut BTreeMap<(String, String), SessionSpend>,
) {
let Ok(projects) = std::fs::read_dir(projects) else {
return;
};
for project in projects.flatten() {
let Ok(entries) = std::fs::read_dir(project.path()) else {
continue;
};
for entry in entries.flatten() {
let path = entry.path();
if path.is_dir() {
let Some(parent) = path.file_name().and_then(|name| name.to_str()) else {
continue;
};
let Ok(subagents) = std::fs::read_dir(path.join("subagents")) else {
continue;
};
for subagent in subagents.flatten() {
let file = subagent.path();
if file.extension().is_some_and(|ext| ext == "jsonl")
&& modified_since(&file, since_ms)
{
let agent = file
.file_stem()
.map(|stem| stem.to_string_lossy().into_owned());
claude_file(&file, since_ms, parent, agent, sessions);
}
}
} else if path.extension().is_some_and(|ext| ext == "jsonl")
&& modified_since(&path, since_ms)
{
if let Some(session) = path.file_stem().and_then(|stem| stem.to_str()) {
claude_file(&path, since_ms, session, None, sessions);
}
}
}
}
}
fn claude_file(
path: &Path,
since_ms: i64,
session_id: &str,
subagent: Option<String>,
sessions: &mut BTreeMap<(String, String), SessionSpend>,
) {
let records = tail_records(path, |record| {
record_ms(record).is_some_and(|at| at < since_ms - OUT_OF_ORDER_MS)
});
let cwd = records
.iter()
.rev()
.find_map(|record| record.get("cwd").and_then(Value::as_str))
.map(str::to_string);
let responses: BTreeMap<String, UseRecord> = records
.iter()
.filter_map(claude_use)
.filter(|(_, used)| used.at_ms.is_some_and(|at| at >= since_ms))
.collect();
if responses.is_empty() {
return;
}
let session = sessions
.entry((crate::HarnessId::CLAUDE_CODE.to_string(), session_id.to_string()))
.or_default();
for used in responses.values() {
session.record(used);
}
if subagent.is_none() || session.cwd.is_none() {
session.cwd = cwd.or(session.cwd.take());
}
if let Some(agent) = subagent {
session.subagents.insert(agent);
}
}
fn codex_spend(
root: &Path,
since_ms: i64,
sessions: &mut BTreeMap<(String, String), SessionSpend>,
) {
let mut stack = vec![(root.to_path_buf(), 0)];
while let Some((directory, depth)) = stack.pop() {
let Ok(entries) = std::fs::read_dir(&directory) else {
continue;
};
for entry in entries.flatten() {
let path = entry.path();
if path.is_dir() {
if depth < 4 {
stack.push((path, depth + 1));
}
continue;
}
let is_rollout = path
.file_name()
.and_then(|name| name.to_str())
.is_some_and(|name| name.starts_with("rollout-") && name.ends_with(".jsonl"));
if is_rollout && modified_since(&path, since_ms) {
codex_file(&path, since_ms, sessions);
}
}
}
}
fn codex_file(path: &Path, since_ms: i64, sessions: &mut BTreeMap<(String, String), SessionSpend>) {
let Some((thread, parent)) = crate::codex_peer::rollout_session(path) else {
return;
};
let (mut total_before, mut model_before) = (false, false);
let records = tail_records(path, |record| {
if record_ms(record).is_some_and(|at| at < since_ms - OUT_OF_ORDER_MS) {
let kind = record.get("type").and_then(Value::as_str);
total_before |= record.pointer("/payload/type").and_then(Value::as_str) == Some("token_count");
model_before |= kind == Some("turn_context");
}
total_before && model_before
});
let mut codex = CodexUse::new(None);
let used: Vec<UseRecord> = records
.iter()
.filter_map(|record| codex.next(record))
.filter(|used| used.at_ms.is_some_and(|at| at >= since_ms))
.collect();
if used.is_empty() {
return;
}
let session_id = parent.clone().unwrap_or_else(|| thread.clone());
let session = sessions
.entry((crate::HarnessId::CODEX.to_string(), session_id))
.or_default();
for record in &used {
session.record(record);
}
if parent.is_some() {
session.subagents.insert(thread);
}
if session.cwd.is_none() || parent.is_none() {
if let Some(cwd) = crate::codex_peer::rollout_cwd(path) {
session.cwd = Some(cwd.to_string_lossy().into_owned());
}
}
}
fn tail_records(path: &Path, mut start: impl FnMut(&Value) -> bool) -> Vec<Value> {
const BLOCK: u64 = 256 * 1024;
let Ok(mut file) = std::fs::File::open(path) else {
return Vec::new();
};
let Ok(length) = file.metadata().map(|meta| meta.len()) else {
return Vec::new();
};
let mut newest_first = Vec::new();
let mut carry: Vec<u8> = Vec::new();
let mut end = length;
while end > 0 {
let begin = end.saturating_sub(BLOCK);
let mut block = vec![0u8; (end - begin) as usize];
if file.seek(SeekFrom::Start(begin)).is_err() || file.read_exact(&mut block).is_err() {
break;
}
block.extend_from_slice(&carry);
let first_newline = if begin == 0 {
None
} else {
match block.iter().position(|byte| *byte == b'\n') {
Some(at) => Some(at),
None => {
carry = block;
end = begin;
continue;
}
}
};
let (head, lines) = match first_newline {
Some(at) => (block[..at].to_vec(), &block[at + 1..]),
None => (Vec::new(), &block[..]),
};
let mut reached = false;
for line in lines.split(|byte| *byte == b'\n').rev() {
let Ok(record) = serde_json::from_slice::<Value>(line) else {
continue;
};
reached = start(&record);
newest_first.push(record);
if reached {
break;
}
}
if reached {
break;
}
carry = head;
end = begin;
}
newest_first.reverse();
newest_first
}