Billing: the agent rate on a workspace's own model key, and tokens by model
Runs on a workspace's own provider were never charged the agent rate: billing kept no session for them, so it never saw their tokens. The run keeps its model session now, the model proxy counts its tokens under it, and the sandbox reports what its harness counted with the cost. The agent rate is charged on the more of the two, once each, on its own meter (agent_tokens_own, $0.25 a million from 2026-10-22 after the notice) and its own line (<run>/agent-own), shown as "Agent rate, your own model key". The model itself stays between the workspace and its provider. Own-provider runs are closed by the cron without the gateway, charging tokens counted late or for a sandbox that died before reporting. Usage gains agent tokens by model. Migration 0041.
| 23 | 23 | mod security; | |
| 24 | 24 | mod tools; | |
| 25 | 25 | ||
| 26 | − | use g1t_contracts::billing::FinishRunArgs; | |
| 26 | + | use g1t_contracts::billing::{FinishRunArgs, RunTokens}; | |
| 27 | 27 | use g1t_contracts::identity::{ | |
| 28 | 28 | DeviceClaim, DeviceClaimArgs, DeviceStart, DeviceStartArgs, TokenArgs, | |
| 29 | 29 | }; | |
| ⋯ | |||
| 457 | 457 | token: body["token"].as_str().unwrap_or_default().to_owned(), | |
| 458 | 458 | cost_usd: body["cost_usd"].as_f64().unwrap_or_default(), | |
| 459 | 459 | turns: body["turns"].as_u64().unwrap_or_default() as u32, | |
| 460 | + | // What the harness counted; the agent rate is charged on no | |
| 461 | + | // fewer, on a workspace's own model key too. | |
| 462 | + | tokens: body.get("tokens").filter(|t| t.is_object()).map(|t| RunTokens { | |
| 463 | + | input: t["input"].as_u64().unwrap_or_default(), | |
| 464 | + | output: t["output"].as_u64().unwrap_or_default(), | |
| 465 | + | cache_read: t["cache_read"].as_u64().unwrap_or_default(), | |
| 466 | + | cache_write: t["cache_write"].as_u64().unwrap_or_default(), | |
| 467 | + | }), | |
| 460 | 468 | }, | |
| 461 | 469 | ) | |
| 462 | 470 | .await?; | |
| 274 | 274 | /// The runner, which is TypeScript, sends it as `billedTo`. | |
| 275 | 275 | #[serde(default = "g1t", alias = "billedTo")] | |
| 276 | 276 | pub billed_to: String, | |
| 277 | − | /// The model session's id, when its requests go through g1t's AI | |
| 278 | − | /// Gateway: settling charges the run what the gateway priced them at. | |
| 277 | + | /// The model session's id. Through g1t's AI Gateway, settling charges | |
| 278 | + | /// the run what the gateway priced its requests at; on the workspace's | |
| 279 | + | /// own provider, it is what the proxy counts the run's tokens under, | |
| 280 | + | /// for the agent rate. | |
| 279 | 281 | #[serde(default)] | |
| 280 | 282 | pub session: Option<String>, | |
| 281 | − | /// `small` or `large`: the tier g1t routed the run to, when g1t pays | |
| 282 | − | /// for its model. None on the workspace's own provider. | |
| 283 | + | /// `small`, `large` or `frontier`: the tier g1t routed the run to. | |
| 284 | + | /// None when the workspace's own provider names its model. | |
| 283 | 285 | #[serde(default)] | |
| 284 | 286 | pub tier: Option<String>, | |
| 285 | 287 | } | |
| ⋯ | |||
| 303 | 305 | pub cost_usd: f64, | |
| 304 | 306 | #[serde(default)] | |
| 305 | 307 | pub turns: u32, | |
| 308 | + | /// The tokens the run used, as the harness counted them from the | |
| 309 | + | /// provider's answers. On the workspace's own provider, the agent rate | |
| 310 | + | /// is charged on no fewer than these. Absent from older sandboxes. | |
| 311 | + | #[serde(default)] | |
| 312 | + | pub tokens: Option<RunTokens>, | |
| 313 | + | } | |
| 314 | + | ||
| 315 | + | /// The tokens one run used, by kind. | |
| 316 | + | #[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)] | |
| 317 | + | #[serde(rename_all = "camelCase")] | |
| 318 | + | pub struct RunTokens { | |
| 319 | + | #[serde(default)] | |
| 320 | + | pub input: u64, | |
| 321 | + | #[serde(default)] | |
| 322 | + | pub output: u64, | |
| 323 | + | #[serde(default)] | |
| 324 | + | pub cache_read: u64, | |
| 325 | + | #[serde(default)] | |
| 326 | + | pub cache_write: u64, | |
| 327 | + | } | |
| 328 | + | ||
| 329 | + | impl RunTokens { | |
| 330 | + | /// Every token, of every kind: what the agent rate is charged on. | |
| 331 | + | pub fn total(&self) -> u64 { | |
| 332 | + | self.input.saturating_add(self.output).saturating_add(self.cache_read).saturating_add(self.cache_write) | |
| 333 | + | } | |
| 306 | 334 | } | |
| 307 | 335 | ||
| 308 | 336 | ||
| ⋯ | |||
| 391 | 419 | #[serde(default)] | |
| 392 | 420 | pub person: Option<String>, | |
| 393 | 421 | pub model: String, | |
| 394 | − | /// On g1t's hosted models: `small` or `large`. | |
| 422 | + | /// The tier g1t routed the run to: `small`, `large` or `frontier`. | |
| 395 | 423 | #[serde(default)] | |
| 396 | 424 | pub tier: Option<String>, | |
| 397 | 425 | #[serde(default)] | |
| ⋯ | |||
| 3221 | 3249 | pub count: u32, | |
| 3222 | 3250 | } | |
| 3223 | 3251 | ||
| 3252 | + | /// The tokens one model used over the range, as the model proxy counted | |
| 3253 | + | /// them: on g1t's models and the workspace's own provider alike. | |
| 3254 | + | #[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)] | |
| 3255 | + | #[serde(rename_all = "camelCase")] | |
| 3256 | + | pub struct ModelTokens { | |
| 3257 | + | /// The model's id, as it ran. | |
| 3258 | + | pub model: String, | |
| 3259 | + | pub input: u64, | |
| 3260 | + | pub output: u64, | |
| 3261 | + | pub cache_read: u64, | |
| 3262 | + | pub cache_write: u64, | |
| 3263 | + | } | |
| 3264 | + | ||
| 3224 | 3265 | /// One product family over the range. | |
| 3225 | 3266 | #[derive(Clone, Debug, PartialEq, Serialize, Deserialize)] | |
| 3226 | 3267 | #[serde(rename_all = "camelCase")] | |
| ⋯ | |||
| 3246 | 3287 | pub products: Vec<ProductUsage>, | |
| 3247 | 3288 | /// Every project with usage in the range, for the filter. | |
| 3248 | 3289 | pub projects: Vec<String>, | |
| 3290 | + | /// Agent tokens by model over the range, most first. | |
| 3291 | + | #[serde(default)] | |
| 3292 | + | pub models: Vec<ModelTokens>, | |
| 3249 | 3293 | /// The plan's included usage this month, when the workspace has it. | |
| 3250 | 3294 | #[serde(default)] | |
| 3251 | 3295 | pub included: Option<Allowance>, | |
| 44 | 44 | } | |
| 45 | 45 | } | |
| 46 | 46 | ||
| 47 | + | /// The tokens a run used, from Claude Code's closing `usage`: input, | |
| 48 | + | /// output, and prompt-cache reads and writes. On a workspace's own model | |
| 49 | + | /// key this is what g1t charges its agent rate on (the model proxy's own | |
| 50 | + | /// count, when it has more, wins), so it is reported with the cost. | |
| 51 | + | fn run_tokens(event: &Value) -> Value { | |
| 52 | + | let usage = &event["usage"]; | |
| 53 | + | let count = |field: &str| usage[field].as_u64().unwrap_or_default(); | |
| 54 | + | serde_json::json!({ | |
| 55 | + | "input": count("input_tokens"), | |
| 56 | + | "output": count("output_tokens"), | |
| 57 | + | "cache_read": count("cache_read_input_tokens"), | |
| 58 | + | "cache_write": count("cache_creation_input_tokens"), | |
| 59 | + | }) | |
| 60 | + | } | |
| 61 | + | ||
| 47 | 62 | /// Tells g1t what the run cost, so that the workspace it was for can be | |
| 48 | 63 | /// charged. The run's own token, given to this sandbox and to nothing | |
| 49 | 64 | /// else, is the credential. Does nothing where runs are not billed. | |
| 50 | − | fn report_cost(cost_usd: f64, turns: u64) { | |
| 65 | + | fn report_cost(cost_usd: f64, turns: u64, tokens: Value) { | |
| 51 | 66 | let (Ok(api), Ok(run), Ok(token)) = ( | |
| 52 | 67 | std::env::var("G1T_API"), | |
| 53 | 68 | std::env::var("BILLING_RUN"), | |
| ⋯ | |||
| 59 | 74 | "token": token, | |
| 60 | 75 | "cost_usd": cost_usd, | |
| 61 | 76 | "turns": turns, | |
| 77 | + | "tokens": tokens, | |
| 62 | 78 | })); | |
| 63 | 79 | if let Err(error) = sent { | |
| 64 | 80 | eprintln!("g1t-runner: could not report what the run cost: {error}"); | |
| ⋯ | |||
| 125 | 141 | // the session so that spend can be read per pull request. | |
| 126 | 142 | if let Some(cost) = event["total_cost_usd"].as_f64() { | |
| 127 | 143 | let turns = event["num_turns"].as_u64().unwrap_or_default(); | |
| 128 | − | report_cost(cost, turns); | |
| 144 | + | report_cost(cost, turns, run_tokens(event)); | |
| 129 | 145 | if let Some(progress) = progress { | |
| 130 | 146 | progress.cost(cost, turns); | |
| 131 | 147 | } | |
| ⋯ | |||
| 316 | 332 | None => bail!("Claude Code exited ({status}) without a result"), | |
| 317 | 333 | } | |
| 318 | 334 | } | |
| 335 | + | ||
| 336 | + | #[cfg(test)] | |
| 337 | + | mod tests { | |
| 338 | + | use super::*; | |
| 339 | + | ||
| 340 | + | #[test] | |
| 341 | + | fn a_runs_tokens_come_from_the_closing_usage_by_kind() { | |
| 342 | + | let event = serde_json::json!({ | |
| 343 | + | "type": "result", | |
| 344 | + | "total_cost_usd": 0.42, | |
| 345 | + | "usage": { | |
| 346 | + | "input_tokens": 1200, | |
| 347 | + | "output_tokens": 340, | |
| 348 | + | "cache_read_input_tokens": 90000, | |
| 349 | + | "cache_creation_input_tokens": 5000, | |
| 350 | + | }, | |
| 351 | + | }); | |
| 352 | + | assert_eq!( | |
| 353 | + | run_tokens(&event), | |
| 354 | + | serde_json::json!({ "input": 1200, "output": 340, "cache_read": 90000, "cache_write": 5000 }) | |
| 355 | + | ); | |
| 356 | + | // An older harness without usage reports nothing, not an error. | |
| 357 | + | assert_eq!( | |
| 358 | + | run_tokens(&serde_json::json!({ "type": "result" })), | |
| 359 | + | serde_json::json!({ "input": 0, "output": 0, "cache_read": 0, "cache_write": 0 }) | |
| 360 | + | ); | |
| 361 | + | } | |
| 362 | + | } | |
| 916 | 916 | /** The person the run is for, by username. */ | |
| 917 | 917 | person?: string | null; | |
| 918 | 918 | model: string; | |
| 919 | − | tier?: "small" | "large" | null; | |
| 919 | + | tier?: "small" | "large" | "frontier" | null; | |
| 920 | 920 | input: number; | |
| 921 | 921 | output: number; | |
| 922 | 922 | cacheRead: number; | |
| ⋯ | |||
| 1082 | 1082 | model: string; | |
| 1083 | 1083 | /** `workspace` when the run uses the workspace's own model provider. */ | |
| 1084 | 1084 | billedTo?: "g1t" | "workspace"; | |
| 1085 | − | /** The model session's id, so the run can be settled at AI Gateway's price. */ | |
| 1085 | + | /** | |
| 1086 | + | * The model session's id: the run is settled at AI Gateway's price by | |
| 1087 | + | * it, and on the workspace's own provider its tokens are counted under | |
| 1088 | + | * it for the agent rate. | |
| 1089 | + | */ | |
| 1086 | 1090 | session?: string | null; | |
| 1087 | − | /** `small` or `large`: the tier g1t routed the run to, on its hosted models. */ | |
| 1088 | − | tier?: "small" | "large" | null; | |
| 1091 | + | /** The tier g1t routed the run to: `small`, `large` or `frontier`. */ | |
| 1092 | + | tier?: "small" | "large" | "frontier" | null; | |
| 1089 | 1093 | }): Promise<Result<RunTicket | null>>; | |
| 1090 | 1094 | /** Usage over a range of days (`YYYY-MM-DD`, both included), at price, by product, meter, project and day. Members only. */ | |
| 1091 | 1095 | usageReport( | |
| ⋯ | |||
| 1470 | 1474 | ||
| 1471 | 1475 | export type ProductUsage = { key: string; label: string; micros: number; meters: MeterLine[]; features?: FeatureUsage[] }; | |
| 1472 | 1476 | ||
| 1477 | + | /** The tokens one model used over a range, as the model proxy counted them. */ | |
| 1478 | + | export type ModelTokens = { model: string; input: number; output: number; cacheRead: number; cacheWrite: number }; | |
| 1479 | + | ||
| 1473 | 1480 | export type UsageReport = { | |
| 1474 | 1481 | from: string; | |
| 1475 | 1482 | until: string; | |
| ⋯ | |||
| 1477 | 1484 | days: UsageDay[]; | |
| 1478 | 1485 | products: ProductUsage[]; | |
| 1479 | 1486 | projects: string[]; | |
| 1487 | + | /** Agent tokens by model over the range, most first. */ | |
| 1488 | + | models?: ModelTokens[]; | |
| 1480 | 1489 | /** The plan's included usage this month, in micros. */ | |
| 1481 | 1490 | included?: UsageAllowance | null; | |
| 1482 | 1491 | discountPercent?: number | null; | |
| 1 | + | -- The g1t agent rate on runs that use a workspace's own model key. | |
| 2 | + | -- | |
| 3 | + | -- Runs on a workspace's own provider pay the provider for the model, and | |
| 4 | + | -- were charged nothing by g1t but their sandbox time: billing never saw | |
| 5 | + | -- their tokens. The model proxy counts them now (under the run's model | |
| 6 | + | -- session, kept on `runs.session_id` for these runs too), and the sandbox | |
| 7 | + | -- reports what its harness counted. The agent rate is charged on them, on | |
| 8 | + | -- a meter of its own so it is priced and shown apart: | |
| 9 | + | -- | |
| 10 | + | -- - `agent_tokens_own`: $0.25 per million tokens (input, output and | |
| 11 | + | -- cached), the same as on g1t's models. The model itself is never | |
| 12 | + | -- charged: the workspace pays its provider. New, so a rise from nothing: | |
| 13 | + | -- it takes effect after the 14 days' notice (2026-10-22), and owners on | |
| 14 | + | -- the plan are emailed once (`tell_owners_of_rises`). | |
| 15 | + | -- | |
| 16 | + | -- Its ledger lines are `<run>/agent-own` (later `<run>/agent-own/<tokens>`), | |
| 17 | + | -- named "Agent rate, your own model key" on Usage and the statement. | |
| 18 | + | ||
| 19 | + | INSERT OR IGNORE INTO prices (meter, title, unit, cost_micros, markup_percent, source, updated_at) VALUES | |
| 20 | + | ('agent_tokens_own', 'g1t agent rate, your own model key', 'million tokens', 0, 0, 'list', '2026-10-08T00:00:00Z'); | |
| 21 | + | ||
| 22 | + | INSERT OR IGNORE INTO price_versions (id, meter, version, cost_micros, markup_percent, effective_at, reason, created_by, created_at, applied_at) VALUES | |
| 23 | + | ('pv_agent_tokens_own_1', 'agent_tokens_own', 1, 0, 0, '2026-10-08T00:00:00Z', | |
| 24 | + | 'The g1t agent rate on your own model key, not charged yet', 'migration', '2026-10-08T00:00:00Z', '2026-10-08T00:00:00Z'), | |
| 25 | + | ('pv_agent_tokens_own_2', 'agent_tokens_own', 2, 250000, 0, '2026-10-22T00:00:00Z', | |
| 26 | + | 'New: $0.25 per million tokens an agent''s run uses on your own model key (input, output and cached), for context, memory, routing and orchestration; the model itself stays between you and your provider', 'migration', '2026-10-08T00:00:00Z', NULL); |
| 804 | 804 | ||
| 805 | 805 | /// Charges a run's agent rate for the tokens counted since it was last | |
| 806 | 806 | /// charged, once each: when the run reports and again when it is | |
| 807 | − | /// settled, so tokens counted late are charged too. | |
| 808 | − | pub(crate) async fn charge_agent_rate(&self, run_id: &str, run: &RunRow) -> Result<()> { | |
| 807 | + | /// settled, so tokens counted late are charged too. `reported` is what | |
| 808 | + | /// the run's harness counted; the rate is charged on no fewer. | |
| 809 | + | /// | |
| 810 | + | /// On the workspace's own provider the model is not g1t's to charge, | |
| 811 | + | /// but the agent rate is, on its own meter (`agent_tokens_own`) and its | |
| 812 | + | /// own line (`<run>/agent-own`), so Usage and the statement show it as | |
| 813 | + | /// the agent rate on the workspace's own model key. | |
| 814 | + | pub(crate) async fn charge_agent_rate(&self, run_id: &str, run: &RunRow, reported: u64) -> Result<()> { | |
| 809 | 815 | if self.stripe.is_none() { | |
| 810 | 816 | return Ok(()); | |
| 811 | 817 | } | |
| ⋯ | |||
| 824 | 830 | else { | |
| 825 | 831 | return Ok(()); | |
| 826 | 832 | }; | |
| 827 | − | let Some(session) = row.session_id else { return Ok(()) }; | |
| 828 | − | let counted = self | |
| 829 | − | .db | |
| 830 | − | .prepare( | |
| 831 | − | "SELECT SUM(input + output + cache_read + cache_write) AS micros FROM token_usage | |
| 832 | − | WHERE workspace = ? AND session = ? AND day >= ?", | |
| 833 | − | ) | |
| 834 | − | .bind(&[run.workspace.as_str().into(), session.as_str().into(), row.created_at[..10].into()])? | |
| 835 | − | .first::<Sum>(None) | |
| 836 | − | .await? | |
| 837 | − | .and_then(|s| s.micros) | |
| 838 | − | .unwrap_or(0.0) as u64; | |
| 833 | + | let counted = match &row.session_id { | |
| 834 | + | Some(session) => self | |
| 835 | + | .db | |
| 836 | + | .prepare( | |
| 837 | + | "SELECT SUM(input + output + cache_read + cache_write) AS micros FROM token_usage | |
| 838 | + | WHERE workspace = ? AND session = ? AND day >= ?", | |
| 839 | + | ) | |
| 840 | + | .bind(&[run.workspace.as_str().into(), session.as_str().into(), row.created_at[..10].into()])? | |
| 841 | + | .first::<Sum>(None) | |
| 842 | + | .await? | |
| 843 | + | .and_then(|s| s.micros) | |
| 844 | + | .unwrap_or(0.0) as u64, | |
| 845 | + | None => 0, | |
| 846 | + | }; | |
| 839 | 847 | let charged = row.agent_tokens.unwrap_or(0.0) as u64; | |
| 840 | − | if counted <= charged { | |
| 848 | + | let Some(total) = tokens_to_charge(counted, reported, charged) else { | |
| 841 | 849 | return Ok(()); | |
| 842 | − | } | |
| 850 | + | }; | |
| 843 | 851 | // Claimed first: two callers never charge the same tokens. | |
| 844 | 852 | let claimed = self | |
| 845 | 853 | .db | |
| 846 | 854 | .prepare("UPDATE runs SET agent_tokens = ?1 WHERE id = ?2 AND agent_tokens = ?3 RETURNING id") | |
| 847 | − | .bind(&[(counted as f64).into(), run_id.into(), (charged as f64).into()])? | |
| 855 | + | .bind(&[(total as f64).into(), run_id.into(), (charged as f64).into()])? | |
| 848 | 856 | .first::<serde_json::Value>(None) | |
| 849 | 857 | .await?; | |
| 850 | 858 | if claimed.is_none() { | |
| 851 | 859 | return Ok(()); | |
| 852 | 860 | } | |
| 853 | − | let tokens = counted - charged; | |
| 854 | − | let base = agent_rate_micros(tokens, self.agent_rate().await?); | |
| 861 | + | let own = run.own_provider(); | |
| 862 | + | let meter = agent_rate_meter(own); | |
| 863 | + | let tokens = total - charged; | |
| 864 | + | let per_million = self.price(meter).await?.map_or(0.0, |(_, price)| price); | |
| 865 | + | let base = agent_rate_micros(tokens, per_million); | |
| 855 | 866 | // Before the rate takes effect: counted, and nothing charged. | |
| 856 | 867 | if base == 0 { | |
| 857 | 868 | return Ok(()); | |
| ⋯ | |||
| 860 | 871 | let now = rfc3339(now_ms()); | |
| 861 | 872 | let eligible = crate::credits::eligible_for(Some(g1t_contracts::billing::ComputeKind::Agent), None); | |
| 862 | 873 | let drawn = self.draw(&run.workspace, charge, &crate::credits::month_of(&now), &eligible).await?; | |
| 863 | − | let reference = if charged == 0 { format!("{run_id}/agent") } else { format!("{run_id}/agent/{counted}") }; | |
| 874 | + | let reference = agent_rate_reference(run_id, own, charged, total); | |
| 864 | 875 | let what = match run.task.as_str() { | |
| 865 | 876 | "plan" => format!("planning for {}", run.repo), | |
| 866 | 877 | "review" => format!("the review of {}#{}", run.repo, run.number), | |
| 867 | 878 | "update" => format!("catching up {}#{}", run.repo, run.number), | |
| 868 | 879 | _ => format!("work on {}#{}", run.repo, run.number), | |
| 869 | 880 | }; | |
| 870 | − | let description = format!( | |
| 871 | − | "g1t agent rate: {} tokens for {what}{terms_note}{}", | |
| 872 | − | crate::features::thousands(tokens), | |
| 873 | − | drawn.note() | |
| 874 | − | ); | |
| 875 | − | self.enter(&run.workspace, EntryKind::Usage, -(charge - drawn.total()), &description, &reference, Some(run), Some(0), None, None) | |
| 881 | + | let label = if own { "g1t agent rate, your own model key" } else { "g1t agent rate" }; | |
| 882 | + | let description = format!("{label}: {} tokens for {what}{terms_note}{}", crate::features::thousands(tokens), drawn.note()); | |
| 883 | + | // g1t's own charge, even on the workspace's provider: it counts | |
| 884 | + | // toward limits and spend like any other. | |
| 885 | + | let line = RunRow { billed_to: None, ..run.clone() }; | |
| 886 | + | self.enter(&run.workspace, EntryKind::Usage, -(charge - drawn.total()), &description, &reference, Some(&line), Some(0), None, None) | |
| 876 | 887 | .await?; | |
| 877 | 888 | self.db | |
| 878 | 889 | .prepare("UPDATE ledger SET quantity = ?, price_version = ? WHERE reference = ?") | |
| 879 | − | .bind(&[(tokens as f64).into(), optional(self.version_now("agent_tokens").await?.as_deref()), reference.as_str().into()])? | |
| 890 | + | .bind(&[(tokens as f64).into(), optional(self.version_now(meter).await?.as_deref()), reference.as_str().into()])? | |
| 880 | 891 | .run() | |
| 881 | 892 | .await?; | |
| 882 | 893 | self.record_drawn(&reference, &drawn).await?; | |
| ⋯ | |||
| 886 | 897 | } | |
| 887 | 898 | } | |
| 888 | 899 | ||
| 900 | + | /// The price-book meter a run's agent rate is on: its own for runs on the | |
| 901 | + | /// workspace's own model key, so it can be priced and shown apart. | |
| 902 | + | pub(crate) fn agent_rate_meter(own_provider: bool) -> &'static str { | |
| 903 | + | if own_provider { "agent_tokens_own" } else { "agent_tokens" } | |
| 904 | + | } | |
| 905 | + | ||
| 906 | + | /// The tokens a run's agent rate covers now: the more of what the proxy | |
| 907 | + | /// counted and what the harness reported, when that is more than was | |
| 908 | + | /// charged already. None when there is nothing new. | |
| 909 | + | pub(crate) fn tokens_to_charge(counted: u64, reported: u64, charged: u64) -> Option<u64> { | |
| 910 | + | let total = counted.max(reported); | |
| 911 | + | (total > charged).then_some(total) | |
| 912 | + | } | |
| 913 | + | ||
| 914 | + | /// The ledger reference of an agent-rate line: `<run>/agent` the first | |
| 915 | + | /// time, `<run>/agent/<tokens>` for tokens counted later; `agent-own` on | |
| 916 | + | /// the workspace's own model key. | |
| 917 | + | pub(crate) fn agent_rate_reference(run_id: &str, own_provider: bool, charged: u64, total: u64) -> String { | |
| 918 | + | let part = if own_provider { "agent-own" } else { "agent" }; | |
| 919 | + | if charged == 0 { format!("{run_id}/{part}") } else { format!("{run_id}/{part}/{total}") } | |
| 920 | + | } | |
| 921 | + | ||
| 889 | 922 | #[cfg(test)] | |
| 890 | 923 | mod tests { | |
| 891 | 924 | use super::*; | |
| ⋯ | |||
| 968 | 1001 | } | |
| 969 | 1002 | ||
| 970 | 1003 | #[test] | |
| 1004 | + | fn own_key_runs_are_charged_the_agent_rate_on_their_own_meter_and_line() { | |
| 1005 | + | assert_eq!(agent_rate_meter(true), "agent_tokens_own"); | |
| 1006 | + | assert_eq!(agent_rate_meter(false), "agent_tokens"); | |
| 1007 | + | assert_eq!(agent_rate_reference("run_1", true, 0, 900), "run_1/agent-own"); | |
| 1008 | + | assert_eq!(agent_rate_reference("run_1", true, 900, 1_200), "run_1/agent-own/1200"); | |
| 1009 | + | assert_eq!(agent_rate_reference("run_1", false, 0, 900), "run_1/agent"); | |
| 1010 | + | assert_eq!(agent_rate_reference("run_1", false, 900, 1_200), "run_1/agent/1200"); | |
| 1011 | + | } | |
| 1012 | + | ||
| 1013 | + | #[test] | |
| 1014 | + | fn the_rate_covers_the_more_of_what_was_counted_and_reported_once() { | |
| 1015 | + | // The proxy counted nothing (no session, or its reports were lost): | |
| 1016 | + | // the harness's count is charged. | |
| 1017 | + | assert_eq!(tokens_to_charge(0, 5_000, 0), Some(5_000)); | |
| 1018 | + | // The proxy counted more: its count. | |
| 1019 | + | assert_eq!(tokens_to_charge(6_000, 5_000, 0), Some(6_000)); | |
| 1020 | + | // Charged already: only what is new, and nothing twice. | |
| 1021 | + | assert_eq!(tokens_to_charge(6_000, 5_000, 6_000), None); | |
| 1022 | + | assert_eq!(tokens_to_charge(7_000, 0, 6_000), Some(7_000)); | |
| 1023 | + | assert_eq!(tokens_to_charge(0, 0, 0), None); | |
| 1024 | + | } | |
| 1025 | + | ||
| 1026 | + | #[test] | |
| 971 | 1027 | fn only_workspaces_paying_on_the_plan_need_ai_credit() { | |
| 972 | 1028 | assert!(needs_credit(PlanKind::Paid)); | |
| 973 | 1029 | assert!(!needs_credit(PlanKind::Internal)); | |
| 580 | 580 | .prepare( | |
| 581 | 581 | "SELECT id, workspace, repo, number, task, model, token_hash, billed_to, session_id, created_at, finished_at | |
| 582 | 582 | FROM runs | |
| 583 | − | WHERE session_id IS NOT NULL AND settled_at IS NULL | |
| 583 | + | WHERE session_id IS NOT NULL AND settled_at IS NULL AND COALESCE(billed_to, 'g1t') = 'g1t' | |
| 584 | 584 | AND ((finished_at IS NOT NULL AND finished_at < ?1) OR created_at < ?2) | |
| 585 | 585 | ORDER BY created_at LIMIT 10", | |
| 586 | 586 | ) | |
| ⋯ | |||
| 605 | 605 | Ok(()) | |
| 606 | 606 | } | |
| 607 | 607 | ||
| 608 | + | /// Closes runs on a workspace's own model provider: none is on g1t's | |
| 609 | + | /// gateway, so nothing is corrected, but tokens the proxy counted after | |
| 610 | + | /// the run reported, or for a sandbox that died before reporting, are | |
| 611 | + | /// charged their agent rate now. Needs no gateway token. | |
| 612 | + | pub(crate) async fn settle_own_runs(&self) -> Result<()> { | |
| 613 | + | #[derive(Deserialize)] | |
| 614 | + | struct Own { | |
| 615 | + | id: String, | |
| 616 | + | workspace: String, | |
| 617 | + | repo: String, | |
| 618 | + | number: u32, | |
| 619 | + | task: String, | |
| 620 | + | model: String, | |
| 621 | + | token_hash: String, | |
| 622 | + | billed_to: Option<String>, | |
| 623 | + | } | |
| 624 | + | let now = now_ms(); | |
| 625 | + | let runs = self | |
| 626 | + | .db | |
| 627 | + | .prepare( | |
| 628 | + | "SELECT id, workspace, repo, number, task, model, token_hash, billed_to | |
| 629 | + | FROM runs | |
| 630 | + | WHERE billed_to = 'workspace' AND session_id IS NOT NULL AND settled_at IS NULL | |
| 631 | + | AND ((finished_at IS NOT NULL AND finished_at < ?1) OR created_at < ?2) | |
| 632 | + | ORDER BY created_at LIMIT 25", | |
| 633 | + | ) | |
| 634 | + | .bind(&[rfc3339(now - SETTLE_AFTER_MS).into(), rfc3339(now - ABANDONED_AFTER_MS).into()])? | |
| 635 | + | .all() | |
| 636 | + | .await? | |
| 637 | + | .results::<Own>()?; | |
| 638 | + | for run in runs { | |
| 639 | + | let settled_at = rfc3339(now_ms()); | |
| 640 | + | let claimed = self | |
| 641 | + | .db | |
| 642 | + | .prepare( | |
| 643 | + | "UPDATE runs SET settled_at = ?1, finished_at = COALESCE(finished_at, ?1) | |
| 644 | + | WHERE id = ?2 AND settled_at IS NULL RETURNING id", | |
| 645 | + | ) | |
| 646 | + | .bind(&[settled_at.as_str().into(), run.id.as_str().into()])? | |
| 647 | + | .first::<Value>(None) | |
| 648 | + | .await?; | |
| 649 | + | if claimed.is_none() { | |
| 650 | + | continue; | |
| 651 | + | } | |
| 652 | + | let row = RunRow { | |
| 653 | + | workspace: run.workspace, | |
| 654 | + | repo: run.repo, | |
| 655 | + | number: run.number, | |
| 656 | + | task: run.task, | |
| 657 | + | model: run.model, | |
| 658 | + | token_hash: run.token_hash, | |
| 659 | + | billed_to: run.billed_to, | |
| 660 | + | }; | |
| 661 | + | self.charge_agent_rate(&run.id, &row, 0).await?; | |
| 662 | + | } | |
| 663 | + | Ok(()) | |
| 664 | + | } | |
| 665 | + | ||
| 608 | 666 | async fn settle(&self, run: &Unsettled, gateway: &SessionCost) -> Result<()> { | |
| 609 | 667 | let requests = gateway.requests; | |
| 610 | 668 | let row = RunRow { | |
| ⋯ | |||
| 657 | 715 | } | |
| 658 | 716 | // Tokens counted after the run reported are charged their agent | |
| 659 | 717 | // rate now (ai.rs). | |
| 660 | − | self.charge_agent_rate(&run.id, &row).await?; | |
| 718 | + | self.charge_agent_rate(&run.id, &row, 0).await?; | |
| 661 | 719 | if let Some(why) = &short { | |
| 662 | 720 | worker::console_warn!("run {} settled at no less than reported: {why}", run.id); | |
| 663 | 721 | } | |
| 157 | 157 | } | |
| 158 | 158 | } | |
| 159 | 159 | ||
| 160 | − | #[derive(Deserialize)] | |
| 160 | + | #[derive(Clone, Deserialize)] | |
| 161 | 161 | struct RunRow { | |
| 162 | 162 | workspace: String, | |
| 163 | 163 | repo: String, | |
| ⋯ | |||
| 779 | 779 | hash(&token).into(), | |
| 780 | 780 | rfc3339(now).into(), | |
| 781 | 781 | if a.billed_to == "workspace" { "workspace" } else { "g1t" }.into(), | |
| 782 | − | optional(a.session.as_deref().filter(|_| a.billed_to != "workspace")), | |
| 783 | − | optional( | |
| 784 | − | a.tier | |
| 785 | − | .as_deref() | |
| 786 | − | .filter(|tier| a.billed_to != "workspace" && matches!(*tier, "small" | "large")), | |
| 787 | − | ), | |
| 782 | + | // On the workspace's own provider too: the proxy counts the | |
| 783 | + | // run's tokens under it, and the agent rate is charged on them. | |
| 784 | + | optional(a.session.as_deref()), | |
| 785 | + | optional(a.tier.as_deref().filter(|tier| tokens::is_tier(tier))), | |
| 788 | 786 | ])? | |
| 789 | 787 | .run() | |
| 790 | 788 | .await?; | |
| ⋯ | |||
| 819 | 817 | if claimed.is_none() { | |
| 820 | 818 | return Ok(Outcome::Ok(false)); | |
| 821 | 819 | } | |
| 822 | − | // On the workspace's own provider, the model was paid for there, | |
| 823 | − | // and the run's sandbox time is recorded on its own: nothing more | |
| 824 | − | // to charge. | |
| 820 | + | // What the harness counted, the floor of what the agent rate is | |
| 821 | + | // charged on when the proxy counted fewer (or none). | |
| 822 | + | let reported = a.tokens.map_or(0, |tokens| tokens.total()); | |
| 823 | + | // On the workspace's own provider, the model was paid for there and | |
| 824 | + | // the run's sandbox time is recorded on its own: only the agent rate | |
| 825 | + | // is charged, on its own meter. | |
| 825 | 826 | if run.own_provider() { | |
| 827 | + | self.charge_agent_rate(&a.run_id, &run, reported).await?; | |
| 826 | 828 | return Ok(Outcome::Ok(true)); | |
| 827 | 829 | } | |
| 828 | 830 | // Its cost plus the margin, on the account's terms; then the plan's | |
| ⋯ | |||
| 860 | 862 | self.record_drawn(&a.run_id, &drawn).await?; | |
| 861 | 863 | self.record_discount(&a.run_id, discount).await?; | |
| 862 | 864 | self.count_spend(&run.workspace, charge_micros(a.cost_usd, 0), charge - drawn.total(), &drawn).await; | |
| 863 | − | self.charge_agent_rate(&a.run_id, &run).await?; | |
| 865 | + | self.charge_agent_rate(&a.run_id, &run, reported).await?; | |
| 864 | 866 | Ok(Outcome::Ok(true)) | |
| 865 | 867 | } | |
| 866 | 868 | } | |
| ⋯ | |||
| 1138 | 1140 | if let Err(error) = billing.settle_runs(&keeper).await { | |
| 1139 | 1141 | worker::console_error!("settling runs failed: {error}"); | |
| 1140 | 1142 | } | |
| 1143 | + | if let Err(error) = billing.settle_own_runs().await { | |
| 1144 | + | worker::console_error!("closing runs on own providers failed: {error}"); | |
| 1145 | + | } | |
| 1141 | 1146 | // Stripe events billing never received, handled now. | |
| 1142 | 1147 | match billing.replay_events().await { | |
| 1143 | 1148 | Ok(done) => worker::console_log!("stripe replay: {done}"), | |
| ⋯ | |||
| 1456 | 1461 | include_str!("../migrations/0038_staff_credits.sql"), | |
| 1457 | 1462 | include_str!("../migrations/0039_discounts_not_comped.sql"), | |
| 1458 | 1463 | include_str!("../migrations/0040_ai_credit.sql"), | |
| 1464 | + | include_str!("../migrations/0041_agent_rate_own_key.sql"), | |
| 1459 | 1465 | ]; | |
| 1460 | 1466 | ||
| 1461 | 1467 | /// The columns of `table` after the migrations: each with whether an | |
| ⋯ | |||
| 1510 | 1516 | assert!(row("('card_fee_fixed'").contains("300000")); | |
| 1511 | 1517 | assert!(row("('card_fee', 'on'").contains("'on'")); | |
| 1512 | 1518 | } | |
| 1519 | + | ||
| 1520 | + | #[test] | |
| 1521 | + | fn the_agent_rate_on_an_own_model_key_is_its_own_dated_meter() { | |
| 1522 | + | let sql = include_str!("../migrations/0041_agent_rate_own_key.sql"); | |
| 1523 | + | let row = |needle: &str| sql.lines().find(|l| l.contains(needle)).unwrap_or_else(|| panic!("no {needle}")).to_owned(); | |
| 1524 | + | // Nothing until the notice has run… | |
| 1525 | + | assert!(row("('pv_agent_tokens_own_1'").contains("0, 0, '2026-10-08")); | |
| 1526 | + | // …then $0.25 a million tokens, as on g1t's models, 14 days after. | |
| 1527 | + | assert!(row("('pv_agent_tokens_own_2'").contains("250000, 0, '2026-10-22")); | |
| 1528 | + | assert!(row("('agent_tokens_own',").contains("'million tokens', 0, 0, 'list'")); | |
| 1529 | + | assert_eq!(ai::agent_rate_meter(true), "agent_tokens_own"); | |
| 1530 | + | } | |
| 1513 | 1531 | #[test] | |
| 1514 | 1532 | fn every_checkout_insert_fills_the_table() { | |
| 1515 | 1533 | // The plan's, the activation's, a card check's, a prepayment's and | |
| 14 | 14 | ||
| 15 | 15 | use futures_util::future::{try_join, try_join5}; | |
| 16 | 16 | use g1t_contracts::billing::{ | |
| 17 | − | Allowance, FeatureUsage, MeterLine, PlanKind, ProductUsage, ProjectUsage, UsageDay, UsageReport, UsageReportArgs, UsageTotals, | |
| 17 | + | Allowance, FeatureUsage, MeterLine, ModelTokens, PlanKind, ProductUsage, ProjectUsage, UsageDay, UsageReport, UsageReportArgs, UsageTotals, | |
| 18 | 18 | PRODUCTS, | |
| 19 | 19 | }; | |
| 20 | 20 | use g1t_contracts::time::{parse_rfc3339, rfc3339}; | |
| ⋯ | |||
| 46 | 46 | WHEN task = 'security' THEN 'security' | |
| 47 | 47 | WHEN task = 'context' THEN 'context' | |
| 48 | 48 | WHEN task = 'gateway' THEN 'gateway' | |
| 49 | + | WHEN reference LIKE '%/agent-own%' THEN 'agent_rate_own' | |
| 49 | 50 | WHEN reference LIKE '%/agent%' THEN 'agent_rate' | |
| 50 | 51 | ELSE 'agent_models' END"; | |
| 51 | 52 | ||
| 52 | 53 | /// Every meter: key, name, product family and the unit of its quantity. | |
| 53 | − | pub(crate) const METERS: [(&str, &str, &str, &str); 16] = [ | |
| 54 | + | pub(crate) const METERS: [(&str, &str, &str, &str); 17] = [ | |
| 54 | 55 | ("agent_models", "Model tokens", "agent", "tokens"), | |
| 55 | 56 | ("agent_rate", "Agent rate", "agent", "tokens"), | |
| 57 | + | ("agent_rate_own", "Agent rate, your own model key", "agent", "tokens"), | |
| 56 | 58 | ("agent_sandbox", "Agent sandbox time", "agent", "seconds"), | |
| 57 | 59 | ("sandbox", "Sandbox time", "sandboxes", "seconds"), | |
| 58 | 60 | ("self_hosted", "Self-hosted runner time", "sandboxes", "seconds"), | |
| ⋯ | |||
| 330 | 332 | .prepare(format!( | |
| 331 | 333 | "SELECT COALESCE(task, 'implement') AS task, SUM({measure}) AS price, SUM(CASE WHEN {RUN_SQL} THEN 1 ELSE 0 END) AS runs | |
| 332 | 334 | FROM ledger WHERE workspace = ?1 AND kind = 'usage' AND created_at >= ?2 AND created_at < ?3 | |
| 333 | − | AND ({meter}) IN ('agent_models', 'agent_rate', 'agent_sandbox') | |
| 335 | + | AND ({meter}) IN ('agent_models', 'agent_rate', 'agent_rate_own', 'agent_sandbox') | |
| 334 | 336 | GROUP BY 1", | |
| 335 | 337 | meter = METER_KEY_SQL | |
| 336 | 338 | )) | |
| ⋯ | |||
| 425 | 427 | } | |
| 426 | 428 | } | |
| 427 | 429 | let ai: i64 = credits.grants.iter().filter(|g| g.scope == "models").map(|g| g.left_micros).sum(); | |
| 430 | + | let models = self.tokens_by_model(&workspace, &from, &until).await?; | |
| 428 | 431 | Ok(Outcome::Ok(UsageReport { | |
| 429 | 432 | from, | |
| 430 | 433 | until, | |
| ⋯ | |||
| 432 | 435 | days: days_out, | |
| 433 | 436 | products, | |
| 434 | 437 | projects: all_projects, | |
| 438 | + | models, | |
| 435 | 439 | included, | |
| 436 | 440 | discount_percent: (percent > 0).then_some(percent), | |
| 437 | 441 | ai_credit_micros: ai, | |
| ⋯ | |||
| 443 | 447 | } | |
| 444 | 448 | } | |
| 445 | 449 | ||
| 450 | + | impl Billing { | |
| 451 | + | /// Agent tokens by model over the days `from` to `until`, most first: | |
| 452 | + | /// what the model proxy counted, on g1t's models and the workspace's | |
| 453 | + | /// own provider alike. | |
| 454 | + | async fn tokens_by_model(&self, workspace: &str, from: &str, until: &str) -> Result<Vec<ModelTokens>> { | |
| 455 | + | #[derive(Deserialize)] | |
| 456 | + | struct Row { | |
| 457 | + | model: String, | |
| 458 | + | input: Option<f64>, | |
| 459 | + | output: Option<f64>, | |
| 460 | + | cache_read: Option<f64>, | |
| 461 | + | cache_write: Option<f64>, | |
| 462 | + | } | |
| 463 | + | let rows = self | |
| 464 | + | .db | |
| 465 | + | .prepare( | |
| 466 | + | "SELECT model, SUM(input) AS input, SUM(output) AS output, SUM(cache_read) AS cache_read, SUM(cache_write) AS cache_write | |
| 467 | + | FROM token_usage WHERE workspace = ?1 AND day >= ?2 AND day <= ?3 | |
| 468 | + | GROUP BY model ORDER BY SUM(input + output + cache_read + cache_write) DESC LIMIT 20", | |
| 469 | + | ) | |
| 470 | + | .bind(&[workspace.into(), from.into(), until.into()])? | |
| 471 | + | .all() | |
| 472 | + | .await? | |
| 473 | + | .results::<Row>()?; | |
| 474 | + | let n = |v: Option<f64>| v.unwrap_or(0.0).max(0.0) as u64; | |
| 475 | + | Ok(rows | |
| 476 | + | .into_iter() | |
| 477 | + | .map(|r| ModelTokens { model: r.model, input: n(r.input), output: n(r.output), cache_read: n(r.cache_read), cache_write: n(r.cache_write) }) | |
| 478 | + | .collect()) | |
| 479 | + | } | |
| 480 | + | } | |
| 481 | + | ||
| 446 | 482 | #[cfg(test)] | |
| 447 | 483 | mod tests { | |
| 448 | 484 | use super::*; | |
| 35 | 35 | WHEN task = 'storage' THEN 'Private storage' | |
| 36 | 36 | WHEN task = 'cache' THEN 'Actions cache storage' | |
| 37 | 37 | WHEN task = 'git' THEN 'Git operations' | |
| 38 | + | WHEN reference LIKE '%/agent-own%' THEN 'Agent rate, your own model key' | |
| 38 | 39 | WHEN billed_to = 'workspace' THEN 'Runs on your own model provider' | |
| 39 | 40 | WHEN reference LIKE '%/agent%' THEN 'Agent rate' | |
| 40 | 41 | WHEN task = 'gateway' THEN 'AI Gateway' | |
| ⋯ | |||
| 58 | 59 | match kind { | |
| 59 | 60 | "Agent runs" => 0, | |
| 60 | 61 | "Agent rate" => 0, | |
| 62 | + | "Agent rate, your own model key" => 0, | |
| 61 | 63 | "AI Gateway" => 1, | |
| 62 | 64 | "Runs on your own model provider" => 1, | |
| 63 | 65 | "Sandbox time" => 2, | |
| 35 | 35 | pub counts: [u64; 4], | |
| 36 | 36 | } | |
| 37 | 37 | ||
| 38 | + | /// Whether `tier` is one g1t routes to: `small`, `large` or `frontier`. | |
| 39 | + | pub(crate) fn is_tier(tier: &str) -> bool { | |
| 40 | + | matches!(tier, "small" | "large" | "frontier") | |
| 41 | + | } | |
| 42 | + | ||
| 38 | 43 | /// The UTC day of a time, `YYYY-MM-DD`. | |
| 39 | 44 | fn day_of(ms: u64) -> String { | |
| 40 | 45 | rfc3339(ms)[..10].to_owned() | |
| ⋯ | |||
| 63 | 68 | person: person_of(a.person.as_deref()), | |
| 64 | 69 | session: session.chars().take(64).collect(), | |
| 65 | 70 | model: if model.is_empty() { "unknown".to_owned() } else { model.chars().take(200).collect() }, | |
| 66 | − | tier: a.tier.as_deref().filter(|tier| matches!(*tier, "small" | "large")).map(str::to_owned), | |
| 71 | + | tier: a.tier.as_deref().filter(|tier| is_tier(tier)).map(str::to_owned), | |
| 67 | 72 | counts, | |
| 68 | 73 | }) | |
| 69 | 74 | } | |
| ⋯ | |||
| 276 | 281 | assert_eq!(row.tier, None); | |
| 277 | 282 | assert_eq!(row.model, "unknown"); | |
| 278 | 283 | assert_eq!(person_of(None), ""); | |
| 284 | + | // The most capable tier is a tier too. | |
| 285 | + | let frontier = token_row(&RecordTokensArgs { tier: Some("frontier".into()), ..args() }, NOW).unwrap(); | |
| 286 | + | assert_eq!(frontier.tier.as_deref(), Some("frontier")); | |
| 279 | 287 | } | |
| 280 | 288 | ||
| 281 | 289 | #[test] | |