| 1 | //! What Cloudflare charges g1t, day by day and product by product. |
| 2 | //! |
| 3 | //! Prices are what g1t pays plus 20%, so g1t has to know what it pays, as |
| 4 | //! Cloudflare counts it, not as g1t assumes it. Some of what g1t runs on |
| 5 | //! is in beta (Artifacts bills "operations" from 2026-10-14 without saying |
| 6 | //! exactly which calls are one), so nothing here assumes a product list: |
| 7 | //! every line Cloudflare bills is kept, whatever it is, and how a line |
| 8 | //! maps to what g1t sells is data (`cost_map`), changed without a deploy. |
| 9 | //! |
| 10 | //! Once a day (`keeper::DAILY`) billing reads: |
| 11 | //! |
| 12 | //! - **Billable usage** (`GET /accounts/{id}/billable-usage`, FOCUS |
| 13 | //! columns): one row per service per day, with its quantity and what it |
| 14 | //! cost g1t. Every Cloudflare product g1t uses appears here once it is |
| 15 | //! used: Workers, Workers for Platforms, D1, KV, R2, Queues, Containers, |
| 16 | //! Durable Objects, Artifacts, Browser Rendering, Workers AI, Vectorize, |
| 17 | //! Cloudflare for SaaS and Email. |
| 18 | //! - **Artifacts events** (GraphQL `artifactsEventsAdaptiveGroups`), by |
| 19 | //! event type and day: what Artifacts itself counted, to compare with |
| 20 | //! what g1t counted (`margin`). |
| 21 | //! |
| 22 | //! Each becomes cost lines in `cost_lines`, one per (day, source, |
| 23 | //! product, meter), upserted, so reading a day again replaces it. The |
| 24 | //! first run reads the last 31 days; after that the last few, since |
| 25 | //! Cloudflare restates recent days as usage settles. |
| 26 | //! |
| 27 | //! The token is `CLOUDFLARE_BILLING_TOKEN` (Account: Billing Read and |
| 28 | //! Account Analytics Read), or the keeper's `CLOUDFLARE_USAGE_TOKEN`, |
| 29 | //! which has both. Without either, nothing is read and nothing fails. |
| 30 | |
| 31 | use std::collections::{BTreeMap, BTreeSet}; |
| 32 | |
| 33 | use g1t_contracts::repos::{GitOperationsArgs, WorkspaceGitOperations}; |
| 34 | use g1t_contracts::time::rfc3339; |
| 35 | use g1t_kit::now_ms; |
| 36 | use serde::Deserialize; |
| 37 | use serde_json::{Value, json}; |
| 38 | use worker::Result; |
| 39 | |
| 40 | use crate::Billing; |
| 41 | use crate::keeper::Keeper; |
| 42 | |
| 43 | /// Where a cost line came from. |
| 44 | pub(crate) const SOURCE_BILLABLE: &str = "billable_usage"; |
| 45 | pub(crate) const SOURCE_ARTIFACTS: &str = "artifacts_events"; |
| 46 | |
| 47 | /// How far back the first run reads: the 31 days GraphQL keeps. |
| 48 | pub(crate) const BACKFILL_DAYS: u64 = 31; |
| 49 | /// How many recent days every later run reads again. |
| 50 | pub(crate) const RESTATE_DAYS: u64 = 4; |
| 51 | |
| 52 | pub(crate) const DAY_MS: u64 = 24 * 60 * 60 * 1000; |
| 53 | |
| 54 | /// One day of one meter of one Cloudflare product. |
| 55 | #[derive(Clone, Debug, PartialEq)] |
| 56 | pub(crate) struct CostLine { |
| 57 | /// YYYY-MM-DD, UTC. |
| 58 | pub day: String, |
| 59 | pub source: &'static str, |
| 60 | /// `containers`, `workers`, `workers_kv`, `artifacts`, … from |
| 61 | /// Cloudflare's own family name. |
| 62 | pub product: String, |
| 63 | /// The service within it, such as `container_memory_per_gib_second` |
| 64 | /// or `events_push`. |
| 65 | pub meter: String, |
| 66 | pub unit: String, |
| 67 | pub quantity: f64, |
| 68 | /// What g1t pays, in dollars: contracted, billed, or list. |
| 69 | pub cost_usd: f64, |
| 70 | /// The name as Cloudflare gave it, for people. |
| 71 | pub raw_name: String, |
| 72 | } |
| 73 | |
| 74 | /// `Workers for Platforms CPU ms (First 60M ms are included)` → |
| 75 | /// `workers_for_platforms_cpu_ms`: lower case, words joined by `_`, and |
| 76 | /// what is in parentheses (the included amount, which changes) left out. |
| 77 | pub(crate) fn slug(text: &str) -> String { |
| 78 | let mut out = String::new(); |
| 79 | let mut depth = 0u32; |
| 80 | let mut gap = false; |
| 81 | for c in text.chars() { |
| 82 | match c { |
| 83 | '(' => depth += 1, |
| 84 | ')' => depth = depth.saturating_sub(1), |
| 85 | _ if depth > 0 => {} |
| 86 | c if c.is_ascii_alphanumeric() => { |
| 87 | if gap && !out.is_empty() { |
| 88 | out.push('_'); |
| 89 | } |
| 90 | gap = false; |
| 91 | out.push(c.to_ascii_lowercase()); |
| 92 | } |
| 93 | _ => gap = true, |
| 94 | } |
| 95 | } |
| 96 | out |
| 97 | } |
| 98 | |
| 99 | /// The product and meter for a family and service name, such as |
| 100 | /// (`D1`, `D1 - Rows Read (first 25 billion included)`) → (`d1`, |
| 101 | /// `d1_rows_read`). Without a family, the service's first word. |
| 102 | pub(crate) fn product_and_meter(family: &str, service: &str) -> (String, String) { |
| 103 | let (family, service) = if family.trim().is_empty() { |
| 104 | match service.split_once(" / ") { |
| 105 | Some((family, service)) => (family, service), |
| 106 | None => (service.split_whitespace().next().unwrap_or("other"), service), |
| 107 | } |
| 108 | } else { |
| 109 | (family, service) |
| 110 | }; |
| 111 | let product = slug(family); |
| 112 | let product = if product.is_empty() { "other".to_owned() } else { product }; |
| 113 | // Kept whole: `Workers for Platforms Requests` under `Workers` must stay |
| 114 | // `workers_for_platforms_requests`. |
| 115 | let meter = slug(service); |
| 116 | let meter = if meter.is_empty() { "usage".to_owned() } else { meter }; |
| 117 | (product, meter) |
| 118 | } |
| 119 | |
| 120 | fn text<'a>(row: &'a Value, keys: &[&str]) -> &'a str { |
| 121 | keys.iter().find_map(|k| row[*k].as_str()).unwrap_or_default() |
| 122 | } |
| 123 | |
| 124 | fn number(row: &Value, keys: &[&str]) -> Option<f64> { |
| 125 | keys.iter() |
| 126 | .find_map(|k| row[*k].as_f64().or_else(|| row[*k].as_str().and_then(|s| s.trim().parse().ok()))) |
| 127 | } |
| 128 | |
| 129 | /// One row of billable usage, read leniently: the API is new, and its |
| 130 | /// field names are FOCUS's (in either case). None without a service or a |
| 131 | /// day. |
| 132 | pub(crate) fn line_from_focus(row: &Value) -> Option<CostLine> { |
| 133 | let service = text(row, &["ServiceName", "service_name", "service"]); |
| 134 | let day = text(row, &["ChargePeriodStart", "charge_period_start", "UsageDate", "date"]); |
| 135 | if service.is_empty() || day.len() < 10 { |
| 136 | return None; |
| 137 | } |
| 138 | let family = text(row, &["ServiceFamilyName", "service_family_name", "ServiceCategory"]); |
| 139 | let (product, meter) = product_and_meter(family, service); |
| 140 | // What g1t pays: contracted, else billed, else list. |
| 141 | let cost = [ |
| 142 | &["ContractedCost", "contracted_cost"][..], |
| 143 | &["BilledCost", "billed_cost"][..], |
| 144 | &["EffectiveCost", "effective_cost"][..], |
| 145 | &["ListCost", "list_cost"][..], |
| 146 | ] |
| 147 | .iter() |
| 148 | .find_map(|keys| number(row, keys).filter(|c| *c > 0.0)) |
| 149 | .unwrap_or(0.0); |
| 150 | Some(CostLine { |
| 151 | day: day[..10].to_owned(), |
| 152 | source: SOURCE_BILLABLE, |
| 153 | product, |
| 154 | meter, |
| 155 | unit: text(row, &["PricingUnit", "pricing_unit", "ConsumedUnit", "consumed_unit"]).to_owned(), |
| 156 | quantity: number(row, &["PricingQuantity", "pricing_quantity", "ConsumedQuantity", "consumed_quantity"]).unwrap_or(0.0), |
| 157 | cost_usd: cost, |
| 158 | raw_name: if family.is_empty() { service.to_owned() } else { format!("{family} / {service}") }, |
| 159 | }) |
| 160 | } |
| 161 | |
| 162 | /// Lines with the same key added together, in key order. Cloudflare can |
| 163 | /// give one service several rows a day (regions, tiers); an upsert of |
| 164 | /// each would keep only the last. |
| 165 | pub(crate) fn aggregate(lines: Vec<CostLine>) -> Vec<CostLine> { |
| 166 | let mut out: Vec<CostLine> = Vec::new(); |
| 167 | for line in lines { |
| 168 | match out |
| 169 | .iter_mut() |
| 170 | .find(|l| l.day == line.day && l.source == line.source && l.product == line.product && l.meter == line.meter) |
| 171 | { |
| 172 | Some(existing) => { |
| 173 | existing.quantity += line.quantity; |
| 174 | existing.cost_usd += line.cost_usd; |
| 175 | } |
| 176 | None => out.push(line), |
| 177 | } |
| 178 | } |
| 179 | out.sort_by(|a, b| (&a.day, a.source, &a.product, &a.meter).cmp(&(&b.day, b.source, &b.product, &b.meter))); |
| 180 | out |
| 181 | } |
| 182 | |
| 183 | /// Every line of a billable-usage answer, aggregated. Errors when |
| 184 | /// Cloudflare says it failed. |
| 185 | pub(crate) fn lines_from_billable(body: &Value) -> std::result::Result<Vec<CostLine>, String> { |
| 186 | if body["success"] == Value::Bool(false) { |
| 187 | return Err(format!("billable usage failed: {}", body["errors"])); |
| 188 | } |
| 189 | let rows = body["result"].as_array().cloned().unwrap_or_default(); |
| 190 | Ok(aggregate(rows.iter().filter_map(line_from_focus).collect())) |
| 191 | } |
| 192 | |
| 193 | /// The GraphQL query for Artifacts' events by type and day. |
| 194 | pub(crate) const ARTIFACTS_QUERY: &str = "query ($account: String!, $since: Date!, $until: Date!) { |
| 195 | viewer { accounts(filter: { accountTag: $account }) { |
| 196 | artifactsEventsAdaptiveGroups(limit: 10000, filter: { date_geq: $since, date_leq: $until }) { |
| 197 | count |
| 198 | dimensions { date eventType repositoryName } |
| 199 | } |
| 200 | } } |
| 201 | }"; |
| 202 | |
| 203 | pub(crate) fn artifacts_variables(account: &str, since: &str, until: &str) -> Value { |
| 204 | json!({ "query": ARTIFACTS_QUERY, "variables": { "account": account, "since": since, "until": until } }) |
| 205 | } |
| 206 | |
| 207 | /// Artifacts' events as lines (`events_push`, `events_pull`, …), with no |
| 208 | /// cost: billable usage carries the cost. Errors when GraphQL does. |
| 209 | pub(crate) fn lines_from_artifacts(body: &Value) -> std::result::Result<Vec<CostLine>, String> { |
| 210 | if let Some(errors) = body["errors"].as_array().filter(|e| !e.is_empty()) { |
| 211 | return Err(format!("Artifacts events failed: {}", Value::Array(errors.clone()))); |
| 212 | } |
| 213 | let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"] |
| 214 | .as_array() |
| 215 | .cloned() |
| 216 | .unwrap_or_default(); |
| 217 | Ok(aggregate( |
| 218 | groups |
| 219 | .iter() |
| 220 | .filter_map(|g| { |
| 221 | let day = g["dimensions"]["date"].as_str()?; |
| 222 | let kind = g["dimensions"]["eventType"].as_str()?; |
| 223 | if day.len() < 10 { |
| 224 | return None; |
| 225 | } |
| 226 | Some(CostLine { |
| 227 | day: day[..10].to_owned(), |
| 228 | source: SOURCE_ARTIFACTS, |
| 229 | product: "artifacts".to_owned(), |
| 230 | meter: format!("events_{}", slug(kind)), |
| 231 | unit: "events".to_owned(), |
| 232 | quantity: g["count"].as_f64().unwrap_or(0.0), |
| 233 | cost_usd: 0.0, |
| 234 | raw_name: format!("Artifacts / {kind}"), |
| 235 | }) |
| 236 | }) |
| 237 | .collect(), |
| 238 | )) |
| 239 | } |
| 240 | |
| 241 | /// Where a cost line came from: AI Gateway's analytics, what it priced |
| 242 | /// g1t's own provider traffic at. |
| 243 | pub(crate) const SOURCE_GATEWAY: &str = "ai_gateway"; |
| 244 | /// The product of AI Gateway's lines. Not `ai_gateway`, which is what a |
| 245 | /// billable-usage line from Cloudflare for AI Gateway itself would slug to. |
| 246 | pub(crate) const GATEWAY_PRODUCT: &str = "ai_gateway_requests"; |
| 247 | /// Meter suffixes of a model's token lines (no cost; for the drift). |
| 248 | pub(crate) const GATEWAY_TOKENS: &str = "__tokens"; |
| 249 | pub(crate) const GATEWAY_CACHE_READ: &str = "__cache_read_tokens"; |
| 250 | pub(crate) const GATEWAY_CACHE_WRITE: &str = "__cache_write_tokens"; |
| 251 | /// The meter prefix of requests Cloudflare billed itself (unified billing), |
| 252 | /// which are on Cloudflare's bill as well as here. |
| 253 | pub(crate) const GATEWAY_WHOLESALE: &str = "wholesale__"; |
| 254 | |
| 255 | /// What AI Gateway priced each day's requests at, per provider and model, |
| 256 | /// for g1t's gateway only: GraphQL `aiGatewayRequestsAdaptiveGroups`, with |
| 257 | /// `sum.cost` (dollars), the tokens it priced and whether Cloudflare billed |
| 258 | /// the request itself (`wholesale`). Field names checked against the |
| 259 | /// schema (`AccountAiGatewayRequestsAdaptiveGroupsSum` and `…Dimensions`). |
| 260 | pub(crate) const GATEWAY_QUERY: &str = "query ($account: String!, $gateway: String!, $since: Date!, $until: Date!) { |
| 261 | viewer { accounts(filter: { accountTag: $account }) { |
| 262 | aiGatewayRequestsAdaptiveGroups(limit: 10000, filter: { date_geq: $since, date_leq: $until, gateway: $gateway }) { |
| 263 | count |
| 264 | sum { cost tokensIn tokensOut cacheReadTokens cacheWriteTokens } |
| 265 | dimensions { date provider model wholesale } |
| 266 | } |
| 267 | } } |
| 268 | }"; |
| 269 | |
| 270 | pub(crate) fn gateway_variables(account: &str, gateway: &str, since: &str, until: &str) -> Value { |
| 271 | json!({ "query": GATEWAY_QUERY, "variables": { "account": account, "gateway": gateway, "since": since, "until": until } }) |
| 272 | } |
| 273 | |
| 274 | /// AI Gateway's analytics as lines: per day and model, the requests at |
| 275 | /// what the gateway priced them (meter `<provider>_<model>`), and their |
| 276 | /// tokens, cache reads and cache writes at no cost (`…__tokens`, …). |
| 277 | /// Errors when GraphQL does. |
| 278 | pub(crate) fn lines_from_gateway(body: &Value) -> std::result::Result<Vec<CostLine>, String> { |
| 279 | if let Some(errors) = body["errors"].as_array().filter(|e| !e.is_empty()) { |
| 280 | return Err(format!("AI Gateway analytics failed: {}", Value::Array(errors.clone()))); |
| 281 | } |
| 282 | let groups = body["data"]["viewer"]["accounts"][0]["aiGatewayRequestsAdaptiveGroups"].as_array().cloned().unwrap_or_default(); |
| 283 | let mut lines = Vec::new(); |
| 284 | for g in &groups { |
| 285 | let d = &g["dimensions"]; |
| 286 | let Some(day) = d["date"].as_str().filter(|day| day.len() >= 10) else { continue }; |
| 287 | let provider = d["provider"].as_str().unwrap_or("unknown"); |
| 288 | let model = d["model"].as_str().unwrap_or("unknown"); |
| 289 | let wholesale = d["wholesale"].as_u64().unwrap_or(0) == 1; |
| 290 | let name = format!("{}{}", if wholesale { GATEWAY_WHOLESALE } else { "" }, slug(&format!("{provider} {model}"))); |
| 291 | let sum = |key: &str| g["sum"][key].as_f64().unwrap_or(0.0); |
| 292 | let line = |meter: String, unit: &str, quantity: f64, cost_usd: f64| CostLine { |
| 293 | day: day[..10].to_owned(), |
| 294 | source: SOURCE_GATEWAY, |
| 295 | product: GATEWAY_PRODUCT.to_owned(), |
| 296 | meter, |
| 297 | unit: unit.to_owned(), |
| 298 | quantity, |
| 299 | cost_usd, |
| 300 | raw_name: format!("AI Gateway / {provider} / {model}{}", if wholesale { " (billed by Cloudflare)" } else { "" }), |
| 301 | }; |
| 302 | lines.push(line(name.clone(), "requests", g["count"].as_f64().unwrap_or(0.0), sum("cost").max(0.0))); |
| 303 | lines.push(line(format!("{name}{GATEWAY_TOKENS}"), "tokens", sum("tokensIn") + sum("tokensOut"), 0.0)); |
| 304 | for (suffix, key) in [(GATEWAY_CACHE_READ, "cacheReadTokens"), (GATEWAY_CACHE_WRITE, "cacheWriteTokens")] { |
| 305 | if sum(key) > 0.0 { |
| 306 | lines.push(line(format!("{name}{suffix}"), "tokens", sum(key), 0.0)); |
| 307 | } |
| 308 | } |
| 309 | } |
| 310 | Ok(aggregate(lines)) |
| 311 | } |
| 312 | |
| 313 | /// What the gateway's lines over a window say about whether its cost can |
| 314 | /// be taken as what the providers bill g1t. |
| 315 | #[derive(Clone, Debug, Default, PartialEq)] |
| 316 | pub(crate) struct GatewayCaveats { |
| 317 | /// Models the gateway put no price on although they used tokens. |
| 318 | pub unpriced: Vec<String>, |
| 319 | pub cache_read_tokens: f64, |
| 320 | pub cache_write_tokens: f64, |
| 321 | /// What Cloudflare billed itself (unified billing), in dollars: on its |
| 322 | /// bill too, so not a cost paid to a provider. |
| 323 | pub wholesale_usd: f64, |
| 324 | /// Runs settled with the gateway's figure short (see `keeper::settled_cost`). |
| 325 | pub short_runs: u32, |
| 326 | } |
| 327 | |
| 328 | /// The caveats in AI Gateway's lines (any days, any order). |
| 329 | pub(crate) fn gateway_caveats(lines: &[(String, f64, f64)]) -> GatewayCaveats { |
| 330 | let mut out = GatewayCaveats::default(); |
| 331 | let mut cost: BTreeMap<&str, f64> = BTreeMap::new(); |
| 332 | let mut tokens: BTreeMap<&str, f64> = BTreeMap::new(); |
| 333 | for (meter, quantity, cost_usd) in lines { |
| 334 | if let Some(model) = meter.strip_suffix(GATEWAY_TOKENS) { |
| 335 | *tokens.entry(model).or_default() += quantity; |
| 336 | } else if meter.ends_with(GATEWAY_CACHE_READ) { |
| 337 | out.cache_read_tokens += quantity; |
| 338 | } else if meter.ends_with(GATEWAY_CACHE_WRITE) { |
| 339 | out.cache_write_tokens += quantity; |
| 340 | } else { |
| 341 | *cost.entry(meter.as_str()).or_default() += cost_usd; |
| 342 | if meter.starts_with(GATEWAY_WHOLESALE) { |
| 343 | out.wholesale_usd += cost_usd; |
| 344 | } |
| 345 | } |
| 346 | } |
| 347 | for (model, used) in tokens { |
| 348 | if used > 0.0 && cost.get(model).copied().unwrap_or(0.0) <= 0.0 { |
| 349 | out.unpriced.push(model.to_owned()); |
| 350 | } |
| 351 | } |
| 352 | out |
| 353 | } |
| 354 | |
| 355 | /// The workspace a repository in the store belongs to: keys are |
| 356 | /// `<workspace>--<repo>`; a pull request's working copy (`pulls--<id>`) is |
| 357 | /// its repository's workspace's, from `owners` (repos' `pull_owners`), else |
| 358 | /// no one's. |
| 359 | pub(crate) fn workspace_of_store_key(key: &str, owners: &BTreeMap<String, String>) -> Option<String> { |
| 360 | let (workspace, rest) = key.split_once("--")?; |
| 361 | if workspace.is_empty() || rest.is_empty() { |
| 362 | return None; |
| 363 | } |
| 364 | if workspace == "pulls" { |
| 365 | return owners.get(rest).map(|w| w.to_lowercase()); |
| 366 | } |
| 367 | Some(workspace.to_lowercase()) |
| 368 | } |
| 369 | |
| 370 | /// The pull request ids whose working copies Artifacts counted events for. |
| 371 | pub(crate) fn pull_ids(body: &Value) -> Vec<String> { |
| 372 | let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"].as_array().cloned().unwrap_or_default(); |
| 373 | let ids: BTreeSet<String> = groups |
| 374 | .iter() |
| 375 | .filter_map(|g| g["dimensions"]["repositoryName"].as_str()?.strip_prefix("pulls--").map(str::to_owned)) |
| 376 | .filter(|id| !id.is_empty()) |
| 377 | .collect(); |
| 378 | ids.into_iter().collect() |
| 379 | } |
| 380 | |
| 381 | /// Artifacts' billable operations per workspace and day, by repository |
| 382 | /// name: how Cloudflare's own count shares out. The meter is |
| 383 | /// `cloudflare_git`, which shares out the git bucket's cost (`margin`). |
| 384 | pub(crate) fn artifacts_by_workspace(body: &Value, owners: &BTreeMap<String, String>) -> Vec<(String, String, f64)> { |
| 385 | let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"].as_array().cloned().unwrap_or_default(); |
| 386 | let mut out: BTreeMap<(String, String), f64> = BTreeMap::new(); |
| 387 | for g in &groups { |
| 388 | let (Some(day), Some(kind)) = (g["dimensions"]["date"].as_str(), g["dimensions"]["eventType"].as_str()) else { continue }; |
| 389 | if day.len() < 10 || !ARTIFACTS_OPERATIONS.contains(&format!("events_{}", slug(kind)).as_str()) { |
| 390 | continue; |
| 391 | } |
| 392 | let Some(workspace) = g["dimensions"]["repositoryName"].as_str().and_then(|key| workspace_of_store_key(key, owners)) else { continue }; |
| 393 | *out.entry((day[..10].to_owned(), workspace)).or_default() += g["count"].as_f64().unwrap_or(0.0); |
| 394 | } |
| 395 | out.into_iter().map(|((day, workspace), count)| (day, workspace, count)).collect() |
| 396 | } |
| 397 | |
| 398 | /// Artifacts' event types that are operations it bills: what its |
| 399 | /// pricing names (create, push, pull, clone) and the metrics list. |
| 400 | /// Errors such as `rateLimited` are not. |
| 401 | pub(crate) const ARTIFACTS_OPERATIONS: [&str; 5] = ["events_create", "events_fork", "events_push", "events_pull", "events_delete"]; |
| 402 | |
| 403 | /// The days to read: the 31 before today on the first run (`last` is |
| 404 | /// None), else the last few, never further back than the backfill. |
| 405 | pub(crate) fn window(last_fetched_day: Option<&str>, now_ms: u64) -> (String, String) { |
| 406 | let day = |ms: u64| g1t_contracts::time::rfc3339(ms)[..10].to_owned(); |
| 407 | let today = day(now_ms); |
| 408 | let since = match last_fetched_day { |
| 409 | None => day(now_ms.saturating_sub(BACKFILL_DAYS * DAY_MS)), |
| 410 | Some(last) => { |
| 411 | let restate = day(now_ms.saturating_sub(RESTATE_DAYS * DAY_MS)); |
| 412 | let floor = day(now_ms.saturating_sub(BACKFILL_DAYS * DAY_MS)); |
| 413 | // A gap since the last run is read too, up to the backfill. |
| 414 | let from = if last < restate.as_str() { last.to_owned() } else { restate }; |
| 415 | if from < floor { floor } else { from } |
| 416 | } |
| 417 | }; |
| 418 | (since, today) |
| 419 | } |
| 420 | |
| 421 | /// Every day from `since` to `until`, inclusive. |
| 422 | pub(crate) fn days_between(since: &str, until: &str) -> Vec<String> { |
| 423 | let start = g1t_contracts::time::parse_rfc3339(&format!("{since}T00:00:00Z")); |
| 424 | let end = g1t_contracts::time::parse_rfc3339(&format!("{until}T00:00:00Z")); |
| 425 | match (start, end) { |
| 426 | (Some(start), Some(end)) if start <= end => (0..=((end - start) / DAY_MS)) |
| 427 | .map(|n| g1t_contracts::time::rfc3339(start + n * DAY_MS)[..10].to_owned()) |
| 428 | .collect(), |
| 429 | _ => Vec::new(), |
| 430 | } |
| 431 | } |
| 432 | |
| 433 | /// A rule from `cost_map`: which Cloudflare lines feed which of g1t's |
| 434 | /// products. |
| 435 | #[derive(Clone, Debug, PartialEq, serde::Deserialize)] |
| 436 | pub(crate) struct Rule { |
| 437 | /// Cloudflare's product, as slugged here. |
| 438 | pub product: String, |
| 439 | /// A meter prefix, or `*` for any meter of the product. |
| 440 | pub meter: String, |
| 441 | /// g1t's product it is a cost of: `sandboxes`, `git`, `platform`, … |
| 442 | pub bucket: String, |
| 443 | /// The price book meter whose cost it measures, if any. |
| 444 | pub price_meter: Option<String>, |
| 445 | /// g1t's own count of the same units, to compare quantities. |
| 446 | pub own_meter: Option<String>, |
| 447 | /// How far g1t's count may be from Cloudflare's before it is drift. |
| 448 | pub drift_percent: f64, |
| 449 | } |
| 450 | |
| 451 | /// The rule for a line: its product's rule with the longest matching |
| 452 | /// meter prefix, `*` last. None means no one decided what pays for it: |
| 453 | /// a leak until someone does. |
| 454 | pub(crate) fn classify<'a>(rules: &'a [Rule], product: &str, meter: &str) -> Option<&'a Rule> { |
| 455 | rules |
| 456 | .iter() |
| 457 | .filter(|r| r.product == product && (r.meter == "*" || meter.starts_with(r.meter.as_str()))) |
| 458 | .max_by_key(|r| if r.meter == "*" { 0 } else { r.meter.len() + 1 }) |
| 459 | } |
| 460 | |
| 461 | /// What a bucket is called in sudo. |
| 462 | pub(crate) fn bucket_title(bucket: &str) -> String { |
| 463 | match bucket { |
| 464 | "sandboxes" => "Sandboxes and builds".into(), |
| 465 | "deployments" => "Deployments".into(), |
| 466 | "git" => "Git operations".into(), |
| 467 | "repo_storage" => "Repository storage".into(), |
| 468 | "actions_cache" => "Actions cache".into(), |
| 469 | "embeddings" => "Search embeddings".into(), |
| 470 | "security" => "Security scans".into(), |
| 471 | "domains" => "Custom domains".into(), |
| 472 | "models" => "Models".into(), |
| 473 | "platform" => "Running g1t (paid by the plan)".into(), |
| 474 | UNMAPPED => "Not mapped".into(), |
| 475 | other => other.replace('_', " "), |
| 476 | } |
| 477 | } |
| 478 | |
| 479 | /// The bucket of a line no rule claims. |
| 480 | pub(crate) const UNMAPPED: &str = "unmapped"; |
| 481 | |
| 482 | /// A day's count per workspace from counts "from this day to the end of |
| 483 | /// its month" (what the repos service's `git_operations` answers with a |
| 484 | /// `since`): each day's is its own less the next day's, within a month. |
| 485 | /// The last day (today, so far) is its own. |
| 486 | pub(crate) fn daily_from_cumulative(days: &[String], cumulative: &[BTreeMap<String, u64>]) -> Vec<(String, String, u64)> { |
| 487 | let mut out = Vec::new(); |
| 488 | for (index, day) in days.iter().enumerate() { |
| 489 | let Some(today) = cumulative.get(index) else { break }; |
| 490 | let next = days |
| 491 | .get(index + 1) |
| 492 | .filter(|next| next[..7] == day[..7]) |
| 493 | .and_then(|_| cumulative.get(index + 1)); |
| 494 | for (workspace, count) in today { |
| 495 | let later = next.and_then(|n| n.get(workspace)).copied().unwrap_or(0); |
| 496 | let own = count.saturating_sub(later); |
| 497 | if own > 0 { |
| 498 | out.push((day.clone(), workspace.clone(), own)); |
| 499 | } |
| 500 | } |
| 501 | } |
| 502 | out |
| 503 | } |
| 504 | |
| 505 | /// One row of the repos service's `artifacts_usage`, read leniently while |
| 506 | /// its shape settles: a day, a workspace, a raw meter and a count. |
| 507 | pub(crate) fn raw_usage_row(row: &Value) -> Option<(String, String, String, f64)> { |
| 508 | let day = text(row, &["day", "date"]); |
| 509 | let workspace = text(row, &["namespace", "workspace"]); |
| 510 | let meter = text(row, &["meter", "kind", "event", "operation"]); |
| 511 | let count = number(row, &["count", "quantity", "operations", "value"])?; |
| 512 | (day.len() >= 10 && !workspace.is_empty() && !meter.is_empty()) |
| 513 | .then(|| (day[..10].to_owned(), workspace.to_lowercase(), slug(meter), count)) |
| 514 | } |
| 515 | |
| 516 | /// The repos service's operation mapping, as `artifacts_usage` returns it |
| 517 | /// beside the rows: for each raw meter (slugged), how many operations it |
| 518 | /// is to Cloudflare (`cost_operations`) and to the customer |
| 519 | /// (`billable_operations`). Repos owns this mapping |
| 520 | /// (`set_operation_mapping`); billing only reads it. |
| 521 | pub(crate) fn operation_mapping(body: &Value) -> BTreeMap<String, (f64, f64)> { |
| 522 | body["mapping"] |
| 523 | .as_array() |
| 524 | .map(|rows| { |
| 525 | rows.iter() |
| 526 | .filter_map(|r| { |
| 527 | let meter = r["meter"].as_str()?; |
| 528 | Some((slug(meter), (r["cost_operations"].as_f64().unwrap_or(0.0), r["billable_operations"].as_f64().unwrap_or(0.0)))) |
| 529 | }) |
| 530 | .collect() |
| 531 | }) |
| 532 | .unwrap_or_default() |
| 533 | } |
| 534 | |
| 535 | /// Raw counts weighted by one column of the mapping: what Cloudflare |
| 536 | /// should count (`cost`), or what customers are charged for. |
| 537 | pub(crate) fn weighted(raw: &[(String, f64)], mapping: &BTreeMap<String, (f64, f64)>, cost: bool) -> f64 { |
| 538 | raw.iter() |
| 539 | .map(|(meter, count)| count * mapping.get(meter).map_or(0.0, |(c, b)| if cost { *c } else { *b })) |
| 540 | .sum() |
| 541 | } |
| 542 | |
| 543 | impl Billing { |
| 544 | /// Reads Cloudflare's bill for the days due (see `window`) into |
| 545 | /// `cost_lines`: the days read and how many lines. None without a |
| 546 | /// token. What could not be read is added to `problems`. |
| 547 | pub(crate) async fn read_cloudflare(&self, keeper: &Keeper, problems: &mut Vec<String>) -> Result<Option<(String, String, u32)>> { |
| 548 | if !keeper.can_read_bill() { |
| 549 | return Ok(None); |
| 550 | } |
| 551 | #[derive(Deserialize)] |
| 552 | struct Last { |
| 553 | day: Option<String>, |
| 554 | } |
| 555 | let last = self |
| 556 | .db |
| 557 | .prepare("SELECT MAX(day) AS day FROM cost_lines WHERE source = ?") |
| 558 | .bind(&[SOURCE_BILLABLE.into()])? |
| 559 | .first::<Last>(None) |
| 560 | .await? |
| 561 | .and_then(|l| l.day); |
| 562 | let (since, until) = window(last.as_deref(), now_ms()); |
| 563 | let fetched_at = rfc3339(now_ms()); |
| 564 | let mut written = 0; |
| 565 | match keeper.billable_usage_body(&since, &until).await.map_err(|e| e.to_string()).and_then(|b| lines_from_billable(&b)) { |
| 566 | Ok(lines) => written += self.upsert_lines(&lines, &fetched_at).await?, |
| 567 | Err(error) => problems.push(format!("Cloudflare's billable usage could not be read: {error}")), |
| 568 | } |
| 569 | match keeper.graphql(artifacts_variables(keeper.account(), &since, &until)).await.map_err(|e| e.to_string()) { |
| 570 | Ok(body) => match lines_from_artifacts(&body) { |
| 571 | Ok(lines) => { |
| 572 | written += self.upsert_lines(&lines, &fetched_at).await?; |
| 573 | let owners = self.pull_owners(&pull_ids(&body)).await; |
| 574 | self.keep_cloudflare_counts(&since, &until, &artifacts_by_workspace(&body, &owners), &fetched_at).await?; |
| 575 | } |
| 576 | Err(error) => problems.push(format!("Artifacts events could not be read: {error}")), |
| 577 | }, |
| 578 | Err(error) => problems.push(format!("Artifacts events could not be read: {error}")), |
| 579 | } |
| 580 | // What AI Gateway priced g1t's own provider traffic at, each day: |
| 581 | // the total the ledger's model cost is checked against (`margin`). |
| 582 | if !keeper.gateway().is_empty() { |
| 583 | match keeper |
| 584 | .graphql_either(gateway_variables(keeper.account(), keeper.gateway(), &since, &until)) |
| 585 | .await |
| 586 | .map_err(|e| e.to_string()) |
| 587 | .and_then(|body| lines_from_gateway(&body)) |
| 588 | { |
| 589 | Ok(lines) => { |
| 590 | // A day's models are read whole: one that is gone from a |
| 591 | // re-read day must not keep its old line. |
| 592 | self.db |
| 593 | .prepare("DELETE FROM cost_lines WHERE source = ?1 AND day >= ?2 AND day <= ?3") |
| 594 | .bind(&[SOURCE_GATEWAY.into(), since.as_str().into(), until.as_str().into()])? |
| 595 | .run() |
| 596 | .await?; |
| 597 | written += self.upsert_lines(&lines, &fetched_at).await?; |
| 598 | } |
| 599 | Err(error) => problems.push(format!("AI Gateway's analytics could not be read: {error}")), |
| 600 | } |
| 601 | } |
| 602 | Ok(Some((since, until, written))) |
| 603 | } |
| 604 | |
| 605 | /// Cloudflare's own per-workspace counts for the days, replacing what |
| 606 | /// was kept for them. |
| 607 | /// The workspace of each pull request's working copy, from repos; |
| 608 | /// nothing while repos does not answer (those events stay no one's). |
| 609 | async fn pull_owners(&self, pulls: &[String]) -> BTreeMap<String, String> { |
| 610 | let Some(repos) = &self.repos else { return BTreeMap::new() }; |
| 611 | if pulls.is_empty() { |
| 612 | return BTreeMap::new(); |
| 613 | } |
| 614 | match g1t_kit::call::<_, Value>(repos, "pull_owners", &json!({ "pulls": pulls })).await { |
| 615 | Ok(body) => body["owners"] |
| 616 | .as_object() |
| 617 | .map(|o| o.iter().filter_map(|(k, v)| Some((k.clone(), v.as_str()?.to_owned()))).collect()) |
| 618 | .unwrap_or_default(), |
| 619 | Err(error) => { |
| 620 | worker::console_error!("pull_owners: {error}"); |
| 621 | BTreeMap::new() |
| 622 | } |
| 623 | } |
| 624 | } |
| 625 | |
| 626 | async fn keep_cloudflare_counts(&self, since: &str, until: &str, counts: &[(String, String, f64)], fetched_at: &str) -> Result<()> { |
| 627 | self.db |
| 628 | .prepare("DELETE FROM own_counts WHERE meter = 'cloudflare_git' AND day >= ?1 AND day <= ?2") |
| 629 | .bind(&[since.into(), until.into()])? |
| 630 | .run() |
| 631 | .await?; |
| 632 | for chunk in counts.chunks(50) { |
| 633 | let mut statements = Vec::with_capacity(chunk.len()); |
| 634 | for (day, workspace, count) in chunk { |
| 635 | statements.push( |
| 636 | self.db |
| 637 | .prepare("INSERT OR REPLACE INTO own_counts (day, meter, workspace, quantity, fetched_at) VALUES (?, 'cloudflare_git', ?, ?, ?)") |
| 638 | .bind(&[day.as_str().into(), workspace.as_str().into(), (*count).into(), fetched_at.into()])?, |
| 639 | ); |
| 640 | } |
| 641 | self.db.batch(statements).await?; |
| 642 | } |
| 643 | Ok(()) |
| 644 | } |
| 645 | |
| 646 | /// Upserts lines, a day's line replacing what was read for it before. |
| 647 | async fn upsert_lines(&self, lines: &[CostLine], fetched_at: &str) -> Result<u32> { |
| 648 | for chunk in lines.chunks(50) { |
| 649 | let mut statements = Vec::with_capacity(chunk.len()); |
| 650 | for line in chunk { |
| 651 | statements.push( |
| 652 | self.db |
| 653 | .prepare( |
| 654 | "INSERT INTO cost_lines (day, source, product, meter, unit, quantity, cost_usd, raw_name, fetched_at) |
| 655 | VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9) |
| 656 | ON CONFLICT (day, source, product, meter) DO UPDATE SET |
| 657 | unit = ?5, quantity = ?6, cost_usd = ?7, raw_name = ?8, fetched_at = ?9", |
| 658 | ) |
| 659 | .bind(&[ |
| 660 | line.day.as_str().into(), |
| 661 | line.source.into(), |
| 662 | line.product.as_str().into(), |
| 663 | line.meter.as_str().into(), |
| 664 | line.unit.as_str().into(), |
| 665 | line.quantity.into(), |
| 666 | line.cost_usd.into(), |
| 667 | line.raw_name.as_str().into(), |
| 668 | fetched_at.into(), |
| 669 | ])?, |
| 670 | ); |
| 671 | } |
| 672 | self.db.batch(statements).await?; |
| 673 | } |
| 674 | Ok(lines.len() as u32) |
| 675 | } |
| 676 | |
| 677 | /// The mapping from Cloudflare's meters to g1t's products. |
| 678 | pub(crate) async fn rules(&self) -> Result<Vec<Rule>> { |
| 679 | #[derive(Deserialize)] |
| 680 | struct Row { |
| 681 | product: String, |
| 682 | meter: String, |
| 683 | bucket: String, |
| 684 | price_meter: Option<String>, |
| 685 | own_meter: Option<String>, |
| 686 | drift_percent: f64, |
| 687 | } |
| 688 | Ok(self |
| 689 | .db |
| 690 | .prepare("SELECT product, meter, bucket, price_meter, own_meter, drift_percent FROM cost_map") |
| 691 | .all() |
| 692 | .await? |
| 693 | .results::<Row>()? |
| 694 | .into_iter() |
| 695 | .map(|r| Rule { |
| 696 | product: r.product, |
| 697 | meter: r.meter, |
| 698 | bucket: r.bucket, |
| 699 | price_meter: r.price_meter, |
| 700 | own_meter: r.own_meter, |
| 701 | drift_percent: r.drift_percent, |
| 702 | }) |
| 703 | .collect()) |
| 704 | } |
| 705 | |
| 706 | /// g1t's own counts for the days, all from the repos service, which |
| 707 | /// owns the mapping from raw meters to operations (`operation_mapping`): |
| 708 | /// |
| 709 | /// - `git_operations`: what customers are charged for, as repos counts |
| 710 | /// it (`git_operations`, already through its mapping). |
| 711 | /// - `cost_operations`: what g1t expects Cloudflare to bill, from the |
| 712 | /// raw meters (`artifacts_usage`) and the mapping's cost column. |
| 713 | /// - `artifacts_<meter>`: each raw meter. |
| 714 | pub(crate) async fn count_own(&self, since: &str, until: &str) -> Result<()> { |
| 715 | let Some(repos) = &self.repos else { return Ok(()) }; |
| 716 | let days = days_between(since, until); |
| 717 | let fetched_at = rfc3339(now_ms()); |
| 718 | let mut rows: Vec<(String, String, String, f64)> = Vec::new(); |
| 719 | |
| 720 | // Raw meters and repos' mapping; skipped while repos does not answer. |
| 721 | if let Ok(body) = g1t_kit::call::<_, Value>(repos, "artifacts_usage", &json!({ "from": since, "to": until })).await { |
| 722 | let mapping = operation_mapping(&body); |
| 723 | let mut by: BTreeMap<(String, String), Vec<(String, f64)>> = BTreeMap::new(); |
| 724 | for (day, workspace, meter, count) in body["rows"].as_array().map(|l| l.iter().filter_map(raw_usage_row).collect::<Vec<_>>()).unwrap_or_default() { |
| 725 | by.entry((day, workspace)).or_default().push((meter, count)); |
| 726 | } |
| 727 | for ((day, workspace), counts) in by { |
| 728 | // Several stores (namespaces) can give the same meter. |
| 729 | let mut merged: BTreeMap<String, f64> = BTreeMap::new(); |
| 730 | for (meter, count) in &counts { |
| 731 | *merged.entry(meter.clone()).or_default() += count; |
| 732 | } |
| 733 | for (meter, count) in &merged { |
| 734 | rows.push((day.clone(), format!("artifacts_{meter}"), workspace.clone(), *count)); |
| 735 | } |
| 736 | if !mapping.is_empty() { |
| 737 | rows.push((day, "cost_operations".to_owned(), workspace, weighted(&counts, &mapping, true))); |
| 738 | } |
| 739 | } |
| 740 | } |
| 741 | // What customers are charged for, as repos counts it through its mapping. |
| 742 | let mut cumulative = Vec::with_capacity(days.len()); |
| 743 | for day in &days { |
| 744 | let list: Vec<WorkspaceGitOperations> = g1t_kit::call( |
| 745 | repos, |
| 746 | "git_operations", |
| 747 | &GitOperationsArgs { month: day[..7].to_owned(), since: Some(format!("{day}T00")), namespace: None }, |
| 748 | ) |
| 749 | .await?; |
| 750 | cumulative.push(list.into_iter().map(|w| (w.namespace.to_lowercase(), w.operations)).collect::<BTreeMap<_, _>>()); |
| 751 | } |
| 752 | for (day, workspace, count) in daily_from_cumulative(&days, &cumulative) { |
| 753 | rows.push((day, "git_operations".to_owned(), workspace, count as f64)); |
| 754 | } |
| 755 | |
| 756 | // Each day's counts replace what was there. |
| 757 | self.db |
| 758 | .prepare("DELETE FROM own_counts WHERE day >= ?1 AND day <= ?2 AND meter NOT LIKE 'cloudflare_%'") |
| 759 | .bind(&[since.into(), until.into()])? |
| 760 | .run() |
| 761 | .await?; |
| 762 | for chunk in rows.chunks(50) { |
| 763 | let mut statements = Vec::with_capacity(chunk.len()); |
| 764 | for (day, meter, workspace, quantity) in chunk { |
| 765 | statements.push( |
| 766 | self.db |
| 767 | .prepare( |
| 768 | "INSERT INTO own_counts (day, meter, workspace, quantity, fetched_at) VALUES (?1, ?2, ?3, ?4, ?5) |
| 769 | ON CONFLICT (day, meter, workspace) DO UPDATE SET quantity = ?4, fetched_at = ?5", |
| 770 | ) |
| 771 | .bind(&[day.as_str().into(), meter.as_str().into(), workspace.as_str().into(), (*quantity).into(), fetched_at.as_str().into()])?, |
| 772 | ); |
| 773 | } |
| 774 | self.db.batch(statements).await?; |
| 775 | } |
| 776 | Ok(()) |
| 777 | } |
| 778 | } |
| 779 | |
| 780 | #[cfg(test)] |
| 781 | mod tests { |
| 782 | use super::*; |
| 783 | |
| 784 | /// AI Gateway's analytics in the shape of the schema |
| 785 | /// (`aiGatewayRequestsAdaptiveGroups`: `count`, `sum`, `dimensions`). |
| 786 | fn gateway_fixture() -> Value { |
| 787 | let group = |day: &str, provider: &str, model: &str, wholesale: u8, count: u64, cost: f64, tokens: (f64, f64, f64, f64)| { |
| 788 | json!({ |
| 789 | "count": count, |
| 790 | "sum": { "cost": cost, "tokensIn": tokens.0, "tokensOut": tokens.1, "cacheReadTokens": tokens.2, "cacheWriteTokens": tokens.3 }, |
| 791 | "dimensions": { "date": day, "provider": provider, "model": model, "wholesale": wholesale } |
| 792 | }) |
| 793 | }; |
| 794 | json!({ "data": { "viewer": { "accounts": [{ "aiGatewayRequestsAdaptiveGroups": [ |
| 795 | group("2026-10-05", "anthropic", "claude-sonnet-5-5", 0, 120, 4.25, (900_000.0, 40_000.0, 3_000_000.0, 200_000.0)), |
| 796 | group("2026-10-05", "anthropic", "claude-haiku-4-5-20251001", 0, 300, 0.40, (400_000.0, 20_000.0, 0.0, 0.0)), |
| 797 | group("2026-10-06", "anthropic", "claude-new-1", 0, 12, 0.0, (80_000.0, 4_000.0, 0.0, 0.0)), |
| 798 | group("2026-10-06", "openai", "gpt-x", 1, 5, 0.10, (1_000.0, 100.0, 0.0, 0.0)), |
| 799 | ] }] } }, "errors": null }) |
| 800 | } |
| 801 | |
| 802 | #[test] |
| 803 | fn the_gateway_query_names_the_fields_its_schema_has() { |
| 804 | // As checked against Cloudflare's GraphQL schema (introspection of |
| 805 | // AccountAiGatewayRequestsAdaptiveGroups{,Sum,Dimensions,Filter}). |
| 806 | for field in ["aiGatewayRequestsAdaptiveGroups", "date_geq", "date_leq", "gateway: $gateway", "count", "cost", "tokensIn", "tokensOut", "cacheReadTokens", "cacheWriteTokens", "date", "provider", "model", "wholesale"] { |
| 807 | assert!(GATEWAY_QUERY.contains(field), "{field}"); |
| 808 | } |
| 809 | let body = gateway_variables("acct", "g1t", "2026-10-01", "2026-10-07"); |
| 810 | assert_eq!(body["variables"]["gateway"], "g1t"); |
| 811 | } |
| 812 | |
| 813 | #[test] |
| 814 | fn what_the_gateway_priced_is_kept_per_day_and_model_with_its_tokens() { |
| 815 | let lines = lines_from_gateway(&gateway_fixture()).unwrap(); |
| 816 | let get = |day: &str, meter: &str| lines.iter().find(|l| l.day == day && l.meter == meter).unwrap_or_else(|| panic!("{day} {meter}")); |
| 817 | let sonnet = get("2026-10-05", "anthropic_claude_sonnet_5_5"); |
| 818 | assert_eq!((sonnet.source, sonnet.product.as_str(), sonnet.quantity, sonnet.cost_usd), (SOURCE_GATEWAY, GATEWAY_PRODUCT, 120.0, 4.25)); |
| 819 | assert_eq!(get("2026-10-05", "anthropic_claude_sonnet_5_5__tokens").quantity, 940_000.0); |
| 820 | assert_eq!(get("2026-10-05", "anthropic_claude_sonnet_5_5__cache_read_tokens").quantity, 3_000_000.0); |
| 821 | assert_eq!(get("2026-10-05", "anthropic_claude_sonnet_5_5__cache_write_tokens").cost_usd, 0.0); |
| 822 | assert_eq!(get("2026-10-06", "wholesale__openai_gpt_x").cost_usd, 0.10); |
| 823 | // Only the request lines carry cost: the day's total is the gateway's. |
| 824 | let day5: f64 = lines.iter().filter(|l| l.day == "2026-10-05").map(|l| l.cost_usd).sum(); |
| 825 | assert!((day5 - 4.65).abs() < 1e-9); |
| 826 | // Errors are a problem for the run, not lines. |
| 827 | assert!(lines_from_gateway(&json!({ "errors": [{ "message": "not authorized for that account" }] })).is_err()); |
| 828 | } |
| 829 | |
| 830 | #[test] |
| 831 | fn the_gateway_lines_say_when_its_cost_cannot_be_the_providers() { |
| 832 | let rows: Vec<(String, f64, f64)> = |
| 833 | lines_from_gateway(&gateway_fixture()).unwrap().into_iter().map(|l| (l.meter, l.quantity, l.cost_usd)).collect(); |
| 834 | let caveats = gateway_caveats(&rows); |
| 835 | assert_eq!(caveats.unpriced, vec!["anthropic_claude_new_1".to_owned()]); |
| 836 | assert_eq!((caveats.cache_read_tokens, caveats.cache_write_tokens), (3_000_000.0, 200_000.0)); |
| 837 | assert!((caveats.wholesale_usd - 0.10).abs() < 1e-9); |
| 838 | // A model priced on one day and not another is priced. |
| 839 | let mixed = vec![ |
| 840 | ("m".to_owned(), 3.0, 0.5), |
| 841 | ("m__tokens".to_owned(), 10.0, 0.0), |
| 842 | ("m".to_owned(), 3.0, 0.0), |
| 843 | ("m__tokens".to_owned(), 10.0, 0.0), |
| 844 | ]; |
| 845 | assert!(gateway_caveats(&mixed).unpriced.is_empty()); |
| 846 | } |
| 847 | |
| 848 | /// Billable usage as Cloudflare answered g1t on 2026-10-06 (two days, |
| 849 | /// trimmed), plus rows past the included amounts, which cost money. |
| 850 | fn billable_fixture() -> Value { |
| 851 | json!({ |
| 852 | "success": true, |
| 853 | "errors": [], |
| 854 | "result": [ |
| 855 | { |
| 856 | "ChargePeriodStart": "2026-10-03T00:00:00Z", |
| 857 | "ChargePeriodEnd": "2026-10-04T00:00:00Z", |
| 858 | "ServiceFamilyName": "Containers", |
| 859 | "ServiceName": "Container Memory, per GiB-Second (First 25 GiB-hours included)", |
| 860 | "PricingUnit": "Count", |
| 861 | "PricingQuantity": "14000", |
| 862 | "ContractedCost": 0, "BilledCost": 0, "ListCost": 0 |
| 863 | }, |
| 864 | { |
| 865 | "ChargePeriodStart": "2026-10-03T00:00:00Z", |
| 866 | "ServiceFamilyName": "D1", |
| 867 | "ServiceName": "D1 - Rows Read (first 25 billion included)", |
| 868 | "PricingUnit": "Count", |
| 869 | "PricingQuantity": 1200000, |
| 870 | "BilledCost": 0 |
| 871 | }, |
| 872 | { |
| 873 | "ChargePeriodStart": "2026-10-15T00:00:00Z", |
| 874 | "ServiceFamilyName": "Artifacts", |
| 875 | "ServiceName": "Artifacts Operations (First 10,000 included)", |
| 876 | "PricingUnit": "Count", |
| 877 | "PricingQuantity": "30000", |
| 878 | "BilledCost": "3.00", |
| 879 | "ListCost": 4.5 |
| 880 | }, |
| 881 | { |
| 882 | "ChargePeriodStart": "2026-10-15T00:00:00Z", |
| 883 | "ServiceFamilyName": "Artifacts", |
| 884 | "ServiceName": "Artifacts Operations (First 10,000 included)", |
| 885 | "PricingUnit": "Count", |
| 886 | "PricingQuantity": "10000", |
| 887 | "BilledCost": "1.50" |
| 888 | }, |
| 889 | { |
| 890 | "ChargePeriodStart": "2026-10-15T00:00:00Z", |
| 891 | "ServiceFamilyName": "Workers", |
| 892 | "ServiceName": "Workers for Platforms Requests (First 20M are included)", |
| 893 | "PricingUnit": "Count", |
| 894 | "PricingQuantity": 25000000, |
| 895 | "ListCost": 1.5 |
| 896 | }, |
| 897 | { "ServiceName": "No day, no line" } |
| 898 | ] |
| 899 | }) |
| 900 | } |
| 901 | |
| 902 | #[test] |
| 903 | fn names_become_products_and_meters() { |
| 904 | assert_eq!(slug("Workers for Platforms CPU ms (First 60M ms are included)"), "workers_for_platforms_cpu_ms"); |
| 905 | assert_eq!(product_and_meter("D1", "D1 - Rows Read (first 25 billion included)"), ("d1".into(), "d1_rows_read".into())); |
| 906 | assert_eq!(product_and_meter("Workers KV", "KV Read Operations (First 10M is included)"), ("workers_kv".into(), "kv_read_operations".into())); |
| 907 | assert_eq!(product_and_meter("", "Containers / Container vCPU"), ("containers".into(), "container_vcpu".into())); |
| 908 | assert_eq!(product_and_meter("Email", ""), ("email".into(), "usage".into())); |
| 909 | } |
| 910 | |
| 911 | #[test] |
| 912 | fn billable_usage_is_parsed_and_rows_of_one_meter_are_added() { |
| 913 | let lines = lines_from_billable(&billable_fixture()).unwrap(); |
| 914 | assert_eq!(lines.len(), 4, "{lines:#?}"); |
| 915 | let artifacts = lines.iter().find(|l| l.product == "artifacts").unwrap(); |
| 916 | // Two rows of the same day and meter: one line, added up. |
| 917 | assert_eq!(artifacts.meter, "artifacts_operations"); |
| 918 | assert_eq!(artifacts.quantity, 40_000.0); |
| 919 | assert!((artifacts.cost_usd - 4.5).abs() < 1e-9, "billed, not list: {}", artifacts.cost_usd); |
| 920 | let wfp = lines.iter().find(|l| l.meter.starts_with("workers_for_platforms")).unwrap(); |
| 921 | assert_eq!(wfp.product, "workers"); |
| 922 | // No billed cost: the list cost. |
| 923 | assert_eq!(wfp.cost_usd, 1.5); |
| 924 | let memory = lines.iter().find(|l| l.product == "containers").unwrap(); |
| 925 | assert_eq!((memory.day.as_str(), memory.quantity, memory.cost_usd), ("2026-10-03", 14_000.0, 0.0)); |
| 926 | assert!(lines_from_billable(&json!({ "success": false, "errors": [{ "code": 10000 }] })).is_err()); |
| 927 | } |
| 928 | |
| 929 | #[test] |
| 930 | fn reading_the_same_days_twice_gives_the_same_lines() { |
| 931 | // Idempotent: the same answer aggregates to the same keys and |
| 932 | // amounts, so the upsert replaces rather than adds. |
| 933 | let once = lines_from_billable(&billable_fixture()).unwrap(); |
| 934 | let twice = lines_from_billable(&billable_fixture()).unwrap(); |
| 935 | assert_eq!(once, twice); |
| 936 | let again = aggregate(once.clone()); |
| 937 | assert_eq!(again, once); |
| 938 | } |
| 939 | |
| 940 | #[test] |
| 941 | fn artifacts_events_are_counted_by_type_and_day() { |
| 942 | let body = json!({ |
| 943 | "data": { "viewer": { "accounts": [{ "artifactsEventsAdaptiveGroups": [ |
| 944 | { "count": 120, "dimensions": { "date": "2026-10-05", "eventType": "pull" } }, |
| 945 | { "count": 30, "dimensions": { "date": "2026-10-05", "eventType": "push" } }, |
| 946 | { "count": 2, "dimensions": { "date": "2026-10-05", "eventType": "rateLimited" } }, |
| 947 | { "count": 5, "dimensions": { "date": "2026-10-06", "eventType": "pull" } } |
| 948 | ] }] } }, |
| 949 | "errors": null |
| 950 | }); |
| 951 | let lines = lines_from_artifacts(&body).unwrap(); |
| 952 | assert_eq!(lines.len(), 4); |
| 953 | assert_eq!(lines[0].meter, "events_pull"); |
| 954 | assert_eq!(lines[0].quantity, 120.0); |
| 955 | assert!(lines.iter().any(|l| l.meter == "events_ratelimited")); |
| 956 | // By workspace, from the repository's store key; forks say none. |
| 957 | let by_repo = json!({ "data": { "viewer": { "accounts": [{ "artifactsEventsAdaptiveGroups": [ |
| 958 | { "count": 100, "dimensions": { "date": "2026-10-05", "eventType": "pull", "repositoryName": "acme--api" } }, |
| 959 | { "count": 20, "dimensions": { "date": "2026-10-05", "eventType": "push", "repositoryName": "acme--web" } }, |
| 960 | { "count": 7, "dimensions": { "date": "2026-10-05", "eventType": "fork", "repositoryName": "pulls--123" } }, |
| 961 | { "count": 3, "dimensions": { "date": "2026-10-05", "eventType": "serverError", "repositoryName": "beta--x" } } |
| 962 | ] }] } } }); |
| 963 | assert_eq!(artifacts_by_workspace(&by_repo, &BTreeMap::new()), vec![("2026-10-05".to_string(), "acme".to_string(), 120.0)]); |
| 964 | let none = BTreeMap::new(); |
| 965 | assert_eq!(workspace_of_store_key("Acme--api", &none), Some("acme".into())); |
| 966 | assert_eq!(workspace_of_store_key("pulls--9", &none), None); |
| 967 | assert_eq!(workspace_of_store_key("plain", &none), None); |
| 968 | let owners = BTreeMap::from([("9".to_string(), "Acme".to_string())]); |
| 969 | assert_eq!(workspace_of_store_key("pulls--9", &owners), Some("acme".into())); |
| 970 | assert_eq!(workspace_of_store_key("pulls--10", &owners), None); |
| 971 | let operations: f64 = lines |
| 972 | .iter() |
| 973 | .filter(|l| l.day == "2026-10-05" && ARTIFACTS_OPERATIONS.contains(&l.meter.as_str())) |
| 974 | .map(|l| l.quantity) |
| 975 | .sum(); |
| 976 | assert_eq!(operations, 150.0); |
| 977 | let failed = json!({ "data": null, "errors": [{ "message": "unknown field" }] }); |
| 978 | assert!(lines_from_artifacts(&failed).is_err()); |
| 979 | } |
| 980 | |
| 981 | #[test] |
| 982 | fn the_first_run_backfills_31_days_and_later_runs_restate_a_few() { |
| 983 | let now = g1t_contracts::time::parse_rfc3339("2026-10-06T04:17:00Z").unwrap(); |
| 984 | assert_eq!(window(None, now), ("2026-09-05".into(), "2026-10-06".into())); |
| 985 | assert_eq!(window(Some("2026-10-06"), now), ("2026-10-02".into(), "2026-10-06".into())); |
| 986 | // A week without a run: from the last day read. |
| 987 | assert_eq!(window(Some("2026-09-28"), now), ("2026-09-28".into(), "2026-10-06".into())); |
| 988 | // Never past the backfill. |
| 989 | assert_eq!(window(Some("2026-01-01"), now), ("2026-09-05".into(), "2026-10-06".into())); |
| 990 | assert_eq!(days_between("2026-09-29", "2026-10-02"), vec!["2026-09-29", "2026-09-30", "2026-10-01", "2026-10-02"]); |
| 991 | assert!(days_between("2026-10-02", "2026-10-01").is_empty()); |
| 992 | } |
| 993 | |
| 994 | #[test] |
| 995 | fn a_days_git_operations_are_its_count_less_the_next_days() { |
| 996 | let days: Vec<String> = ["2026-09-29", "2026-09-30", "2026-10-01", "2026-10-02"].iter().map(|d| d.to_string()).collect(); |
| 997 | let map = |pairs: &[(&str, u64)]| pairs.iter().map(|(w, n)| (w.to_string(), *n)).collect::<BTreeMap<_, _>>(); |
| 998 | // From each day to the end of its month. |
| 999 | let cumulative = vec![map(&[("acme", 30), ("beta", 4)]), map(&[("acme", 10)]), map(&[("acme", 7)]), map(&[("acme", 2)])]; |
| 1000 | let daily = daily_from_cumulative(&days, &cumulative); |
| 1001 | assert_eq!( |
| 1002 | daily, |
| 1003 | vec![ |
| 1004 | ("2026-09-29".into(), "acme".into(), 20), |
| 1005 | ("2026-09-29".into(), "beta".into(), 4), |
| 1006 | // The month's last day is its own. |
| 1007 | ("2026-09-30".into(), "acme".into(), 10), |
| 1008 | ("2026-10-01".into(), "acme".into(), 5), |
| 1009 | // Today so far. |
| 1010 | ("2026-10-02".into(), "acme".into(), 2), |
| 1011 | ] |
| 1012 | ); |
| 1013 | } |
| 1014 | |
| 1015 | #[test] |
| 1016 | fn raw_meters_are_weighted_by_the_repos_mapping() { |
| 1017 | // As repos' artifacts_usage answers: rows, and its operation_mapping. |
| 1018 | let body = json!({ |
| 1019 | "rows": [ |
| 1020 | { "day": "2026-10-15", "store": "g1t", "workspace": "Acme", "meter": "git.fetch", "count": 12, "bytes_in": 0, "bytes_out": 0 }, |
| 1021 | { "day": "2026-10-15", "store": "g1t", "workspace": "acme", "meter": "git.receive_pack", "count": 3, "bytes_in": 0, "bytes_out": 0 }, |
| 1022 | { "day": "2026-10-15", "store": "g1t", "workspace": "acme", "meter": "binding.read_blob", "count": 400, "bytes_in": 0, "bytes_out": 0 } |
| 1023 | ], |
| 1024 | "mapping": [ |
| 1025 | { "meter": "git.fetch", "cost_operations": 1, "billable_operations": 1 }, |
| 1026 | { "meter": "git.receive_pack", "cost_operations": 1, "billable_operations": 1 }, |
| 1027 | { "meter": "binding.read_blob", "cost_operations": 0, "billable_operations": 0 } |
| 1028 | ], |
| 1029 | "truncated": false |
| 1030 | }); |
| 1031 | let row = raw_usage_row(&body["rows"][0]).unwrap(); |
| 1032 | assert_eq!(row, ("2026-10-15".into(), "acme".into(), "git_fetch".into(), 12.0)); |
| 1033 | assert!(raw_usage_row(&json!({ "day": "2026-10-15", "count": 1 })).is_none()); |
| 1034 | let mut mapping = operation_mapping(&body); |
| 1035 | let raw: Vec<(String, f64)> = body["rows"].as_array().unwrap().iter().filter_map(raw_usage_row).map(|r| (r.2, r.3)).collect(); |
| 1036 | assert_eq!(weighted(&raw, &mapping, true), 15.0); |
| 1037 | assert_eq!(weighted(&raw, &mapping, false), 15.0); |
| 1038 | // Cloudflare turns out to bill binding reads: repos changes one row |
| 1039 | // (set_operation_mapping), and the bill g1t expects follows. |
| 1040 | mapping.insert("binding_read_blob".into(), (1.0, 0.0)); |
| 1041 | assert_eq!(weighted(&raw, &mapping, true), 415.0); |
| 1042 | assert_eq!(weighted(&raw, &mapping, false), 15.0); |
| 1043 | assert!(operation_mapping(&json!({})).is_empty()); |
| 1044 | } |
| 1045 | |
| 1046 | #[test] |
| 1047 | fn the_most_specific_rule_claims_a_line() { |
| 1048 | let rule = |product: &str, meter: &str, bucket: &str| Rule { |
| 1049 | product: product.into(), |
| 1050 | meter: meter.into(), |
| 1051 | bucket: bucket.into(), |
| 1052 | price_meter: None, |
| 1053 | own_meter: None, |
| 1054 | drift_percent: 10.0, |
| 1055 | }; |
| 1056 | let rules = vec![rule("workers", "*", "platform"), rule("workers", "workers_for_platforms", "deployments"), rule("artifacts", "*", "git")]; |
| 1057 | assert_eq!(classify(&rules, "workers", "workers_for_platforms_cpu_ms").unwrap().bucket, "deployments"); |
| 1058 | assert_eq!(classify(&rules, "workers", "workers_cpu_ms").unwrap().bucket, "platform"); |
| 1059 | assert_eq!(classify(&rules, "artifacts", "artifacts_operations").unwrap().bucket, "git"); |
| 1060 | assert!(classify(&rules, "browser_rendering", "browser_hours").is_none()); |
| 1061 | } |
| 1062 | } |