Skip to main content

supercode_harness/
token_spend.rs

1//! Token spend: what each harness session used, read from the harness's own records (Claude Code's `message.usage`
2//! on each response, Codex's cumulative `token_count` totals) and priced at the model's list rates
3//! ([`crate::pricing::built_in`]). A model with no list price is unrecorded, never zero.
4//!
5//! [`session_spend`] answers one loaded session per UTC day and model (`harness.v1.sessions.usage`), from a moment
6//! on when asked. [`machine_spend`] answers every Claude Code and Codex session on this machine with use in a window
7//! (`harness.v1.sessions.spend`, `supercode sessions usage`): a subagent's use counts on its parent session. It reads
8//! only transcripts modified inside the window, and of each only its tail back to the window's start.
9
10use 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/// One response's (Claude Code) or one turn's (Codex) recorded use.
21#[derive(Debug, Clone)]
22pub struct UseRecord {
23    /// When the harness recorded it.
24    pub at: Option<String>,
25    /// The same, in unix milliseconds.
26    pub at_ms: Option<i64>,
27    /// The model that answered.
28    pub model: String,
29    /// What it used.
30    pub tokens: RecordedTokens,
31    /// Whether it ran in fast mode (`usage.speed: "fast"`).
32    pub fast: bool,
33    /// Whether it ran US-only (`usage.inference_geo: "us"`).
34    pub us_only: bool,
35}
36
37/// Use summed, and what it cost at list rates where its models are priced.
38#[derive(Debug, Clone, Default)]
39pub struct Spend {
40    /// Tokens of every kind.
41    pub tokens: RecordedTokens,
42    /// List-rate cost of the priced part, in US dollars.
43    pub cost_usd: f64,
44    /// Whether any of it was priced.
45    pub priced: bool,
46    /// Whether any of it was not: its model has no list price.
47    pub unpriced: bool,
48}
49
50impl Spend {
51    /// Add one record, priced at its model's list rates.
52    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    /// Add another sum.
65    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    /// The fields every reading answers: each token kind (`cache_write_tokens` is both writes), `cost` only where some
73    /// use was priced, and `cost_unrecorded` when some was not.
74    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
99/// The use one Claude Code transcript line records, with its response's message id: Claude writes a response's
100/// `usage` on each of its lines, so a reader keeps one per id (the last). A synthetic message records none.
101/// `cache_creation_input_tokens` is split by TTL where Claude records the split (`usage.cache_creation`); a write
102/// the split does not account for is at the default 5-minute TTL.
103pub 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/// Codex's cumulative `token_count` totals, read in order: each change is one turn's use. Codex counts cached prompt
151/// tokens inside `input_tokens`; a use separates them, as Claude's `usage` does. The model is the latest
152/// `turn_context`'s.
153#[derive(Debug, Clone)]
154pub struct CodexUse {
155    model: String,
156    previous: [u64; 3],
157}
158
159impl CodexUse {
160    /// A reader starting from no totals, with the session's recorded model until a `turn_context` names one.
161    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    /// The use `record` adds, when it is a changed total.
169    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
212/// One loaded session's use per UTC day and model, each record counted from `since_ms` on (all of it without one).
213pub 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
264/// A window's start: a duration back from `now_ms` (`90s`, `10m`, `1h`, `2d`) or an RFC 3339 UTC moment.
265pub 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/// One session's use in a window, its subagents' included.
289#[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
315/// Every Claude Code and Codex session on this machine with recorded use from `since_ms` to now, a subagent's use on
316/// its parent. Sorted by `sort`: `cost` (the default) puts priced sessions first, most cost first, then the rest by
317/// tokens; `tokens` sorts all by tokens. `harnesses` limits the reading (empty: both); `limit` keeps the first rows.
318pub 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
412/// How far before a window's start a tail is read: records are mostly, not strictly, in time order (a response's
413/// lines, a subagent's result), so the read stops only at a record this much older than the window, and what it
414/// reads in between is still counted by each record's own time. The limit: a record inside the window that sits in
415/// the file before one more than ten minutes older than the window's start is not read, so not counted.
416const OUT_OF_ORDER_MS: i64 = 10 * 60 * 1000;
417
418/// Whether `path` was written at or after `since_ms`.
419fn 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
434/// Claude Code's transcripts under `projects`: `<project>/<session>.jsonl`, and a subagent's
435/// `<project>/<session>/subagents/*.jsonl`, counted on `<session>`.
436fn 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    // A session's folder is its own transcript's; a subagent's stands in until that is read.
509    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
517/// Codex's rollouts under `sessions` (`YYYY/MM/DD/rollout-*.jsonl`); a subagent thread's use counts on its parent.
518fn 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    // Read back to a total recorded before the window (each later total is a delta from it) and the model then.
552    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
587/// The records of a JSONL file from the newest back to the first one `start` accepts (included), oldest first; the
588/// whole file when none does. Reads back from the end in blocks, so a long transcript costs only its tail. `start`
589/// sees each record newest first.
590fn 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    // The bytes before the first newline of the block read last: a line whose start is further back.
600    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}