g1t/services/billing/src/keeper.rs

970 lines43,428 bytesCodeBlame
1//! Keeps every price current with what g1t actually pays.
2//!
3//! g1t passes its own costs through, so a price is only right while the
4//! cost under it is. Two jobs keep them right, on the billing service's
5//! cron:
6//!
7//! - **Settling runs** (every 15 minutes). A model run is charged when it
8//! finishes at what the sandbox reported. Each of g1t's hosted runs goes
9//! through its AI Gateway, which prices every request at the provider's
10//! current rates and logs it with the run's session. Settling sums those
11//! logs and corrects the charge to the gateway's figure, with a
12//! correction on the statement. A run whose sandbox died before
13//! reporting is charged here instead of never.
14//! - **Checking costs** (daily). What Cloudflare billed g1t's account, from
15//! its usage API, is measured against how much was used: Containers
16//! against the seconds containers ran, Workers for Platforms per request
17//! and per CPU millisecond. When a measured cost moves, the price book
18//! moves with it, since each price is its cost plus a set markup, and
19//! the change is recorded where anyone can see it. A measurement goes
20//! through `pricing` as a proposal: a small move is applied on its own (a
21//! rise only after customers have had notice), a large or suspect one
22//! waits for staff in sudo, so one odd day of data cannot reprice
23//! everything.
24
25use g1t_contracts::billing::{EntryKind, MICROS_PER_DOLLAR, Price, PriceBook, PriceChange};
26
27use g1t_contracts::time::rfc3339;
28use g1t_kit::now_ms;
29use serde::Deserialize;
30use serde_json::{Value, json};
31use worker::{Env, Fetch, Headers, Method, Request, RequestInit, Result};
32
33use crate::{Billing, RunRow};
34
35/// The cron that also checks costs against Cloudflare's bill.
36pub(crate) const DAILY: &str = "17 4 * * *";
37
38/// A run is settled once its logs have had time to land.
39const SETTLE_AFTER_MS: u64 = 5 * 60 * 1000;
40/// A run with no gateway logs after this is left as reported.
41const GIVE_UP_AFTER_MS: u64 = 3 * 60 * 60 * 1000;
42/// A run never finished after this died without reporting.
43const ABANDONED_AFTER_MS: u64 = 3 * 60 * 60 * 1000;
44
45/// Where the keeper reads what g1t pays.
46pub(crate) struct Keeper {
47 /// `CLOUDFLARE_USAGE_TOKEN`: Billing, Account Analytics and AI Gateway,
48 /// read only.
49 token: Option<String>,
50 /// What reads the bill for `costs`: `CLOUDFLARE_BILLING_TOKEN`
51 /// (Account: Billing Read and Account Analytics Read), or the usage
52 /// token, which has both.
53 billing_token: Option<String>,
54 account: String,
55 gateway: String,
56}
57
58impl Keeper {
59 pub(crate) fn from_env(env: &Env) -> Self {
60 let var = |name: &str| env.var(name).map(|v| v.to_string()).unwrap_or_default();
61 let secret = |name: &str| env.secret(name).ok().map(|v| v.to_string()).filter(|v| !v.is_empty());
62 let token = secret("CLOUDFLARE_USAGE_TOKEN");
63 Keeper {
64 billing_token: secret("CLOUDFLARE_BILLING_TOKEN").or_else(|| token.clone()),
65 token,
66 account: var("CLOUDFLARE_ACCOUNT_ID"),
67 gateway: var("AI_GATEWAY_ID"),
68 }
69 }
70
71 async fn send(&self, method: Method, url: &str, body: Option<Value>) -> Result<Value> {
72 let Some(token) = &self.token else {
73 return Err(worker::Error::RustError("no CLOUDFLARE_USAGE_TOKEN".into()));
74 };
75 send_with(token, method, url, body).await
76 }
77
78 /// Whether Cloudflare's bill can be read.
79 pub(crate) fn can_read_bill(&self) -> bool {
80 self.billing_token.is_some() && !self.account.is_empty()
81 }
82
83 fn billing_token(&self) -> Result<&str> {
84 self.billing_token
85 .as_deref()
86 .ok_or_else(|| worker::Error::RustError("no CLOUDFLARE_BILLING_TOKEN or CLOUDFLARE_USAGE_TOKEN".into()))
87 }
88
89 /// Billable usage from `from` to `to` (dates), as Cloudflare answers it.
90 pub(crate) async fn billable_usage_body(&self, from: &str, to: &str) -> Result<Value> {
91 send_with(self.billing_token()?, Method::Get, &self.api(&format!("/billable-usage?from={from}&to={to}")), None).await
92 }
93
94 /// The account's subscriptions (Workers Paid, add-ons), as Cloudflare
95 /// answers them. Needs Account: Billing Read.
96 pub(crate) async fn subscriptions_body(&self) -> Result<Value> {
97 send_with(self.billing_token()?, Method::Get, &self.api("/subscriptions"), None).await
98 }
99
100 /// A GraphQL Analytics query, as Cloudflare answers it, errors and all.
101 pub(crate) async fn graphql(&self, body: Value) -> Result<Value> {
102 send_with(self.billing_token()?, Method::Post, "https://api.cloudflare.com/client/v4/graphql", Some(body)).await
103 }
104
105 pub(crate) fn account(&self) -> &str {
106 &self.account
107 }
108
109 fn api(&self, path: &str) -> String {
110 format!("https://api.cloudflare.com/client/v4/accounts/{}{path}", self.account)
111 }
112
113 /// AI Gateway's id, empty when there is none.
114 pub(crate) fn gateway(&self) -> &str {
115 &self.gateway
116 }
117
118 /// A GraphQL query with the bill's token, and on failure with the
119 /// keeper's (AI Gateway Read), when that is a different token.
120 pub(crate) async fn graphql_either(&self, body: Value) -> Result<Value> {
121 match self.graphql(body.clone()).await {
122 Ok(answer) => Ok(answer),
123 Err(error) => match &self.token {
124 Some(token) if Some(token) != self.billing_token.as_ref() => {
125 send_with(token, Method::Post, "https://api.cloudflare.com/client/v4/graphql", Some(body)).await
126 }
127 _ => Err(error),
128 },
129 }
130 }
131
132 /// What AI Gateway priced a session's requests at, and how many there
133 /// were, with the requests it had no price for.
134 async fn session_cost(&self, session: &str) -> Result<SessionCost> {
135 let mut total = SessionCost { complete: true, ..SessionCost::default() };
136 for page in 1..=MAX_LOG_PAGES {
137 // The filter goes as URL-encoded JSON; the bracket form is
138 // ignored, and would sum every log there is. Session ids are
139 // [a-z0-9_], which need no escaping inside it.
140 let filter = format!(
141 "%5B%7B%22key%22%3A%22metadata.value%22%2C%22operator%22%3A%22eq%22%2C%22value%22%3A%5B%22{session}%22%5D%7D%5D"
142 );
143 let url = self.api(&format!(
144 "/ai-gateway/gateways/{}/logs?per_page=50&page={page}&filters={filter}",
145 self.gateway
146 ));
147 let body = self.send(Method::Get, &url, None).await?;
148 let logs = body["result"].as_array().cloned().unwrap_or_default();
149 for log in &logs {
150 total.add(log);
151 }
152 if logs.len() < 50 {
153 return Ok(total);
154 }
155 }
156 // More logs than were read: what was read is less than the run.
157 total.complete = false;
158 Ok(total)
159 }
160
161 /// The account's billable usage, one row per service per day, as
162 /// Cloudflare reports it.
163 async fn billable_usage(&self, from: &str, to: &str) -> Result<Vec<UsageRow>> {
164 let body = self
165 .send(Method::Get, &self.api(&format!("/billable-usage?from={from}&to={to}")), None)
166 .await?;
167 let rows = body["result"].as_array().cloned().unwrap_or_default();
168 Ok(rows.iter().filter_map(UsageRow::from_value).collect())
169 }
170
171 /// What g1t's containers used from `since` to `until` (dates), as
172 /// Cloudflare bills it: memory in byte-seconds, and CPU seconds.
173 async fn container_usage(&self, since: &str, until: &str) -> Result<ContainerUsage> {
174 let query = "query ($account: String!, $since: Date!, $until: Date!) {
175 viewer { accounts(filter: { accountTag: $account }) {
176 containersUsageAdaptiveGroups(limit: 1000, filter: { date_geq: $since, date_leq: $until }) {
177 sum { cpuTimeSec allocatedMemory }
178 }
179 } }
180 }";
181 let body = self
182 .send(
183 Method::Post,
184 "https://api.cloudflare.com/client/v4/graphql",
185 Some(json!({ "query": query, "variables": { "account": self.account, "since": since, "until": until } })),
186 )
187 .await?;
188 let groups = body["data"]["viewer"]["accounts"][0]["containersUsageAdaptiveGroups"]
189 .as_array()
190 .cloned()
191 .unwrap_or_default();
192 Ok(groups.iter().fold(ContainerUsage::default(), |total, g| ContainerUsage {
193 cpu_seconds: total.cpu_seconds + g["sum"]["cpuTimeSec"].as_f64().unwrap_or(0.0),
194 memory_byte_seconds: total.memory_byte_seconds + g["sum"]["allocatedMemory"].as_f64().unwrap_or(0.0),
195 }))
196 }
197}
198
199/// A request to Cloudflare's API with a bearer token; anything but 200 is
200/// an error with what Cloudflare said.
201async fn send_with(token: &str, method: Method, url: &str, body: Option<Value>) -> Result<Value> {
202 let headers = Headers::new();
203 headers.set("authorization", &format!("Bearer {token}"))?;
204 headers.set("content-type", "application/json")?;
205 let mut init = RequestInit::new();
206 init.with_method(method).with_headers(headers);
207 if let Some(body) = body {
208 init.with_body(Some(body.to_string().into()));
209 }
210 let mut response = Fetch::Request(Request::new_with_init(url, &init)?).send().await?;
211 let status = response.status_code();
212 let value: Value = response.json().await.unwrap_or(Value::Null);
213 if status != 200 {
214 return Err(worker::Error::RustError(format!("Cloudflare answered {status}: {value}")));
215 }
216 Ok(value)
217}
218
219#[derive(Debug, Default, Clone, Copy)]
220pub(crate) struct ContainerUsage {
221 cpu_seconds: f64,
222 memory_byte_seconds: f64,
223}
224
225/// g1t's sandboxes: Containers' standard-1, half a vCPU, 4 GiB, 8 GB disk.
226const SANDBOX_GIB: f64 = 4.0;
227const SANDBOX_DISK_GB: f64 = 8.0;
228const GIB: f64 = 1024.0 * 1024.0 * 1024.0;
229/// The Durable Object behind each container is billed for as long as the
230/// container runs, at 128 MB.
231const SANDBOX_DO_GB: f64 = 0.125;
232
233/// Cloudflare's published Containers rates, in dollars, used for any rate
234/// the bill does not show yet (while usage is inside the included amount).
235const LIST_MEMORY_GIB_SECOND: f64 = 0.000_002_5;
236const LIST_DISK_GB_SECOND: f64 = 0.000_000_07;
237const LIST_VCPU_SECOND: f64 = 0.000_02;
238/// Durable Objects duration: $12.50 per million GB-seconds.
239const LIST_DO_GB_SECOND: f64 = 0.000_012_5;
240
241/// What one second of a sandbox costs whatever it does, in millionths of a
242/// dollar: its memory and disk, and the Durable Object behind it, for the
243/// whole second. CPU is billed only while busy, on top.
244pub(crate) fn sandbox_base_micros(memory: f64, disk: f64, durable_object: f64) -> f64 {
245 (SANDBOX_GIB * memory + SANDBOX_DISK_GB * disk + SANDBOX_DO_GB * durable_object) * MICROS_PER_DOLLAR as f64
246}
247
248/// What one second of a sandbox costs on average, in millionths of a
249/// dollar: its base, and the CPU sandboxes actually use per second of
250/// running. Runs that report their own CPU are priced on it instead (see
251/// `run_cost`).
252pub(crate) fn sandbox_second_micros(usage: ContainerUsage, memory: f64, disk: f64, vcpu: f64, durable_object: f64) -> Option<f64> {
253 let instance_seconds = usage.memory_byte_seconds / (SANDBOX_GIB * GIB);
254 if instance_seconds < 3600.0 {
255 return None;
256 }
257 let cpu_share = usage.cpu_seconds / instance_seconds;
258 Some(sandbox_base_micros(memory, disk, durable_object) + cpu_share * vcpu * MICROS_PER_DOLLAR as f64)
259}
260
261/// How much more a second of a larger machine's memory and disk (and the
262/// Durable Object behind it) costs than the standard sandbox's, at
263/// Cloudflare's list rates: 1 for the standard machine.
264pub(crate) fn base_scale(memory_gib: f64, disk_gb: f64) -> f64 {
265 let base = |memory: f64, disk: f64| memory * LIST_MEMORY_GIB_SECOND + disk * LIST_DISK_GB_SECOND + SANDBOX_DO_GB * LIST_DO_GB_SECOND;
266 base(memory_gib, disk_gb) / base(SANDBOX_GIB, SANDBOX_DISK_GB)
267}
268
269/// What a run that reported its own CPU cost g1t: its base for every
270/// second, and its vCPU-seconds at the vCPU rate.
271pub(crate) fn run_cost(seconds: i64, cpu_seconds: f64, base_per_second: f64, per_vcpu_second: f64) -> f64 {
272 seconds.max(0) as f64 * base_per_second + cpu_seconds.max(0.0) * per_vcpu_second
273}
274
275/// A unit's marginal rate from the bill: the median, over the days that
276/// were charged, of cost over quantity. None while nothing was charged.
277pub(crate) fn billed_rate(rows: &[&UsageRow]) -> Option<f64> {
278 let mut rates: Vec<f64> = rows
279 .iter()
280 .filter(|r| r.cost > 0.0 && r.quantity > 0.0)
281 .map(|r| r.cost / r.quantity)
282 .collect();
283 if rates.is_empty() {
284 return None;
285 }
286 rates.sort_by(f64::total_cmp);
287 Some(rates[rates.len() / 2])
288}
289
290/// One line of Cloudflare's billable usage.
291#[derive(Debug, Clone)]
292pub(crate) struct UsageRow {
293 period_start: String,
294 period_end: String,
295 service: String,
296 unit: String,
297 quantity: f64,
298 cost: f64,
299}
300
301impl UsageRow {
302 /// Read leniently: the API is new, and its field names are FOCUS's.
303 fn from_value(row: &Value) -> Option<Self> {
304 let text = |keys: &[&str]| keys.iter().find_map(|k| row[*k].as_str()).unwrap_or_default().to_owned();
305 let number = |keys: &[&str]| {
306 keys.iter()
307 .find_map(|k| row[*k].as_f64().or_else(|| row[*k].as_str().and_then(|s| s.parse().ok())))
308 .unwrap_or(0.0)
309 };
310 let service = text(&["ServiceName", "service_name", "service"]);
311 if service.is_empty() {
312 return None;
313 }
314 let family = text(&["ServiceFamilyName", "service_family_name"]);
315 Some(UsageRow {
316 period_start: text(&["ChargePeriodStart", "charge_period_start"]),
317 period_end: text(&["ChargePeriodEnd", "charge_period_end"]),
318 service: if family.is_empty() { service } else { format!("{family} / {service}") },
319 unit: text(&["PricingUnit", "ConsumedUnit", "consumed_unit"]),
320 quantity: number(&["PricingQuantity", "ConsumedQuantity", "pricing_quantity"]),
321 // What g1t pays; list price if nothing was contracted.
322 cost: Some(number(&["ContractedCost", "BilledCost", "contracted_cost"]))
323 .filter(|cost| *cost > 0.0)
324 .unwrap_or_else(|| number(&["ListCost", "list_cost"])),
325 })
326 }
327}
328
329#[derive(Deserialize)]
330struct PriceRow {
331 meter: String,
332 title: String,
333 unit: String,
334 cost_micros: f64,
335 markup_percent: u32,
336 source: String,
337 checked_at: Option<String>,
338 updated_at: String,
339}
340
341#[derive(Deserialize)]
342struct ChangeRow {
343 meter: String,
344 old_cost_micros: f64,
345 new_cost_micros: f64,
346 markup_percent: u32,
347 old_markup_percent: Option<u32>,
348 reason: String,
349 created_at: String,
350}
351
352#[derive(Deserialize)]
353struct Unsettled {
354 id: String,
355 workspace: String,
356 repo: String,
357 number: u32,
358 task: String,
359 model: String,
360 token_hash: String,
361 billed_to: Option<String>,
362 session_id: String,
363 created_at: String,
364 finished_at: Option<String>,
365}
366
367#[derive(Deserialize)]
368struct Charged {
369 cost_micros: Option<i64>,
370 description: String,
371 amount_micros: i64,
372}
373
374/// What a settled run's correction comes to, from the charges at the
375/// reported and the gateway's cost and what the workspace was charged at
376/// first. A charge up is drawn down like any charge; a charge down is
377/// given back only up to what the workspace paid, since what a credit or
378/// a pool paid was never the workspace's money.
379pub(crate) fn correction(reported_charge: i64, gateway_charge: i64, first_charged: i64) -> i64 {
380 let delta = gateway_charge - reported_charge;
381 if delta >= 0 { delta } else { delta.max(-first_charged.max(0)) }
382}
383
384/// A session's logs are read 50 at a time, up to this many pages.
385const MAX_LOG_PAGES: u32 = 40;
386
387/// What AI Gateway's logs say a session cost.
388#[derive(Clone, Debug, Default, PartialEq)]
389pub(crate) struct SessionCost {
390 /// What the gateway priced the requests at, in dollars.
391 pub cost_usd: f64,
392 pub requests: u32,
393 /// Requests that used tokens but that the gateway put no price on: a
394 /// model it has no price for. Their cost is not in `cost_usd`.
395 pub unpriced: u32,
396 /// The models of those, for the statement and the drift.
397 pub unpriced_models: Vec<String>,
398 /// False when there were more logs than were read.
399 pub complete: bool,
400}
401
402impl SessionCost {
403 /// Adds one log, read leniently: `cost` in dollars, `tokens_in` and
404 /// `tokens_out`, `cached` for an answer from the gateway's own cache
405 /// (which costs nothing).
406 pub(crate) fn add(&mut self, log: &Value) {
407 self.requests += 1;
408 let number = |key: &str| log[key].as_f64().or_else(|| log[key].as_str().and_then(|s| s.parse().ok()));
409 let cost = number("cost").filter(|c| c.is_finite() && *c > 0.0);
410 let tokens = number("tokens_in").unwrap_or(0.0) + number("tokens_out").unwrap_or(0.0);
411 let cached = log["cached"].as_bool().unwrap_or(false);
412 match cost {
413 Some(cost) => self.cost_usd += cost,
414 None if tokens > 0.0 && !cached => {
415 self.unpriced += 1;
416 let model = log["model"].as_str().unwrap_or("an unnamed model").to_owned();
417 if !self.unpriced_models.contains(&model) {
418 self.unpriced_models.push(model);
419 }
420 }
421 None => {}
422 }
423 }
424
425 /// Whether the gateway's figure is the whole of what the run cost.
426 pub(crate) fn whole(&self) -> bool {
427 self.complete && self.unpriced == 0
428 }
429}
430
431/// The cost a run is settled at, in millionths: the gateway's figure when
432/// it priced every request; otherwise (a model it has no price for, or
433/// more logs than were read) never less than the sandbox reported, since
434/// the gateway's sum is then short of what the provider bills. With a
435/// reason for the statement and the drift when it is not the gateway's
436/// figure alone.
437pub(crate) fn settled_cost(reported_micros: i64, gateway: &SessionCost) -> (i64, Option<String>) {
438 // Not held to MAX_RUN_COST_USD: the gateway's figure is trusted.
439 let priced = if gateway.cost_usd.is_finite() { (gateway.cost_usd.max(0.0) * MICROS_PER_DOLLAR as f64).ceil() as i64 } else { 0 };
440 if gateway.whole() {
441 return (priced, None);
442 }
443 let mut why = Vec::new();
444 if gateway.unpriced > 0 {
445 why.push(format!(
446 "AI Gateway has no price for {} of its {} requests ({})",
447 gateway.unpriced,
448 gateway.requests,
449 gateway.unpriced_models.join(", ")
450 ));
451 }
452 if !gateway.complete {
453 why.push(format!("more than {} of its requests were logged", gateway.requests));
454 }
455 (priced.max(reported_micros.max(0)), Some(why.join("; ")))
456}
457
458fn ms(timestamp: &str) -> u64 {
459 // RFC 3339 in UTC, as g1t writes them.
460 worker::js_sys::Date::parse(timestamp) as u64
461}
462
463impl Billing {
464 pub(crate) async fn prices(&self) -> Result<PriceBook> {
465 // Three reads at once; the plans are priced from the rows already
466 // read, not one query per meter (status.g1t.sh times this call).
467 let (prices, changes, coming) = futures_util::future::join3(
468 async { self.db.prepare("SELECT * FROM prices ORDER BY rowid").all().await?.results::<PriceRow>() },
469 async {
470 self.db
471 .prepare("SELECT * FROM price_changes ORDER BY created_at DESC LIMIT 20")
472 .all()
473 .await?
474 .results::<ChangeRow>()
475 },
476 // Changes still to come first, so a rise is seen before it is charged.
477 self.coming_changes(),
478 )
479 .await;
480 let (prices, changes, coming) = (prices?, changes?, coming?);
481 let book: std::collections::BTreeMap<&str, f64> = prices
482 .iter()
483 .map(|row| (row.meter.as_str(), Price::price_for(row.cost_micros, row.markup_percent)))
484 .collect();
485 let plans: Vec<_> = g1t_contracts::billing::Feature::ALL.iter().map(|_| self.plan_at(&book)).collect();
486 Ok(PriceBook {
487 prices: prices
488 .into_iter()
489 .map(|row| Price {
490 price_micros: Price::price_for(row.cost_micros, row.markup_percent),
491 meter: row.meter,
492 title: row.title,
493 unit: row.unit,
494 cost_micros: row.cost_micros,
495 markup_percent: row.markup_percent,
496 source: row.source,
497 checked_at: row.checked_at,
498 updated_at: row.updated_at,
499 })
500 .collect(),
501 changes: coming
502 .into_iter()
503 .chain(changes.into_iter().map(|row| PriceChange {
504 meter: row.meter,
505 old_cost_micros: row.old_cost_micros,
506 new_cost_micros: row.new_cost_micros,
507 markup_percent: row.markup_percent,
508 old_markup_percent: row.old_markup_percent,
509 reason: row.reason,
510 created_at: row.created_at,
511 effective_at: None,
512 }))
513 .collect(),
514 model_margin_percent: self.margin_percent,
515 plans,
516 free: Some(g1t_contracts::billing::FreeTier {
517 trial_workspace_micros: if self.trials_on { self.plans.trial_workspace_micros } else { 0 },
518 trial_monthly_pool_micros: if self.trials_on { self.plans.trial_monthly_pool_micros } else { 0 },
519 oss_pool_micros: self.plans.oss_pool_micros,
520 oss_repo_micros: self.plans.oss_repo_micros,
521 free_private_storage_bytes: self.plans.free_storage_bytes,
522 audit_retention_days: self.plans.free_audit_days,
523 plan_audit_retention_days: self.plans.audit_days,
524 min_charge_micros: self.plans.min_charge_micros,
525 git_operations_included: self.plans.git_included,
526 paid_start_ceiling_micros: self.plans.paid_start_micros,
527 overage_forgive_cost_micros: self.plans.forgive_cost_micros,
528 }),
529 })
530 }
531
532 /// Whether the costs have never been checked against Cloudflare's bill.
533 pub(crate) async fn never_checked(&self) -> Result<bool> {
534 Ok(self
535 .db
536 .prepare("SELECT meter FROM prices WHERE checked_at IS NOT NULL LIMIT 1")
537 .first::<Value>(None)
538 .await?
539 .is_none())
540 }
541
542 /// A meter's cost and price per unit, from the book.
543 pub(crate) async fn price(&self, meter: &str) -> Result<Option<(f64, f64)>> {
544 #[derive(Deserialize)]
545 struct Row {
546 cost_micros: f64,
547 markup_percent: u32,
548 }
549 Ok(self
550 .db
551 .prepare("SELECT cost_micros, markup_percent FROM prices WHERE meter = ?")
552 .bind(&[meter.into()])?
553 .first::<Row>(None)
554 .await?
555 .map(|row| (row.cost_micros, Price::price_for(row.cost_micros, row.markup_percent))))
556 }
557
558 /// Corrects finished runs to what AI Gateway priced them at, and
559 /// charges runs whose sandbox died before reporting.
560 pub(crate) async fn settle_runs(&self, keeper: &Keeper) -> Result<()> {
561 if keeper.token.is_none() || keeper.gateway.is_empty() {
562 return Ok(());
563 }
564 let now = now_ms();
565 let runs = self
566 .db
567 .prepare(
568 "SELECT id, workspace, repo, number, task, model, token_hash, billed_to, session_id, created_at, finished_at
569 FROM runs
570 WHERE session_id IS NOT NULL AND settled_at IS NULL
571 AND ((finished_at IS NOT NULL AND finished_at < ?1) OR created_at < ?2)
572 ORDER BY created_at LIMIT 10",
573 )
574 .bind(&[rfc3339(now - SETTLE_AFTER_MS).into(), rfc3339(now - ABANDONED_AFTER_MS).into()])?
575 .all()
576 .await?
577 .results::<Unsettled>()?;
578 for run in runs {
579 let gateway = match keeper.session_cost(&run.session_id).await {
580 Ok(found) => found,
581 Err(error) => {
582 worker::console_error!("could not read gateway logs for {}: {error}", run.id);
583 continue;
584 }
585 };
586 let since = ms(run.finished_at.as_deref().unwrap_or(&run.created_at));
587 if gateway.requests == 0 && now.saturating_sub(since) < GIVE_UP_AFTER_MS {
588 continue;
589 }
590 self.settle(&run, &gateway).await?;
591 }
592 Ok(())
593 }
594
595 async fn settle(&self, run: &Unsettled, gateway: &SessionCost) -> Result<()> {
596 let requests = gateway.requests;
597 let row = RunRow {
598 workspace: run.workspace.clone(),
599 repo: run.repo.clone(),
600 number: run.number,
601 task: run.task.clone(),
602 model: run.model.clone(),
603 token_hash: run.token_hash.clone(),
604 billed_to: run.billed_to.clone(),
605 };
606 let charged = self
607 .db
608 .prepare("SELECT cost_micros, description, amount_micros FROM ledger WHERE reference = ?")
609 .bind(&[run.id.as_str().into()])?
610 .first::<Charged>(None)
611 .await?;
612 let reported = charged.as_ref().and_then(|c| c.cost_micros).unwrap_or(0);
613 // The gateway's figure; never under what the sandbox reported when
614 // the gateway could not price all of it (see `settled_cost`).
615 let (gateway_micros, short) = settled_cost(reported, gateway);
616 let terms = self.terms_of(&run.workspace).await?;
617 // A cost's charge on the account's terms, and what a discount gave
618 // below cost plus the margin (counted as given, see `charged`).
619 let charge_for = |micros: i64| {
620 if self.free { (0, 0) } else { terms.discounted(crate::margin_on(micros, self.margin_percent)) }
621 };
622 let settled_at = rfc3339(now_ms());
623 // Claim it, so two crons never settle it twice.
624 let claimed = self
625 .db
626 .prepare(
627 "UPDATE runs SET settled_at = ?, gateway_cost_micros = ?, gateway_note = ?, finished_at = COALESCE(finished_at, ?)
628 WHERE id = ? AND settled_at IS NULL RETURNING id",
629 )
630 .bind(&[
631 settled_at.as_str().into(),
632 (gateway_micros as f64).into(),
633 short.as_deref().map_or(worker::wasm_bindgen::JsValue::NULL, Into::into),
634 settled_at.as_str().into(),
635 run.id.as_str().into(),
636 ])?
637 .first::<Value>(None)
638 .await?;
639 if claimed.is_none() || requests == 0 {
640 return Ok(());
641 }
642 if let Some(why) = &short {
643 worker::console_warn!("run {} settled at no less than reported: {why}", run.id);
644 }
645 let short_note = short.as_ref().map_or(String::new(), |why| format!(" ({why}; charged at no less than the sandbox reported)"));
646 let free_note = if self.free { " (free while g1t is being built out)" } else { "" };
647 match charged {
648 // Never reported: charged now, from the gateway's figure.
649 None => {
650 let (charge, discount) = charge_for(gateway_micros);
651 let eligible = crate::credits::eligible_for(Some(g1t_contracts::billing::ComputeKind::Agent), None);
652 let drawn = self.draw(&run.workspace, charge, &settled_at[..7], &eligible).await?;
653 let description = format!(
654 "Work on {}#{}, settled from AI Gateway after the sandbox stopped without reporting{short_note}{free_note}{}",
655 run.repo,
656 run.number,
657 drawn.note()
658 );
659 self.enter(&run.workspace, EntryKind::Usage, -(charge - drawn.total()), &description, &run.id, Some(&row), Some(gateway_micros), None, None)
660 .await?;
661 self.record_drawn(&run.id, &drawn).await?;
662 self.record_discount(&run.id, discount).await?;
663 self.count_spend(&run.workspace, gateway_micros, charge - drawn.total(), &drawn).await;
664 }
665 Some(charged) => {
666 let delta = gateway_micros - reported;
667 if delta == 0 {
668 return Ok(());
669 }
670 let ((was, was_given), (now, now_given)) = (charge_for(reported), charge_for(gateway_micros));
671 let change = correction(was, now, -charged.amount_micros);
672 // What the discount gives moves with the charge.
673 let discount = now_given - was_given;
674 // A charge up is paid for like any other charge.
675 let drawn = if change > 0 {
676 let eligible = crate::credits::eligible_for(Some(g1t_contracts::billing::ComputeKind::Agent), None);
677 self.draw(&run.workspace, change, &settled_at[..7], &eligible).await?
678 } else {
679 crate::credits::Drawn::default()
680 };
681 let description = format!(
682 "Correction to “{}”: AI Gateway priced its {requests} model requests at {}, not {}{short_note}{}",
683 charged.description,
684 crate::features::dollars(gateway_micros),
685 crate::features::dollars(reported),
686 drawn.note(),
687 );
688 let reference = format!("{}/settled", run.id);
689 self.enter(
690 &run.workspace,
691 EntryKind::Usage,
692 -(change - drawn.total()),
693 &description,
694 &reference,
695 Some(&row),
696 Some(delta),
697 None,
698 None,
699 )
700 .await?;
701 self.record_drawn(&reference, &drawn).await?;
702 self.record_discount(&reference, discount).await?;
703 self.count_spend(&run.workspace, delta, change - drawn.total(), &drawn).await;
704 }
705 }
706 Ok(())
707 }
708
709 /// Checks each cost against what Cloudflare billed this month, and
710 /// moves the ones that changed.
711 pub(crate) async fn reconcile(&self, keeper: &Keeper) -> Result<()> {
712 if keeper.token.is_none() {
713 return Ok(());
714 }
715 let now = rfc3339(now_ms());
716 let today = &now[..10];
717 let since = rfc3339(now_ms() - 30 * 24 * 60 * 60 * 1000);
718 let rows = keeper.billable_usage(&since[..10], today).await?;
719 for row in &rows {
720 self.db
721 .prepare(
722 "INSERT INTO cloudflare_usage (period_start, period_end, service, unit, quantity, cost_usd, fetched_at)
723 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
724 ON CONFLICT (period_start, service, unit) DO UPDATE SET
725 period_end = ?2, quantity = ?5, cost_usd = ?6, fetched_at = ?7",
726 )
727 .bind(&[
728 row.period_start.as_str().into(),
729 row.period_end.as_str().into(),
730 row.service.as_str().into(),
731 row.unit.as_str().into(),
732 row.quantity.into(),
733 row.cost.into(),
734 now.as_str().into(),
735 ])?
736 .run()
737 .await?;
738 }
739
740 let named = |words: &[&str]| -> Vec<&UsageRow> {
741 rows.iter()
742 .filter(|r| {
743 let service = r.service.to_lowercase();
744 words.iter().all(|word| service.contains(word))
745 })
746 .collect()
747 };
748
749 // Containers: each resource at what the bill shows it costs, or
750 // the published rate while the included amount still covers it,
751 // over how much CPU g1t's sandboxes really use per second.
752 let memory = billed_rate(&named(&["container memory"]));
753 let disk = billed_rate(&named(&["container disk"]));
754 let vcpu = billed_rate(&named(&["container vcpu"]));
755 let durable_object = billed_rate(&named(&["durable objects", "duration"]));
756 let usage = keeper.container_usage(&since[..10], today).await?;
757 let rates = (
758 memory.unwrap_or(LIST_MEMORY_GIB_SECOND),
759 disk.unwrap_or(LIST_DISK_GB_SECOND),
760 vcpu.unwrap_or(LIST_VCPU_SECOND),
761 durable_object.unwrap_or(LIST_DO_GB_SECOND),
762 );
763 // The parts, for runs that report their own CPU.
764 let parts_reason = "Cloudflare's Containers and Durable Objects rates, as billed or published";
765 self.measure("sandbox_base_second", sandbox_base_micros(rates.0, rates.1, rates.3), parts_reason).await?;
766 self.measure("sandbox_cpu_second", rates.2 * MICROS_PER_DOLLAR as f64, parts_reason).await?;
767 if let Some(per_second) = sandbox_second_micros(usage, rates.0, rates.1, rates.2, rates.3) {
768 let instance_seconds = usage.memory_byte_seconds / (SANDBOX_GIB * GIB);
769 let billed = [("memory", memory), ("disk", disk), ("vCPU", vcpu), ("Durable Object duration", durable_object)]
770 .iter()
771 .filter(|(_, rate)| rate.is_some())
772 .map(|(name, _)| *name)
773 .collect::<Vec<_>>();
774 let reason = format!(
775 "Sandboxes used {:.2} vCPU per second over {:.0} hours of Cloudflare Containers in the last 30 days, with the Durable Object behind each; {}",
776 usage.cpu_seconds / instance_seconds,
777 instance_seconds / 3600.0,
778 if billed.is_empty() {
779 "rates are Cloudflare's published ones".to_owned()
780 } else {
781 format!("{} at what Cloudflare billed", billed.join(", "))
782 },
783 );
784 for meter in ["sandbox_second", "build_second"] {
785 self.measure(meter, per_second, &reason).await?;
786 }
787 }
788 // Apps run as Workers: per million requests and CPU milliseconds,
789 // once the bill shows them charged.
790 // Security scans' CPU follows the same Workers CPU rate.
791 let app_meters: [(&str, &[&str], &str); 3] = [
792 ("app_requests", &["workers", "requests"], "requests"),
793 ("app_cpu", &["workers cpu"], "CPU ms"),
794 ("scan_cpu", &["workers cpu"], "CPU ms"),
795 ];
796 for (meter, words, unit) in app_meters {
797 if let Some(rate) = billed_rate(&named(words)) {
798 let reason = format!("Cloudflare billed Workers {unit} at ${:.2} per million", rate * 1e6);
799 self.measure(meter, rate * 1e6 * MICROS_PER_DOLLAR as f64, &reason).await?;
800 }
801 }
802 self.db
803 .prepare("UPDATE prices SET checked_at = ?")
804 .bind(&[now.as_str().into()])?
805 .run()
806 .await?;
807 Ok(())
808 }
809
810 /// Proposes moving a meter's cost to a measurement (see `pricing`):
811 /// applied on its own when small, after notice when a rise; left for
812 /// staff when large or suspect.
813 async fn measure(&self, meter: &str, measured: f64, reason: &str) -> Result<()> {
814 if let Some(outcome) = self.propose(meter, measured, reason, "keeper").await? {
815 worker::console_log!("{meter}: {outcome}");
816 }
817 Ok(())
818 }
819}
820
821#[cfg(test)]
822mod tests {
823 use super::*;
824
825 #[test]
826 fn a_correction_never_gives_back_what_the_workspace_did_not_pay() {
827 // Up by 2 cents: charged in full (then drawn down like any charge).
828 assert_eq!(correction(100_000, 120_000, 100_000), 20_000);
829 // Down by 2 cents, all of it paid by the workspace: given back.
830 assert_eq!(correction(120_000, 100_000, 120_000), -20_000);
831 // Down, but the open-source pool paid all but a cent: a cent back.
832 assert_eq!(correction(120_000, 100_000, 10_000), -10_000);
833 // Paid entirely by a credit or pool: nothing back.
834 assert_eq!(correction(120_000, 100_000, 0), 0);
835 }
836
837 fn logs(entries: &[Value]) -> SessionCost {
838 let mut total = SessionCost { complete: true, ..SessionCost::default() };
839 for log in entries {
840 total.add(log);
841 }
842 total
843 }
844
845 #[test]
846 fn a_run_is_settled_at_the_gateways_figure_when_it_priced_every_request() {
847 let gateway = logs(&[
848 json!({ "cost": 0.012, "tokens_in": 4000, "tokens_out": 300, "model": "claude-sonnet-5-5" }),
849 json!({ "cost": "0.003", "tokens_in": 900, "tokens_out": 40, "model": "claude-haiku-4-5" }),
850 // Served from the gateway's own cache: no cost, and none owed.
851 json!({ "cost": 0, "tokens_in": 900, "tokens_out": 40, "cached": true }),
852 // An error with no tokens costs nothing either.
853 json!({ "cost": null, "tokens_in": 0, "tokens_out": 0 }),
854 ]);
855 assert!(gateway.whole());
856 assert_eq!(gateway.requests, 4);
857 // Down from what the sandbox said, or up: the gateway's figure.
858 assert_eq!(settled_cost(20_000, &gateway), (15_000, None));
859 assert_eq!(settled_cost(9_000, &gateway), (15_000, None));
860 }
861
862 #[test]
863 fn a_model_the_gateway_cannot_price_is_never_settled_down_to_nothing() {
864 let gateway = logs(&[
865 json!({ "cost": 0.002, "tokens_in": 100, "tokens_out": 10, "model": "claude-haiku-4-5" }),
866 json!({ "cost": 0, "tokens_in": 50_000, "tokens_out": 2_000, "model": "claude-new-1" }),
867 json!({ "tokens_in": 50_000, "tokens_out": 2_000, "model": "claude-new-1" }),
868 ]);
869 assert!(!gateway.whole());
870 assert_eq!((gateway.unpriced, gateway.unpriced_models.clone()), (2, vec!["claude-new-1".to_owned()]));
871 // The sandbox said $0.90: kept, not cut to the gateway's $0.002.
872 let (cost, why) = settled_cost(900_000, &gateway);
873 assert_eq!(cost, 900_000);
874 assert!(why.unwrap().contains("no price for 2 of its 3 requests (claude-new-1)"));
875 // A sandbox that reported less than the gateway priced: the gateway's.
876 assert_eq!(settled_cost(1_000, &gateway).0, 2_000);
877 }
878
879 #[test]
880 fn more_logs_than_were_read_never_settle_a_run_down() {
881 let mut gateway = logs(&[json!({ "cost": 1.0, "tokens_in": 1, "tokens_out": 1 })]);
882 gateway.complete = false;
883 let (cost, why) = settled_cost(3_000_000, &gateway);
884 assert_eq!(cost, 3_000_000);
885 assert!(why.unwrap().contains("more than 1 of its requests"));
886 }
887
888 #[test]
889 fn a_gateway_figure_over_the_report_cap_is_charged_in_full() {
890 // A sandbox's report is believed up to $100; the gateway's is not capped.
891 let gateway = logs(&[json!({ "cost": 140.0, "tokens_in": 1, "tokens_out": 1 })]);
892 assert_eq!(settled_cost(100_000_000, &gateway).0, 140_000_000);
893 assert_eq!(crate::charge_micros(140.0, 20), 120_000_000);
894 assert_eq!(crate::margin_on(140_000_000, 20), 168_000_000);
895 // Exactly cost plus the margin, rounded up, in whole micros (dollars
896 // as floats can come out a micro high), never under it.
897 for cost in [0_i64, 1, 7, 999, 123_457, 99_999_999] {
898 let exact = (cost * 120 + 99) / 100;
899 assert_eq!(crate::margin_on(cost, 20), exact, "{cost}");
900 assert!(crate::margin_on(cost, 20) * 100 >= cost * 120, "{cost}");
901 assert!(crate::charge_micros(cost as f64 / 1e6, 20) >= exact, "{cost}");
902 }
903 }
904
905 #[test]
906 fn a_sandbox_second_is_its_memory_and_disk_and_the_cpu_it_uses() {
907 // An hour of sandboxes that kept a fifth of a vCPU busy.
908 let usage = ContainerUsage { cpu_seconds: 720.0, memory_byte_seconds: 3600.0 * 4.0 * GIB };
909 let micros =
910 sandbox_second_micros(usage, LIST_MEMORY_GIB_SECOND, LIST_DISK_GB_SECOND, LIST_VCPU_SECOND, LIST_DO_GB_SECOND).unwrap();
911 // 4 x 2.5 + 8 x 0.07 + 0.125 x 12.5 + 0.2 x 20 = 16.1225
912 assert!((micros - 16.1225).abs() < 1e-9, "{micros}");
913 // The Durable Object adds about 11% to the second it left out.
914 assert!((sandbox_base_micros(LIST_MEMORY_GIB_SECOND, LIST_DISK_GB_SECOND, LIST_DO_GB_SECOND) - 12.1225).abs() < 1e-9);
915 // Too little use to say anything.
916 assert!(sandbox_second_micros(ContainerUsage { cpu_seconds: 1.0, memory_byte_seconds: GIB }, 1.0, 1.0, 1.0, 1.0).is_none());
917 }
918
919 #[test]
920 fn a_run_that_reports_its_cpu_is_priced_on_it() {
921 let base = sandbox_base_micros(LIST_MEMORY_GIB_SECOND, LIST_DISK_GB_SECOND, LIST_DO_GB_SECOND);
922 let vcpu = LIST_VCPU_SECOND * MICROS_PER_DOLLAR as f64;
923 // A 10-minute cargo build that kept its half vCPU busy throughout.
924 let heavy = run_cost(600, 300.0, base, vcpu);
925 assert!((heavy - (600.0 * 12.1225 + 300.0 * 20.0)).abs() < 1e-6);
926 // The same ten minutes, mostly idle, costs less.
927 let light = run_cost(600, 30.0, base, vcpu);
928 assert!(light < heavy);
929 // The average would have under-priced the heavy one.
930 let average = 600.0 * (base + 0.195 * vcpu);
931 assert!(average < heavy && average > light);
932 assert_eq!(run_cost(0, -1.0, base, vcpu), 0.0);
933 // A larger machine's base: its memory and disk, not its CPU.
934 assert!((base_scale(SANDBOX_GIB, SANDBOX_DISK_GB) - 1.0).abs() < 1e-12);
935 assert!((base_scale(12.0, 20.0) - 2.72).abs() < 0.01, "{}", base_scale(12.0, 20.0));
936 assert!((base_scale(8.0, 16.0) - 1.87).abs() < 0.01, "{}", base_scale(8.0, 16.0));
937 }
938
939 #[test]
940 fn a_billed_rate_is_the_median_of_the_charged_days() {
941 let row = |quantity: f64, cost: f64| UsageRow {
942 period_start: String::new(),
943 period_end: String::new(),
944 service: "Containers / Container Memory".into(),
945 unit: "Count".into(),
946 quantity,
947 cost,
948 };
949 let rows = [row(100.0, 0.0), row(100.0, 0.0002), row(100.0, 0.00025), row(100.0, 0.00025)];
950 assert_eq!(billed_rate(&rows.iter().collect::<Vec<_>>()), Some(0.000_002_5));
951 assert_eq!(billed_rate(&[&row(5.0, 0.0)]), None);
952 }
953
954 #[test]
955 fn usage_rows_are_read_by_their_focus_names() {
956 let row = UsageRow::from_value(&json!({
957 "ServiceFamilyName": "Containers",
958 "ServiceName": "Memory",
959 "PricingUnit": "GiB-seconds",
960 "PricingQuantity": "1200.5",
961 "ContractedCost": 0.003,
962 "ChargePeriodStart": "2026-10-01",
963 }))
964 .unwrap();
965 assert_eq!(row.service, "Containers / Memory");
966 assert_eq!(row.quantity, 1200.5);
967 assert_eq!(row.cost, 0.003);
968 assert!(UsageRow::from_value(&json!({ "nothing": 1 })).is_none());
969 }
970}