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