g1t/services/billing/src/keeper.rs
Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.
| Prices keep themselves current with what g1t pays | 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 { | |
| The keeper reads Cloudflare as it really answers | 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 | ); | |
| Prices keep themselves current with what g1t pays | 103 | let url = self.api(&format!( |
| The keeper reads Cloudflare as it really answers | 104 | "/ai-gateway/gateways/{}/logs?per_page=50&page={page}&filters={filter}", |
| Prices keep themselves current with what g1t pays | 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 | ||
| The keeper reads Cloudflare as it really answers | 120 | /// The account's billable usage, one row per service per day, as |
| 121 | /// Cloudflare reports it. | |
| Prices keep themselves current with what g1t pays | 122 | async fn billable_usage(&self, from: &str, to: &str) -> Result<Vec<UsageRow>> { |
| 123 | let body = self | |
| The keeper reads Cloudflare as it really answers | 124 | .send(Method::Get, &self.api(&format!("/billable-usage?from={from}&to={to}")), None) |
| Prices keep themselves current with what g1t pays | 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 | ||
| The keeper reads Cloudflare as it really answers | 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!) { | |
| Prices keep themselves current with what g1t pays | 134 | viewer { accounts(filter: { accountTag: $account }) { |
| The keeper reads Cloudflare as it really answers | 135 | containersUsageAdaptiveGroups(limit: 1000, filter: { date_geq: $since, date_leq: $until }) { |
| 136 | sum { cpuTimeSec allocatedMemory } | |
| Prices keep themselves current with what g1t pays | 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?; | |
| The keeper reads Cloudflare as it really answers | 147 | let groups = body["data"]["viewer"]["accounts"][0]["containersUsageAdaptiveGroups"] |
| Prices keep themselves current with what g1t pays | 148 | .as_array() |
| 149 | .cloned() | |
| 150 | .unwrap_or_default(); | |
| The keeper reads Cloudflare as it really answers | 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 | })) | |
| Prices keep themselves current with what g1t pays | 155 | } |
| 156 | } | |
| 157 | ||
| The keeper reads Cloudflare as it really answers | 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 | ||
| 169 | /// Cloudflare's published Containers rates, in dollars, used for any rate | |
| 170 | /// the bill does not show yet (while usage is inside the included amount). | |
| 171 | const LIST_MEMORY_GIB_SECOND: f64 = 0.000_002_5; | |
| 172 | const LIST_DISK_GB_SECOND: f64 = 0.000_000_07; | |
| 173 | const LIST_VCPU_SECOND: f64 = 0.000_02; | |
| 174 | ||
| 175 | /// What one second of a sandbox costs, in millionths of a dollar: its | |
| 176 | /// memory and disk for the whole second, and the CPU sandboxes actually | |
| 177 | /// use per second of running, which is billed only while busy. | |
| 178 | pub(crate) fn sandbox_second_micros(usage: ContainerUsage, memory: f64, disk: f64, vcpu: f64) -> Option<f64> { | |
| 179 | let instance_seconds = usage.memory_byte_seconds / (SANDBOX_GIB * GIB); | |
| 180 | if instance_seconds < 3600.0 { | |
| 181 | return None; | |
| 182 | } | |
| 183 | let cpu_share = usage.cpu_seconds / instance_seconds; | |
| 184 | Some((SANDBOX_GIB * memory + SANDBOX_DISK_GB * disk + cpu_share * vcpu) * MICROS_PER_DOLLAR as f64) | |
| 185 | } | |
| 186 | ||
| 187 | /// A unit's marginal rate from the bill: the median, over the days that | |
| 188 | /// were charged, of cost over quantity. None while nothing was charged. | |
| 189 | pub(crate) fn billed_rate(rows: &[&UsageRow]) -> Option<f64> { | |
| 190 | let mut rates: Vec<f64> = rows | |
| 191 | .iter() | |
| 192 | .filter(|r| r.cost > 0.0 && r.quantity > 0.0) | |
| 193 | .map(|r| r.cost / r.quantity) | |
| 194 | .collect(); | |
| 195 | if rates.is_empty() { | |
| 196 | return None; | |
| 197 | } | |
| 198 | rates.sort_by(f64::total_cmp); | |
| 199 | Some(rates[rates.len() / 2]) | |
| 200 | } | |
| 201 | ||
| Prices keep themselves current with what g1t pays | 202 | /// One line of Cloudflare's billable usage. |
| 203 | #[derive(Debug, Clone)] | |
| 204 | pub(crate) struct UsageRow { | |
| 205 | period_start: String, | |
| 206 | period_end: String, | |
| 207 | service: String, | |
| 208 | unit: String, | |
| 209 | quantity: f64, | |
| 210 | cost: f64, | |
| 211 | } | |
| 212 | ||
| 213 | impl UsageRow { | |
| 214 | /// Read leniently: the API is new, and its field names are FOCUS's. | |
| 215 | fn from_value(row: &Value) -> Option<Self> { | |
| 216 | let text = |keys: &[&str]| keys.iter().find_map(|k| row[*k].as_str()).unwrap_or_default().to_owned(); | |
| 217 | let number = |keys: &[&str]| { | |
| 218 | keys.iter() | |
| 219 | .find_map(|k| row[*k].as_f64().or_else(|| row[*k].as_str().and_then(|s| s.parse().ok()))) | |
| 220 | .unwrap_or(0.0) | |
| 221 | }; | |
| 222 | let service = text(&["ServiceName", "service_name", "service"]); | |
| 223 | if service.is_empty() { | |
| 224 | return None; | |
| 225 | } | |
| 226 | let family = text(&["ServiceFamilyName", "service_family_name"]); | |
| 227 | Some(UsageRow { | |
| 228 | period_start: text(&["ChargePeriodStart", "charge_period_start"]), | |
| 229 | period_end: text(&["ChargePeriodEnd", "charge_period_end"]), | |
| 230 | service: if family.is_empty() { service } else { format!("{family} / {service}") }, | |
| The keeper reads Cloudflare as it really answers | 231 | unit: text(&["PricingUnit", "ConsumedUnit", "consumed_unit"]), |
| Prices keep themselves current with what g1t pays | 232 | quantity: number(&["PricingQuantity", "ConsumedQuantity", "pricing_quantity"]), |
| The keeper reads Cloudflare as it really answers | 233 | // What g1t pays; list price if nothing was contracted. |
| 234 | cost: Some(number(&["ContractedCost", "BilledCost", "contracted_cost"])) | |
| 235 | .filter(|cost| *cost > 0.0) | |
| 236 | .unwrap_or_else(|| number(&["ListCost", "list_cost"])), | |
| Prices keep themselves current with what g1t pays | 237 | }) |
| 238 | } | |
| 239 | } | |
| 240 | ||
| 241 | /// What a cost should become from a measurement, or why not. | |
| 242 | pub(crate) fn adopt(current: f64, measured: f64) -> std::result::Result<Option<f64>, String> { | |
| 243 | if !measured.is_finite() || measured <= 0.0 { | |
| 244 | return Err("nothing to measure".into()); | |
| 245 | } | |
| 246 | let ratio = measured / current; | |
| 247 | if !(1.0 / MAX_FACTOR..=MAX_FACTOR).contains(&ratio) { | |
| 248 | return Err(format!("measured {measured:.4} against {current:.4}, too far off to adopt")); | |
| 249 | } | |
| 250 | Ok(((ratio - 1.0).abs() >= MIN_CHANGE).then_some(measured)) | |
| 251 | } | |
| 252 | ||
| 253 | #[derive(Deserialize)] | |
| 254 | struct PriceRow { | |
| 255 | meter: String, | |
| 256 | title: String, | |
| 257 | unit: String, | |
| 258 | cost_micros: f64, | |
| 259 | markup_percent: u32, | |
| 260 | source: String, | |
| 261 | checked_at: Option<String>, | |
| 262 | updated_at: String, | |
| 263 | } | |
| 264 | ||
| 265 | #[derive(Deserialize)] | |
| 266 | struct ChangeRow { | |
| 267 | meter: String, | |
| 268 | old_cost_micros: f64, | |
| 269 | new_cost_micros: f64, | |
| 270 | markup_percent: u32, | |
| 271 | reason: String, | |
| 272 | created_at: String, | |
| 273 | } | |
| 274 | ||
| 275 | #[derive(Deserialize)] | |
| 276 | struct Unsettled { | |
| 277 | id: String, | |
| 278 | workspace: String, | |
| 279 | repo: String, | |
| 280 | number: u32, | |
| 281 | task: String, | |
| 282 | model: String, | |
| 283 | token_hash: String, | |
| 284 | billed_to: Option<String>, | |
| 285 | session_id: String, | |
| 286 | created_at: String, | |
| 287 | finished_at: Option<String>, | |
| 288 | } | |
| 289 | ||
| 290 | #[derive(Deserialize)] | |
| 291 | struct Charged { | |
| 292 | cost_micros: Option<i64>, | |
| 293 | description: String, | |
| 294 | } | |
| 295 | ||
| 296 | fn ms(timestamp: &str) -> u64 { | |
| 297 | // RFC 3339 in UTC, as g1t writes them. | |
| 298 | worker::js_sys::Date::parse(timestamp) as u64 | |
| 299 | } | |
| 300 | ||
| 301 | impl Billing { | |
| 302 | pub(crate) async fn prices(&self) -> Result<PriceBook> { | |
| 303 | let prices = self | |
| 304 | .db | |
| 305 | .prepare("SELECT * FROM prices ORDER BY rowid") | |
| 306 | .all() | |
| 307 | .await? | |
| 308 | .results::<PriceRow>()?; | |
| 309 | let changes = self | |
| 310 | .db | |
| 311 | .prepare("SELECT * FROM price_changes ORDER BY created_at DESC LIMIT 20") | |
| 312 | .all() | |
| 313 | .await? | |
| 314 | .results::<ChangeRow>()?; | |
| 315 | Ok(PriceBook { | |
| 316 | prices: prices | |
| 317 | .into_iter() | |
| 318 | .map(|row| Price { | |
| 319 | price_micros: Price::price_for(row.cost_micros, row.markup_percent), | |
| 320 | meter: row.meter, | |
| 321 | title: row.title, | |
| 322 | unit: row.unit, | |
| 323 | cost_micros: row.cost_micros, | |
| 324 | markup_percent: row.markup_percent, | |
| 325 | source: row.source, | |
| 326 | checked_at: row.checked_at, | |
| 327 | updated_at: row.updated_at, | |
| 328 | }) | |
| 329 | .collect(), | |
| 330 | changes: changes | |
| 331 | .into_iter() | |
| 332 | .map(|row| PriceChange { | |
| 333 | meter: row.meter, | |
| 334 | old_cost_micros: row.old_cost_micros, | |
| 335 | new_cost_micros: row.new_cost_micros, | |
| 336 | markup_percent: row.markup_percent, | |
| 337 | reason: row.reason, | |
| 338 | created_at: row.created_at, | |
| 339 | }) | |
| 340 | .collect(), | |
| 341 | model_margin_percent: self.margin_percent, | |
| 342 | }) | |
| 343 | } | |
| 344 | ||
| The keeper reads Cloudflare as it really answers | 345 | /// Whether the costs have never been checked against Cloudflare's bill. |
| 346 | pub(crate) async fn never_checked(&self) -> Result<bool> { | |
| 347 | Ok(self | |
| 348 | .db | |
| 349 | .prepare("SELECT meter FROM prices WHERE checked_at IS NOT NULL LIMIT 1") | |
| 350 | .first::<Value>(None) | |
| 351 | .await? | |
| 352 | .is_none()) | |
| 353 | } | |
| 354 | ||
| Prices keep themselves current with what g1t pays | 355 | /// A meter's cost and price per unit, from the book. |
| 356 | pub(crate) async fn price(&self, meter: &str) -> Result<Option<(f64, f64)>> { | |
| 357 | #[derive(Deserialize)] | |
| 358 | struct Row { | |
| 359 | cost_micros: f64, | |
| 360 | markup_percent: u32, | |
| 361 | } | |
| 362 | Ok(self | |
| 363 | .db | |
| 364 | .prepare("SELECT cost_micros, markup_percent FROM prices WHERE meter = ?") | |
| 365 | .bind(&[meter.into()])? | |
| 366 | .first::<Row>(None) | |
| 367 | .await? | |
| 368 | .map(|row| (row.cost_micros, Price::price_for(row.cost_micros, row.markup_percent)))) | |
| 369 | } | |
| 370 | ||
| 371 | /// Corrects finished runs to what AI Gateway priced them at, and | |
| 372 | /// charges runs whose sandbox died before reporting. | |
| 373 | pub(crate) async fn settle_runs(&self, keeper: &Keeper) -> Result<()> { | |
| 374 | if keeper.token.is_none() || keeper.gateway.is_empty() { | |
| 375 | return Ok(()); | |
| 376 | } | |
| 377 | let now = now_ms(); | |
| 378 | let runs = self | |
| 379 | .db | |
| 380 | .prepare( | |
| 381 | "SELECT id, workspace, repo, number, task, model, token_hash, billed_to, session_id, created_at, finished_at | |
| 382 | FROM runs | |
| 383 | WHERE session_id IS NOT NULL AND settled_at IS NULL | |
| 384 | AND ((finished_at IS NOT NULL AND finished_at < ?1) OR created_at < ?2) | |
| 385 | ORDER BY created_at LIMIT 10", | |
| 386 | ) | |
| 387 | .bind(&[rfc3339(now - SETTLE_AFTER_MS).into(), rfc3339(now - ABANDONED_AFTER_MS).into()])? | |
| 388 | .all() | |
| 389 | .await? | |
| 390 | .results::<Unsettled>()?; | |
| 391 | for run in runs { | |
| 392 | let (cost_usd, requests) = match keeper.session_cost(&run.session_id).await { | |
| 393 | Ok(found) => found, | |
| 394 | Err(error) => { | |
| 395 | worker::console_error!("could not read gateway logs for {}: {error}", run.id); | |
| 396 | continue; | |
| 397 | } | |
| 398 | }; | |
| 399 | let since = ms(run.finished_at.as_deref().unwrap_or(&run.created_at)); | |
| 400 | if requests == 0 && now.saturating_sub(since) < GIVE_UP_AFTER_MS { | |
| 401 | continue; | |
| 402 | } | |
| 403 | self.settle(&run, cost_usd, requests).await?; | |
| 404 | } | |
| 405 | Ok(()) | |
| 406 | } | |
| 407 | ||
| 408 | async fn settle(&self, run: &Unsettled, cost_usd: f64, requests: u32) -> Result<()> { | |
| 409 | let row = RunRow { | |
| 410 | workspace: run.workspace.clone(), | |
| 411 | repo: run.repo.clone(), | |
| 412 | number: run.number, | |
| 413 | task: run.task.clone(), | |
| 414 | model: run.model.clone(), | |
| 415 | token_hash: run.token_hash.clone(), | |
| 416 | billed_to: run.billed_to.clone(), | |
| 417 | }; | |
| 418 | let gateway_micros = charge_micros(cost_usd, 0); | |
| 419 | let charged = self | |
| 420 | .db | |
| 421 | .prepare("SELECT cost_micros, description FROM ledger WHERE reference = ?") | |
| 422 | .bind(&[run.id.as_str().into()])? | |
| 423 | .first::<Charged>(None) | |
| 424 | .await?; | |
| Billing accounts, terms and enterprises; g1t is no longer free | 425 | let terms = self.terms_of(&run.workspace).await?; |
| Prices keep themselves current with what g1t pays | 426 | let charge_for = |micros: i64| { |
| 427 | if self.free { | |
| 428 | 0 | |
| 429 | } else { | |
| Billing accounts, terms and enterprises; g1t is no longer free | 430 | terms.apply(charge_micros(micros as f64 / MICROS_PER_DOLLAR as f64, self.margin_percent)) |
| Prices keep themselves current with what g1t pays | 431 | } |
| 432 | }; | |
| 433 | let settled_at = rfc3339(now_ms()); | |
| 434 | // Claim it, so two crons never settle it twice. | |
| 435 | let claimed = self | |
| 436 | .db | |
| 437 | .prepare("UPDATE runs SET settled_at = ?, gateway_cost_micros = ?, finished_at = COALESCE(finished_at, ?) WHERE id = ? AND settled_at IS NULL RETURNING id") | |
| 438 | .bind(&[ | |
| 439 | settled_at.as_str().into(), | |
| 440 | (gateway_micros as f64).into(), | |
| 441 | settled_at.as_str().into(), | |
| 442 | run.id.as_str().into(), | |
| 443 | ])? | |
| 444 | .first::<Value>(None) | |
| 445 | .await?; | |
| 446 | if claimed.is_none() || requests == 0 { | |
| 447 | return Ok(()); | |
| 448 | } | |
| 449 | let free_note = if self.free { " (free while g1t is being built out)" } else { "" }; | |
| 450 | match charged { | |
| 451 | // Never reported: charged now, from the gateway's figure. | |
| 452 | None => { | |
| 453 | let description = format!( | |
| 454 | "Work on {}#{}, settled from AI Gateway after the sandbox stopped without reporting{free_note}", | |
| 455 | run.repo, run.number | |
| 456 | ); | |
| 457 | self.enter(&run.workspace, EntryKind::Usage, -charge_for(gateway_micros), &description, &run.id, Some(&row), Some(gateway_micros), None, None) | |
| 458 | .await?; | |
| 459 | } | |
| 460 | Some(charged) => { | |
| 461 | let reported = charged.cost_micros.unwrap_or(0); | |
| 462 | let delta = gateway_micros - reported; | |
| 463 | if delta == 0 { | |
| 464 | return Ok(()); | |
| 465 | } | |
| 466 | let amount = -(charge_for(gateway_micros) - charge_for(reported)); | |
| 467 | let description = format!( | |
| 468 | "Correction to “{}”: AI Gateway priced its {requests} model requests at {}, not {}", | |
| 469 | charged.description, | |
| 470 | crate::features::dollars(gateway_micros), | |
| 471 | crate::features::dollars(reported), | |
| 472 | ); | |
| 473 | self.enter( | |
| 474 | &run.workspace, | |
| 475 | EntryKind::Usage, | |
| 476 | amount, | |
| 477 | &description, | |
| 478 | &format!("{}/settled", run.id), | |
| 479 | Some(&row), | |
| 480 | Some(delta), | |
| 481 | None, | |
| 482 | None, | |
| 483 | ) | |
| 484 | .await?; | |
| 485 | } | |
| 486 | } | |
| 487 | Ok(()) | |
| 488 | } | |
| 489 | ||
| 490 | /// Checks each cost against what Cloudflare billed this month, and | |
| 491 | /// moves the ones that changed. | |
| 492 | pub(crate) async fn reconcile(&self, keeper: &Keeper) -> Result<()> { | |
| 493 | if keeper.token.is_none() { | |
| 494 | return Ok(()); | |
| 495 | } | |
| 496 | let now = rfc3339(now_ms()); | |
| 497 | let today = &now[..10]; | |
| The keeper reads Cloudflare as it really answers | 498 | let since = rfc3339(now_ms() - 30 * 24 * 60 * 60 * 1000); |
| 499 | let rows = keeper.billable_usage(&since[..10], today).await?; | |
| Prices keep themselves current with what g1t pays | 500 | for row in &rows { |
| 501 | self.db | |
| 502 | .prepare( | |
| 503 | "INSERT INTO cloudflare_usage (period_start, period_end, service, unit, quantity, cost_usd, fetched_at) | |
| 504 | VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7) | |
| 505 | ON CONFLICT (period_start, service, unit) DO UPDATE SET | |
| 506 | period_end = ?2, quantity = ?5, cost_usd = ?6, fetched_at = ?7", | |
| 507 | ) | |
| 508 | .bind(&[ | |
| 509 | row.period_start.as_str().into(), | |
| 510 | row.period_end.as_str().into(), | |
| 511 | row.service.as_str().into(), | |
| 512 | row.unit.as_str().into(), | |
| 513 | row.quantity.into(), | |
| 514 | row.cost.into(), | |
| 515 | now.as_str().into(), | |
| 516 | ])? | |
| 517 | .run() | |
| 518 | .await?; | |
| 519 | } | |
| 520 | ||
| The keeper reads Cloudflare as it really answers | 521 | let named = |words: &[&str]| -> Vec<&UsageRow> { |
| Prices keep themselves current with what g1t pays | 522 | rows.iter() |
| The keeper reads Cloudflare as it really answers | 523 | .filter(|r| { |
| 524 | let service = r.service.to_lowercase(); | |
| 525 | words.iter().all(|word| service.contains(word)) | |
| 526 | }) | |
| 527 | .collect() | |
| Prices keep themselves current with what g1t pays | 528 | }; |
| 529 | ||
| The keeper reads Cloudflare as it really answers | 530 | // Containers: each resource at what the bill shows it costs, or |
| 531 | // the published rate while the included amount still covers it, | |
| 532 | // over how much CPU g1t's sandboxes really use per second. | |
| 533 | let memory = billed_rate(&named(&["container memory"])); | |
| 534 | let disk = billed_rate(&named(&["container disk"])); | |
| 535 | let vcpu = billed_rate(&named(&["container vcpu"])); | |
| 536 | let usage = keeper.container_usage(&since[..10], today).await?; | |
| 537 | if let Some(per_second) = sandbox_second_micros( | |
| 538 | usage, | |
| 539 | memory.unwrap_or(LIST_MEMORY_GIB_SECOND), | |
| 540 | disk.unwrap_or(LIST_DISK_GB_SECOND), | |
| 541 | vcpu.unwrap_or(LIST_VCPU_SECOND), | |
| 542 | ) { | |
| 543 | let instance_seconds = usage.memory_byte_seconds / (SANDBOX_GIB * GIB); | |
| 544 | let billed = [("memory", memory), ("disk", disk), ("vCPU", vcpu)] | |
| 545 | .iter() | |
| 546 | .filter(|(_, rate)| rate.is_some()) | |
| 547 | .map(|(name, _)| *name) | |
| 548 | .collect::<Vec<_>>(); | |
| 549 | let reason = format!( | |
| 550 | "Sandboxes used {:.2} vCPU per second over {:.0} hours of Cloudflare Containers in the last 30 days; {}", | |
| 551 | usage.cpu_seconds / instance_seconds, | |
| 552 | instance_seconds / 3600.0, | |
| 553 | if billed.is_empty() { | |
| 554 | "rates are Cloudflare's published ones".to_owned() | |
| 555 | } else { | |
| 556 | format!("{} at what Cloudflare billed", billed.join(", ")) | |
| 557 | }, | |
| 558 | ); | |
| 559 | for meter in ["sandbox_second", "build_second"] { | |
| 560 | self.measure(meter, per_second, &reason).await?; | |
| Prices keep themselves current with what g1t pays | 561 | } |
| 562 | } | |
| The keeper reads Cloudflare as it really answers | 563 | // Apps run as Workers: per million requests and CPU milliseconds, |
| 564 | // once the bill shows them charged. | |
| 565 | let app_meters: [(&str, &[&str], &str); 2] = [ | |
| 566 | ("app_requests", &["workers", "requests"], "requests"), | |
| 567 | ("app_cpu", &["workers cpu"], "CPU ms"), | |
| 568 | ]; | |
| 569 | for (meter, words, unit) in app_meters { | |
| 570 | if let Some(rate) = billed_rate(&named(words)) { | |
| 571 | let reason = format!("Cloudflare billed Workers {unit} at ${:.2} per million", rate * 1e6); | |
| 572 | self.measure(meter, rate * 1e6 * MICROS_PER_DOLLAR as f64, &reason).await?; | |
| Prices keep themselves current with what g1t pays | 573 | } |
| 574 | } | |
| 575 | self.db | |
| 576 | .prepare("UPDATE prices SET checked_at = ?") | |
| 577 | .bind(&[now.as_str().into()])? | |
| 578 | .run() | |
| 579 | .await?; | |
| 580 | Ok(()) | |
| 581 | } | |
| 582 | ||
| 583 | /// Moves a meter's cost to a measurement, if it is sound and different. | |
| 584 | async fn measure(&self, meter: &str, measured: f64, reason: &str) -> Result<()> { | |
| 585 | let Some((current, _)) = self.price(meter).await? else { | |
| 586 | return Ok(()); | |
| 587 | }; | |
| 588 | match adopt(current, measured) { | |
| 589 | Err(why) => worker::console_log!("{meter}: {why}"), | |
| 590 | Ok(None) => {} | |
| 591 | Ok(Some(cost)) => { | |
| 592 | let now = now_ms(); | |
| 593 | self.db | |
| 594 | .batch(vec![ | |
| 595 | self.db | |
| 596 | .prepare("UPDATE prices SET cost_micros = ?, source = 'cloudflare', updated_at = ? WHERE meter = ?") | |
| 597 | .bind(&[cost.into(), rfc3339(now).into(), meter.into()])?, | |
| 598 | self.db | |
| 599 | .prepare( | |
| 600 | "INSERT INTO price_changes (id, meter, old_cost_micros, new_cost_micros, markup_percent, reason, created_at) | |
| 601 | SELECT ?, meter, ?, ?, markup_percent, ?, ? FROM prices WHERE meter = ?", | |
| 602 | ) | |
| 603 | .bind(&[ | |
| 604 | new_id("prc", now).into(), | |
| 605 | current.into(), | |
| 606 | cost.into(), | |
| 607 | reason.into(), | |
| 608 | rfc3339(now).into(), | |
| 609 | meter.into(), | |
| 610 | ])?, | |
| 611 | ]) | |
| 612 | .await?; | |
| 613 | } | |
| 614 | } | |
| 615 | Ok(()) | |
| 616 | } | |
| 617 | } | |
| 618 | ||
| 619 | #[cfg(test)] | |
| 620 | mod tests { | |
| 621 | use super::*; | |
| 622 | ||
| 623 | #[test] | |
| 624 | fn small_moves_are_noise_and_wild_ones_are_not_believed() { | |
| 625 | assert_eq!(adopt(21.0, 21.2), Ok(None)); | |
| 626 | assert_eq!(adopt(21.0, 25.0), Ok(Some(25.0))); | |
| 627 | assert_eq!(adopt(21.0, 15.0), Ok(Some(15.0))); | |
| 628 | assert!(adopt(21.0, 200.0).is_err()); | |
| 629 | assert!(adopt(21.0, 0.0).is_err()); | |
| 630 | } | |
| 631 | ||
| 632 | #[test] | |
| The keeper reads Cloudflare as it really answers | 633 | fn a_sandbox_second_is_its_memory_and_disk_and_the_cpu_it_uses() { |
| 634 | // An hour of sandboxes that kept a fifth of a vCPU busy. | |
| 635 | let usage = ContainerUsage { cpu_seconds: 720.0, memory_byte_seconds: 3600.0 * 4.0 * GIB }; | |
| 636 | let micros = sandbox_second_micros(usage, LIST_MEMORY_GIB_SECOND, LIST_DISK_GB_SECOND, LIST_VCPU_SECOND).unwrap(); | |
| 637 | // 4 x 2.5 + 8 x 0.07 + 0.2 x 20 = 14.56 | |
| 638 | assert!((micros - 14.56).abs() < 1e-9, "{micros}"); | |
| 639 | // Too little use to say anything. | |
| 640 | assert!(sandbox_second_micros(ContainerUsage { cpu_seconds: 1.0, memory_byte_seconds: GIB }, 1.0, 1.0, 1.0).is_none()); | |
| 641 | } | |
| 642 | ||
| 643 | #[test] | |
| 644 | fn a_billed_rate_is_the_median_of_the_charged_days() { | |
| 645 | let row = |quantity: f64, cost: f64| UsageRow { | |
| 646 | period_start: String::new(), | |
| 647 | period_end: String::new(), | |
| 648 | service: "Containers / Container Memory".into(), | |
| 649 | unit: "Count".into(), | |
| 650 | quantity, | |
| 651 | cost, | |
| 652 | }; | |
| 653 | let rows = [row(100.0, 0.0), row(100.0, 0.0002), row(100.0, 0.00025), row(100.0, 0.00025)]; | |
| 654 | assert_eq!(billed_rate(&rows.iter().collect::<Vec<_>>()), Some(0.000_002_5)); | |
| 655 | assert_eq!(billed_rate(&[&row(5.0, 0.0)]), None); | |
| 656 | } | |
| 657 | ||
| 658 | #[test] | |
| Prices keep themselves current with what g1t pays | 659 | fn usage_rows_are_read_by_their_focus_names() { |
| 660 | let row = UsageRow::from_value(&json!({ | |
| 661 | "ServiceFamilyName": "Containers", | |
| 662 | "ServiceName": "Memory", | |
| The keeper reads Cloudflare as it really answers | 663 | "PricingUnit": "GiB-seconds", |
| Prices keep themselves current with what g1t pays | 664 | "PricingQuantity": "1200.5", |
| 665 | "ContractedCost": 0.003, | |
| 666 | "ChargePeriodStart": "2026-10-01", | |
| 667 | })) | |
| 668 | .unwrap(); | |
| 669 | assert_eq!(row.service, "Containers / Memory"); | |
| 670 | assert_eq!(row.quantity, 1200.5); | |
| 671 | assert_eq!(row.cost, 0.003); | |
| 672 | assert!(UsageRow::from_value(&json!({ "nothing": 1 })).is_none()); | |
| 673 | } | |
| 674 | } |