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 tokens | 1 | //! The AI Gateway: a workspace's own model requests, sent with one of its |
| AI Gateway: OpenAI's format, open models, and your own providers | 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). | |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 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 | |
| AI Gateway: OpenAI's format, open models, and your own providers | 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 | |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 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. | |
| AI Gateway: OpenAI's format, open models, and your own providers | 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. | |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 27 | //! - **The log.** Every request is kept for `RETENTION_DAYS`, with its |
| AI Gateway: OpenAI's format, open models, and your own providers | 28 | //! 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 tokens | 30 | //! its prompt or answer. |
| 31 | ||
| 32 | use g1t_contracts::billing::{ | |
| 33 | GatewayAdmitArgs, GatewayModel, GatewayRequest, GatewayRequests, GatewayRequestsArgs, LimitState, PlanKind, RecordGatewayArgs, | |
| 34 | }; | |
| 35 | use g1t_contracts::time::rfc3339; | |
| 36 | use g1t_contracts::{FailureCode, Outcome, new_id}; | |
| 37 | use g1t_kit::now_ms; | |
| 38 | use serde::Deserialize; | |
| 39 | use worker::Result; | |
| 40 | use worker::wasm_bindgen::JsValue; | |
| 41 | ||
| 42 | use crate::credits::{Eligible, month_of}; | |
| 43 | use crate::features::thousands; | |
| 44 | use crate::{Billing, margin_on, members_only, optional}; | |
| 45 | ||
| 46 | /// How long the log keeps a request. | |
| 47 | pub(crate) const RETENTION_DAYS: u32 = 30; | |
| 48 | /// A page of the log: this many when not asked, and at most. | |
| 49 | const PAGE: u32 = 50; | |
| 50 | const MAX_PAGE: u32 = 200; | |
| 51 | const DAY_MS: u64 = 86_400_000; | |
| 52 | ||
| AI Gateway: OpenAI's format, open models, and your own providers | 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)] | |
| 56 | pub(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 | ||
| 64 | impl 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. | |
| 77 | pub(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. | |
| 84 | pub(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 tokens | 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)] | |
| 111 | pub(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. | |
| 125 | pub(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 providers | 139 | pub(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 tokens | 140 | match model { |
| AI Gateway: OpenAI's format, open models, and your own providers | 141 | Some(model) if !own_key => cost_micros(model, used), |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 142 | _ => 0, |
| 143 | } | |
| 144 | } | |
| 145 | ||
| 146 | /// Why a request is refused, in words for whoever sent it, or None. | |
| 147 | pub(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 providers | 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." |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 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. | |
| AI Gateway: OpenAI's format, open models, and your own providers | 164 | pub(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 tokens | 165 | 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 providers | 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`. | |
| 171 | pub(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 tokens | 173 | } |
| 174 | ||
| 175 | /// Whether a request id is one the proxy makes: `gw_` and up to 64 letters, | |
| 176 | /// digits, `_` and `-`. | |
| 177 | pub(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-artifacts | 181 | /// A row of `gateway_models`, the model catalogue (catalogue.rs reads the |
| 182 | /// rest of its columns). | |
| 183 | #[derive(Clone, Deserialize)] | |
| 184 | pub(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 providers | 188 | #[serde(default)] |
| Merge branch 'main' into actions-toolkit-oidc-artifacts | 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, | |
| AI Gateway: OpenAI's format, open models, and your own providers | 194 | #[serde(default)] |
| Merge branch 'main' into actions-toolkit-oidc-artifacts | 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>, | |
| AI Gateway: OpenAI's format, open models, and your own providers | 206 | #[serde(default)] |
| Merge branch 'main' into actions-toolkit-oidc-artifacts | 207 | pub over_cache_write_1h_micros: Option<i64>, |
| AI Gateway: OpenAI's format, open models, and your own providers | 208 | #[serde(default)] |
| Merge branch 'main' into actions-toolkit-oidc-artifacts | 209 | pub aliases: Option<String>, |
| AI Gateway: OpenAI's format, open models, and your own providers | 210 | #[serde(default)] |
| Merge branch 'main' into actions-toolkit-oidc-artifacts | 211 | pub family: Option<String>, |
| AI Gateway: OpenAI's format, open models, and your own providers | 212 | #[serde(default)] |
| Merge branch 'main' into actions-toolkit-oidc-artifacts | 213 | pub tier_hint: Option<String>, |
| AI Gateway: OpenAI's format, open models, and your own providers | 214 | #[serde(default)] |
| Merge branch 'main' into actions-toolkit-oidc-artifacts | 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>, | |
| AI Gateway: OpenAI's format, open models, and your own providers | 238 | #[serde(default)] |
| Merge branch 'main' into actions-toolkit-oidc-artifacts | 239 | pub note: Option<String>, |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 240 | } |
| 241 | ||
| Merge branch 'main' into actions-toolkit-oidc-artifacts | 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. | |
| 245 | pub(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 tokens | 247 | impl 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 providers | 253 | kind: row.kind.unwrap_or_else(|| "chat".to_owned()), |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 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, | |
| AI Gateway: OpenAI's format, open models, and your own providers | 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), | |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 265 | } |
| 266 | } | |
| 267 | } | |
| 268 | ||
| 269 | #[derive(Deserialize)] | |
| 270 | struct 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 providers | 280 | #[serde(default)] |
| 281 | cache_write_1h: Option<f64>, | |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 282 | cost_micros: f64, |
| 283 | charged_micros: f64, | |
| 284 | status: f64, | |
| 285 | own_key: f64, | |
| AI Gateway: OpenAI's format, open models, and your own providers | 286 | #[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 tokens | 292 | streamed: f64, |
| 293 | duration_ms: f64, | |
| 294 | error: Option<String>, | |
| 295 | } | |
| 296 | ||
| 297 | impl 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 providers | 310 | cache_write_hour: n(row.cache_write_1h.unwrap_or(0.0)), |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 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, | |
| AI Gateway: OpenAI's format, open models, and your own providers | 315 | 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 tokens | 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. | |
| 326 | fn number(n: u64) -> JsValue { | |
| 327 | JsValue::from_f64(n as f64) | |
| 328 | } | |
| 329 | ||
| 330 | impl Billing { | |
| Merge branch 'main' into actions-toolkit-oidc-artifacts | 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. | |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 334 | pub(crate) async fn gateway_models(&self) -> Result<Vec<GatewayModel>> { |
| Merge branch 'main' into actions-toolkit-oidc-artifacts | 335 | let mut models: Vec<GatewayModel> = self |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 336 | .db |
| Merge branch 'main' into actions-toolkit-oidc-artifacts | 337 | .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 tokens | 338 | .all() |
| 339 | .await? | |
| 340 | .results::<ModelRow>()? | |
| 341 | .into_iter() | |
| 342 | .map(GatewayModel::from) | |
| Merge branch 'main' into actions-toolkit-oidc-artifacts | 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) | |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 348 | } |
| 349 | ||
| 350 | async fn gateway_model(&self, model: &str) -> Result<Option<GatewayModel>> { | |
| 351 | Ok(self | |
| 352 | .db | |
| Merge branch 'main' into actions-toolkit-oidc-artifacts | 353 | .prepare(format!("SELECT * FROM gateway_models WHERE model = ? AND {OFFERED_SQL}")) |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 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 | } | |
| AI Gateway: OpenAI's format, open models, and your own providers | 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 | }; | |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 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>()); | |
| AI Gateway: OpenAI's format, open models, and your own providers | 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>()); | |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 413 | let priced = if a.own_key { None } else { self.gateway_model(&model_name).await? }; |
| AI Gateway: OpenAI's format, open models, and your own providers | 414 | let cost = request_cost(a.own_key, priced.as_ref(), &used); |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 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, | |
| AI Gateway: OpenAI's format, open models, and your own providers | 423 | 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 tokens | 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), | |
| AI Gateway: OpenAI's format, open models, and your own providers | 439 | number(used.cache_write_1h), |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 440 | (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 providers | 443 | 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 tokens | 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(×tamp), &eligible).await?; | |
| 465 | let owed = charge - drawn.total(); | |
| AI Gateway: OpenAI's format, open models, and your own providers | 466 | 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 tokens | 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)] | |
| AI Gateway: OpenAI's format, open models, and your own providers | 555 | #[path = "gateway_tests.rs"] |
| 556 | mod tests; | |
| Merge the AI Gateway: Anthropic's Messages API on a workspace's tokens | 557 |
This file's history is long; its oldest lines are credited to the oldest commit read.