g1t/services/billing/src/tokens.rs

300 lines11,771 bytesCodeBlame

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 mix1//! 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/// The UTC day of a time, `YYYY-MM-DD`.
39fn 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.
45pub(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.
52pub(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.
72pub(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.
78pub(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.
88fn number(n: u64) -> JsValue {
89 JsValue::from_f64(n as f64)
90}
91
92impl 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 };
Usage's cost shows comped runs at what they cost the provider, not $0173 // What the window's runs cost, whoever paid: at g1t's price, what was
Mission control shows model usage, yours and the workspace's: tokens, cost, active days, cache share, each day, and the mix174 // left to pay plus what included usage, credit, the trial, the
Usage's cost shows comped runs at what they cost the provider, not $0175 // 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.
Mission control shows model usage, yours and the workspace's: tokens, cost, active days, cache share, each day, and the mix178 let measure = if self.free {
179 "COALESCE(l.cost_micros, 0)"
180 } else {
Usage's cost shows comped runs at what they cost the provider, not $0181 "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"
Mission control shows model usage, yours and the workspace's: tokens, cost, active days, cache share, each day, and the mix184 };
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)]
231mod 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}