g1t/services/billing/src/keeper.rs

555 lines22,022 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 far off
20//! the current cost is not adopted, only logged, so one odd day of data
21//! cannot reprice everything.
22
23use g1t_contracts::billing::{EntryKind, MICROS_PER_DOLLAR, Price, PriceBook, PriceChange};
24use g1t_contracts::new_id;
25use g1t_contracts::time::rfc3339;
26use g1t_kit::now_ms;
27use serde::Deserialize;
28use serde_json::{Value, json};
29use worker::{Env, Fetch, Headers, Method, Request, RequestInit, Result};
30
31use crate::{Billing, RunRow, charge_micros};
32
33/// The cron that also checks costs against Cloudflare's bill.
34pub(crate) const DAILY: &str = "17 4 * * *";
35
36/// A run is settled once its logs have had time to land.
37const SETTLE_AFTER_MS: u64 = 5 * 60 * 1000;
38/// A run with no gateway logs after this is left as reported.
39const GIVE_UP_AFTER_MS: u64 = 3 * 60 * 60 * 1000;
40/// A run never finished after this died without reporting.
41const ABANDONED_AFTER_MS: u64 = 3 * 60 * 60 * 1000;
42/// Too little spend to measure a cost from.
43const MIN_MEASURED_USD: f64 = 5.0;
44/// Smaller moves are noise.
45const MIN_CHANGE: f64 = 0.02;
46/// A measurement outside this factor of the current cost is suspect.
47const MAX_FACTOR: f64 = 4.0;
48
49/// Where the keeper reads what g1t pays.
50pub(crate) struct Keeper {
51 /// `CLOUDFLARE_USAGE_TOKEN`: Billing, Account Analytics and AI Gateway,
52 /// read only.
53 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 Keeper {
62 token: env.secret("CLOUDFLARE_USAGE_TOKEN").ok().map(|v| v.to_string()).filter(|v| !v.is_empty()),
63 account: var("CLOUDFLARE_ACCOUNT_ID"),
64 gateway: var("AI_GATEWAY_ID"),
65 }
66 }
67
68 async fn send(&self, method: Method, url: &str, body: Option<Value>) -> Result<Value> {
69 let Some(token) = &self.token else {
70 return Err(worker::Error::RustError("no CLOUDFLARE_USAGE_TOKEN".into()));
71 };
72 let headers = Headers::new();
73 headers.set("authorization", &format!("Bearer {token}"))?;
74 headers.set("content-type", "application/json")?;
75 let mut init = RequestInit::new();
76 init.with_method(method).with_headers(headers);
77 if let Some(body) = body {
78 init.with_body(Some(body.to_string().into()));
79 }
80 let mut response = Fetch::Request(Request::new_with_init(url, &init)?).send().await?;
81 let status = response.status_code();
82 let value: Value = response.json().await.unwrap_or(Value::Null);
83 if status != 200 {
84 return Err(worker::Error::RustError(format!("Cloudflare answered {status}: {value}")));
85 }
86 Ok(value)
87 }
88
89 fn api(&self, path: &str) -> String {
90 format!("https://api.cloudflare.com/client/v4/accounts/{}{path}", self.account)
91 }
92
93 /// What AI Gateway priced a session's requests at, in dollars, and how
94 /// many there were.
95 async fn session_cost(&self, session: &str) -> Result<(f64, u32)> {
96 let mut cost = 0.0;
97 let mut count = 0;
98 for page in 1..=40 {
99 let url = self.api(&format!(
100 "/ai-gateway/gateways/{}/logs?per_page=50&page={page}\
101 &filters[0][key]=metadata.value&filters[0][operator]=eq&filters[0][value][0]={session}",
102 self.gateway
103 ));
104 let body = self.send(Method::Get, &url, None).await?;
105 let logs = body["result"].as_array().cloned().unwrap_or_default();
106 for log in &logs {
107 cost += log["cost"].as_f64().unwrap_or(0.0);
108 count += 1;
109 }
110 if logs.len() < 50 {
111 break;
112 }
113 }
114 Ok((cost, count))
115 }
116
117 /// The account's billable usage this month, as Cloudflare reports it.
118 async fn billable_usage(&self, from: &str, to: &str) -> Result<Vec<UsageRow>> {
119 let body = self
120 .send(Method::Get, &self.api(&format!("/billing/usage/paygo?from={from}&to={to}")), None)
121 .await?;
122 let rows = body["result"].as_array().cloned().unwrap_or_default();
123 Ok(rows.iter().filter_map(UsageRow::from_value).collect())
124 }
125
126 /// Seconds g1t's containers ran since `since` (RFC 3339), across the
127 /// account.
128 async fn container_seconds(&self, since: &str, until: &str) -> Result<f64> {
129 let query = "query ($account: String!, $since: Time!, $until: Time!) {
130 viewer { accounts(filter: { accountTag: $account }) {
131 containersMetricsAdaptiveGroups(limit: 10000, filter: { datetime_geq: $since, datetime_leq: $until }) {
132 sum { containerUptime }
133 }
134 } }
135 }";
136 let body = self
137 .send(
138 Method::Post,
139 "https://api.cloudflare.com/client/v4/graphql",
140 Some(json!({ "query": query, "variables": { "account": self.account, "since": since, "until": until } })),
141 )
142 .await?;
143 let groups = body["data"]["viewer"]["accounts"][0]["containersMetricsAdaptiveGroups"]
144 .as_array()
145 .cloned()
146 .unwrap_or_default();
147 Ok(groups.iter().map(|g| g["sum"]["containerUptime"].as_f64().unwrap_or(0.0)).sum())
148 }
149}
150
151/// One line of Cloudflare's billable usage.
152#[derive(Debug, Clone)]
153pub(crate) struct UsageRow {
154 period_start: String,
155 period_end: String,
156 service: String,
157 unit: String,
158 quantity: f64,
159 cost: f64,
160}
161
162impl UsageRow {
163 /// Read leniently: the API is new, and its field names are FOCUS's.
164 fn from_value(row: &Value) -> Option<Self> {
165 let text = |keys: &[&str]| keys.iter().find_map(|k| row[*k].as_str()).unwrap_or_default().to_owned();
166 let number = |keys: &[&str]| {
167 keys.iter()
168 .find_map(|k| row[*k].as_f64().or_else(|| row[*k].as_str().and_then(|s| s.parse().ok())))
169 .unwrap_or(0.0)
170 };
171 let service = text(&["ServiceName", "service_name", "service"]);
172 if service.is_empty() {
173 return None;
174 }
175 let family = text(&["ServiceFamilyName", "service_family_name"]);
176 Some(UsageRow {
177 period_start: text(&["ChargePeriodStart", "charge_period_start"]),
178 period_end: text(&["ChargePeriodEnd", "charge_period_end"]),
179 service: if family.is_empty() { service } else { format!("{family} / {service}") },
180 unit: text(&["ConsumedUnit", "PricingUnit", "consumed_unit"]),
181 quantity: number(&["PricingQuantity", "ConsumedQuantity", "pricing_quantity"]),
182 cost: number(&["ContractedCost", "BilledCost", "contracted_cost"]),
183 })
184 }
185}
186
187/// What a cost should become from a measurement, or why not.
188pub(crate) fn adopt(current: f64, measured: f64) -> std::result::Result<Option<f64>, String> {
189 if !measured.is_finite() || measured <= 0.0 {
190 return Err("nothing to measure".into());
191 }
192 let ratio = measured / current;
193 if !(1.0 / MAX_FACTOR..=MAX_FACTOR).contains(&ratio) {
194 return Err(format!("measured {measured:.4} against {current:.4}, too far off to adopt"));
195 }
196 Ok(((ratio - 1.0).abs() >= MIN_CHANGE).then_some(measured))
197}
198
199#[derive(Deserialize)]
200struct PriceRow {
201 meter: String,
202 title: String,
203 unit: String,
204 cost_micros: f64,
205 markup_percent: u32,
206 source: String,
207 checked_at: Option<String>,
208 updated_at: String,
209}
210
211#[derive(Deserialize)]
212struct ChangeRow {
213 meter: String,
214 old_cost_micros: f64,
215 new_cost_micros: f64,
216 markup_percent: u32,
217 reason: String,
218 created_at: String,
219}
220
221#[derive(Deserialize)]
222struct Unsettled {
223 id: String,
224 workspace: String,
225 repo: String,
226 number: u32,
227 task: String,
228 model: String,
229 token_hash: String,
230 billed_to: Option<String>,
231 session_id: String,
232 created_at: String,
233 finished_at: Option<String>,
234}
235
236#[derive(Deserialize)]
237struct Charged {
238 cost_micros: Option<i64>,
239 description: String,
240}
241
242fn ms(timestamp: &str) -> u64 {
243 // RFC 3339 in UTC, as g1t writes them.
244 worker::js_sys::Date::parse(timestamp) as u64
245}
246
247impl Billing {
248 pub(crate) async fn prices(&self) -> Result<PriceBook> {
249 let prices = self
250 .db
251 .prepare("SELECT * FROM prices ORDER BY rowid")
252 .all()
253 .await?
254 .results::<PriceRow>()?;
255 let changes = self
256 .db
257 .prepare("SELECT * FROM price_changes ORDER BY created_at DESC LIMIT 20")
258 .all()
259 .await?
260 .results::<ChangeRow>()?;
261 Ok(PriceBook {
262 prices: prices
263 .into_iter()
264 .map(|row| Price {
265 price_micros: Price::price_for(row.cost_micros, row.markup_percent),
266 meter: row.meter,
267 title: row.title,
268 unit: row.unit,
269 cost_micros: row.cost_micros,
270 markup_percent: row.markup_percent,
271 source: row.source,
272 checked_at: row.checked_at,
273 updated_at: row.updated_at,
274 })
275 .collect(),
276 changes: changes
277 .into_iter()
278 .map(|row| PriceChange {
279 meter: row.meter,
280 old_cost_micros: row.old_cost_micros,
281 new_cost_micros: row.new_cost_micros,
282 markup_percent: row.markup_percent,
283 reason: row.reason,
284 created_at: row.created_at,
285 })
286 .collect(),
287 model_margin_percent: self.margin_percent,
288 })
289 }
290
291 /// A meter's cost and price per unit, from the book.
292 pub(crate) async fn price(&self, meter: &str) -> Result<Option<(f64, f64)>> {
293 #[derive(Deserialize)]
294 struct Row {
295 cost_micros: f64,
296 markup_percent: u32,
297 }
298 Ok(self
299 .db
300 .prepare("SELECT cost_micros, markup_percent FROM prices WHERE meter = ?")
301 .bind(&[meter.into()])?
302 .first::<Row>(None)
303 .await?
304 .map(|row| (row.cost_micros, Price::price_for(row.cost_micros, row.markup_percent))))
305 }
306
307 /// Corrects finished runs to what AI Gateway priced them at, and
308 /// charges runs whose sandbox died before reporting.
309 pub(crate) async fn settle_runs(&self, keeper: &Keeper) -> Result<()> {
310 if keeper.token.is_none() || keeper.gateway.is_empty() {
311 return Ok(());
312 }
313 let now = now_ms();
314 let runs = self
315 .db
316 .prepare(
317 "SELECT id, workspace, repo, number, task, model, token_hash, billed_to, session_id, created_at, finished_at
318 FROM runs
319 WHERE session_id IS NOT NULL AND settled_at IS NULL
320 AND ((finished_at IS NOT NULL AND finished_at < ?1) OR created_at < ?2)
321 ORDER BY created_at LIMIT 10",
322 )
323 .bind(&[rfc3339(now - SETTLE_AFTER_MS).into(), rfc3339(now - ABANDONED_AFTER_MS).into()])?
324 .all()
325 .await?
326 .results::<Unsettled>()?;
327 for run in runs {
328 let (cost_usd, requests) = match keeper.session_cost(&run.session_id).await {
329 Ok(found) => found,
330 Err(error) => {
331 worker::console_error!("could not read gateway logs for {}: {error}", run.id);
332 continue;
333 }
334 };
335 let since = ms(run.finished_at.as_deref().unwrap_or(&run.created_at));
336 if requests == 0 && now.saturating_sub(since) < GIVE_UP_AFTER_MS {
337 continue;
338 }
339 self.settle(&run, cost_usd, requests).await?;
340 }
341 Ok(())
342 }
343
344 async fn settle(&self, run: &Unsettled, cost_usd: f64, requests: u32) -> Result<()> {
345 let row = RunRow {
346 workspace: run.workspace.clone(),
347 repo: run.repo.clone(),
348 number: run.number,
349 task: run.task.clone(),
350 model: run.model.clone(),
351 token_hash: run.token_hash.clone(),
352 billed_to: run.billed_to.clone(),
353 };
354 let gateway_micros = charge_micros(cost_usd, 0);
355 let charged = self
356 .db
357 .prepare("SELECT cost_micros, description FROM ledger WHERE reference = ?")
358 .bind(&[run.id.as_str().into()])?
359 .first::<Charged>(None)
360 .await?;
361 let charge_for = |micros: i64| {
362 if self.free {
363 0
364 } else {
365 charge_micros(micros as f64 / MICROS_PER_DOLLAR as f64, self.margin_percent)
366 }
367 };
368 let settled_at = rfc3339(now_ms());
369 // Claim it, so two crons never settle it twice.
370 let claimed = self
371 .db
372 .prepare("UPDATE runs SET settled_at = ?, gateway_cost_micros = ?, finished_at = COALESCE(finished_at, ?) WHERE id = ? AND settled_at IS NULL RETURNING id")
373 .bind(&[
374 settled_at.as_str().into(),
375 (gateway_micros as f64).into(),
376 settled_at.as_str().into(),
377 run.id.as_str().into(),
378 ])?
379 .first::<Value>(None)
380 .await?;
381 if claimed.is_none() || requests == 0 {
382 return Ok(());
383 }
384 let free_note = if self.free { " (free while g1t is being built out)" } else { "" };
385 match charged {
386 // Never reported: charged now, from the gateway's figure.
387 None => {
388 let description = format!(
389 "Work on {}#{}, settled from AI Gateway after the sandbox stopped without reporting{free_note}",
390 run.repo, run.number
391 );
392 self.enter(&run.workspace, EntryKind::Usage, -charge_for(gateway_micros), &description, &run.id, Some(&row), Some(gateway_micros), None, None)
393 .await?;
394 }
395 Some(charged) => {
396 let reported = charged.cost_micros.unwrap_or(0);
397 let delta = gateway_micros - reported;
398 if delta == 0 {
399 return Ok(());
400 }
401 let amount = -(charge_for(gateway_micros) - charge_for(reported));
402 let description = format!(
403 "Correction to “{}”: AI Gateway priced its {requests} model requests at {}, not {}",
404 charged.description,
405 crate::features::dollars(gateway_micros),
406 crate::features::dollars(reported),
407 );
408 self.enter(
409 &run.workspace,
410 EntryKind::Usage,
411 amount,
412 &description,
413 &format!("{}/settled", run.id),
414 Some(&row),
415 Some(delta),
416 None,
417 None,
418 )
419 .await?;
420 }
421 }
422 Ok(())
423 }
424
425 /// Checks each cost against what Cloudflare billed this month, and
426 /// moves the ones that changed.
427 pub(crate) async fn reconcile(&self, keeper: &Keeper) -> Result<()> {
428 if keeper.token.is_none() {
429 return Ok(());
430 }
431 let now = rfc3339(now_ms());
432 let month_start = format!("{}-01", &now[..7]);
433 let today = &now[..10];
434 let rows = keeper.billable_usage(&month_start, today).await?;
435 for row in &rows {
436 self.db
437 .prepare(
438 "INSERT INTO cloudflare_usage (period_start, period_end, service, unit, quantity, cost_usd, fetched_at)
439 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
440 ON CONFLICT (period_start, service, unit) DO UPDATE SET
441 period_end = ?2, quantity = ?5, cost_usd = ?6, fetched_at = ?7",
442 )
443 .bind(&[
444 row.period_start.as_str().into(),
445 row.period_end.as_str().into(),
446 row.service.as_str().into(),
447 row.unit.as_str().into(),
448 row.quantity.into(),
449 row.cost.into(),
450 now.as_str().into(),
451 ])?
452 .run()
453 .await?;
454 }
455
456 let matching = |service: &str, unit: Option<&str>| -> (f64, f64) {
457 rows.iter()
458 .filter(|r| r.service.to_lowercase().contains(service))
459 .filter(|r| unit.is_none_or(|u| r.unit.to_lowercase().contains(u)))
460 .fold((0.0, 0.0), |(q, c), r| (q + r.quantity, c + r.cost))
461 };
462
463 // Containers: what they cost, over the seconds they ran.
464 let (_, container_cost) = matching("container", None);
465 if container_cost >= MIN_MEASURED_USD {
466 let seconds = keeper.container_seconds(&format!("{month_start}T00:00:00Z"), &now).await?;
467 if seconds > 0.0 {
468 let per_second = container_cost * MICROS_PER_DOLLAR as f64 / seconds;
469 for meter in ["sandbox_second", "build_second"] {
470 self.measure(meter, per_second, &format!("Cloudflare billed ${container_cost:.2} for {seconds:.0} container-seconds this month")).await?;
471 }
472 }
473 }
474 // Workers for Platforms: per million requests and CPU milliseconds.
475 for (meter, unit, scale) in [("app_requests", "request", 1e6), ("app_cpu", "ms", 1e6)] {
476 let (quantity, cost) = matching("workers for platforms", Some(unit));
477 if cost >= MIN_MEASURED_USD && quantity > 0.0 {
478 let per = cost * MICROS_PER_DOLLAR as f64 / quantity * scale;
479 self.measure(meter, per, &format!("Cloudflare billed ${cost:.2} for {quantity:.0} {unit}s this month")).await?;
480 }
481 }
482 self.db
483 .prepare("UPDATE prices SET checked_at = ?")
484 .bind(&[now.as_str().into()])?
485 .run()
486 .await?;
487 Ok(())
488 }
489
490 /// Moves a meter's cost to a measurement, if it is sound and different.
491 async fn measure(&self, meter: &str, measured: f64, reason: &str) -> Result<()> {
492 let Some((current, _)) = self.price(meter).await? else {
493 return Ok(());
494 };
495 match adopt(current, measured) {
496 Err(why) => worker::console_log!("{meter}: {why}"),
497 Ok(None) => {}
498 Ok(Some(cost)) => {
499 let now = now_ms();
500 self.db
501 .batch(vec![
502 self.db
503 .prepare("UPDATE prices SET cost_micros = ?, source = 'cloudflare', updated_at = ? WHERE meter = ?")
504 .bind(&[cost.into(), rfc3339(now).into(), meter.into()])?,
505 self.db
506 .prepare(
507 "INSERT INTO price_changes (id, meter, old_cost_micros, new_cost_micros, markup_percent, reason, created_at)
508 SELECT ?, meter, ?, ?, markup_percent, ?, ? FROM prices WHERE meter = ?",
509 )
510 .bind(&[
511 new_id("prc", now).into(),
512 current.into(),
513 cost.into(),
514 reason.into(),
515 rfc3339(now).into(),
516 meter.into(),
517 ])?,
518 ])
519 .await?;
520 }
521 }
522 Ok(())
523 }
524}
525
526#[cfg(test)]
527mod tests {
528 use super::*;
529
530 #[test]
531 fn small_moves_are_noise_and_wild_ones_are_not_believed() {
532 assert_eq!(adopt(21.0, 21.2), Ok(None));
533 assert_eq!(adopt(21.0, 25.0), Ok(Some(25.0)));
534 assert_eq!(adopt(21.0, 15.0), Ok(Some(15.0)));
535 assert!(adopt(21.0, 200.0).is_err());
536 assert!(adopt(21.0, 0.0).is_err());
537 }
538
539 #[test]
540 fn usage_rows_are_read_by_their_focus_names() {
541 let row = UsageRow::from_value(&json!({
542 "ServiceFamilyName": "Containers",
543 "ServiceName": "Memory",
544 "ConsumedUnit": "GiB-seconds",
545 "PricingQuantity": "1200.5",
546 "ContractedCost": 0.003,
547 "ChargePeriodStart": "2026-10-01",
548 }))
549 .unwrap();
550 assert_eq!(row.service, "Containers / Memory");
551 assert_eq!(row.quantity, 1200.5);
552 assert_eq!(row.cost, 0.003);
553 assert!(UsageRow::from_value(&json!({ "nothing": 1 })).is_none());
554 }
555}