Skip to content

g1t/services/billing/src/gateway.rs

557 lines24,064 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.

Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens1//! The AI Gateway: a workspace's own model requests, sent with one of its
AI Gateway: OpenAI's format, open models, and your own providers2//! 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).
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens5//!
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
AI Gateway: OpenAI's format, open models, and your own providers15//! 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
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens19//! 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.
AI Gateway: OpenAI's format, open models, and your own providers24//! - **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.
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens27//! - **The log.** Every request is kept for `RETENTION_DAYS`, with its
AI Gateway: OpenAI's format, open models, and your own providers28//! model, format, provider, tokens by kind, cost, status and the token
29//! that sent it; never
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens30//! 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
AI Gateway: OpenAI's format, open models, and your own providers53/// 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 ];
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens101 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.
AI Gateway: OpenAI's format, open models, and your own providers139pub(crate) fn request_cost(own_key: bool, model: Option<&GatewayModel>, used: &Used) -> i64 {
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens140 match model {
AI Gateway: OpenAI's format, open models, and your own providers141 Some(model) if !own_key => cost_micros(model, used),
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens142 _ => 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!(
AI Gateway: OpenAI's format, open models, and your own providers152 "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."
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens153 )),
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.
AI Gateway: OpenAI's format, open models, and your own providers164pub(crate) fn describe(model_name: &str, used: &Used, over: bool, token_name: Option<&str>) -> String {
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens165 let by = token_name.map(str::trim).filter(|name| !name.is_empty()).map_or(String::new(), |name| format!(", token {name}"));
AI Gateway: OpenAI's format, open models, and your own providers166 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" }
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens173}
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
Merge branch 'main' into actions-toolkit-oidc-artifacts181/// 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,
AI Gateway: OpenAI's format, open models, and your own providers188 #[serde(default)]
Merge branch 'main' into actions-toolkit-oidc-artifacts189 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,
AI Gateway: OpenAI's format, open models, and your own providers194 #[serde(default)]
Merge branch 'main' into actions-toolkit-oidc-artifacts195 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>,
AI Gateway: OpenAI's format, open models, and your own providers206 #[serde(default)]
Merge branch 'main' into actions-toolkit-oidc-artifacts207 pub over_cache_write_1h_micros: Option<i64>,
AI Gateway: OpenAI's format, open models, and your own providers208 #[serde(default)]
Merge branch 'main' into actions-toolkit-oidc-artifacts209 pub aliases: Option<String>,
AI Gateway: OpenAI's format, open models, and your own providers210 #[serde(default)]
Merge branch 'main' into actions-toolkit-oidc-artifacts211 pub family: Option<String>,
AI Gateway: OpenAI's format, open models, and your own providers212 #[serde(default)]
Merge branch 'main' into actions-toolkit-oidc-artifacts213 pub tier_hint: Option<String>,
AI Gateway: OpenAI's format, open models, and your own providers214 #[serde(default)]
Merge branch 'main' into actions-toolkit-oidc-artifacts215 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>,
AI Gateway: OpenAI's format, open models, and your own providers238 #[serde(default)]
Merge branch 'main' into actions-toolkit-oidc-artifacts239 pub note: Option<String>,
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens240}
241
Merge branch 'main' into actions-toolkit-oidc-artifacts242/// 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
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens247impl From<ModelRow> for GatewayModel {
248 fn from(row: ModelRow) -> Self {
249 GatewayModel {
250 model: row.model,
251 name: row.name,
252 provider: row.provider,
AI Gateway: OpenAI's format, open models, and your own providers253 kind: row.kind.unwrap_or_else(|| "chat".to_owned()),
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens254 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,
AI Gateway: OpenAI's format, open models, and your own providers258 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),
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens265 }
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,
AI Gateway: OpenAI's format, open models, and your own providers280 #[serde(default)]
281 cache_write_1h: Option<f64>,
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens282 cost_micros: f64,
283 charged_micros: f64,
284 status: f64,
285 own_key: f64,
AI Gateway: OpenAI's format, open models, and your own providers286 #[serde(default)]
287 format: Option<String>,
288 #[serde(default)]
289 provider: Option<String>,
290 #[serde(default)]
291 connection: Option<String>,
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens292 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),
AI Gateway: OpenAI's format, open models, and your own providers310 cache_write_hour: n(row.cache_write_1h.unwrap_or(0.0)),
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens311 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,
AI Gateway: OpenAI's format, open models, and your own providers315 format: format_of(row.format.as_deref().unwrap_or_default()).to_owned(),
316 provider: row.provider.unwrap_or_default(),
317 connection: row.connection,
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens318 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 {
Merge branch 'main' into actions-toolkit-oidc-artifacts331 /// `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.
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens334 pub(crate) async fn gateway_models(&self) -> Result<Vec<GatewayModel>> {
Merge branch 'main' into actions-toolkit-oidc-artifacts335 let mut models: Vec<GatewayModel> = self
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens336 .db
Merge branch 'main' into actions-toolkit-oidc-artifacts337 .prepare(format!("SELECT * FROM gateway_models WHERE {OFFERED_SQL} ORDER BY position, model"))
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens338 .all()
339 .await?
340 .results::<ModelRow>()?
341 .into_iter()
342 .map(GatewayModel::from)
Merge branch 'main' into actions-toolkit-oidc-artifacts343 .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)
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens348 }
349
350 async fn gateway_model(&self, model: &str) -> Result<Option<GatewayModel>> {
351 Ok(self
352 .db
Merge branch 'main' into actions-toolkit-oidc-artifacts353 .prepare(format!("SELECT * FROM gateway_models WHERE model = ? AND {OFFERED_SQL}"))
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens354 .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 }
AI Gateway: OpenAI's format, open models, and your own providers402 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 };
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens409 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>());
AI Gateway: OpenAI's format, open models, and your own providers411 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>());
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens413 let priced = if a.own_key { None } else { self.gateway_model(&model_name).await? };
AI Gateway: OpenAI's format, open models, and your own providers414 let cost = request_cost(a.own_key, priced.as_ref(), &used);
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens415 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,
AI Gateway: OpenAI's format, open models, and your own providers423 cache_write_1h, cost_micros, charged_micros, status, own_key, format, provider, connection,
424 streamed, duration_ms, error)
425 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 0, ?, ?, ?, ?, ?, ?, ?, ?)
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens426 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),
AI Gateway: OpenAI's format, open models, and your own providers439 number(used.cache_write_1h),
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens440 (cost as f64).into(),
441 f64::from(a.status).into(),
442 f64::from(u8::from(a.own_key)).into(),
AI Gateway: OpenAI's format, open models, and your own providers443 format_of(&a.format).into(),
444 provider.as_str().into(),
445 optional(connection.as_deref()),
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens446 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();
AI Gateway: OpenAI's format, open models, and your own providers466 let over = over_threshold(&model, &used);
467 let description = format!("{}{terms_note}{}", describe(&model.name, &used, over, token_name.as_deref()), drawn.note());
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens468 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)]
AI Gateway: OpenAI's format, open models, and your own providers555#[path = "gateway_tests.rs"]
556mod tests;
Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens557

This file's history is long; its oldest lines are credited to the oldest commit read.