Skip to content

g1t/services/billing/src/tokens.rs

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