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.
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 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 | ||
| Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners) | 31 | use std::collections::{BTreeMap, BTreeSet}; |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 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 | ||
| Merge branch 'worktree-agent-a633ac0f7f66d419d' | 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 | ||
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 355 | /// The workspace a repository in the store belongs to: keys are |
| Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners) | 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> { | |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 360 | let (workspace, rest) = key.split_once("--")?; |
| Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners) | 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()) | |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 368 | } |
| 369 | ||
| Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners) | 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 | ||
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 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`). | |
| Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners) | 384 | pub(crate) fn artifacts_by_workspace(body: &Value, owners: &BTreeMap<String, String>) -> Vec<(String, String, f64)> { |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 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 | } | |
| Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners) | 392 | let Some(workspace) = g["dimensions"]["repositoryName"].as_str().and_then(|key| workspace_of_store_key(key, owners)) else { continue }; |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 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 | ||
| One operation mapping, owned by repos; billing reads it instead of keeping its own | 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() | |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 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?; | |
| Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners) | 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?; | |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 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 | } | |
| Merge branch 'worktree-agent-a633ac0f7f66d419d' | 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 | } | |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 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. | |
| Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners) | 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 | ||
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 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 | ||
| One operation mapping, owned by repos; billing reads it instead of keeping its own | 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. | |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 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 | ||
| One operation mapping, owned by repos; billing reads it instead of keeping its own | 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); | |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 723 | let mut by: BTreeMap<(String, String), Vec<(String, f64)>> = BTreeMap::new(); |
| One operation mapping, owned by repos; billing reads it instead of keeping its own | 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)); | |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 726 | } |
| 727 | for ((day, workspace), counts) in by { | |
| One operation mapping, owned by repos; billing reads it instead of keeping its own | 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 | } | |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 739 | } |
| 740 | } | |
| One operation mapping, owned by repos; billing reads it instead of keeping its own | 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 | } | |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 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 | ||
| Merge branch 'worktree-agent-a633ac0f7f66d419d' | 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 | ||
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 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 | ] }] } } }); | |
| Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners) | 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); | |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 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] | |
| One operation mapping, owned by repos; billing reads it instead of keeping its own | 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)); | |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 1033 | assert!(raw_usage_row(&json!({ "day": "2026-10-15", "count": 1 })).is_none()); |
| One operation mapping, owned by repos; billing reads it instead of keeping its own | 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()); | |
| Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily | 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 | } |
This file's history is long; its oldest lines are credited to the oldest commit read.