Skip to content

g1t/services/billing/src/gateway.rs

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