| 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, whoever paid: at g1t's price, what was |
| 174 | // left to pay plus what included usage, credit, the trial, the |
| 175 | // open-source pool or a comp covered; a run charged nothing at all |
| 176 | // (comped, or while g1t charges nothing) at the provider's cost. |
| 177 | // 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 | "CASE WHEN (-l.amount_micros + l.credit_micros + l.trial_micros + l.oss_micros + l.given_micros) > 0 |
| 182 | THEN (-l.amount_micros + l.credit_micros + l.trial_micros + l.oss_micros + l.given_micros) |
| 183 | ELSE COALESCE(l.cost_micros, 0) END" |
| 184 | }; |
| 185 | let sessions = if person.is_some() { |
| 186 | " AND r.session_id IN (SELECT session FROM token_usage WHERE workspace = ?1 AND person = ?3 AND day >= ?2)" |
| 187 | } else { |
| 188 | "" |
| 189 | }; |
| 190 | let charged = async { |
| 191 | self.db |
| 192 | .prepare(format!( |
| 193 | "SELECT SUM({measure}) AS micros FROM ledger l JOIN runs r ON r.id = l.reference |
| 194 | WHERE l.workspace = ?1 AND l.kind = 'usage' AND l.created_at >= ?2{sessions}" |
| 195 | )) |
| 196 | .bind(&binds)? |
| 197 | .first::<Charged>(None) |
| 198 | .await |
| 199 | }; |
| 200 | // Both read independently, so they go to D1 at once. |
| 201 | let (rows, charged) = try_join(by_day, charged).await?; |
| 202 | |
| 203 | let count = |n: Option<f64>| n.unwrap_or_default().max(0.0) as u64; |
| 204 | let mut totals = [0u64; 4]; |
| 205 | let mut counted = Vec::with_capacity(rows.len()); |
| 206 | for row in &rows { |
| 207 | let row_counts = [count(row.input), count(row.output), count(row.cache_read), count(row.cache_write)]; |
| 208 | for (total, n) in totals.iter_mut().zip(row_counts) { |
| 209 | *total += n; |
| 210 | } |
| 211 | counted.push((row.day.clone(), row_counts.iter().sum::<u64>())); |
| 212 | } |
| 213 | let by_day = fill(&days, &counted); |
| 214 | Ok(Outcome::Ok(TokenUsage { |
| 215 | since, |
| 216 | days: days.len() as u32, |
| 217 | person, |
| 218 | total_tokens: totals.iter().sum(), |
| 219 | input_tokens: totals[0], |
| 220 | output_tokens: totals[1], |
| 221 | cache_read_tokens: totals[2], |
| 222 | cache_write_tokens: totals[3], |
| 223 | cost_micros: charged.and_then(|c| c.micros).unwrap_or_default().round() as i64, |
| 224 | active_days: by_day.iter().filter(|day| day.tokens > 0).count() as u32, |
| 225 | by_day, |
| 226 | })) |
| 227 | } |
| 228 | } |
| 229 | |
| 230 | #[cfg(test)] |
| 231 | mod tests { |
| 232 | use super::*; |
| 233 | |
| 234 | // 2026-10-06T12:00:00Z. |
| 235 | const NOW: u64 = 1_791_288_000_000; |
| 236 | |
| 237 | fn args() -> RecordTokensArgs { |
| 238 | RecordTokensArgs { |
| 239 | workspace: " Acme ".into(), |
| 240 | session: "ms_abc".into(), |
| 241 | person: Some("Ada".into()), |
| 242 | model: "claude-opus".into(), |
| 243 | tier: Some("large".into()), |
| 244 | input: 10, |
| 245 | output: 5, |
| 246 | cache_read: 0, |
| 247 | cache_write: 0, |
| 248 | } |
| 249 | } |
| 250 | |
| 251 | #[test] |
| 252 | fn a_report_is_added_to_today_under_its_person() { |
| 253 | let row = token_row(&args(), NOW).unwrap(); |
| 254 | assert_eq!(row.day, "2026-10-06"); |
| 255 | assert_eq!(row.workspace, "acme"); |
| 256 | assert_eq!(row.person, "ada"); |
| 257 | assert_eq!(row.tier.as_deref(), Some("large")); |
| 258 | assert_eq!(row.counts, [10, 5, 0, 0]); |
| 259 | } |
| 260 | |
| 261 | #[test] |
| 262 | fn nothing_used_or_nowhere_to_put_it_is_not_counted() { |
| 263 | assert!(token_row(&RecordTokensArgs { input: 0, output: 0, ..args() }, NOW).is_none()); |
| 264 | assert!(token_row(&RecordTokensArgs { session: " ".into(), ..args() }, NOW).is_none()); |
| 265 | assert!(token_row(&RecordTokensArgs { workspace: String::new(), ..args() }, NOW).is_none()); |
| 266 | } |
| 267 | |
| 268 | #[test] |
| 269 | fn the_agent_is_nobody_and_odd_tiers_and_models_are_tidied() { |
| 270 | let row = token_row( |
| 271 | &RecordTokensArgs { person: Some("g1t".into()), tier: Some("huge".into()), model: " ".into(), ..args() }, |
| 272 | NOW, |
| 273 | ) |
| 274 | .unwrap(); |
| 275 | assert_eq!(row.person, ""); |
| 276 | assert_eq!(row.tier, None); |
| 277 | assert_eq!(row.model, "unknown"); |
| 278 | assert_eq!(person_of(None), ""); |
| 279 | } |
| 280 | |
| 281 | #[test] |
| 282 | fn the_window_ends_today_oldest_first_and_is_bounded() { |
| 283 | let days = window(NOW, Some(3)); |
| 284 | assert_eq!(days, vec!["2026-10-04", "2026-10-05", "2026-10-06"]); |
| 285 | assert_eq!(window(NOW, None).len(), DEFAULT_DAYS as usize); |
| 286 | assert_eq!(window(NOW, Some(0)).len(), 1); |
| 287 | assert_eq!(window(NOW, Some(5000)).len(), MAX_DAYS as usize); |
| 288 | // Across a month's end. |
| 289 | assert_eq!(window(NOW, Some(7))[0], "2026-09-30"); |
| 290 | } |
| 291 | |
| 292 | #[test] |
| 293 | fn every_day_is_filled_with_zeros_where_nothing_ran() { |
| 294 | let days = window(NOW, Some(3)); |
| 295 | let filled = fill(&days, &[("2026-10-05".into(), 40), ("2026-09-01".into(), 9)]); |
| 296 | let tokens: Vec<u64> = filled.iter().map(|d| d.tokens).collect(); |
| 297 | assert_eq!(tokens, vec![0, 40, 0]); |
| 298 | assert_eq!(filled[2].day, "2026-10-06"); |
| 299 | } |
| 300 | } |