Skip to content

g1t/services/billing/src/gateway.rs

504 lines22,295 bytesCodeBlame
1//! The AI Gateway: a workspace's own model requests, sent with one of its
2//! access tokens to the model proxy at `models.g1t.sh/anthropic`.
3//!
4//! - **Admission.** Before a request goes to g1t's models the proxy asks
5//! `gateway_admit`. A workspace over its spend limit is refused, as one is
6//! for anything else. On the plan it needs AI credit or included usage
7//! left, as an agent run does (auto-reload is tried first). A workspace
8//! with no plan is refused: the gateway on g1t's key is paid for from AI
9//! credit, which comes with the plan. A 100% discount and an enterprise
10//! need nothing more. On the workspace's own provider key nothing is
11//! asked: those requests cost g1t nothing.
12//! - **Charging.** Each request that used tokens on g1t's models is
13//! charged its tokens at the model's list price (`gateway_models`), plus
14//! the price book's `gateway_models` markup (0 while the gateway is in
15//! beta), on a ledger line of its own (task `gateway`, one request each).
16//! The plan's included usage pays first, then AI credit (`grants::replay`
17//! counts gateway lines as model usage). Not an agent run, so never the
18//! agent rate, and never trial credit or g1t's pools.
19//! - **On the workspace's own key.** Logged with its tokens and charged
20//! nothing.
21//! - **The log.** Every request is kept for `RETENTION_DAYS`, with its
22//! model, tokens by kind, cost, status and the token that sent it; never
23//! its prompt or answer.
24
25use g1t_contracts::billing::{
26 GatewayAdmitArgs, GatewayModel, GatewayRequest, GatewayRequests, GatewayRequestsArgs, LimitState, PlanKind, RecordGatewayArgs,
27};
28use g1t_contracts::time::rfc3339;
29use g1t_contracts::{FailureCode, Outcome, new_id};
30use g1t_kit::now_ms;
31use serde::Deserialize;
32use worker::Result;
33use worker::wasm_bindgen::JsValue;
34
35use crate::credits::{Eligible, month_of};
36use crate::features::thousands;
37use crate::{Billing, margin_on, members_only, optional};
38
39/// How long the log keeps a request.
40pub(crate) const RETENTION_DAYS: u32 = 30;
41/// A page of the log: this many when not asked, and at most.
42const PAGE: u32 = 50;
43const MAX_PAGE: u32 = 200;
44const DAY_MS: u64 = 86_400_000;
45
46/// What `tokens` (input, output, cache reads, cache writes) cost at a
47/// model's prices per million, rounded up to a whole millionth of a dollar.
48pub(crate) fn cost_micros(model: &GatewayModel, tokens: [u64; 4]) -> i64 {
49 let prices = [model.input_micros, model.output_micros, model.cache_read_micros, model.cache_write_micros];
50 let millionths: u128 = tokens
51 .iter()
52 .zip(prices)
53 .map(|(n, price)| u128::from(*n) * u128::from(price.max(0).unsigned_abs()))
54 .sum();
55 i64::try_from(millionths.div_ceil(1_000_000)).unwrap_or(i64::MAX)
56}
57
58/// Where a workspace stands for a request on g1t's models.
59#[derive(Clone, Debug, PartialEq)]
60pub(crate) enum Standing {
61 /// Over its spend limit, with the limit's own message.
62 Stopped(String),
63 /// No plan.
64 NoPlan,
65 /// On the plan with no AI credit or included usage left.
66 OutOfCredit { reload_failed: bool },
67 Admitted,
68}
69
70/// Where a workspace within its limit stands, by its plan: `exhausted` is
71/// `credit_exhausted`'s answer (asked only on the plan). A workspace with
72/// no plan has no AI credit to spend; a 100% discount and an enterprise
73/// need none.
74pub(crate) fn plan_standing(plan: PlanKind, exhausted: Option<bool>) -> Standing {
75 match plan {
76 PlanKind::Free => Standing::NoPlan,
77 PlanKind::Internal | PlanKind::Enterprise => Standing::Admitted,
78 PlanKind::Paid => match exhausted {
79 Some(reload_failed) => Standing::OutOfCredit { reload_failed },
80 None => Standing::Admitted,
81 },
82 }
83}
84
85/// What one request costs g1t, to charge: its tokens at its model's prices
86/// on g1t's key; nothing on the workspace's own key, or for a model g1t
87/// has no price for.
88pub(crate) fn request_cost(own_key: bool, model: Option<&GatewayModel>, tokens: [u64; 4]) -> i64 {
89 match model {
90 Some(model) if !own_key => cost_micros(model, tokens),
91 _ => 0,
92 }
93}
94
95/// Why a request is refused, in words for whoever sent it, or None.
96pub(crate) fn refusal(workspace: &str, standing: &Standing) -> Option<String> {
97 match standing {
98 Standing::Admitted => None,
99 Standing::Stopped(message) => Some(message.clone()),
100 Standing::NoPlan => Some(format!(
101 "The AI Gateway on g1t's models is paid for from AI credit, which comes with the g1t plan. An owner can start the plan for {workspace} at /{workspace}/-/billing, or connect the workspace's own Anthropic key under Integrations to use the gateway at no charge."
102 )),
103 Standing::OutOfCredit { reload_failed } => {
104 let reload = if *reload_failed { " Auto-reload was turned off after its last charge failed." } else { "" };
105 Some(format!(
106 "The {workspace} workspace is out of AI credit and has used this month's included usage, so the AI Gateway refuses requests to g1t's models.{reload} An owner can buy AI credit or turn on auto-reload at /{workspace}/-/billing#ai-credit."
107 ))
108 }
109 }
110}
111
112/// What a ledger line for one request says.
113pub(crate) fn describe(model_name: &str, tokens: [u64; 4], token_name: Option<&str>) -> String {
114 let used: u64 = tokens.iter().sum();
115 let by = token_name.map(str::trim).filter(|name| !name.is_empty()).map_or(String::new(), |name| format!(", token {name}"));
116 format!("AI Gateway: {model_name}, {} tokens{by}", thousands(used))
117}
118
119/// Whether a request id is one the proxy makes: `gw_` and up to 64 letters,
120/// digits, `_` and `-`.
121pub(crate) fn valid_id(id: &str) -> bool {
122 id.starts_with("gw_") && id.len() <= 64 && id.chars().all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')
123}
124
125#[derive(Deserialize)]
126struct ModelRow {
127 model: String,
128 name: String,
129 provider: String,
130 input_micros: i64,
131 output_micros: i64,
132 cache_read_micros: i64,
133 cache_write_micros: i64,
134}
135
136impl From<ModelRow> for GatewayModel {
137 fn from(row: ModelRow) -> Self {
138 GatewayModel {
139 model: row.model,
140 name: row.name,
141 provider: row.provider,
142 input_micros: row.input_micros,
143 output_micros: row.output_micros,
144 cache_read_micros: row.cache_read_micros,
145 cache_write_micros: row.cache_write_micros,
146 }
147 }
148}
149
150#[derive(Deserialize)]
151struct RequestRow {
152 id: String,
153 created_at: String,
154 model: String,
155 token_id: String,
156 token_name: Option<String>,
157 input: f64,
158 output: f64,
159 cache_read: f64,
160 cache_write: f64,
161 cost_micros: f64,
162 charged_micros: f64,
163 status: f64,
164 own_key: f64,
165 streamed: f64,
166 duration_ms: f64,
167 error: Option<String>,
168}
169
170impl From<RequestRow> for GatewayRequest {
171 fn from(row: RequestRow) -> Self {
172 let n = |v: f64| v.max(0.0) as u64;
173 GatewayRequest {
174 id: row.id,
175 created_at: row.created_at,
176 model: row.model,
177 token_id: row.token_id,
178 token_name: row.token_name,
179 input: n(row.input),
180 output: n(row.output),
181 cache_read: n(row.cache_read),
182 cache_write: n(row.cache_write),
183 cost_micros: row.cost_micros as i64,
184 charged_micros: row.charged_micros as i64,
185 status: row.status as u16,
186 own_key: row.own_key != 0.0,
187 streamed: row.streamed != 0.0,
188 duration_ms: n(row.duration_ms),
189 error: row.error,
190 }
191 }
192}
193
194/// D1 takes numbers as doubles; counts and amounts here fit exactly.
195fn number(n: u64) -> JsValue {
196 JsValue::from_f64(n as f64)
197}
198
199impl Billing {
200 /// `gateway_models`: what the gateway offers on g1t's key, with prices.
201 pub(crate) async fn gateway_models(&self) -> Result<Vec<GatewayModel>> {
202 Ok(self
203 .db
204 .prepare("SELECT * FROM gateway_models ORDER BY position, model")
205 .all()
206 .await?
207 .results::<ModelRow>()?
208 .into_iter()
209 .map(GatewayModel::from)
210 .collect())
211 }
212
213 async fn gateway_model(&self, model: &str) -> Result<Option<GatewayModel>> {
214 Ok(self
215 .db
216 .prepare("SELECT * FROM gateway_models WHERE model = ?")
217 .bind(&[model.into()])?
218 .first::<ModelRow>(None)
219 .await?
220 .map(GatewayModel::from))
221 }
222
223 /// Where a workspace stands for a request on g1t's models.
224 pub(crate) async fn gateway_standing(&self, workspace: &str) -> Result<Standing> {
225 if self.stripe.is_none() || self.free {
226 return Ok(Standing::Admitted);
227 }
228 let limit = self.limit_of(workspace).await?;
229 if limit.state == LimitState::Stopped {
230 return Ok(Standing::Stopped(limit.message.unwrap_or_else(|| "This workspace is over its limit.".to_owned())));
231 }
232 if let Some(Outcome::Fail(failure)) = self.out_of_credit::<bool>(workspace).await? {
233 return Ok(Standing::Stopped(failure.message));
234 }
235 let account = self.account_of(workspace).await?;
236 let plan = self.plan_kind_for(workspace, &account).await?;
237 // Only a workspace paying on the plan needs credit to spend.
238 let exhausted = if plan == PlanKind::Paid { self.credit_exhausted(workspace).await? } else { None };
239 Ok(plan_standing(plan, exhausted))
240 }
241
242 /// `gateway_admit`.
243 pub(crate) async fn gateway_admit(&self, a: GatewayAdmitArgs) -> Result<Outcome<bool>> {
244 let workspace = a.workspace.trim().to_lowercase();
245 if workspace.is_empty() {
246 return Ok(Outcome::fail(FailureCode::Invalid, "Name the workspace."));
247 }
248 let standing = self.gateway_standing(&workspace).await?;
249 Ok(match refusal(&workspace, &standing) {
250 Some(why) => Outcome::fail(FailureCode::PaymentRequired, why),
251 None => Outcome::Ok(true),
252 })
253 }
254
255 /// `record_gateway`: logs a request once, and charges it once when it
256 /// used tokens on g1t's models.
257 pub(crate) async fn record_gateway(&self, a: RecordGatewayArgs) -> Result<Outcome<bool>> {
258 if !valid_id(&a.id) {
259 return Ok(Outcome::fail(FailureCode::Invalid, "A gateway request's id is gw_ and up to 64 letters and digits."));
260 }
261 let workspace = a.workspace.trim().to_lowercase();
262 if workspace.is_empty() {
263 return Ok(Outcome::fail(FailureCode::Invalid, "Name the workspace."));
264 }
265 let tokens = [a.input, a.output, a.cache_read, a.cache_write];
266 let model_name: String = a.model.trim().chars().take(200).collect();
267 let token_name = a.token_name.as_deref().map(|name| name.trim().chars().take(100).collect::<String>());
268 let priced = if a.own_key { None } else { self.gateway_model(&model_name).await? };
269 let cost = request_cost(a.own_key, priced.as_ref(), tokens);
270 let now = now_ms();
271 let timestamp = rfc3339(now);
272 // Claimed first: the same request recorded twice is one row and one charge.
273 let claimed = self
274 .db
275 .prepare(
276 "INSERT OR IGNORE INTO gateway_requests
277 (id, workspace, created_at, token_id, token_name, model, input, output, cache_read, cache_write,
278 cost_micros, charged_micros, status, own_key, streamed, duration_ms, error)
279 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 0, ?, ?, ?, ?, ?)
280 RETURNING id",
281 )
282 .bind(&[
283 a.id.as_str().into(),
284 workspace.as_str().into(),
285 timestamp.as_str().into(),
286 a.token_id.chars().take(64).collect::<String>().into(),
287 optional(token_name.as_deref()),
288 model_name.as_str().into(),
289 number(a.input),
290 number(a.output),
291 number(a.cache_read),
292 number(a.cache_write),
293 (cost as f64).into(),
294 f64::from(a.status).into(),
295 f64::from(u8::from(a.own_key)).into(),
296 f64::from(u8::from(a.streamed)).into(),
297 number(a.duration_ms),
298 optional(a.error.as_deref().map(|e| e.chars().take(500).collect::<String>()).as_deref()),
299 ])?
300 .first::<serde_json::Value>(None)
301 .await?;
302 if claimed.is_none() {
303 return Ok(Outcome::Ok(false));
304 }
305 // Without a card processor there is no bill to put it on.
306 let Some(model) = priced.filter(|_| cost > 0 && self.stripe.is_some()) else {
307 return Ok(Outcome::Ok(true));
308 };
309 let base = margin_on(cost, self.gateway_markup().await?);
310 let (charge, terms_note, discount) = self.charged(&workspace, base).await?;
311 // Included usage pays first, then AI credit. Never the trial or
312 // g1t's pools, and g1t never covers the rest.
313 let eligible = Eligible { trial: false, repo: None, cover_rest: false };
314 let drawn = self.draw(&workspace, charge, &month_of(&timestamp), &eligible).await?;
315 let owed = charge - drawn.total();
316 let description = format!("{}{terms_note}{}", describe(&model.name, tokens, token_name.as_deref()), drawn.note());
317 self.db
318 .batch(vec![
319 self.db
320 .prepare(
321 "INSERT INTO ledger
322 (id, workspace, kind, amount_micros, description, task, model, cost_micros, reference, created_at,
323 billed_to, credit_micros, trial_micros, oss_micros, given_micros, price_version, quantity)
324 VALUES (?, ?, 'usage', ?, ?, 'gateway', ?, ?, ?, ?, 'g1t', ?, ?, ?, ?, ?, 1)",
325 )
326 .bind(&[
327 new_id("led", now).into(),
328 workspace.as_str().into(),
329 (-(owed as f64)).into(),
330 description.as_str().into(),
331 model.model.as_str().into(),
332 (cost as f64).into(),
333 a.id.as_str().into(),
334 timestamp.as_str().into(),
335 (drawn.credit as f64).into(),
336 (drawn.trial as f64).into(),
337 (drawn.oss as f64).into(),
338 (drawn.given as f64).into(),
339 optional(self.version_now("gateway_models").await?.as_deref()),
340 ])?,
341 self.db
342 .prepare(
343 "INSERT INTO accounts (workspace, balance_micros, created_at)
344 VALUES (?1, ?2, ?3)
345 ON CONFLICT (workspace) DO UPDATE SET balance_micros = balance_micros + ?2",
346 )
347 .bind(&[workspace.as_str().into(), (-(owed as f64)).into(), timestamp.as_str().into()])?,
348 self.db
349 .prepare("UPDATE gateway_requests SET charged_micros = ? WHERE id = ?")
350 .bind(&[(charge as f64).into(), a.id.as_str().into()])?,
351 ])
352 .await?;
353 self.record_discount(&a.id, discount).await?;
354 self.count_spend(&workspace, cost, owed, &drawn).await;
355 Ok(Outcome::Ok(true))
356 }
357
358 /// `gateway_requests`: the log, newest first, for members.
359 pub(crate) async fn gateway_requests(&self, a: GatewayRequestsArgs) -> Result<Outcome<GatewayRequests>> {
360 let workspace = a.workspace.to_lowercase();
361 if !a.viewer.is_some_and(|viewer| viewer.is_member(&workspace)) {
362 return Ok(members_only());
363 }
364 let limit = a.limit.unwrap_or(PAGE).clamp(1, MAX_PAGE);
365 let mut binds: Vec<JsValue> = vec![workspace.as_str().into()];
366 let older = match a.before.as_deref().map(str::trim).filter(|id| !id.is_empty()) {
367 Some(before) => {
368 binds.push(before.into());
369 " AND (created_at, id) < (SELECT created_at, id FROM gateway_requests WHERE id = ?2)"
370 }
371 None => "",
372 };
373 binds.push(f64::from(limit + 1).into());
374 let at = binds.len();
375 let mut rows = self
376 .db
377 .prepare(format!(
378 "SELECT * FROM gateway_requests WHERE workspace = ?1{older} ORDER BY created_at DESC, id DESC LIMIT ?{at}"
379 ))
380 .bind(&binds)?
381 .all()
382 .await?
383 .results::<RequestRow>()?;
384 let more = rows.len() > limit as usize;
385 rows.truncate(limit as usize);
386 let requests: Vec<GatewayRequest> = rows.into_iter().map(GatewayRequest::from).collect();
387 let next = if more { requests.last().map(|r| r.id.clone()) } else { None };
388 Ok(Outcome::Ok(GatewayRequests { requests, next, retention_days: RETENTION_DAYS }))
389 }
390
391 /// Daily: requests older than the log keeps are deleted.
392 pub(crate) async fn forget_gateway_requests(&self) -> Result<()> {
393 let cutoff = rfc3339(now_ms().saturating_sub(u64::from(RETENTION_DAYS) * DAY_MS));
394 self.db
395 .prepare("DELETE FROM gateway_requests WHERE created_at < ?")
396 .bind(&[cutoff.into()])?
397 .run()
398 .await?;
399 Ok(())
400 }
401}
402
403#[cfg(test)]
404mod tests {
405 use super::*;
406
407 fn sonnet() -> GatewayModel {
408 GatewayModel {
409 model: "claude-sonnet-5-5".into(),
410 name: "Claude Sonnet 5.5".into(),
411 provider: "anthropic".into(),
412 input_micros: 2_000_000,
413 output_micros: 10_000_000,
414 cache_read_micros: 200_000,
415 cache_write_micros: 2_500_000,
416 }
417 }
418
419 #[test]
420 fn a_request_costs_its_tokens_at_the_models_prices() {
421 // A million of each: $2 + $10 + $0.20 + $2.50.
422 assert_eq!(cost_micros(&sonnet(), [1_000_000; 4]), 14_700_000);
423 // 1,000 input and 500 output: $0.002 + $0.005.
424 assert_eq!(cost_micros(&sonnet(), [1_000, 500, 0, 0]), 7_000);
425 // Nothing used costs nothing.
426 assert_eq!(cost_micros(&sonnet(), [0; 4]), 0);
427 }
428
429 #[test]
430 fn a_fraction_of_a_millionth_rounds_up() {
431 // One cache-read token: 0.2 millionths.
432 assert_eq!(cost_micros(&sonnet(), [0, 0, 1, 0]), 1);
433 // A negative price in the table is never a credit: it counts as nothing.
434 let odd = GatewayModel { input_micros: -5, ..sonnet() };
435 assert_eq!(cost_micros(&odd, [10, 0, 0, 0]), 0);
436 }
437
438 #[test]
439 fn the_markup_is_the_price_books_and_zero_in_beta_charges_the_cost() {
440 let cost = cost_micros(&sonnet(), [12_000, 800, 40_000, 0]);
441 assert_eq!(margin_on(cost, 0), cost);
442 assert_eq!(margin_on(1_000, 20), 1_200);
443 }
444
445 #[test]
446 fn the_workspaces_own_key_is_counted_and_never_charged() {
447 let tokens = [50_000, 2_000, 0, 0];
448 assert_eq!(request_cost(false, Some(&sonnet()), tokens), 120_000);
449 assert_eq!(request_cost(true, Some(&sonnet()), tokens), 0);
450 // A model g1t has no price for is never charged a guess.
451 assert_eq!(request_cost(false, None, tokens), 0);
452 }
453
454 #[test]
455 fn out_of_credit_on_the_plan_is_refused_and_the_rest_by_plan() {
456 assert_eq!(plan_standing(PlanKind::Paid, None), Standing::Admitted);
457 assert_eq!(plan_standing(PlanKind::Paid, Some(false)), Standing::OutOfCredit { reload_failed: false });
458 assert_eq!(plan_standing(PlanKind::Paid, Some(true)), Standing::OutOfCredit { reload_failed: true });
459 assert_eq!(plan_standing(PlanKind::Free, None), Standing::NoPlan);
460 // Comped and invoiced workspaces need no credit.
461 assert_eq!(plan_standing(PlanKind::Internal, Some(false)), Standing::Admitted);
462 assert_eq!(plan_standing(PlanKind::Enterprise, Some(false)), Standing::Admitted);
463 assert!(refusal("acme", &plan_standing(PlanKind::Paid, Some(false))).is_some());
464 }
465
466 #[test]
467 fn only_an_admitted_workspace_is_let_through() {
468 assert_eq!(refusal("acme", &Standing::Admitted), None);
469 let out = refusal("acme", &Standing::OutOfCredit { reload_failed: false }).unwrap();
470 assert!(out.contains("out of AI credit") && out.contains("/acme/-/billing#ai-credit"), "{out}");
471 assert!(!out.contains("Auto-reload"));
472 let failed = refusal("acme", &Standing::OutOfCredit { reload_failed: true }).unwrap();
473 assert!(failed.contains("Auto-reload was turned off"));
474 let no_plan = refusal("acme", &Standing::NoPlan).unwrap();
475 assert!(no_plan.contains("g1t plan") && no_plan.contains("own Anthropic key"), "{no_plan}");
476 assert_eq!(refusal("acme", &Standing::Stopped("Over the limit.".into())).as_deref(), Some("Over the limit."));
477 }
478
479 #[test]
480 fn a_ledger_line_names_the_model_tokens_and_token() {
481 assert_eq!(describe("Claude Sonnet 5.5", [12_000, 800, 0, 0], Some("ci")), "AI Gateway: Claude Sonnet 5.5, 12,800 tokens, token ci");
482 assert_eq!(describe("Claude Haiku 4.5", [5, 0, 0, 0], None), "AI Gateway: Claude Haiku 4.5, 5 tokens");
483 assert_eq!(describe("Claude Haiku 4.5", [5, 0, 0, 0], Some(" ")), "AI Gateway: Claude Haiku 4.5, 5 tokens");
484 }
485
486 #[test]
487 fn request_ids_are_the_proxys() {
488 assert!(valid_id("gw_01kkr2m4c8f1t7qh3d6n9w5p0x"));
489 assert!(valid_id("gw_a-b_c"));
490 assert!(!valid_id("run_1"));
491 assert!(!valid_id("gw_x'; DROP TABLE ledger"));
492 assert!(!valid_id(&format!("gw_{}", "a".repeat(80))));
493 }
494
495 #[test]
496 fn the_migration_prices_every_model_it_offers() {
497 let sql = include_str!("../migrations/0045_gateway.sql");
498 for model in ["claude-opus-5-5", "claude-sonnet-5-5", "claude-haiku-4-5"] {
499 assert!(sql.contains(&format!("('{model}', ")), "{model}");
500 }
501 // Sonnet 5.5 at $2 / $10, cache reads $0.20, five-minute cache writes $2.50.
502 assert!(sql.contains("('claude-sonnet-5-5', 'Claude Sonnet 5.5', 'anthropic', 2000000, 10000000, 200000, 2500000"));
503 }
504}