Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.
| Mission control shows model usage, yours and the workspace's: tokens, cost, active days, cache share, each day, and the mix | 1 | //! Model tokens, counted per run for usage views. |
| 2 | //! | |
| 3 | //! The model proxy reports what each answer used (`record_tokens`), and the | |
| 4 | //! day's row for that run adds it up: one row per day, workspace, person, | |
| 5 | //! session and model. `token_usage` reads a window of those back, for the | |
| 6 | //! whole workspace or for one person, with what the window's runs were | |
| 7 | //! charged. Billing still prices runs from AI Gateway, never from these. | |
| 8 | ||
| 9 | use futures_util::future::try_join; | |
| 10 | use g1t_contracts::billing::{DayTokens, RecordTokensArgs, TokenUsage, TokenUsageArgs}; | |
| 11 | use g1t_contracts::identity::AGENT_NAME; | |
| 12 | use g1t_contracts::time::rfc3339; | |
| 13 | use g1t_contracts::{FailureCode, Outcome, Role}; | |
| 14 | use g1t_kit::now_ms; | |
| 15 | use serde::Deserialize; | |
| 16 | use worker::wasm_bindgen::JsValue; | |
| 17 | use worker::Result; | |
| 18 | ||
| 19 | use crate::{Billing, members_only}; | |
| 20 | ||
| 21 | /// The window `token_usage` shows when asked for none, and the longest. | |
| 22 | pub(crate) const DEFAULT_DAYS: u32 = 42; | |
| 23 | pub(crate) const MAX_DAYS: u32 = 366; | |
| 24 | const DAY_MS: u64 = 86_400_000; | |
| 25 | ||
| 26 | /// One answer's tokens, ready to add to its day's row. | |
| 27 | #[derive(Debug, PartialEq)] | |
| 28 | pub(crate) struct TokenRow { | |
| 29 | pub day: String, | |
| 30 | pub workspace: String, | |
| 31 | pub person: String, | |
| 32 | pub session: String, | |
| 33 | pub model: String, | |
| 34 | pub tier: Option<String>, | |
| 35 | pub counts: [u64; 4], | |
| 36 | } | |
| 37 | ||
| 38 | /// The UTC day of a time, `YYYY-MM-DD`. | |
| 39 | fn day_of(ms: u64) -> String { | |
| 40 | rfc3339(ms)[..10].to_owned() | |
| 41 | } | |
| 42 | ||
| 43 | /// Who a run's tokens count for: a username, lowercased, or empty for | |
| 44 | /// nobody. The agent is not a person. | |
| 45 | pub(crate) fn person_of(username: Option<&str>) -> String { | |
| 46 | let name = username.unwrap_or_default().trim().to_lowercase(); | |
| 47 | if name == AGENT_NAME { String::new() } else { name } | |
| 48 | } | |
| 49 | ||
| 50 | /// What to add for one report, or None when there is nothing to count or | |
| 51 | /// nothing to count it under. | |
| 52 | pub(crate) fn token_row(a: &RecordTokensArgs, now: u64) -> Option<TokenRow> { | |
| 53 | let counts = [a.input, a.output, a.cache_read, a.cache_write]; | |
| 54 | let workspace = a.workspace.trim().to_lowercase(); | |
| 55 | let session = a.session.trim(); | |
| 56 | if counts.iter().all(|n| *n == 0) || workspace.is_empty() || session.is_empty() { | |
| 57 | return None; | |
| 58 | } | |
| 59 | let model = a.model.trim(); | |
| 60 | Some(TokenRow { | |
| 61 | day: day_of(now), | |
| 62 | workspace, | |
| 63 | person: person_of(a.person.as_deref()), | |
| 64 | session: session.chars().take(64).collect(), | |
| 65 | model: if model.is_empty() { "unknown".to_owned() } else { model.chars().take(200).collect() }, | |
| 66 | tier: a.tier.as_deref().filter(|tier| matches!(*tier, "small" | "large")).map(str::to_owned), | |
| 67 | counts, | |
| 68 | }) | |
| 69 | } | |
| 70 | ||
| 71 | /// The days of a window that ends today, oldest first, and how many. | |
| 72 | pub(crate) fn window(now: u64, days: Option<u32>) -> Vec<String> { | |
| 73 | let days = days.unwrap_or(DEFAULT_DAYS).clamp(1, MAX_DAYS); | |
| 74 | (0..u64::from(days)).rev().map(|back| day_of(now.saturating_sub(back * DAY_MS))).collect() | |
| 75 | } | |
| 76 | ||
| 77 | /// Every day of the window with its tokens, zeros included. | |
| 78 | pub(crate) fn fill(days: &[String], counted: &[(String, u64)]) -> Vec<DayTokens> { | |
| 79 | days.iter() | |
| 80 | .map(|day| DayTokens { | |
| 81 | day: day.clone(), | |
| 82 | tokens: counted.iter().filter(|(d, _)| d == day).map(|(_, n)| n).sum(), | |
| 83 | }) | |
| 84 | .collect() | |
| 85 | } | |
| 86 | ||
| 87 | /// D1 takes numbers as doubles; a count of tokens fits exactly. | |
| 88 | fn number(n: u64) -> JsValue { | |
| 89 | JsValue::from_f64(n as f64) | |
| 90 | } | |
| 91 | ||
| 92 | impl Billing { | |
| 93 | pub(crate) async fn record_tokens(&self, a: RecordTokensArgs) -> Result<Outcome<bool>> { | |
| 94 | let Some(row) = token_row(&a, now_ms()) else { | |
| 95 | return Ok(Outcome::Ok(false)); | |
| 96 | }; | |
| 97 | let [input, output, cache_read, cache_write] = row.counts; | |
| 98 | self.db | |
| 99 | .prepare( | |
| 100 | "INSERT INTO token_usage (day, workspace, person, session, model, tier, input, output, cache_read, cache_write, requests) | |
| 101 | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 1) | |
| 102 | ON CONFLICT (day, workspace, person, session, model) DO UPDATE SET | |
| 103 | input = input + excluded.input, | |
| 104 | output = output + excluded.output, | |
| 105 | cache_read = cache_read + excluded.cache_read, | |
| 106 | cache_write = cache_write + excluded.cache_write, | |
| 107 | requests = requests + 1, | |
| 108 | tier = COALESCE(excluded.tier, tier)", | |
| 109 | ) | |
| 110 | .bind(&[ | |
| 111 | row.day.into(), | |
| 112 | row.workspace.into(), | |
| 113 | row.person.into(), | |
| 114 | row.session.into(), | |
| 115 | row.model.into(), | |
| 116 | row.tier.map_or(JsValue::NULL, JsValue::from), | |
| 117 | number(input), | |
| 118 | number(output), | |
| 119 | number(cache_read), | |
| 120 | number(cache_write), | |
| 121 | ])? | |
| 122 | .run() | |
| 123 | .await?; | |
| 124 | Ok(Outcome::Ok(true)) | |
| 125 | } | |
| 126 | ||
| 127 | pub(crate) async fn token_usage(&self, a: TokenUsageArgs) -> Result<Outcome<TokenUsage>> { | |
| 128 | let workspace = a.workspace.to_lowercase(); | |
| 129 | let Some(viewer) = a.viewer.filter(|viewer| viewer.is_member(&workspace)) else { | |
| 130 | return Ok(members_only()); | |
| 131 | }; | |
| 132 | let person = a.person.as_deref().map(|name| person_of(Some(name))).filter(|name| !name.is_empty()); | |
| 133 | if let Some(person) = &person | |
| 134 | && *person != viewer.username.to_lowercase() | |
| 135 | && viewer.role_in(&workspace) != Some(Role::Owner) | |
| 136 | { | |
| 137 | return Ok(Outcome::fail( | |
| 138 | FailureCode::Forbidden, | |
| 139 | "Only owners can see another person's usage.", | |
| 140 | )); | |
| 141 | } | |
| 142 | let days = window(now_ms(), a.days); | |
| 143 | let since = days[0].clone(); | |
| 144 | let mut binds: Vec<JsValue> = vec![workspace.as_str().into(), since.as_str().into()]; | |
| 145 | if let Some(person) = &person { | |
| 146 | binds.push(person.as_str().into()); | |
| 147 | } | |
| 148 | let for_person = if person.is_some() { " AND person = ?3" } else { "" }; | |
| 149 | ||
| 150 | #[derive(Deserialize)] | |
| 151 | struct DayRow { | |
| 152 | day: String, | |
| 153 | input: Option<f64>, | |
| 154 | output: Option<f64>, | |
| 155 | cache_read: Option<f64>, | |
| 156 | cache_write: Option<f64>, | |
| 157 | } | |
| 158 | #[derive(Deserialize)] | |
| 159 | struct Charged { | |
| 160 | micros: Option<f64>, | |
| 161 | } | |
| 162 | let by_day = async { | |
| 163 | self.db | |
| 164 | .prepare(format!( | |
| 165 | "SELECT day, SUM(input) AS input, SUM(output) AS output, SUM(cache_read) AS cache_read, SUM(cache_write) AS cache_write | |
| 166 | FROM token_usage WHERE workspace = ?1 AND day >= ?2{for_person} GROUP BY day" | |
| 167 | )) | |
| 168 | .bind(&binds)? | |
| 169 | .all() | |
| 170 | .await? | |
| 171 | .results::<DayRow>() | |
| 172 | }; | |
| 173 | // What the window's runs cost at g1t's price, whoever paid: what was | |
| 174 | // left to pay plus what included usage, credit, the trial, the | |
| 175 | // open-source pool or a comp covered, so a comped workspace's runs | |
| 176 | // still show their cost. While g1t charges nothing, the provider's | |
| 177 | // cost. For one person, the runs whose sessions counted their tokens. | |
| 178 | let measure = if self.free { | |
| 179 | "COALESCE(l.cost_micros, 0)" | |
| 180 | } else { | |
| 181 | "(-l.amount_micros + l.credit_micros + l.trial_micros + l.oss_micros + l.given_micros)" | |
| 182 | }; | |
| 183 | let sessions = if person.is_some() { | |
| 184 | " AND r.session_id IN (SELECT session FROM token_usage WHERE workspace = ?1 AND person = ?3 AND day >= ?2)" | |
| 185 | } else { | |
| 186 | "" | |
| 187 | }; | |
| 188 | let charged = async { | |
| 189 | self.db | |
| 190 | .prepare(format!( | |
| 191 | "SELECT SUM({measure}) AS micros FROM ledger l JOIN runs r ON r.id = l.reference | |
| 192 | WHERE l.workspace = ?1 AND l.kind = 'usage' AND l.created_at >= ?2{sessions}" | |
| 193 | )) | |
| 194 | .bind(&binds)? | |
| 195 | .first::<Charged>(None) | |
| 196 | .await | |
| 197 | }; | |
| 198 | // Both read independently, so they go to D1 at once. | |
| 199 | let (rows, charged) = try_join(by_day, charged).await?; | |
| 200 | ||
| 201 | let count = |n: Option<f64>| n.unwrap_or_default().max(0.0) as u64; | |
| 202 | let mut totals = [0u64; 4]; | |
| 203 | let mut counted = Vec::with_capacity(rows.len()); | |
| 204 | for row in &rows { | |
| 205 | let row_counts = [count(row.input), count(row.output), count(row.cache_read), count(row.cache_write)]; | |
| 206 | for (total, n) in totals.iter_mut().zip(row_counts) { | |
| 207 | *total += n; | |
| 208 | } | |
| 209 | counted.push((row.day.clone(), row_counts.iter().sum::<u64>())); | |
| 210 | } | |
| 211 | let by_day = fill(&days, &counted); | |
| 212 | Ok(Outcome::Ok(TokenUsage { | |
| 213 | since, | |
| 214 | days: days.len() as u32, | |
| 215 | person, | |
| 216 | total_tokens: totals.iter().sum(), | |
| 217 | input_tokens: totals[0], | |
| 218 | output_tokens: totals[1], | |
| 219 | cache_read_tokens: totals[2], | |
| 220 | cache_write_tokens: totals[3], | |
| 221 | cost_micros: charged.and_then(|c| c.micros).unwrap_or_default().round() as i64, | |
| 222 | active_days: by_day.iter().filter(|day| day.tokens > 0).count() as u32, | |
| 223 | by_day, | |
| 224 | })) | |
| 225 | } | |
| 226 | } | |
| 227 | ||
| 228 | #[cfg(test)] | |
| 229 | mod tests { | |
| 230 | use super::*; | |
| 231 | ||
| 232 | // 2026-10-06T12:00:00Z. | |
| 233 | const NOW: u64 = 1_791_288_000_000; | |
| 234 | ||
| 235 | fn args() -> RecordTokensArgs { | |
| 236 | RecordTokensArgs { | |
| 237 | workspace: " Acme ".into(), | |
| 238 | session: "ms_abc".into(), | |
| 239 | person: Some("Ada".into()), | |
| 240 | model: "claude-opus".into(), | |
| 241 | tier: Some("large".into()), | |
| 242 | input: 10, | |
| 243 | output: 5, | |
| 244 | cache_read: 0, | |
| 245 | cache_write: 0, | |
| 246 | } | |
| 247 | } | |
| 248 | ||
| 249 | #[test] | |
| 250 | fn a_report_is_added_to_today_under_its_person() { | |
| 251 | let row = token_row(&args(), NOW).unwrap(); | |
| 252 | assert_eq!(row.day, "2026-10-06"); | |
| 253 | assert_eq!(row.workspace, "acme"); | |
| 254 | assert_eq!(row.person, "ada"); | |
| 255 | assert_eq!(row.tier.as_deref(), Some("large")); | |
| 256 | assert_eq!(row.counts, [10, 5, 0, 0]); | |
| 257 | } | |
| 258 | ||
| 259 | #[test] | |
| 260 | fn nothing_used_or_nowhere_to_put_it_is_not_counted() { | |
| 261 | assert!(token_row(&RecordTokensArgs { input: 0, output: 0, ..args() }, NOW).is_none()); | |
| 262 | assert!(token_row(&RecordTokensArgs { session: " ".into(), ..args() }, NOW).is_none()); | |
| 263 | assert!(token_row(&RecordTokensArgs { workspace: String::new(), ..args() }, NOW).is_none()); | |
| 264 | } | |
| 265 | ||
| 266 | #[test] | |
| 267 | fn the_agent_is_nobody_and_odd_tiers_and_models_are_tidied() { | |
| 268 | let row = token_row( | |
| 269 | &RecordTokensArgs { person: Some("g1t".into()), tier: Some("huge".into()), model: " ".into(), ..args() }, | |
| 270 | NOW, | |
| 271 | ) | |
| 272 | .unwrap(); | |
| 273 | assert_eq!(row.person, ""); | |
| 274 | assert_eq!(row.tier, None); | |
| 275 | assert_eq!(row.model, "unknown"); | |
| 276 | assert_eq!(person_of(None), ""); | |
| 277 | } | |
| 278 | ||
| 279 | #[test] | |
| 280 | fn the_window_ends_today_oldest_first_and_is_bounded() { | |
| 281 | let days = window(NOW, Some(3)); | |
| 282 | assert_eq!(days, vec!["2026-10-04", "2026-10-05", "2026-10-06"]); | |
| 283 | assert_eq!(window(NOW, None).len(), DEFAULT_DAYS as usize); | |
| 284 | assert_eq!(window(NOW, Some(0)).len(), 1); | |
| 285 | assert_eq!(window(NOW, Some(5000)).len(), MAX_DAYS as usize); | |
| 286 | // Across a month's end. | |
| 287 | assert_eq!(window(NOW, Some(7))[0], "2026-09-30"); | |
| 288 | } | |
| 289 | ||
| 290 | #[test] | |
| 291 | fn every_day_is_filled_with_zeros_where_nothing_ran() { | |
| 292 | let days = window(NOW, Some(3)); | |
| 293 | let filled = fill(&days, &[("2026-10-05".into(), 40), ("2026-09-01".into(), 9)]); | |
| 294 | let tokens: Vec<u64> = filled.iter().map(|d| d.tokens).collect(); | |
| 295 | assert_eq!(tokens, vec![0, 40, 0]); | |
| 296 | assert_eq!(filled[2].day, "2026-10-06"); | |
| 297 | } | |
| 298 | } |