1use std::collections::{BTreeMap, BTreeSet, HashMap};
11use std::io::{Read, Seek, SeekFrom};
12use std::path::Path;
13
14use serde_json::{json, Map, Value};
15use supercode_interchange::sidecar::{ms_to_rfc3339, rfc3339_to_ms};
16
17use crate::pricing::{built_in, RecordedTokens};
18use crate::HarnessHomes;
19
20#[derive(Debug, Clone)]
22pub struct UseRecord {
23 pub at: Option<String>,
25 pub at_ms: Option<i64>,
27 pub model: String,
29 pub tokens: RecordedTokens,
31 pub fast: bool,
33 pub us_only: bool,
35}
36
37#[derive(Debug, Clone, Default)]
39pub struct Spend {
40 pub tokens: RecordedTokens,
42 pub cost_usd: f64,
44 pub priced: bool,
46 pub unpriced: bool,
48}
49
50impl Spend {
51 pub fn record(&mut self, record: &UseRecord) {
53 self.tokens.add(&record.tokens);
54 match built_in(&record.model) {
55 Some(price) => {
56 self.cost_usd +=
57 price.recorded_cost_usd(&record.tokens, record.fast, record.us_only);
58 self.priced = true;
59 }
60 None => self.unpriced = true,
61 }
62 }
63
64 pub fn add(&mut self, other: &Spend) {
66 self.tokens.add(&other.tokens);
67 self.cost_usd += other.cost_usd;
68 self.priced |= other.priced;
69 self.unpriced |= other.unpriced;
70 }
71
72 pub fn fields(&self) -> Map<String, Value> {
75 let t = &self.tokens;
76 let mut fields = Map::new();
77 fields.insert("input_tokens".into(), json!(t.input));
78 fields.insert("output_tokens".into(), json!(t.output));
79 fields.insert("cache_read_tokens".into(), json!(t.cache_read));
80 fields.insert(
81 "cache_write_tokens".into(),
82 json!(t.cache_write_5m + t.cache_write_1h),
83 );
84 fields.insert("cache_write_5m_tokens".into(), json!(t.cache_write_5m));
85 fields.insert("cache_write_1h_tokens".into(), json!(t.cache_write_1h));
86 if self.priced {
87 fields.insert(
88 "cost".into(),
89 json!({"amount": (self.cost_usd * 1e6).round() / 1e6, "currency": "USD", "source": "rate_table"}),
90 );
91 }
92 if self.unpriced {
93 fields.insert("cost_unrecorded".into(), json!(true));
94 }
95 fields
96 }
97}
98
99pub fn claude_use(record: &Value) -> Option<(String, UseRecord)> {
104 let message = record.get("message")?;
105 let id = message.get("id")?.as_str()?;
106 let usage = message.get("usage")?;
107 let model = message
108 .get("model")
109 .and_then(Value::as_str)
110 .unwrap_or("unknown");
111 if model == "<synthetic>" {
112 return None;
113 }
114 let n = |value: Option<&Value>| value.and_then(Value::as_u64).unwrap_or(0);
115 let written = n(usage.get("cache_creation_input_tokens"));
116 let (write_5m, write_1h) = match usage.get("cache_creation") {
117 Some(split) => {
118 let write_5m = n(split.get("ephemeral_5m_input_tokens"));
119 let write_1h = n(split.get("ephemeral_1h_input_tokens"));
120 (
121 write_5m + written.saturating_sub(write_5m + write_1h),
122 write_1h,
123 )
124 }
125 None => (written, 0),
126 };
127 let at = record
128 .get("timestamp")
129 .and_then(Value::as_str)
130 .map(str::to_string);
131 Some((
132 id.to_string(),
133 UseRecord {
134 at_ms: at.as_deref().and_then(rfc3339_to_ms),
135 at,
136 model: model.to_string(),
137 tokens: RecordedTokens {
138 input: n(usage.get("input_tokens")),
139 output: n(usage.get("output_tokens")),
140 cache_read: n(usage.get("cache_read_input_tokens")),
141 cache_write_5m: write_5m,
142 cache_write_1h: write_1h,
143 },
144 fast: usage.get("speed").and_then(Value::as_str) == Some("fast"),
145 us_only: usage.get("inference_geo").and_then(Value::as_str) == Some("us"),
146 },
147 ))
148}
149
150#[derive(Debug, Clone)]
154pub struct CodexUse {
155 model: String,
156 previous: [u64; 3],
157}
158
159impl CodexUse {
160 pub fn new(model: Option<String>) -> Self {
162 CodexUse {
163 model: model.unwrap_or_else(|| "unknown".into()),
164 previous: [0; 3],
165 }
166 }
167
168 pub fn next(&mut self, record: &Value) -> Option<UseRecord> {
170 let payload = record.get("payload").unwrap_or(&Value::Null);
171 if record.get("type").and_then(Value::as_str) == Some("turn_context") {
172 if let Some(model) = payload.get("model").and_then(Value::as_str) {
173 self.model = model.to_string();
174 }
175 return None;
176 }
177 if payload.get("type").and_then(Value::as_str) != Some("token_count") {
178 return None;
179 }
180 let total = payload.get("info")?.get("total_token_usage")?;
181 let n = |value: Option<&Value>| value.and_then(Value::as_u64).unwrap_or(0);
182 let now = [
183 n(total.get("input_tokens")),
184 n(total.get("output_tokens")),
185 n(total.get("cached_input_tokens")),
186 ];
187 if now == self.previous {
188 return None;
189 }
190 let [input, output, cached] = [0, 1, 2].map(|i| now[i].saturating_sub(self.previous[i]));
191 self.previous = now;
192 let at = record
193 .get("timestamp")
194 .and_then(Value::as_str)
195 .map(str::to_string);
196 Some(UseRecord {
197 at_ms: at.as_deref().and_then(rfc3339_to_ms),
198 at,
199 model: self.model.clone(),
200 tokens: RecordedTokens {
201 input: input.saturating_sub(cached),
202 output,
203 cache_read: cached,
204 ..RecordedTokens::default()
205 },
206 fast: false,
207 us_only: false,
208 })
209 }
210}
211
212pub fn session_spend(
214 source: &crate::SessionSource,
215 model: Option<String>,
216 raw: &[String],
217 since_ms: Option<i64>,
218) -> Vec<Value> {
219 let mut days: BTreeMap<(String, String), Spend> = BTreeMap::new();
220 let mut count = |record: &UseRecord| {
221 if since_ms.is_some_and(|since| record.at_ms.is_none_or(|at| at < since)) {
222 return;
223 }
224 let day = record
225 .at
226 .as_deref()
227 .and_then(|at| at.get(..10))
228 .unwrap_or("unknown")
229 .to_string();
230 days.entry((day, record.model.clone()))
231 .or_default()
232 .record(record);
233 };
234 let records = raw
235 .iter()
236 .filter_map(|line| serde_json::from_str::<Value>(line).ok());
237 match source {
238 crate::SessionSource::ClaudeCode => {
239 let responses: BTreeMap<String, UseRecord> =
240 records.filter_map(|record| claude_use(&record)).collect();
241 responses.values().for_each(&mut count);
242 }
243 crate::SessionSource::Codex => {
244 let mut codex = CodexUse::new(model);
245 for record in records {
246 if let Some(used) = codex.next(&record) {
247 count(&used);
248 }
249 }
250 }
251 _ => {}
252 }
253 days.into_iter()
254 .map(|((day, model), spend)| {
255 let mut row = Map::new();
256 row.insert("at".into(), json!(format!("{day}T00:00:00Z")));
257 row.insert("model".into(), json!(model));
258 row.extend(spend.fields());
259 Value::Object(row)
260 })
261 .collect()
262}
263
264pub fn parse_since(text: &str, now_ms: i64) -> Result<i64, String> {
266 let text = text.trim();
267 if let Some(at) = rfc3339_to_ms(text) {
268 return Ok(at);
269 }
270 let (number, unit) = text.split_at(text.find(|c: char| !c.is_ascii_digit()).unwrap_or(text.len()));
271 let amount: i64 = number
272 .parse()
273 .map_err(|_| format!("`{text}` is neither a duration (10m, 1h, 2d) nor an RFC 3339 moment"))?;
274 let unit_ms = match unit {
275 "s" => 1_000,
276 "m" => 60_000,
277 "h" => 3_600_000,
278 "d" => 86_400_000,
279 _ => {
280 return Err(format!(
281 "`{text}`: a duration's unit is s, m, h or d (10m, 1h, 2d)"
282 ))
283 }
284 };
285 Ok(now_ms - amount * unit_ms)
286}
287
288#[derive(Debug, Default)]
290struct SessionSpend {
291 by_model: BTreeMap<String, Spend>,
292 cwd: Option<String>,
293 subagents: BTreeSet<String>,
294 last_ms: i64,
295}
296
297impl SessionSpend {
298 fn record(&mut self, record: &UseRecord) {
299 self.by_model
300 .entry(record.model.clone())
301 .or_default()
302 .record(record);
303 self.last_ms = self.last_ms.max(record.at_ms.unwrap_or_default());
304 }
305
306 fn total(&self) -> Spend {
307 let mut total = Spend::default();
308 for spend in self.by_model.values() {
309 total.add(spend);
310 }
311 total
312 }
313}
314
315pub fn machine_spend(
319 homes: &HarnessHomes,
320 since_ms: i64,
321 now_ms: i64,
322 harnesses: &[String],
323 sort: &str,
324 limit: Option<usize>,
325) -> Value {
326 let wants = |harness: &str| harnesses.is_empty() || harnesses.iter().any(|h| h == harness);
327 let mut sessions: BTreeMap<(String, String), SessionSpend> = BTreeMap::new();
328 if wants(crate::HarnessId::CLAUDE_CODE) {
329 claude_spend(&homes.claude_code, since_ms, &mut sessions);
330 }
331 if wants(crate::HarnessId::CODEX) {
332 codex_spend(&homes.codex, since_ms, &mut sessions);
333 }
334 let names: HashMap<String, String> =
335 crate::claude_peer::read_registry(&crate::claude_peer::registry_dir(homes))
336 .into_iter()
337 .map(|session| (session.session_id, session.name))
338 .collect();
339 let machine = crate::mailbox::local_machine_name();
340 let mut rows: Vec<(Spend, Value)> = sessions
341 .into_iter()
342 .map(|((harness, session_id), spend)| {
343 let total = spend.total();
344 let name = match harness.as_str() {
345 crate::HarnessId::CODEX => Some(crate::mail_route::codex_name(&session_id)),
346 _ => names.get(&session_id).cloned(),
347 };
348 let mut row = Map::new();
349 row.insert("harness".into(), json!(harness));
350 row.insert("session_id".into(), json!(session_id));
351 row.insert(
352 "address".into(),
353 json!(crate::mailbox::MailAddress::new(&machine, &harness, &session_id)
354 .ok()
355 .map(|address| address.to_string())),
356 );
357 row.insert("name".into(), json!(name));
358 row.insert("cwd".into(), json!(spend.cwd));
359 row.insert("subagents".into(), json!(spend.subagents.len()));
360 row.insert(
361 "last_at".into(),
362 json!((spend.last_ms > 0).then(|| ms_to_rfc3339(spend.last_ms))),
363 );
364 row.extend(total.fields());
365 row.insert(
366 "models".into(),
367 Value::Array(
368 spend
369 .by_model
370 .iter()
371 .map(|(model, spend)| {
372 let mut entry = Map::new();
373 entry.insert("model".into(), json!(model));
374 entry.extend(spend.fields());
375 Value::Object(entry)
376 })
377 .collect(),
378 ),
379 );
380 (total, Value::Object(row))
381 })
382 .collect();
383 let by_tokens = sort == "tokens";
384 rows.sort_by(|(a, _), (b, _)| {
385 if by_tokens || a.priced != b.priced || !a.priced {
386 (b.priced && !by_tokens)
387 .cmp(&(a.priced && !by_tokens))
388 .then(b.tokens.total().cmp(&a.tokens.total()))
389 } else {
390 b.cost_usd.total_cmp(&a.cost_usd)
391 }
392 });
393 let mut total = Spend::default();
394 for (spend, _) in &rows {
395 total.add(spend);
396 }
397 let count = rows.len();
398 let listed: Vec<Value> = rows
399 .into_iter()
400 .take(limit.unwrap_or(usize::MAX))
401 .map(|(_, row)| row)
402 .collect();
403 let mut answer = Map::new();
404 answer.insert("since".into(), json!(ms_to_rfc3339(since_ms)));
405 answer.insert("until".into(), json!(ms_to_rfc3339(now_ms)));
406 answer.insert("session_count".into(), json!(count));
407 answer.insert("sessions".into(), Value::Array(listed));
408 answer.insert("total".into(), Value::Object(total.fields()));
409 Value::Object(answer)
410}
411
412const OUT_OF_ORDER_MS: i64 = 10 * 60 * 1000;
417
418fn modified_since(path: &Path, since_ms: i64) -> bool {
420 std::fs::metadata(path)
421 .and_then(|meta| meta.modified())
422 .ok()
423 .and_then(|at| at.duration_since(std::time::UNIX_EPOCH).ok())
424 .is_some_and(|at| at.as_millis() as i64 >= since_ms)
425}
426
427fn record_ms(record: &Value) -> Option<i64> {
428 record
429 .get("timestamp")
430 .and_then(Value::as_str)
431 .and_then(rfc3339_to_ms)
432}
433
434fn claude_spend(
437 projects: &Path,
438 since_ms: i64,
439 sessions: &mut BTreeMap<(String, String), SessionSpend>,
440) {
441 let Ok(projects) = std::fs::read_dir(projects) else {
442 return;
443 };
444 for project in projects.flatten() {
445 let Ok(entries) = std::fs::read_dir(project.path()) else {
446 continue;
447 };
448 for entry in entries.flatten() {
449 let path = entry.path();
450 if path.is_dir() {
451 let Some(parent) = path.file_name().and_then(|name| name.to_str()) else {
452 continue;
453 };
454 let Ok(subagents) = std::fs::read_dir(path.join("subagents")) else {
455 continue;
456 };
457 for subagent in subagents.flatten() {
458 let file = subagent.path();
459 if file.extension().is_some_and(|ext| ext == "jsonl")
460 && modified_since(&file, since_ms)
461 {
462 let agent = file
463 .file_stem()
464 .map(|stem| stem.to_string_lossy().into_owned());
465 claude_file(&file, since_ms, parent, agent, sessions);
466 }
467 }
468 } else if path.extension().is_some_and(|ext| ext == "jsonl")
469 && modified_since(&path, since_ms)
470 {
471 if let Some(session) = path.file_stem().and_then(|stem| stem.to_str()) {
472 claude_file(&path, since_ms, session, None, sessions);
473 }
474 }
475 }
476 }
477}
478
479fn claude_file(
480 path: &Path,
481 since_ms: i64,
482 session_id: &str,
483 subagent: Option<String>,
484 sessions: &mut BTreeMap<(String, String), SessionSpend>,
485) {
486 let records = tail_records(path, |record| {
487 record_ms(record).is_some_and(|at| at < since_ms - OUT_OF_ORDER_MS)
488 });
489 let cwd = records
490 .iter()
491 .rev()
492 .find_map(|record| record.get("cwd").and_then(Value::as_str))
493 .map(str::to_string);
494 let responses: BTreeMap<String, UseRecord> = records
495 .iter()
496 .filter_map(claude_use)
497 .filter(|(_, used)| used.at_ms.is_some_and(|at| at >= since_ms))
498 .collect();
499 if responses.is_empty() {
500 return;
501 }
502 let session = sessions
503 .entry((crate::HarnessId::CLAUDE_CODE.to_string(), session_id.to_string()))
504 .or_default();
505 for used in responses.values() {
506 session.record(used);
507 }
508 if subagent.is_none() || session.cwd.is_none() {
510 session.cwd = cwd.or(session.cwd.take());
511 }
512 if let Some(agent) = subagent {
513 session.subagents.insert(agent);
514 }
515}
516
517fn codex_spend(
519 root: &Path,
520 since_ms: i64,
521 sessions: &mut BTreeMap<(String, String), SessionSpend>,
522) {
523 let mut stack = vec![(root.to_path_buf(), 0)];
524 while let Some((directory, depth)) = stack.pop() {
525 let Ok(entries) = std::fs::read_dir(&directory) else {
526 continue;
527 };
528 for entry in entries.flatten() {
529 let path = entry.path();
530 if path.is_dir() {
531 if depth < 4 {
532 stack.push((path, depth + 1));
533 }
534 continue;
535 }
536 let is_rollout = path
537 .file_name()
538 .and_then(|name| name.to_str())
539 .is_some_and(|name| name.starts_with("rollout-") && name.ends_with(".jsonl"));
540 if is_rollout && modified_since(&path, since_ms) {
541 codex_file(&path, since_ms, sessions);
542 }
543 }
544 }
545}
546
547fn codex_file(path: &Path, since_ms: i64, sessions: &mut BTreeMap<(String, String), SessionSpend>) {
548 let Some((thread, parent)) = crate::codex_peer::rollout_session(path) else {
549 return;
550 };
551 let (mut total_before, mut model_before) = (false, false);
553 let records = tail_records(path, |record| {
554 if record_ms(record).is_some_and(|at| at < since_ms - OUT_OF_ORDER_MS) {
555 let kind = record.get("type").and_then(Value::as_str);
556 total_before |= record.pointer("/payload/type").and_then(Value::as_str) == Some("token_count");
557 model_before |= kind == Some("turn_context");
558 }
559 total_before && model_before
560 });
561 let mut codex = CodexUse::new(None);
562 let used: Vec<UseRecord> = records
563 .iter()
564 .filter_map(|record| codex.next(record))
565 .filter(|used| used.at_ms.is_some_and(|at| at >= since_ms))
566 .collect();
567 if used.is_empty() {
568 return;
569 }
570 let session_id = parent.clone().unwrap_or_else(|| thread.clone());
571 let session = sessions
572 .entry((crate::HarnessId::CODEX.to_string(), session_id))
573 .or_default();
574 for record in &used {
575 session.record(record);
576 }
577 if parent.is_some() {
578 session.subagents.insert(thread);
579 }
580 if session.cwd.is_none() || parent.is_none() {
581 if let Some(cwd) = crate::codex_peer::rollout_cwd(path) {
582 session.cwd = Some(cwd.to_string_lossy().into_owned());
583 }
584 }
585}
586
587fn tail_records(path: &Path, mut start: impl FnMut(&Value) -> bool) -> Vec<Value> {
591 const BLOCK: u64 = 256 * 1024;
592 let Ok(mut file) = std::fs::File::open(path) else {
593 return Vec::new();
594 };
595 let Ok(length) = file.metadata().map(|meta| meta.len()) else {
596 return Vec::new();
597 };
598 let mut newest_first = Vec::new();
599 let mut carry: Vec<u8> = Vec::new();
601 let mut end = length;
602 while end > 0 {
603 let begin = end.saturating_sub(BLOCK);
604 let mut block = vec![0u8; (end - begin) as usize];
605 if file.seek(SeekFrom::Start(begin)).is_err() || file.read_exact(&mut block).is_err() {
606 break;
607 }
608 block.extend_from_slice(&carry);
609 let first_newline = if begin == 0 {
610 None
611 } else {
612 match block.iter().position(|byte| *byte == b'\n') {
613 Some(at) => Some(at),
614 None => {
615 carry = block;
616 end = begin;
617 continue;
618 }
619 }
620 };
621 let (head, lines) = match first_newline {
622 Some(at) => (block[..at].to_vec(), &block[at + 1..]),
623 None => (Vec::new(), &block[..]),
624 };
625 let mut reached = false;
626 for line in lines.split(|byte| *byte == b'\n').rev() {
627 let Ok(record) = serde_json::from_slice::<Value>(line) else {
628 continue;
629 };
630 reached = start(&record);
631 newest_first.push(record);
632 if reached {
633 break;
634 }
635 }
636 if reached {
637 break;
638 }
639 carry = head;
640 end = begin;
641 }
642 newest_first.reverse();
643 newest_first
644}