g1t/services/billing/src/costs.rs

1,062 lines51,745 bytesCodeBlame
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
31use std::collections::{BTreeMap, BTreeSet};
32
33use g1t_contracts::repos::{GitOperationsArgs, WorkspaceGitOperations};
34use g1t_contracts::time::rfc3339;
35use g1t_kit::now_ms;
36use serde::Deserialize;
37use serde_json::{Value, json};
38use worker::Result;
39
40use crate::Billing;
41use crate::keeper::Keeper;
42
43/// Where a cost line came from.
44pub(crate) const SOURCE_BILLABLE: &str = "billable_usage";
45pub(crate) const SOURCE_ARTIFACTS: &str = "artifacts_events";
46
47/// How far back the first run reads: the 31 days GraphQL keeps.
48pub(crate) const BACKFILL_DAYS: u64 = 31;
49/// How many recent days every later run reads again.
50pub(crate) const RESTATE_DAYS: u64 = 4;
51
52pub(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)]
56pub(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.
77pub(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.
102pub(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
120fn text<'a>(row: &'a Value, keys: &[&str]) -> &'a str {
121 keys.iter().find_map(|k| row[*k].as_str()).unwrap_or_default()
122}
123
124fn 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.
132pub(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.
165pub(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.
185pub(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.
194pub(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
203pub(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.
209pub(crate) fn lines_from_artifacts(body: &Value) -> std::result::Result<Vec<CostLine>, String> {
210 if let Some(errors) = body["errors"].as_array().filter(|e| !e.is_empty()) {
211 return Err(format!("Artifacts events failed: {}", Value::Array(errors.clone())));
212 }
213 let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"]
214 .as_array()
215 .cloned()
216 .unwrap_or_default();
217 Ok(aggregate(
218 groups
219 .iter()
220 .filter_map(|g| {
221 let day = g["dimensions"]["date"].as_str()?;
222 let kind = g["dimensions"]["eventType"].as_str()?;
223 if day.len() < 10 {
224 return None;
225 }
226 Some(CostLine {
227 day: day[..10].to_owned(),
228 source: SOURCE_ARTIFACTS,
229 product: "artifacts".to_owned(),
230 meter: format!("events_{}", slug(kind)),
231 unit: "events".to_owned(),
232 quantity: g["count"].as_f64().unwrap_or(0.0),
233 cost_usd: 0.0,
234 raw_name: format!("Artifacts / {kind}"),
235 })
236 })
237 .collect(),
238 ))
239}
240
241/// Where a cost line came from: AI Gateway's analytics, what it priced
242/// g1t's own provider traffic at.
243pub(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.
246pub(crate) const GATEWAY_PRODUCT: &str = "ai_gateway_requests";
247/// Meter suffixes of a model's token lines (no cost; for the drift).
248pub(crate) const GATEWAY_TOKENS: &str = "__tokens";
249pub(crate) const GATEWAY_CACHE_READ: &str = "__cache_read_tokens";
250pub(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.
253pub(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`).
260pub(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
270pub(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.
278pub(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)]
316pub(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).
329pub(crate) fn gateway_caveats(lines: &[(String, f64, f64)]) -> GatewayCaveats {
330 let mut out = GatewayCaveats::default();
331 let mut cost: BTreeMap<&str, f64> = BTreeMap::new();
332 let mut tokens: BTreeMap<&str, f64> = BTreeMap::new();
333 for (meter, quantity, cost_usd) in lines {
334 if let Some(model) = meter.strip_suffix(GATEWAY_TOKENS) {
335 *tokens.entry(model).or_default() += quantity;
336 } else if meter.ends_with(GATEWAY_CACHE_READ) {
337 out.cache_read_tokens += quantity;
338 } else if meter.ends_with(GATEWAY_CACHE_WRITE) {
339 out.cache_write_tokens += quantity;
340 } else {
341 *cost.entry(meter.as_str()).or_default() += cost_usd;
342 if meter.starts_with(GATEWAY_WHOLESALE) {
343 out.wholesale_usd += cost_usd;
344 }
345 }
346 }
347 for (model, used) in tokens {
348 if used > 0.0 && cost.get(model).copied().unwrap_or(0.0) <= 0.0 {
349 out.unpriced.push(model.to_owned());
350 }
351 }
352 out
353}
354
355/// The workspace a repository in the store belongs to: keys are
356/// `<workspace>--<repo>`; a pull request's working copy (`pulls--<id>`) is
357/// its repository's workspace's, from `owners` (repos' `pull_owners`), else
358/// no one's.
359pub(crate) fn workspace_of_store_key(key: &str, owners: &BTreeMap<String, String>) -> Option<String> {
360 let (workspace, rest) = key.split_once("--")?;
361 if workspace.is_empty() || rest.is_empty() {
362 return None;
363 }
364 if workspace == "pulls" {
365 return owners.get(rest).map(|w| w.to_lowercase());
366 }
367 Some(workspace.to_lowercase())
368}
369
370/// The pull request ids whose working copies Artifacts counted events for.
371pub(crate) fn pull_ids(body: &Value) -> Vec<String> {
372 let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"].as_array().cloned().unwrap_or_default();
373 let ids: BTreeSet<String> = groups
374 .iter()
375 .filter_map(|g| g["dimensions"]["repositoryName"].as_str()?.strip_prefix("pulls--").map(str::to_owned))
376 .filter(|id| !id.is_empty())
377 .collect();
378 ids.into_iter().collect()
379}
380
381/// Artifacts' billable operations per workspace and day, by repository
382/// name: how Cloudflare's own count shares out. The meter is
383/// `cloudflare_git`, which shares out the git bucket's cost (`margin`).
384pub(crate) fn artifacts_by_workspace(body: &Value, owners: &BTreeMap<String, String>) -> Vec<(String, String, f64)> {
385 let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"].as_array().cloned().unwrap_or_default();
386 let mut out: BTreeMap<(String, String), f64> = BTreeMap::new();
387 for g in &groups {
388 let (Some(day), Some(kind)) = (g["dimensions"]["date"].as_str(), g["dimensions"]["eventType"].as_str()) else { continue };
389 if day.len() < 10 || !ARTIFACTS_OPERATIONS.contains(&format!("events_{}", slug(kind)).as_str()) {
390 continue;
391 }
392 let Some(workspace) = g["dimensions"]["repositoryName"].as_str().and_then(|key| workspace_of_store_key(key, owners)) else { continue };
393 *out.entry((day[..10].to_owned(), workspace)).or_default() += g["count"].as_f64().unwrap_or(0.0);
394 }
395 out.into_iter().map(|((day, workspace), count)| (day, workspace, count)).collect()
396}
397
398/// Artifacts' event types that are operations it bills: what its
399/// pricing names (create, push, pull, clone) and the metrics list.
400/// Errors such as `rateLimited` are not.
401pub(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.
405pub(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.
422pub(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)]
436pub(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.
454pub(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.
462pub(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.
480pub(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.
486pub(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.
507pub(crate) fn raw_usage_row(row: &Value) -> Option<(String, String, String, f64)> {
508 let day = text(row, &["day", "date"]);
509 let workspace = text(row, &["namespace", "workspace"]);
510 let meter = text(row, &["meter", "kind", "event", "operation"]);
511 let count = number(row, &["count", "quantity", "operations", "value"])?;
512 (day.len() >= 10 && !workspace.is_empty() && !meter.is_empty())
513 .then(|| (day[..10].to_owned(), workspace.to_lowercase(), slug(meter), count))
514}
515
516/// The repos service's operation mapping, as `artifacts_usage` returns it
517/// beside the rows: for each raw meter (slugged), how many operations it
518/// is to Cloudflare (`cost_operations`) and to the customer
519/// (`billable_operations`). Repos owns this mapping
520/// (`set_operation_mapping`); billing only reads it.
521pub(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.
537pub(crate) fn weighted(raw: &[(String, f64)], mapping: &BTreeMap<String, (f64, f64)>, cost: bool) -> f64 {
538 raw.iter()
539 .map(|(meter, count)| count * mapping.get(meter).map_or(0.0, |(c, b)| if cost { *c } else { *b }))
540 .sum()
541}
542
543impl Billing {
544 /// Reads Cloudflare's bill for the days due (see `window`) into
545 /// `cost_lines`: the days read and how many lines. None without a
546 /// token. What could not be read is added to `problems`.
547 pub(crate) async fn read_cloudflare(&self, keeper: &Keeper, problems: &mut Vec<String>) -> Result<Option<(String, String, u32)>> {
548 if !keeper.can_read_bill() {
549 return Ok(None);
550 }
551 #[derive(Deserialize)]
552 struct Last {
553 day: Option<String>,
554 }
555 let last = self
556 .db
557 .prepare("SELECT MAX(day) AS day FROM cost_lines WHERE source = ?")
558 .bind(&[SOURCE_BILLABLE.into()])?
559 .first::<Last>(None)
560 .await?
561 .and_then(|l| l.day);
562 let (since, until) = window(last.as_deref(), now_ms());
563 let fetched_at = rfc3339(now_ms());
564 let mut written = 0;
565 match keeper.billable_usage_body(&since, &until).await.map_err(|e| e.to_string()).and_then(|b| lines_from_billable(&b)) {
566 Ok(lines) => written += self.upsert_lines(&lines, &fetched_at).await?,
567 Err(error) => problems.push(format!("Cloudflare's billable usage could not be read: {error}")),
568 }
569 match keeper.graphql(artifacts_variables(keeper.account(), &since, &until)).await.map_err(|e| e.to_string()) {
570 Ok(body) => match lines_from_artifacts(&body) {
571 Ok(lines) => {
572 written += self.upsert_lines(&lines, &fetched_at).await?;
573 let owners = self.pull_owners(&pull_ids(&body)).await;
574 self.keep_cloudflare_counts(&since, &until, &artifacts_by_workspace(&body, &owners), &fetched_at).await?;
575 }
576 Err(error) => problems.push(format!("Artifacts events could not be read: {error}")),
577 },
578 Err(error) => problems.push(format!("Artifacts events could not be read: {error}")),
579 }
580 // What AI Gateway priced g1t's own provider traffic at, each day:
581 // the total the ledger's model cost is checked against (`margin`).
582 if !keeper.gateway().is_empty() {
583 match keeper
584 .graphql_either(gateway_variables(keeper.account(), keeper.gateway(), &since, &until))
585 .await
586 .map_err(|e| e.to_string())
587 .and_then(|body| lines_from_gateway(&body))
588 {
589 Ok(lines) => {
590 // A day's models are read whole: one that is gone from a
591 // re-read day must not keep its old line.
592 self.db
593 .prepare("DELETE FROM cost_lines WHERE source = ?1 AND day >= ?2 AND day <= ?3")
594 .bind(&[SOURCE_GATEWAY.into(), since.as_str().into(), until.as_str().into()])?
595 .run()
596 .await?;
597 written += self.upsert_lines(&lines, &fetched_at).await?;
598 }
599 Err(error) => problems.push(format!("AI Gateway's analytics could not be read: {error}")),
600 }
601 }
602 Ok(Some((since, until, written)))
603 }
604
605 /// Cloudflare's own per-workspace counts for the days, replacing what
606 /// was kept for them.
607 /// The workspace of each pull request's working copy, from repos;
608 /// nothing while repos does not answer (those events stay no one's).
609 async fn pull_owners(&self, pulls: &[String]) -> BTreeMap<String, String> {
610 let Some(repos) = &self.repos else { return BTreeMap::new() };
611 if pulls.is_empty() {
612 return BTreeMap::new();
613 }
614 match g1t_kit::call::<_, Value>(repos, "pull_owners", &json!({ "pulls": pulls })).await {
615 Ok(body) => body["owners"]
616 .as_object()
617 .map(|o| o.iter().filter_map(|(k, v)| Some((k.clone(), v.as_str()?.to_owned()))).collect())
618 .unwrap_or_default(),
619 Err(error) => {
620 worker::console_error!("pull_owners: {error}");
621 BTreeMap::new()
622 }
623 }
624 }
625
626 async fn keep_cloudflare_counts(&self, since: &str, until: &str, counts: &[(String, String, f64)], fetched_at: &str) -> Result<()> {
627 self.db
628 .prepare("DELETE FROM own_counts WHERE meter = 'cloudflare_git' AND day >= ?1 AND day <= ?2")
629 .bind(&[since.into(), until.into()])?
630 .run()
631 .await?;
632 for chunk in counts.chunks(50) {
633 let mut statements = Vec::with_capacity(chunk.len());
634 for (day, workspace, count) in chunk {
635 statements.push(
636 self.db
637 .prepare("INSERT OR REPLACE INTO own_counts (day, meter, workspace, quantity, fetched_at) VALUES (?, 'cloudflare_git', ?, ?, ?)")
638 .bind(&[day.as_str().into(), workspace.as_str().into(), (*count).into(), fetched_at.into()])?,
639 );
640 }
641 self.db.batch(statements).await?;
642 }
643 Ok(())
644 }
645
646 /// Upserts lines, a day's line replacing what was read for it before.
647 async fn upsert_lines(&self, lines: &[CostLine], fetched_at: &str) -> Result<u32> {
648 for chunk in lines.chunks(50) {
649 let mut statements = Vec::with_capacity(chunk.len());
650 for line in chunk {
651 statements.push(
652 self.db
653 .prepare(
654 "INSERT INTO cost_lines (day, source, product, meter, unit, quantity, cost_usd, raw_name, fetched_at)
655 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)
656 ON CONFLICT (day, source, product, meter) DO UPDATE SET
657 unit = ?5, quantity = ?6, cost_usd = ?7, raw_name = ?8, fetched_at = ?9",
658 )
659 .bind(&[
660 line.day.as_str().into(),
661 line.source.into(),
662 line.product.as_str().into(),
663 line.meter.as_str().into(),
664 line.unit.as_str().into(),
665 line.quantity.into(),
666 line.cost_usd.into(),
667 line.raw_name.as_str().into(),
668 fetched_at.into(),
669 ])?,
670 );
671 }
672 self.db.batch(statements).await?;
673 }
674 Ok(lines.len() as u32)
675 }
676
677 /// The mapping from Cloudflare's meters to g1t's products.
678 pub(crate) async fn rules(&self) -> Result<Vec<Rule>> {
679 #[derive(Deserialize)]
680 struct Row {
681 product: String,
682 meter: String,
683 bucket: String,
684 price_meter: Option<String>,
685 own_meter: Option<String>,
686 drift_percent: f64,
687 }
688 Ok(self
689 .db
690 .prepare("SELECT product, meter, bucket, price_meter, own_meter, drift_percent FROM cost_map")
691 .all()
692 .await?
693 .results::<Row>()?
694 .into_iter()
695 .map(|r| Rule {
696 product: r.product,
697 meter: r.meter,
698 bucket: r.bucket,
699 price_meter: r.price_meter,
700 own_meter: r.own_meter,
701 drift_percent: r.drift_percent,
702 })
703 .collect())
704 }
705
706 /// g1t's own counts for the days, all from the repos service, which
707 /// owns the mapping from raw meters to operations (`operation_mapping`):
708 ///
709 /// - `git_operations`: what customers are charged for, as repos counts
710 /// it (`git_operations`, already through its mapping).
711 /// - `cost_operations`: what g1t expects Cloudflare to bill, from the
712 /// raw meters (`artifacts_usage`) and the mapping's cost column.
713 /// - `artifacts_<meter>`: each raw meter.
714 pub(crate) async fn count_own(&self, since: &str, until: &str) -> Result<()> {
715 let Some(repos) = &self.repos else { return Ok(()) };
716 let days = days_between(since, until);
717 let fetched_at = rfc3339(now_ms());
718 let mut rows: Vec<(String, String, String, f64)> = Vec::new();
719
720 // Raw meters and repos' mapping; skipped while repos does not answer.
721 if let Ok(body) = g1t_kit::call::<_, Value>(repos, "artifacts_usage", &json!({ "from": since, "to": until })).await {
722 let mapping = operation_mapping(&body);
723 let mut by: BTreeMap<(String, String), Vec<(String, f64)>> = BTreeMap::new();
724 for (day, workspace, meter, count) in body["rows"].as_array().map(|l| l.iter().filter_map(raw_usage_row).collect::<Vec<_>>()).unwrap_or_default() {
725 by.entry((day, workspace)).or_default().push((meter, count));
726 }
727 for ((day, workspace), counts) in by {
728 // Several stores (namespaces) can give the same meter.
729 let mut merged: BTreeMap<String, f64> = BTreeMap::new();
730 for (meter, count) in &counts {
731 *merged.entry(meter.clone()).or_default() += count;
732 }
733 for (meter, count) in &merged {
734 rows.push((day.clone(), format!("artifacts_{meter}"), workspace.clone(), *count));
735 }
736 if !mapping.is_empty() {
737 rows.push((day, "cost_operations".to_owned(), workspace, weighted(&counts, &mapping, true)));
738 }
739 }
740 }
741 // What customers are charged for, as repos counts it through its mapping.
742 let mut cumulative = Vec::with_capacity(days.len());
743 for day in &days {
744 let list: Vec<WorkspaceGitOperations> = g1t_kit::call(
745 repos,
746 "git_operations",
747 &GitOperationsArgs { month: day[..7].to_owned(), since: Some(format!("{day}T00")), namespace: None },
748 )
749 .await?;
750 cumulative.push(list.into_iter().map(|w| (w.namespace.to_lowercase(), w.operations)).collect::<BTreeMap<_, _>>());
751 }
752 for (day, workspace, count) in daily_from_cumulative(&days, &cumulative) {
753 rows.push((day, "git_operations".to_owned(), workspace, count as f64));
754 }
755
756 // Each day's counts replace what was there.
757 self.db
758 .prepare("DELETE FROM own_counts WHERE day >= ?1 AND day <= ?2 AND meter NOT LIKE 'cloudflare_%'")
759 .bind(&[since.into(), until.into()])?
760 .run()
761 .await?;
762 for chunk in rows.chunks(50) {
763 let mut statements = Vec::with_capacity(chunk.len());
764 for (day, meter, workspace, quantity) in chunk {
765 statements.push(
766 self.db
767 .prepare(
768 "INSERT INTO own_counts (day, meter, workspace, quantity, fetched_at) VALUES (?1, ?2, ?3, ?4, ?5)
769 ON CONFLICT (day, meter, workspace) DO UPDATE SET quantity = ?4, fetched_at = ?5",
770 )
771 .bind(&[day.as_str().into(), meter.as_str().into(), workspace.as_str().into(), (*quantity).into(), fetched_at.as_str().into()])?,
772 );
773 }
774 self.db.batch(statements).await?;
775 }
776 Ok(())
777 }
778}
779
780#[cfg(test)]
781mod tests {
782 use super::*;
783
784 /// AI Gateway's analytics in the shape of the schema
785 /// (`aiGatewayRequestsAdaptiveGroups`: `count`, `sum`, `dimensions`).
786 fn gateway_fixture() -> Value {
787 let group = |day: &str, provider: &str, model: &str, wholesale: u8, count: u64, cost: f64, tokens: (f64, f64, f64, f64)| {
788 json!({
789 "count": count,
790 "sum": { "cost": cost, "tokensIn": tokens.0, "tokensOut": tokens.1, "cacheReadTokens": tokens.2, "cacheWriteTokens": tokens.3 },
791 "dimensions": { "date": day, "provider": provider, "model": model, "wholesale": wholesale }
792 })
793 };
794 json!({ "data": { "viewer": { "accounts": [{ "aiGatewayRequestsAdaptiveGroups": [
795 group("2026-10-05", "anthropic", "claude-sonnet-5-5", 0, 120, 4.25, (900_000.0, 40_000.0, 3_000_000.0, 200_000.0)),
796 group("2026-10-05", "anthropic", "claude-haiku-4-5-20251001", 0, 300, 0.40, (400_000.0, 20_000.0, 0.0, 0.0)),
797 group("2026-10-06", "anthropic", "claude-new-1", 0, 12, 0.0, (80_000.0, 4_000.0, 0.0, 0.0)),
798 group("2026-10-06", "openai", "gpt-x", 1, 5, 0.10, (1_000.0, 100.0, 0.0, 0.0)),
799 ] }] } }, "errors": null })
800 }
801
802 #[test]
803 fn the_gateway_query_names_the_fields_its_schema_has() {
804 // As checked against Cloudflare's GraphQL schema (introspection of
805 // AccountAiGatewayRequestsAdaptiveGroups{,Sum,Dimensions,Filter}).
806 for field in ["aiGatewayRequestsAdaptiveGroups", "date_geq", "date_leq", "gateway: $gateway", "count", "cost", "tokensIn", "tokensOut", "cacheReadTokens", "cacheWriteTokens", "date", "provider", "model", "wholesale"] {
807 assert!(GATEWAY_QUERY.contains(field), "{field}");
808 }
809 let body = gateway_variables("acct", "g1t", "2026-10-01", "2026-10-07");
810 assert_eq!(body["variables"]["gateway"], "g1t");
811 }
812
813 #[test]
814 fn what_the_gateway_priced_is_kept_per_day_and_model_with_its_tokens() {
815 let lines = lines_from_gateway(&gateway_fixture()).unwrap();
816 let get = |day: &str, meter: &str| lines.iter().find(|l| l.day == day && l.meter == meter).unwrap_or_else(|| panic!("{day} {meter}"));
817 let sonnet = get("2026-10-05", "anthropic_claude_sonnet_5_5");
818 assert_eq!((sonnet.source, sonnet.product.as_str(), sonnet.quantity, sonnet.cost_usd), (SOURCE_GATEWAY, GATEWAY_PRODUCT, 120.0, 4.25));
819 assert_eq!(get("2026-10-05", "anthropic_claude_sonnet_5_5__tokens").quantity, 940_000.0);
820 assert_eq!(get("2026-10-05", "anthropic_claude_sonnet_5_5__cache_read_tokens").quantity, 3_000_000.0);
821 assert_eq!(get("2026-10-05", "anthropic_claude_sonnet_5_5__cache_write_tokens").cost_usd, 0.0);
822 assert_eq!(get("2026-10-06", "wholesale__openai_gpt_x").cost_usd, 0.10);
823 // Only the request lines carry cost: the day's total is the gateway's.
824 let day5: f64 = lines.iter().filter(|l| l.day == "2026-10-05").map(|l| l.cost_usd).sum();
825 assert!((day5 - 4.65).abs() < 1e-9);
826 // Errors are a problem for the run, not lines.
827 assert!(lines_from_gateway(&json!({ "errors": [{ "message": "not authorized for that account" }] })).is_err());
828 }
829
830 #[test]
831 fn the_gateway_lines_say_when_its_cost_cannot_be_the_providers() {
832 let rows: Vec<(String, f64, f64)> =
833 lines_from_gateway(&gateway_fixture()).unwrap().into_iter().map(|l| (l.meter, l.quantity, l.cost_usd)).collect();
834 let caveats = gateway_caveats(&rows);
835 assert_eq!(caveats.unpriced, vec!["anthropic_claude_new_1".to_owned()]);
836 assert_eq!((caveats.cache_read_tokens, caveats.cache_write_tokens), (3_000_000.0, 200_000.0));
837 assert!((caveats.wholesale_usd - 0.10).abs() < 1e-9);
838 // A model priced on one day and not another is priced.
839 let mixed = vec![
840 ("m".to_owned(), 3.0, 0.5),
841 ("m__tokens".to_owned(), 10.0, 0.0),
842 ("m".to_owned(), 3.0, 0.0),
843 ("m__tokens".to_owned(), 10.0, 0.0),
844 ];
845 assert!(gateway_caveats(&mixed).unpriced.is_empty());
846 }
847
848 /// Billable usage as Cloudflare answered g1t on 2026-10-06 (two days,
849 /// trimmed), plus rows past the included amounts, which cost money.
850 fn billable_fixture() -> Value {
851 json!({
852 "success": true,
853 "errors": [],
854 "result": [
855 {
856 "ChargePeriodStart": "2026-10-03T00:00:00Z",
857 "ChargePeriodEnd": "2026-10-04T00:00:00Z",
858 "ServiceFamilyName": "Containers",
859 "ServiceName": "Container Memory, per GiB-Second (First 25 GiB-hours included)",
860 "PricingUnit": "Count",
861 "PricingQuantity": "14000",
862 "ContractedCost": 0, "BilledCost": 0, "ListCost": 0
863 },
864 {
865 "ChargePeriodStart": "2026-10-03T00:00:00Z",
866 "ServiceFamilyName": "D1",
867 "ServiceName": "D1 - Rows Read (first 25 billion included)",
868 "PricingUnit": "Count",
869 "PricingQuantity": 1200000,
870 "BilledCost": 0
871 },
872 {
873 "ChargePeriodStart": "2026-10-15T00:00:00Z",
874 "ServiceFamilyName": "Artifacts",
875 "ServiceName": "Artifacts Operations (First 10,000 included)",
876 "PricingUnit": "Count",
877 "PricingQuantity": "30000",
878 "BilledCost": "3.00",
879 "ListCost": 4.5
880 },
881 {
882 "ChargePeriodStart": "2026-10-15T00:00:00Z",
883 "ServiceFamilyName": "Artifacts",
884 "ServiceName": "Artifacts Operations (First 10,000 included)",
885 "PricingUnit": "Count",
886 "PricingQuantity": "10000",
887 "BilledCost": "1.50"
888 },
889 {
890 "ChargePeriodStart": "2026-10-15T00:00:00Z",
891 "ServiceFamilyName": "Workers",
892 "ServiceName": "Workers for Platforms Requests (First 20M are included)",
893 "PricingUnit": "Count",
894 "PricingQuantity": 25000000,
895 "ListCost": 1.5
896 },
897 { "ServiceName": "No day, no line" }
898 ]
899 })
900 }
901
902 #[test]
903 fn names_become_products_and_meters() {
904 assert_eq!(slug("Workers for Platforms CPU ms (First 60M ms are included)"), "workers_for_platforms_cpu_ms");
905 assert_eq!(product_and_meter("D1", "D1 - Rows Read (first 25 billion included)"), ("d1".into(), "d1_rows_read".into()));
906 assert_eq!(product_and_meter("Workers KV", "KV Read Operations (First 10M is included)"), ("workers_kv".into(), "kv_read_operations".into()));
907 assert_eq!(product_and_meter("", "Containers / Container vCPU"), ("containers".into(), "container_vcpu".into()));
908 assert_eq!(product_and_meter("Email", ""), ("email".into(), "usage".into()));
909 }
910
911 #[test]
912 fn billable_usage_is_parsed_and_rows_of_one_meter_are_added() {
913 let lines = lines_from_billable(&billable_fixture()).unwrap();
914 assert_eq!(lines.len(), 4, "{lines:#?}");
915 let artifacts = lines.iter().find(|l| l.product == "artifacts").unwrap();
916 // Two rows of the same day and meter: one line, added up.
917 assert_eq!(artifacts.meter, "artifacts_operations");
918 assert_eq!(artifacts.quantity, 40_000.0);
919 assert!((artifacts.cost_usd - 4.5).abs() < 1e-9, "billed, not list: {}", artifacts.cost_usd);
920 let wfp = lines.iter().find(|l| l.meter.starts_with("workers_for_platforms")).unwrap();
921 assert_eq!(wfp.product, "workers");
922 // No billed cost: the list cost.
923 assert_eq!(wfp.cost_usd, 1.5);
924 let memory = lines.iter().find(|l| l.product == "containers").unwrap();
925 assert_eq!((memory.day.as_str(), memory.quantity, memory.cost_usd), ("2026-10-03", 14_000.0, 0.0));
926 assert!(lines_from_billable(&json!({ "success": false, "errors": [{ "code": 10000 }] })).is_err());
927 }
928
929 #[test]
930 fn reading_the_same_days_twice_gives_the_same_lines() {
931 // Idempotent: the same answer aggregates to the same keys and
932 // amounts, so the upsert replaces rather than adds.
933 let once = lines_from_billable(&billable_fixture()).unwrap();
934 let twice = lines_from_billable(&billable_fixture()).unwrap();
935 assert_eq!(once, twice);
936 let again = aggregate(once.clone());
937 assert_eq!(again, once);
938 }
939
940 #[test]
941 fn artifacts_events_are_counted_by_type_and_day() {
942 let body = json!({
943 "data": { "viewer": { "accounts": [{ "artifactsEventsAdaptiveGroups": [
944 { "count": 120, "dimensions": { "date": "2026-10-05", "eventType": "pull" } },
945 { "count": 30, "dimensions": { "date": "2026-10-05", "eventType": "push" } },
946 { "count": 2, "dimensions": { "date": "2026-10-05", "eventType": "rateLimited" } },
947 { "count": 5, "dimensions": { "date": "2026-10-06", "eventType": "pull" } }
948 ] }] } },
949 "errors": null
950 });
951 let lines = lines_from_artifacts(&body).unwrap();
952 assert_eq!(lines.len(), 4);
953 assert_eq!(lines[0].meter, "events_pull");
954 assert_eq!(lines[0].quantity, 120.0);
955 assert!(lines.iter().any(|l| l.meter == "events_ratelimited"));
956 // By workspace, from the repository's store key; forks say none.
957 let by_repo = json!({ "data": { "viewer": { "accounts": [{ "artifactsEventsAdaptiveGroups": [
958 { "count": 100, "dimensions": { "date": "2026-10-05", "eventType": "pull", "repositoryName": "acme--api" } },
959 { "count": 20, "dimensions": { "date": "2026-10-05", "eventType": "push", "repositoryName": "acme--web" } },
960 { "count": 7, "dimensions": { "date": "2026-10-05", "eventType": "fork", "repositoryName": "pulls--123" } },
961 { "count": 3, "dimensions": { "date": "2026-10-05", "eventType": "serverError", "repositoryName": "beta--x" } }
962 ] }] } } });
963 assert_eq!(artifacts_by_workspace(&by_repo, &BTreeMap::new()), vec![("2026-10-05".to_string(), "acme".to_string(), 120.0)]);
964 let none = BTreeMap::new();
965 assert_eq!(workspace_of_store_key("Acme--api", &none), Some("acme".into()));
966 assert_eq!(workspace_of_store_key("pulls--9", &none), None);
967 assert_eq!(workspace_of_store_key("plain", &none), None);
968 let owners = BTreeMap::from([("9".to_string(), "Acme".to_string())]);
969 assert_eq!(workspace_of_store_key("pulls--9", &owners), Some("acme".into()));
970 assert_eq!(workspace_of_store_key("pulls--10", &owners), None);
971 let operations: f64 = lines
972 .iter()
973 .filter(|l| l.day == "2026-10-05" && ARTIFACTS_OPERATIONS.contains(&l.meter.as_str()))
974 .map(|l| l.quantity)
975 .sum();
976 assert_eq!(operations, 150.0);
977 let failed = json!({ "data": null, "errors": [{ "message": "unknown field" }] });
978 assert!(lines_from_artifacts(&failed).is_err());
979 }
980
981 #[test]
982 fn the_first_run_backfills_31_days_and_later_runs_restate_a_few() {
983 let now = g1t_contracts::time::parse_rfc3339("2026-10-06T04:17:00Z").unwrap();
984 assert_eq!(window(None, now), ("2026-09-05".into(), "2026-10-06".into()));
985 assert_eq!(window(Some("2026-10-06"), now), ("2026-10-02".into(), "2026-10-06".into()));
986 // A week without a run: from the last day read.
987 assert_eq!(window(Some("2026-09-28"), now), ("2026-09-28".into(), "2026-10-06".into()));
988 // Never past the backfill.
989 assert_eq!(window(Some("2026-01-01"), now), ("2026-09-05".into(), "2026-10-06".into()));
990 assert_eq!(days_between("2026-09-29", "2026-10-02"), vec!["2026-09-29", "2026-09-30", "2026-10-01", "2026-10-02"]);
991 assert!(days_between("2026-10-02", "2026-10-01").is_empty());
992 }
993
994 #[test]
995 fn a_days_git_operations_are_its_count_less_the_next_days() {
996 let days: Vec<String> = ["2026-09-29", "2026-09-30", "2026-10-01", "2026-10-02"].iter().map(|d| d.to_string()).collect();
997 let map = |pairs: &[(&str, u64)]| pairs.iter().map(|(w, n)| (w.to_string(), *n)).collect::<BTreeMap<_, _>>();
998 // From each day to the end of its month.
999 let cumulative = vec![map(&[("acme", 30), ("beta", 4)]), map(&[("acme", 10)]), map(&[("acme", 7)]), map(&[("acme", 2)])];
1000 let daily = daily_from_cumulative(&days, &cumulative);
1001 assert_eq!(
1002 daily,
1003 vec![
1004 ("2026-09-29".into(), "acme".into(), 20),
1005 ("2026-09-29".into(), "beta".into(), 4),
1006 // The month's last day is its own.
1007 ("2026-09-30".into(), "acme".into(), 10),
1008 ("2026-10-01".into(), "acme".into(), 5),
1009 // Today so far.
1010 ("2026-10-02".into(), "acme".into(), 2),
1011 ]
1012 );
1013 }
1014
1015 #[test]
1016 fn raw_meters_are_weighted_by_the_repos_mapping() {
1017 // As repos' artifacts_usage answers: rows, and its operation_mapping.
1018 let body = json!({
1019 "rows": [
1020 { "day": "2026-10-15", "store": "g1t", "workspace": "Acme", "meter": "git.fetch", "count": 12, "bytes_in": 0, "bytes_out": 0 },
1021 { "day": "2026-10-15", "store": "g1t", "workspace": "acme", "meter": "git.receive_pack", "count": 3, "bytes_in": 0, "bytes_out": 0 },
1022 { "day": "2026-10-15", "store": "g1t", "workspace": "acme", "meter": "binding.read_blob", "count": 400, "bytes_in": 0, "bytes_out": 0 }
1023 ],
1024 "mapping": [
1025 { "meter": "git.fetch", "cost_operations": 1, "billable_operations": 1 },
1026 { "meter": "git.receive_pack", "cost_operations": 1, "billable_operations": 1 },
1027 { "meter": "binding.read_blob", "cost_operations": 0, "billable_operations": 0 }
1028 ],
1029 "truncated": false
1030 });
1031 let row = raw_usage_row(&body["rows"][0]).unwrap();
1032 assert_eq!(row, ("2026-10-15".into(), "acme".into(), "git_fetch".into(), 12.0));
1033 assert!(raw_usage_row(&json!({ "day": "2026-10-15", "count": 1 })).is_none());
1034 let mut mapping = operation_mapping(&body);
1035 let raw: Vec<(String, f64)> = body["rows"].as_array().unwrap().iter().filter_map(raw_usage_row).map(|r| (r.2, r.3)).collect();
1036 assert_eq!(weighted(&raw, &mapping, true), 15.0);
1037 assert_eq!(weighted(&raw, &mapping, false), 15.0);
1038 // Cloudflare turns out to bill binding reads: repos changes one row
1039 // (set_operation_mapping), and the bill g1t expects follows.
1040 mapping.insert("binding_read_blob".into(), (1.0, 0.0));
1041 assert_eq!(weighted(&raw, &mapping, true), 415.0);
1042 assert_eq!(weighted(&raw, &mapping, false), 15.0);
1043 assert!(operation_mapping(&json!({})).is_empty());
1044 }
1045
1046 #[test]
1047 fn the_most_specific_rule_claims_a_line() {
1048 let rule = |product: &str, meter: &str, bucket: &str| Rule {
1049 product: product.into(),
1050 meter: meter.into(),
1051 bucket: bucket.into(),
1052 price_meter: None,
1053 own_meter: None,
1054 drift_percent: 10.0,
1055 };
1056 let rules = vec![rule("workers", "*", "platform"), rule("workers", "workers_for_platforms", "deployments"), rule("artifacts", "*", "git")];
1057 assert_eq!(classify(&rules, "workers", "workers_for_platforms_cpu_ms").unwrap().bucket, "deployments");
1058 assert_eq!(classify(&rules, "workers", "workers_cpu_ms").unwrap().bucket, "platform");
1059 assert_eq!(classify(&rules, "artifacts", "artifacts_operations").unwrap().bucket, "git");
1060 assert!(classify(&rules, "browser_rendering", "browser_hours").is_none());
1061 }
1062}