g1t/services/billing/src/tokens.rs

298 lines11,601 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 };
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)]
229mod 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}