pr_01m47d24b0e6n91zwymwxg0vpx/services/billing/src/keeper.rs
| 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 | |
| 23 | use g1t_contracts::billing::{EntryKind, MICROS_PER_DOLLAR, Price, PriceBook, PriceChange}; |
| 24 | use g1t_contracts::new_id; |
| 25 | use g1t_contracts::time::rfc3339; |
| 26 | use g1t_kit::now_ms; |
| 27 | use serde::Deserialize; |
| 28 | use serde_json::{Value, json}; |
| 29 | use worker::{Env, Fetch, Headers, Method, Request, RequestInit, Result}; |
| 30 | |
| 31 | use crate::{Billing, RunRow, charge_micros}; |
| 32 | |
| 33 | /// The cron that also checks costs against Cloudflare's bill. |
| 34 | pub(crate) const DAILY: &str = "17 4 * * *"; |
| 35 | |
| 36 | /// A run is settled once its logs have had time to land. |
| 37 | const SETTLE_AFTER_MS: u64 = 5 * 60 * 1000; |
| 38 | /// A run with no gateway logs after this is left as reported. |
| 39 | const GIVE_UP_AFTER_MS: u64 = 3 * 60 * 60 * 1000; |
| 40 | /// A run never finished after this died without reporting. |
| 41 | const ABANDONED_AFTER_MS: u64 = 3 * 60 * 60 * 1000; |
| 42 | /// Smaller moves are noise. |
| 43 | const MIN_CHANGE: f64 = 0.02; |
| 44 | /// A measurement outside this factor of the current cost is suspect. |
| 45 | const MAX_FACTOR: f64 = 4.0; |
| 46 | |
| 47 | /// Where the keeper reads what g1t pays. |
| 48 | pub(crate) struct Keeper { |
| 49 | /// `CLOUDFLARE_USAGE_TOKEN`: Billing, Account Analytics and AI Gateway, |
| 50 | /// read only. |
| 51 | token: Option<String>, |
| 52 | account: String, |
| 53 | gateway: String, |
| 54 | } |
| 55 | |
| 56 | impl Keeper { |
| 57 | pub(crate) fn from_env(env: &Env) -> Self { |
| 58 | let var = |name: &str| env.var(name).map(|v| v.to_string()).unwrap_or_default(); |
| 59 | Keeper { |
| 60 | token: env.secret("CLOUDFLARE_USAGE_TOKEN").ok().map(|v| v.to_string()).filter(|v| !v.is_empty()), |
| 61 | account: var("CLOUDFLARE_ACCOUNT_ID"), |
| 62 | gateway: var("AI_GATEWAY_ID"), |
| 63 | } |
| 64 | } |
| 65 | |
| 66 | async fn send(&self, method: Method, url: &str, body: Option<Value>) -> Result<Value> { |
| 67 | let Some(token) = &self.token else { |
| 68 | return Err(worker::Error::RustError("no CLOUDFLARE_USAGE_TOKEN".into())); |
| 69 | }; |
| 70 | let headers = Headers::new(); |
| 71 | headers.set("authorization", &format!("Bearer {token}"))?; |
| 72 | headers.set("content-type", "application/json")?; |
| 73 | let mut init = RequestInit::new(); |
| 74 | init.with_method(method).with_headers(headers); |
| 75 | if let Some(body) = body { |
| 76 | init.with_body(Some(body.to_string().into())); |
| 77 | } |
| 78 | let mut response = Fetch::Request(Request::new_with_init(url, &init)?).send().await?; |
| 79 | let status = response.status_code(); |
| 80 | let value: Value = response.json().await.unwrap_or(Value::Null); |
| 81 | if status != 200 { |
| 82 | return Err(worker::Error::RustError(format!("Cloudflare answered {status}: {value}"))); |
| 83 | } |
| 84 | Ok(value) |
| 85 | } |
| 86 | |
| 87 | fn api(&self, path: &str) -> String { |
| 88 | format!("https://api.cloudflare.com/client/v4/accounts/{}{path}", self.account) |
| 89 | } |
| 90 | |
| 91 | /// What AI Gateway priced a session's requests at, in dollars, and how |
| 92 | /// many there were. |
| 93 | async fn session_cost(&self, session: &str) -> Result<(f64, u32)> { |
| 94 | let mut cost = 0.0; |
| 95 | let mut count = 0; |
| 96 | for page in 1..=40 { |
| 97 | // The filter goes as URL-encoded JSON; the bracket form is |
| 98 | // ignored, and would sum every log there is. Session ids are |
| 99 | // [a-z0-9_], which need no escaping inside it. |
| 100 | let filter = format!( |
| 101 | "%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" |
| 102 | ); |
| 103 | let url = self.api(&format!( |
| 104 | "/ai-gateway/gateways/{}/logs?per_page=50&page={page}&filters={filter}", |
| 105 | self.gateway |
| 106 | )); |
| 107 | let body = self.send(Method::Get, &url, None).await?; |
| 108 | let logs = body["result"].as_array().cloned().unwrap_or_default(); |
| 109 | for log in &logs { |
| 110 | cost += log["cost"].as_f64().unwrap_or(0.0); |
| 111 | count += 1; |
| 112 | } |
| 113 | if logs.len() < 50 { |
| 114 | break; |
| 115 | } |
| 116 | } |
| 117 | Ok((cost, count)) |
| 118 | } |
| 119 | |
| 120 | /// The account's billable usage, one row per service per day, as |
| 121 | /// Cloudflare reports it. |
| 122 | async fn billable_usage(&self, from: &str, to: &str) -> Result<Vec<UsageRow>> { |
| 123 | let body = self |
| 124 | .send(Method::Get, &self.api(&format!("/billable-usage?from={from}&to={to}")), None) |
| 125 | .await?; |
| 126 | let rows = body["result"].as_array().cloned().unwrap_or_default(); |
| 127 | Ok(rows.iter().filter_map(UsageRow::from_value).collect()) |
| 128 | } |
| 129 | |
| 130 | /// What g1t's containers used from `since` to `until` (dates), as |
| 131 | /// Cloudflare bills it: memory in byte-seconds, and CPU seconds. |
| 132 | async fn container_usage(&self, since: &str, until: &str) -> Result<ContainerUsage> { |
| 133 | let query = "query ($account: String!, $since: Date!, $until: Date!) { |
| 134 | viewer { accounts(filter: { accountTag: $account }) { |
| 135 | containersUsageAdaptiveGroups(limit: 1000, filter: { date_geq: $since, date_leq: $until }) { |
| 136 | sum { cpuTimeSec allocatedMemory } |
| 137 | } |
| 138 | } } |
| 139 | }"; |
| 140 | let body = self |
| 141 | .send( |
| 142 | Method::Post, |
| 143 | "https://api.cloudflare.com/client/v4/graphql", |
| 144 | Some(json!({ "query": query, "variables": { "account": self.account, "since": since, "until": until } })), |
| 145 | ) |
| 146 | .await?; |
| 147 | let groups = body["data"]["viewer"]["accounts"][0]["containersUsageAdaptiveGroups"] |
| 148 | .as_array() |
| 149 | .cloned() |
| 150 | .unwrap_or_default(); |
| 151 | Ok(groups.iter().fold(ContainerUsage::default(), |total, g| ContainerUsage { |
| 152 | cpu_seconds: total.cpu_seconds + g["sum"]["cpuTimeSec"].as_f64().unwrap_or(0.0), |
| 153 | memory_byte_seconds: total.memory_byte_seconds + g["sum"]["allocatedMemory"].as_f64().unwrap_or(0.0), |
| 154 | })) |
| 155 | } |
| 156 | } |
| 157 | |
| 158 | #[derive(Debug, Default, Clone, Copy)] |
| 159 | pub(crate) struct ContainerUsage { |
| 160 | cpu_seconds: f64, |
| 161 | memory_byte_seconds: f64, |
| 162 | } |
| 163 | |
| 164 | /// g1t's sandboxes: Containers' standard-1, half a vCPU, 4 GiB, 8 GB disk. |
| 165 | const SANDBOX_GIB: f64 = 4.0; |
| 166 | const SANDBOX_DISK_GB: f64 = 8.0; |
| 167 | const GIB: f64 = 1024.0 * 1024.0 * 1024.0; |
| 168 | /// The Durable Object behind each container is billed for as long as the |
| 169 | /// container runs, at 128 MB. |
| 170 | const SANDBOX_DO_GB: f64 = 0.125; |
| 171 | |
| 172 | /// Cloudflare's published Containers rates, in dollars, used for any rate |
| 173 | /// the bill does not show yet (while usage is inside the included amount). |
| 174 | const LIST_MEMORY_GIB_SECOND: f64 = 0.000_002_5; |
| 175 | const LIST_DISK_GB_SECOND: f64 = 0.000_000_07; |
| 176 | const LIST_VCPU_SECOND: f64 = 0.000_02; |
| 177 | /// Durable Objects duration: $12.50 per million GB-seconds. |
| 178 | const LIST_DO_GB_SECOND: f64 = 0.000_012_5; |
| 179 | |
| 180 | /// What one second of a sandbox costs whatever it does, in millionths of a |
| 181 | /// dollar: its memory and disk, and the Durable Object behind it, for the |
| 182 | /// whole second. CPU is billed only while busy, on top. |
| 183 | pub(crate) fn sandbox_base_micros(memory: f64, disk: f64, durable_object: f64) -> f64 { |
| 184 | (SANDBOX_GIB * memory + SANDBOX_DISK_GB * disk + SANDBOX_DO_GB * durable_object) * MICROS_PER_DOLLAR as f64 |
| 185 | } |
| 186 | |
| 187 | /// What one second of a sandbox costs on average, in millionths of a |
| 188 | /// dollar: its base, and the CPU sandboxes actually use per second of |
| 189 | /// running. Runs that report their own CPU are priced on it instead (see |
| 190 | /// `run_cost`). |
| 191 | pub(crate) fn sandbox_second_micros(usage: ContainerUsage, memory: f64, disk: f64, vcpu: f64, durable_object: f64) -> Option<f64> { |
| 192 | let instance_seconds = usage.memory_byte_seconds / (SANDBOX_GIB * GIB); |
| 193 | if instance_seconds < 3600.0 { |
| 194 | return None; |
| 195 | } |
| 196 | let cpu_share = usage.cpu_seconds / instance_seconds; |
| 197 | Some(sandbox_base_micros(memory, disk, durable_object) + cpu_share * vcpu * MICROS_PER_DOLLAR as f64) |
| 198 | } |
| 199 | |
| 200 | /// What a run that reported its own CPU cost g1t: its base for every |
| 201 | /// second, and its vCPU-seconds at the vCPU rate. |
| 202 | pub(crate) fn run_cost(seconds: i64, cpu_seconds: f64, base_per_second: f64, per_vcpu_second: f64) -> f64 { |
| 203 | seconds.max(0) as f64 * base_per_second + cpu_seconds.max(0.0) * per_vcpu_second |
| 204 | } |
| 205 | |
| 206 | /// A unit's marginal rate from the bill: the median, over the days that |
| 207 | /// were charged, of cost over quantity. None while nothing was charged. |
| 208 | pub(crate) fn billed_rate(rows: &[&UsageRow]) -> Option<f64> { |
| 209 | let mut rates: Vec<f64> = rows |
| 210 | .iter() |
| 211 | .filter(|r| r.cost > 0.0 && r.quantity > 0.0) |
| 212 | .map(|r| r.cost / r.quantity) |
| 213 | .collect(); |
| 214 | if rates.is_empty() { |
| 215 | return None; |
| 216 | } |
| 217 | rates.sort_by(f64::total_cmp); |
| 218 | Some(rates[rates.len() / 2]) |
| 219 | } |
| 220 | |
| 221 | /// One line of Cloudflare's billable usage. |
| 222 | #[derive(Debug, Clone)] |
| 223 | pub(crate) struct UsageRow { |
| 224 | period_start: String, |
| 225 | period_end: String, |
| 226 | service: String, |
| 227 | unit: String, |
| 228 | quantity: f64, |
| 229 | cost: f64, |
| 230 | } |
| 231 | |
| 232 | impl UsageRow { |
| 233 | /// Read leniently: the API is new, and its field names are FOCUS's. |
| 234 | fn from_value(row: &Value) -> Option<Self> { |
| 235 | let text = |keys: &[&str]| keys.iter().find_map(|k| row[*k].as_str()).unwrap_or_default().to_owned(); |
| 236 | let number = |keys: &[&str]| { |
| 237 | keys.iter() |
| 238 | .find_map(|k| row[*k].as_f64().or_else(|| row[*k].as_str().and_then(|s| s.parse().ok()))) |
| 239 | .unwrap_or(0.0) |
| 240 | }; |
| 241 | let service = text(&["ServiceName", "service_name", "service"]); |
| 242 | if service.is_empty() { |
| 243 | return None; |
| 244 | } |
| 245 | let family = text(&["ServiceFamilyName", "service_family_name"]); |
| 246 | Some(UsageRow { |
| 247 | period_start: text(&["ChargePeriodStart", "charge_period_start"]), |
| 248 | period_end: text(&["ChargePeriodEnd", "charge_period_end"]), |
| 249 | service: if family.is_empty() { service } else { format!("{family} / {service}") }, |
| 250 | unit: text(&["PricingUnit", "ConsumedUnit", "consumed_unit"]), |
| 251 | quantity: number(&["PricingQuantity", "ConsumedQuantity", "pricing_quantity"]), |
| 252 | // What g1t pays; list price if nothing was contracted. |
| 253 | cost: Some(number(&["ContractedCost", "BilledCost", "contracted_cost"])) |
| 254 | .filter(|cost| *cost > 0.0) |
| 255 | .unwrap_or_else(|| number(&["ListCost", "list_cost"])), |
| 256 | }) |
| 257 | } |
| 258 | } |
| 259 | |
| 260 | /// What a cost should become from a measurement, or why not. |
| 261 | pub(crate) fn adopt(current: f64, measured: f64) -> std::result::Result<Option<f64>, String> { |
| 262 | if !measured.is_finite() || measured <= 0.0 { |
| 263 | return Err("nothing to measure".into()); |
| 264 | } |
| 265 | let ratio = measured / current; |
| 266 | if !(1.0 / MAX_FACTOR..=MAX_FACTOR).contains(&ratio) { |
| 267 | return Err(format!("measured {measured:.4} against {current:.4}, too far off to adopt")); |
| 268 | } |
| 269 | Ok(((ratio - 1.0).abs() >= MIN_CHANGE).then_some(measured)) |
| 270 | } |
| 271 | |
| 272 | #[derive(Deserialize)] |
| 273 | struct PriceRow { |
| 274 | meter: String, |
| 275 | title: String, |
| 276 | unit: String, |
| 277 | cost_micros: f64, |
| 278 | markup_percent: u32, |
| 279 | source: String, |
| 280 | checked_at: Option<String>, |
| 281 | updated_at: String, |
| 282 | } |
| 283 | |
| 284 | #[derive(Deserialize)] |
| 285 | struct ChangeRow { |
| 286 | meter: String, |
| 287 | old_cost_micros: f64, |
| 288 | new_cost_micros: f64, |
| 289 | markup_percent: u32, |
| 290 | old_markup_percent: Option<u32>, |
| 291 | reason: String, |
| 292 | created_at: String, |
| 293 | } |
| 294 | |
| 295 | #[derive(Deserialize)] |
| 296 | struct Unsettled { |
| 297 | id: String, |
| 298 | workspace: String, |
| 299 | repo: String, |
| 300 | number: u32, |
| 301 | task: String, |
| 302 | model: String, |
| 303 | token_hash: String, |
| 304 | billed_to: Option<String>, |
| 305 | session_id: String, |
| 306 | created_at: String, |
| 307 | finished_at: Option<String>, |
| 308 | } |
| 309 | |
| 310 | #[derive(Deserialize)] |
| 311 | struct Charged { |
| 312 | cost_micros: Option<i64>, |
| 313 | description: String, |
| 314 | amount_micros: i64, |
| 315 | } |
| 316 | |
| 317 | /// What a settled run's correction comes to, from the charges at the |
| 318 | /// reported and the gateway's cost and what the workspace was charged at |
| 319 | /// first. A charge up is drawn down like any charge; a charge down is |
| 320 | /// given back only up to what the workspace paid, since what a credit or |
| 321 | /// a pool paid was never the workspace's money. |
| 322 | pub(crate) fn correction(reported_charge: i64, gateway_charge: i64, first_charged: i64) -> i64 { |
| 323 | let delta = gateway_charge - reported_charge; |
| 324 | if delta >= 0 { delta } else { delta.max(-first_charged.max(0)) } |
| 325 | } |
| 326 | |
| 327 | fn ms(timestamp: &str) -> u64 { |
| 328 | // RFC 3339 in UTC, as g1t writes them. |
| 329 | worker::js_sys::Date::parse(timestamp) as u64 |
| 330 | } |
| 331 | |
| 332 | impl Billing { |
| 333 | pub(crate) async fn prices(&self) -> Result<PriceBook> { |
| 334 | let prices = self |
| 335 | .db |
| 336 | .prepare("SELECT * FROM prices ORDER BY rowid") |
| 337 | .all() |
| 338 | .await? |
| 339 | .results::<PriceRow>()?; |
| 340 | let changes = self |
| 341 | .db |
| 342 | .prepare("SELECT * FROM price_changes ORDER BY created_at DESC LIMIT 20") |
| 343 | .all() |
| 344 | .await? |
| 345 | .results::<ChangeRow>()?; |
| 346 | Ok(PriceBook { |
| 347 | prices: prices |
| 348 | .into_iter() |
| 349 | .map(|row| Price { |
| 350 | price_micros: Price::price_for(row.cost_micros, row.markup_percent), |
| 351 | meter: row.meter, |
| 352 | title: row.title, |
| 353 | unit: row.unit, |
| 354 | cost_micros: row.cost_micros, |
| 355 | markup_percent: row.markup_percent, |
| 356 | source: row.source, |
| 357 | checked_at: row.checked_at, |
| 358 | updated_at: row.updated_at, |
| 359 | }) |
| 360 | .collect(), |
| 361 | changes: changes |
| 362 | .into_iter() |
| 363 | .map(|row| PriceChange { |
| 364 | meter: row.meter, |
| 365 | old_cost_micros: row.old_cost_micros, |
| 366 | new_cost_micros: row.new_cost_micros, |
| 367 | markup_percent: row.markup_percent, |
| 368 | old_markup_percent: row.old_markup_percent, |
| 369 | reason: row.reason, |
| 370 | created_at: row.created_at, |
| 371 | }) |
| 372 | .collect(), |
| 373 | model_margin_percent: self.margin_percent, |
| 374 | plans: g1t_contracts::billing::Feature::ALL.iter().map(|feature| self.plan(*feature)).collect(), |
| 375 | free: Some(g1t_contracts::billing::FreeTier { |
| 376 | trial_workspace_micros: if self.trials_on { self.plans.trial_workspace_micros } else { 0 }, |
| 377 | trial_monthly_pool_micros: if self.trials_on { self.plans.trial_monthly_pool_micros } else { 0 }, |
| 378 | oss_pool_micros: self.plans.oss_pool_micros, |
| 379 | oss_repo_micros: self.plans.oss_repo_micros, |
| 380 | free_private_storage_bytes: self.plans.free_storage_bytes, |
| 381 | audit_retention_days: self.plans.audit_days, |
| 382 | min_charge_micros: self.plans.min_charge_micros, |
| 383 | git_operations_included: self.plans.git_included, |
| 384 | git_operations_free_cap: self.plans.git_free_cap, |
| 385 | plan_private_storage_bytes: self.plans.plan_storage_bytes, |
| 386 | paid_start_ceiling_micros: self.plans.paid_start_micros, |
| 387 | overage_forgive_cost_micros: self.plans.forgive_cost_micros, |
| 388 | }), |
| 389 | }) |
| 390 | } |
| 391 | |
| 392 | /// Whether the costs have never been checked against Cloudflare's bill. |
| 393 | pub(crate) async fn never_checked(&self) -> Result<bool> { |
| 394 | Ok(self |
| 395 | .db |
| 396 | .prepare("SELECT meter FROM prices WHERE checked_at IS NOT NULL LIMIT 1") |
| 397 | .first::<Value>(None) |
| 398 | .await? |
| 399 | .is_none()) |
| 400 | } |
| 401 | |
| 402 | /// A meter's cost and price per unit, from the book. |
| 403 | pub(crate) async fn price(&self, meter: &str) -> Result<Option<(f64, f64)>> { |
| 404 | #[derive(Deserialize)] |
| 405 | struct Row { |
| 406 | cost_micros: f64, |
| 407 | markup_percent: u32, |
| 408 | } |
| 409 | Ok(self |
| 410 | .db |
| 411 | .prepare("SELECT cost_micros, markup_percent FROM prices WHERE meter = ?") |
| 412 | .bind(&[meter.into()])? |
| 413 | .first::<Row>(None) |
| 414 | .await? |
| 415 | .map(|row| (row.cost_micros, Price::price_for(row.cost_micros, row.markup_percent)))) |
| 416 | } |
| 417 | |
| 418 | /// Corrects finished runs to what AI Gateway priced them at, and |
| 419 | /// charges runs whose sandbox died before reporting. |
| 420 | pub(crate) async fn settle_runs(&self, keeper: &Keeper) -> Result<()> { |
| 421 | if keeper.token.is_none() || keeper.gateway.is_empty() { |
| 422 | return Ok(()); |
| 423 | } |
| 424 | let now = now_ms(); |
| 425 | let runs = self |
| 426 | .db |
| 427 | .prepare( |
| 428 | "SELECT id, workspace, repo, number, task, model, token_hash, billed_to, session_id, created_at, finished_at |
| 429 | FROM runs |
| 430 | WHERE session_id IS NOT NULL AND settled_at IS NULL |
| 431 | AND ((finished_at IS NOT NULL AND finished_at < ?1) OR created_at < ?2) |
| 432 | ORDER BY created_at LIMIT 10", |
| 433 | ) |
| 434 | .bind(&[rfc3339(now - SETTLE_AFTER_MS).into(), rfc3339(now - ABANDONED_AFTER_MS).into()])? |
| 435 | .all() |
| 436 | .await? |
| 437 | .results::<Unsettled>()?; |
| 438 | for run in runs { |
| 439 | let (cost_usd, requests) = match keeper.session_cost(&run.session_id).await { |
| 440 | Ok(found) => found, |
| 441 | Err(error) => { |
| 442 | worker::console_error!("could not read gateway logs for {}: {error}", run.id); |
| 443 | continue; |
| 444 | } |
| 445 | }; |
| 446 | let since = ms(run.finished_at.as_deref().unwrap_or(&run.created_at)); |
| 447 | if requests == 0 && now.saturating_sub(since) < GIVE_UP_AFTER_MS { |
| 448 | continue; |
| 449 | } |
| 450 | self.settle(&run, cost_usd, requests).await?; |
| 451 | } |
| 452 | Ok(()) |
| 453 | } |
| 454 | |
| 455 | async fn settle(&self, run: &Unsettled, cost_usd: f64, requests: u32) -> Result<()> { |
| 456 | let row = RunRow { |
| 457 | workspace: run.workspace.clone(), |
| 458 | repo: run.repo.clone(), |
| 459 | number: run.number, |
| 460 | task: run.task.clone(), |
| 461 | model: run.model.clone(), |
| 462 | token_hash: run.token_hash.clone(), |
| 463 | billed_to: run.billed_to.clone(), |
| 464 | }; |
| 465 | let gateway_micros = charge_micros(cost_usd, 0); |
| 466 | let charged = self |
| 467 | .db |
| 468 | .prepare("SELECT cost_micros, description, amount_micros FROM ledger WHERE reference = ?") |
| 469 | .bind(&[run.id.as_str().into()])? |
| 470 | .first::<Charged>(None) |
| 471 | .await?; |
| 472 | let terms = self.terms_of(&run.workspace).await?; |
| 473 | let charge_for = |micros: i64| { |
| 474 | if self.free { |
| 475 | 0 |
| 476 | } else { |
| 477 | terms.apply(charge_micros(micros as f64 / MICROS_PER_DOLLAR as f64, self.margin_percent)) |
| 478 | } |
| 479 | }; |
| 480 | let settled_at = rfc3339(now_ms()); |
| 481 | // Claim it, so two crons never settle it twice. |
| 482 | let claimed = self |
| 483 | .db |
| 484 | .prepare("UPDATE runs SET settled_at = ?, gateway_cost_micros = ?, finished_at = COALESCE(finished_at, ?) WHERE id = ? AND settled_at IS NULL RETURNING id") |
| 485 | .bind(&[ |
| 486 | settled_at.as_str().into(), |
| 487 | (gateway_micros as f64).into(), |
| 488 | settled_at.as_str().into(), |
| 489 | run.id.as_str().into(), |
| 490 | ])? |
| 491 | .first::<Value>(None) |
| 492 | .await?; |
| 493 | if claimed.is_none() || requests == 0 { |
| 494 | return Ok(()); |
| 495 | } |
| 496 | let free_note = if self.free { " (free while g1t is being built out)" } else { "" }; |
| 497 | match charged { |
| 498 | // Never reported: charged now, from the gateway's figure. |
| 499 | None => { |
| 500 | let charge = charge_for(gateway_micros); |
| 501 | let eligible = crate::credits::eligible_for(Some(g1t_contracts::billing::ComputeKind::Agent), None); |
| 502 | let drawn = self.draw(&run.workspace, charge, &settled_at[..7], &eligible).await?; |
| 503 | let description = format!( |
| 504 | "Work on {}#{}, settled from AI Gateway after the sandbox stopped without reporting{free_note}{}", |
| 505 | run.repo, |
| 506 | run.number, |
| 507 | drawn.note() |
| 508 | ); |
| 509 | self.enter(&run.workspace, EntryKind::Usage, -(charge - drawn.total()), &description, &run.id, Some(&row), Some(gateway_micros), None, None) |
| 510 | .await?; |
| 511 | self.record_drawn(&run.id, &drawn).await?; |
| 512 | } |
| 513 | Some(charged) => { |
| 514 | let reported = charged.cost_micros.unwrap_or(0); |
| 515 | let delta = gateway_micros - reported; |
| 516 | if delta == 0 { |
| 517 | return Ok(()); |
| 518 | } |
| 519 | let change = correction(charge_for(reported), charge_for(gateway_micros), -charged.amount_micros); |
| 520 | // A charge up is paid for like any other charge. |
| 521 | let drawn = if change > 0 { |
| 522 | let eligible = crate::credits::eligible_for(Some(g1t_contracts::billing::ComputeKind::Agent), None); |
| 523 | self.draw(&run.workspace, change, &settled_at[..7], &eligible).await? |
| 524 | } else { |
| 525 | crate::credits::Drawn::default() |
| 526 | }; |
| 527 | let description = format!( |
| 528 | "Correction to “{}”: AI Gateway priced its {requests} model requests at {}, not {}{}", |
| 529 | charged.description, |
| 530 | crate::features::dollars(gateway_micros), |
| 531 | crate::features::dollars(reported), |
| 532 | drawn.note(), |
| 533 | ); |
| 534 | let reference = format!("{}/settled", run.id); |
| 535 | self.enter( |
| 536 | &run.workspace, |
| 537 | EntryKind::Usage, |
| 538 | -(change - drawn.total()), |
| 539 | &description, |
| 540 | &reference, |
| 541 | Some(&row), |
| 542 | Some(delta), |
| 543 | None, |
| 544 | None, |
| 545 | ) |
| 546 | .await?; |
| 547 | self.record_drawn(&reference, &drawn).await?; |
| 548 | } |
| 549 | } |
| 550 | Ok(()) |
| 551 | } |
| 552 | |
| 553 | /// Checks each cost against what Cloudflare billed this month, and |
| 554 | /// moves the ones that changed. |
| 555 | pub(crate) async fn reconcile(&self, keeper: &Keeper) -> Result<()> { |
| 556 | if keeper.token.is_none() { |
| 557 | return Ok(()); |
| 558 | } |
| 559 | let now = rfc3339(now_ms()); |
| 560 | let today = &now[..10]; |
| 561 | let since = rfc3339(now_ms() - 30 * 24 * 60 * 60 * 1000); |
| 562 | let rows = keeper.billable_usage(&since[..10], today).await?; |
| 563 | for row in &rows { |
| 564 | self.db |
| 565 | .prepare( |
| 566 | "INSERT INTO cloudflare_usage (period_start, period_end, service, unit, quantity, cost_usd, fetched_at) |
| 567 | VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7) |
| 568 | ON CONFLICT (period_start, service, unit) DO UPDATE SET |
| 569 | period_end = ?2, quantity = ?5, cost_usd = ?6, fetched_at = ?7", |
| 570 | ) |
| 571 | .bind(&[ |
| 572 | row.period_start.as_str().into(), |
| 573 | row.period_end.as_str().into(), |
| 574 | row.service.as_str().into(), |
| 575 | row.unit.as_str().into(), |
| 576 | row.quantity.into(), |
| 577 | row.cost.into(), |
| 578 | now.as_str().into(), |
| 579 | ])? |
| 580 | .run() |
| 581 | .await?; |
| 582 | } |
| 583 | |
| 584 | let named = |words: &[&str]| -> Vec<&UsageRow> { |
| 585 | rows.iter() |
| 586 | .filter(|r| { |
| 587 | let service = r.service.to_lowercase(); |
| 588 | words.iter().all(|word| service.contains(word)) |
| 589 | }) |
| 590 | .collect() |
| 591 | }; |
| 592 | |
| 593 | // Containers: each resource at what the bill shows it costs, or |
| 594 | // the published rate while the included amount still covers it, |
| 595 | // over how much CPU g1t's sandboxes really use per second. |
| 596 | let memory = billed_rate(&named(&["container memory"])); |
| 597 | let disk = billed_rate(&named(&["container disk"])); |
| 598 | let vcpu = billed_rate(&named(&["container vcpu"])); |
| 599 | let durable_object = billed_rate(&named(&["durable objects", "duration"])); |
| 600 | let usage = keeper.container_usage(&since[..10], today).await?; |
| 601 | let rates = ( |
| 602 | memory.unwrap_or(LIST_MEMORY_GIB_SECOND), |
| 603 | disk.unwrap_or(LIST_DISK_GB_SECOND), |
| 604 | vcpu.unwrap_or(LIST_VCPU_SECOND), |
| 605 | durable_object.unwrap_or(LIST_DO_GB_SECOND), |
| 606 | ); |
| 607 | // The parts, for runs that report their own CPU. |
| 608 | let parts_reason = "Cloudflare's Containers and Durable Objects rates, as billed or published"; |
| 609 | self.measure("sandbox_base_second", sandbox_base_micros(rates.0, rates.1, rates.3), parts_reason).await?; |
| 610 | self.measure("sandbox_cpu_second", rates.2 * MICROS_PER_DOLLAR as f64, parts_reason).await?; |
| 611 | if let Some(per_second) = sandbox_second_micros(usage, rates.0, rates.1, rates.2, rates.3) { |
| 612 | let instance_seconds = usage.memory_byte_seconds / (SANDBOX_GIB * GIB); |
| 613 | let billed = [("memory", memory), ("disk", disk), ("vCPU", vcpu), ("Durable Object duration", durable_object)] |
| 614 | .iter() |
| 615 | .filter(|(_, rate)| rate.is_some()) |
| 616 | .map(|(name, _)| *name) |
| 617 | .collect::<Vec<_>>(); |
| 618 | let reason = format!( |
| 619 | "Sandboxes used {:.2} vCPU per second over {:.0} hours of Cloudflare Containers in the last 30 days, with the Durable Object behind each; {}", |
| 620 | usage.cpu_seconds / instance_seconds, |
| 621 | instance_seconds / 3600.0, |
| 622 | if billed.is_empty() { |
| 623 | "rates are Cloudflare's published ones".to_owned() |
| 624 | } else { |
| 625 | format!("{} at what Cloudflare billed", billed.join(", ")) |
| 626 | }, |
| 627 | ); |
| 628 | for meter in ["sandbox_second", "build_second"] { |
| 629 | self.measure(meter, per_second, &reason).await?; |
| 630 | } |
| 631 | } |
| 632 | // Apps run as Workers: per million requests and CPU milliseconds, |
| 633 | // once the bill shows them charged. |
| 634 | // Security scans' CPU follows the same Workers CPU rate. |
| 635 | let app_meters: [(&str, &[&str], &str); 3] = [ |
| 636 | ("app_requests", &["workers", "requests"], "requests"), |
| 637 | ("app_cpu", &["workers cpu"], "CPU ms"), |
| 638 | ("scan_cpu", &["workers cpu"], "CPU ms"), |
| 639 | ]; |
| 640 | for (meter, words, unit) in app_meters { |
| 641 | if let Some(rate) = billed_rate(&named(words)) { |
| 642 | let reason = format!("Cloudflare billed Workers {unit} at ${:.2} per million", rate * 1e6); |
| 643 | self.measure(meter, rate * 1e6 * MICROS_PER_DOLLAR as f64, &reason).await?; |
| 644 | } |
| 645 | } |
| 646 | self.db |
| 647 | .prepare("UPDATE prices SET checked_at = ?") |
| 648 | .bind(&[now.as_str().into()])? |
| 649 | .run() |
| 650 | .await?; |
| 651 | Ok(()) |
| 652 | } |
| 653 | |
| 654 | /// Moves a meter's cost to a measurement, if it is sound and different. |
| 655 | async fn measure(&self, meter: &str, measured: f64, reason: &str) -> Result<()> { |
| 656 | let Some((current, _)) = self.price(meter).await? else { |
| 657 | return Ok(()); |
| 658 | }; |
| 659 | match adopt(current, measured) { |
| 660 | Err(why) => worker::console_log!("{meter}: {why}"), |
| 661 | Ok(None) => {} |
| 662 | Ok(Some(cost)) => { |
| 663 | let now = now_ms(); |
| 664 | self.db |
| 665 | .batch(vec![ |
| 666 | self.db |
| 667 | .prepare("UPDATE prices SET cost_micros = ?, source = 'cloudflare', updated_at = ? WHERE meter = ?") |
| 668 | .bind(&[cost.into(), rfc3339(now).into(), meter.into()])?, |
| 669 | self.db |
| 670 | .prepare( |
| 671 | "INSERT INTO price_changes (id, meter, old_cost_micros, new_cost_micros, markup_percent, reason, created_at) |
| 672 | SELECT ?, meter, ?, ?, markup_percent, ?, ? FROM prices WHERE meter = ?", |
| 673 | ) |
| 674 | .bind(&[ |
| 675 | new_id("prc", now).into(), |
| 676 | current.into(), |
| 677 | cost.into(), |
| 678 | reason.into(), |
| 679 | rfc3339(now).into(), |
| 680 | meter.into(), |
| 681 | ])?, |
| 682 | ]) |
| 683 | .await?; |
| 684 | } |
| 685 | } |
| 686 | Ok(()) |
| 687 | } |
| 688 | } |
| 689 | |
| 690 | #[cfg(test)] |
| 691 | mod tests { |
| 692 | use super::*; |
| 693 | |
| 694 | #[test] |
| 695 | fn a_correction_never_gives_back_what_the_workspace_did_not_pay() { |
| 696 | // Up by 2 cents: charged in full (then drawn down like any charge). |
| 697 | assert_eq!(correction(100_000, 120_000, 100_000), 20_000); |
| 698 | // Down by 2 cents, all of it paid by the workspace: given back. |
| 699 | assert_eq!(correction(120_000, 100_000, 120_000), -20_000); |
| 700 | // Down, but the open-source pool paid all but a cent: a cent back. |
| 701 | assert_eq!(correction(120_000, 100_000, 10_000), -10_000); |
| 702 | // Paid entirely by a credit or pool: nothing back. |
| 703 | assert_eq!(correction(120_000, 100_000, 0), 0); |
| 704 | } |
| 705 | |
| 706 | #[test] |
| 707 | fn small_moves_are_noise_and_wild_ones_are_not_believed() { |
| 708 | assert_eq!(adopt(21.0, 21.2), Ok(None)); |
| 709 | assert_eq!(adopt(21.0, 25.0), Ok(Some(25.0))); |
| 710 | assert_eq!(adopt(21.0, 15.0), Ok(Some(15.0))); |
| 711 | assert!(adopt(21.0, 200.0).is_err()); |
| 712 | assert!(adopt(21.0, 0.0).is_err()); |
| 713 | } |
| 714 | |
| 715 | #[test] |
| 716 | fn a_sandbox_second_is_its_memory_and_disk_and_the_cpu_it_uses() { |
| 717 | // An hour of sandboxes that kept a fifth of a vCPU busy. |
| 718 | let usage = ContainerUsage { cpu_seconds: 720.0, memory_byte_seconds: 3600.0 * 4.0 * GIB }; |
| 719 | let micros = |
| 720 | sandbox_second_micros(usage, LIST_MEMORY_GIB_SECOND, LIST_DISK_GB_SECOND, LIST_VCPU_SECOND, LIST_DO_GB_SECOND).unwrap(); |
| 721 | // 4 x 2.5 + 8 x 0.07 + 0.125 x 12.5 + 0.2 x 20 = 16.1225 |
| 722 | assert!((micros - 16.1225).abs() < 1e-9, "{micros}"); |
| 723 | // The Durable Object adds about 11% to the second it left out. |
| 724 | assert!((sandbox_base_micros(LIST_MEMORY_GIB_SECOND, LIST_DISK_GB_SECOND, LIST_DO_GB_SECOND) - 12.1225).abs() < 1e-9); |
| 725 | // Too little use to say anything. |
| 726 | assert!(sandbox_second_micros(ContainerUsage { cpu_seconds: 1.0, memory_byte_seconds: GIB }, 1.0, 1.0, 1.0, 1.0).is_none()); |
| 727 | } |
| 728 | |
| 729 | #[test] |
| 730 | fn a_run_that_reports_its_cpu_is_priced_on_it() { |
| 731 | let base = sandbox_base_micros(LIST_MEMORY_GIB_SECOND, LIST_DISK_GB_SECOND, LIST_DO_GB_SECOND); |
| 732 | let vcpu = LIST_VCPU_SECOND * MICROS_PER_DOLLAR as f64; |
| 733 | // A 10-minute cargo build that kept its half vCPU busy throughout. |
| 734 | let heavy = run_cost(600, 300.0, base, vcpu); |
| 735 | assert!((heavy - (600.0 * 12.1225 + 300.0 * 20.0)).abs() < 1e-6); |
| 736 | // The same ten minutes, mostly idle, costs less. |
| 737 | let light = run_cost(600, 30.0, base, vcpu); |
| 738 | assert!(light < heavy); |
| 739 | // The average would have under-priced the heavy one. |
| 740 | let average = 600.0 * (base + 0.195 * vcpu); |
| 741 | assert!(average < heavy && average > light); |
| 742 | assert_eq!(run_cost(0, -1.0, base, vcpu), 0.0); |
| 743 | } |
| 744 | |
| 745 | #[test] |
| 746 | fn a_billed_rate_is_the_median_of_the_charged_days() { |
| 747 | let row = |quantity: f64, cost: f64| UsageRow { |
| 748 | period_start: String::new(), |
| 749 | period_end: String::new(), |
| 750 | service: "Containers / Container Memory".into(), |
| 751 | unit: "Count".into(), |
| 752 | quantity, |
| 753 | cost, |
| 754 | }; |
| 755 | let rows = [row(100.0, 0.0), row(100.0, 0.0002), row(100.0, 0.00025), row(100.0, 0.00025)]; |
| 756 | assert_eq!(billed_rate(&rows.iter().collect::<Vec<_>>()), Some(0.000_002_5)); |
| 757 | assert_eq!(billed_rate(&[&row(5.0, 0.0)]), None); |
| 758 | } |
| 759 | |
| 760 | #[test] |
| 761 | fn usage_rows_are_read_by_their_focus_names() { |
| 762 | let row = UsageRow::from_value(&json!({ |
| 763 | "ServiceFamilyName": "Containers", |
| 764 | "ServiceName": "Memory", |
| 765 | "PricingUnit": "GiB-seconds", |
| 766 | "PricingQuantity": "1200.5", |
| 767 | "ContractedCost": 0.003, |
| 768 | "ChargePeriodStart": "2026-10-01", |
| 769 | })) |
| 770 | .unwrap(); |
| 771 | assert_eq!(row.service, "Containers / Memory"); |
| 772 | assert_eq!(row.quantity, 1200.5); |
| 773 | assert_eq!(row.cost, 0.003); |
| 774 | assert!(UsageRow::from_value(&json!({ "nothing": 1 })).is_none()); |
| 775 | } |
| 776 | } |