Skip to content
1,126 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
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) and the tokens it priced; once for the requests
258/// g1t's own provider keys paid (`wholesale: 0`) and once for those
259/// Cloudflare billed itself (`wholesale: 1`). Field names checked against
260/// the schema (`AccountAiGatewayRequestsAdaptiveGroupsSum`, `…Dimensions`
261/// and `…Filter_InputObject`).
262///
263/// `wholesale` is a filter here, never a dimension: grouped by it,
264/// Cloudflare answers no rows at all and no error. Seen 2026-10-08: with
265/// `wholesale` in `dimensions` the query returned nothing for 562 requests
266/// that cost $11.11, which it returns without it.
267pub(crate) const GATEWAY_QUERY: &str = "query ($account: String!, $gateway: String!, $since: Date!, $until: Date!) {
268 viewer { accounts(filter: { accountTag: $account }) {
269 own: aiGatewayRequestsAdaptiveGroups(limit: 10000, filter: { date_geq: $since, date_leq: $until, gateway: $gateway, wholesale: 0 }) {
270 count
271 sum { cost tokensIn tokensOut cacheReadTokens cacheWriteTokens }
272 dimensions { date provider model }
273 }
274 wholesale: aiGatewayRequestsAdaptiveGroups(limit: 10000, filter: { date_geq: $since, date_leq: $until, gateway: $gateway, wholesale: 1 }) {
275 count
276 sum { cost tokensIn tokensOut cacheReadTokens cacheWriteTokens }
277 dimensions { date provider model }
278 }
279 } }
280}";
281
282pub(crate) fn gateway_variables(account: &str, gateway: &str, since: &str, until: &str) -> Value {
283 json!({ "query": GATEWAY_QUERY, "variables": { "account": account, "gateway": gateway, "since": since, "until": until } })
284}
285
286/// AI Gateway's analytics as lines: per day and model, the requests at
287/// what the gateway priced them (meter `<provider>_<model>`), and their
288/// tokens, cache reads and cache writes at no cost (`…__tokens`, …).
289/// Errors when GraphQL does.
290pub(crate) fn lines_from_gateway(body: &Value) -> std::result::Result<Vec<CostLine>, String> {
291 if let Some(errors) = body["errors"].as_array().filter(|e| !e.is_empty()) {
292 return Err(format!("AI Gateway analytics failed: {}", Value::Array(errors.clone())));
293 }
294 let account = &body["data"]["viewer"]["accounts"][0];
295 let groups = |alias: &str| account[alias].as_array().cloned().unwrap_or_default();
296 let mut lines = Vec::new();
297 for (wholesale, g) in groups("own").into_iter().map(|g| (false, g)).chain(groups("wholesale").into_iter().map(|g| (true, g))) {
298 let d = &g["dimensions"];
299 let Some(day) = d["date"].as_str().filter(|day| day.len() >= 10) else { continue };
300 let provider = d["provider"].as_str().unwrap_or("unknown");
301 let model = d["model"].as_str().unwrap_or("unknown");
302 let name = format!("{}{}", if wholesale { GATEWAY_WHOLESALE } else { "" }, slug(&format!("{provider} {model}")));
303 let sum = |key: &str| g["sum"][key].as_f64().unwrap_or(0.0);
304 let line = |meter: String, unit: &str, quantity: f64, cost_usd: f64| CostLine {
305 day: day[..10].to_owned(),
306 source: SOURCE_GATEWAY,
307 product: GATEWAY_PRODUCT.to_owned(),
308 meter,
309 unit: unit.to_owned(),
310 quantity,
311 cost_usd,
312 raw_name: format!("AI Gateway / {provider} / {model}{}", if wholesale { " (billed by Cloudflare)" } else { "" }),
313 };
314 lines.push(line(name.clone(), "requests", g["count"].as_f64().unwrap_or(0.0), sum("cost").max(0.0)));
315 lines.push(line(format!("{name}{GATEWAY_TOKENS}"), "tokens", sum("tokensIn") + sum("tokensOut"), 0.0));
316 for (suffix, key) in [(GATEWAY_CACHE_READ, "cacheReadTokens"), (GATEWAY_CACHE_WRITE, "cacheWriteTokens")] {
317 if sum(key) > 0.0 {
318 lines.push(line(format!("{name}{suffix}"), "tokens", sum(key), 0.0));
319 }
320 }
321 }
322 Ok(aggregate(lines))
323}
324
325/// What the gateway's lines over a window say about whether its cost can
326/// be taken as what the providers bill g1t.
327#[derive(Clone, Debug, Default, PartialEq)]
328pub(crate) struct GatewayCaveats {
329 /// Models the gateway put no price on although they used tokens.
330 pub unpriced: Vec<String>,
331 pub cache_read_tokens: f64,
332 pub cache_write_tokens: f64,
333 /// What Cloudflare billed itself (unified billing), in dollars: on its
334 /// bill too, so not a cost paid to a provider.
335 pub wholesale_usd: f64,
336 /// Runs settled with the gateway's figure short (see `keeper::settled_cost`).
337 pub short_runs: u32,
338 /// Requests the gateway logged, priced or not.
339 pub requests: f64,
340}
341
342/// What the last read of AI Gateway's analytics found, so the models drift
343/// can say why the gateway's side is empty rather than guess.
344#[derive(Clone, Debug, Default, PartialEq)]
345pub(crate) enum GatewayRead {
346 /// Not read: no `AI_GATEWAY_ID`, or no token to read it with.
347 #[default]
348 NotRead,
349 /// GraphQL refused it, or did not answer.
350 Failed(String),
351 /// It answered with requests.
352 Rows,
353 /// It answered with no rows. `visible`: whether the token it was read
354 /// with can see the gateway (`GET …/ai-gateway/gateways/{id}`: 403
355 /// without AI Gateway Read, 404 for an id that is not there), to tell
356 /// that from a gateway nothing went through; None when that could not
357 /// be told.
358 Empty { visible: Option<bool> },
359}
360
361/// The caveats in AI Gateway's lines (any days, any order).
362pub(crate) fn gateway_caveats(lines: &[(String, f64, f64)]) -> GatewayCaveats {
363 let mut out = GatewayCaveats::default();
364 let mut cost: BTreeMap<&str, f64> = BTreeMap::new();
365 let mut tokens: BTreeMap<&str, f64> = BTreeMap::new();
366 for (meter, quantity, cost_usd) in lines {
367 if let Some(model) = meter.strip_suffix(GATEWAY_TOKENS) {
368 *tokens.entry(model).or_default() += quantity;
369 } else if meter.ends_with(GATEWAY_CACHE_READ) {
370 out.cache_read_tokens += quantity;
371 } else if meter.ends_with(GATEWAY_CACHE_WRITE) {
372 out.cache_write_tokens += quantity;
373 } else {
374 out.requests += quantity;
375 *cost.entry(meter.as_str()).or_default() += cost_usd;
376 if meter.starts_with(GATEWAY_WHOLESALE) {
377 out.wholesale_usd += cost_usd;
378 }
379 }
380 }
381 for (model, used) in tokens {
382 if used > 0.0 && cost.get(model).copied().unwrap_or(0.0) <= 0.0 {
383 out.unpriced.push(model.to_owned());
384 }
385 }
386 out
387}
388
389/// The workspace a repository in the store belongs to: keys are
390/// `<workspace>--<repo>`; a pull request's working copy (`pulls--<id>`) is
391/// its repository's workspace's, from `owners` (repos' `pull_owners`), else
392/// no one's.
393pub(crate) fn workspace_of_store_key(key: &str, owners: &BTreeMap<String, String>) -> Option<String> {
394 let (workspace, rest) = key.split_once("--")?;
395 if workspace.is_empty() || rest.is_empty() {
396 return None;
397 }
398 if workspace == "pulls" {
399 return owners.get(rest).map(|w| w.to_lowercase());
400 }
401 Some(workspace.to_lowercase())
402}
403
404/// The pull request ids whose working copies Artifacts counted events for.
405pub(crate) fn pull_ids(body: &Value) -> Vec<String> {
406 let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"].as_array().cloned().unwrap_or_default();
407 let ids: BTreeSet<String> = groups
408 .iter()
409 .filter_map(|g| g["dimensions"]["repositoryName"].as_str()?.strip_prefix("pulls--").map(str::to_owned))
410 .filter(|id| !id.is_empty())
411 .collect();
412 ids.into_iter().collect()
413}
414
415/// Artifacts' billable operations per workspace and day, by repository
416/// name: how Cloudflare's own count shares out. The meter is
417/// `cloudflare_git`, which shares out the git bucket's cost (`margin`).
418pub(crate) fn artifacts_by_workspace(body: &Value, owners: &BTreeMap<String, String>) -> Vec<(String, String, f64)> {
419 let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"].as_array().cloned().unwrap_or_default();
420 let mut out: BTreeMap<(String, String), f64> = BTreeMap::new();
421 for g in &groups {
422 let (Some(day), Some(kind)) = (g["dimensions"]["date"].as_str(), g["dimensions"]["eventType"].as_str()) else { continue };
423 if day.len() < 10 || !ARTIFACTS_OPERATIONS.contains(&format!("events_{}", slug(kind)).as_str()) {
424 continue;
425 }
426 let Some(workspace) = g["dimensions"]["repositoryName"].as_str().and_then(|key| workspace_of_store_key(key, owners)) else { continue };
427 *out.entry((day[..10].to_owned(), workspace)).or_default() += g["count"].as_f64().unwrap_or(0.0);
428 }
429 out.into_iter().map(|((day, workspace), count)| (day, workspace, count)).collect()
430}
431
432/// Artifacts' event types that are operations it bills: what its
433/// pricing names (create, push, pull, clone) and the metrics list.
434/// Errors such as `rateLimited` are not.
435pub(crate) const ARTIFACTS_OPERATIONS: [&str; 5] = ["events_create", "events_fork", "events_push", "events_pull", "events_delete"];
436
437/// The days to read: the 31 before today on the first run (`last` is
438/// None), else the last few, never further back than the backfill.
439pub(crate) fn window(last_fetched_day: Option<&str>, now_ms: u64) -> (String, String) {
440 let day = |ms: u64| g1t_contracts::time::rfc3339(ms)[..10].to_owned();
441 let today = day(now_ms);
442 let since = match last_fetched_day {
443 None => day(now_ms.saturating_sub(BACKFILL_DAYS * DAY_MS)),
444 Some(last) => {
445 let restate = day(now_ms.saturating_sub(RESTATE_DAYS * DAY_MS));
446 let floor = day(now_ms.saturating_sub(BACKFILL_DAYS * DAY_MS));
447 // A gap since the last run is read too, up to the backfill.
448 let from = if last < restate.as_str() { last.to_owned() } else { restate };
449 if from < floor { floor } else { from }
450 }
451 };
452 (since, today)
453}
454
455/// Every day from `since` to `until`, inclusive.
456pub(crate) fn days_between(since: &str, until: &str) -> Vec<String> {
457 let start = g1t_contracts::time::parse_rfc3339(&format!("{since}T00:00:00Z"));
458 let end = g1t_contracts::time::parse_rfc3339(&format!("{until}T00:00:00Z"));
459 match (start, end) {
460 (Some(start), Some(end)) if start <= end => (0..=((end - start) / DAY_MS))
461 .map(|n| g1t_contracts::time::rfc3339(start + n * DAY_MS)[..10].to_owned())
462 .collect(),
463 _ => Vec::new(),
464 }
465}
466
467/// A rule from `cost_map`: which Cloudflare lines feed which of g1t's
468/// products.
469#[derive(Clone, Debug, PartialEq, serde::Deserialize)]
470pub(crate) struct Rule {
471 /// Cloudflare's product, as slugged here.
472 pub product: String,
473 /// A meter prefix, or `*` for any meter of the product.
474 pub meter: String,
475 /// g1t's product it is a cost of: `sandboxes`, `git`, `platform`, …
476 pub bucket: String,
477 /// The price book meter whose cost it measures, if any.
478 pub price_meter: Option<String>,
479 /// g1t's own count of the same units, to compare quantities.
480 pub own_meter: Option<String>,
481 /// How far g1t's count may be from Cloudflare's before it is drift.
482 pub drift_percent: f64,
483}
484
485/// The rule for a line: its product's rule with the longest matching
486/// meter prefix, `*` last. None means no one decided what pays for it:
487/// a leak until someone does.
488pub(crate) fn classify<'a>(rules: &'a [Rule], product: &str, meter: &str) -> Option<&'a Rule> {
489 rules
490 .iter()
491 .filter(|r| r.product == product && (r.meter == "*" || meter.starts_with(r.meter.as_str())))
492 .max_by_key(|r| if r.meter == "*" { 0 } else { r.meter.len() + 1 })
493}
494
495/// What a bucket is called in sudo.
496pub(crate) fn bucket_title(bucket: &str) -> String {
497 match bucket {
498 "sandboxes" => "Sandboxes and builds".into(),
499 "deployments" => "Deployments".into(),
500 "git" => "Git operations".into(),
501 "repo_storage" => "Repository storage".into(),
502 "actions_cache" => "Actions cache".into(),
503 "embeddings" => "Search embeddings".into(),
504 "security" => "Security scans".into(),
505 "domains" => "Custom domains".into(),
506 "models" => "Models".into(),
507 "platform" => "Running g1t (paid by the plan)".into(),
508 UNMAPPED => "Not mapped".into(),
509 other => other.replace('_', " "),
510 }
511}
512
513/// The bucket of a line no rule claims.
514pub(crate) const UNMAPPED: &str = "unmapped";
515
516/// A day's count per workspace from counts "from this day to the end of
517/// its month" (what the repos service's `git_operations` answers with a
518/// `since`): each day's is its own less the next day's, within a month.
519/// The last day (today, so far) is its own.
520pub(crate) fn daily_from_cumulative(days: &[String], cumulative: &[BTreeMap<String, u64>]) -> Vec<(String, String, u64)> {
521 let mut out = Vec::new();
522 for (index, day) in days.iter().enumerate() {
523 let Some(today) = cumulative.get(index) else { break };
524 let next = days
525 .get(index + 1)
526 .filter(|next| next[..7] == day[..7])
527 .and_then(|_| cumulative.get(index + 1));
528 for (workspace, count) in today {
529 let later = next.and_then(|n| n.get(workspace)).copied().unwrap_or(0);
530 let own = count.saturating_sub(later);
531 if own > 0 {
532 out.push((day.clone(), workspace.clone(), own));
533 }
534 }
535 }
536 out
537}
538
539/// One row of the repos service's `artifacts_usage`, read leniently while
540/// its shape settles: a day, a workspace, a raw meter and a count.
541pub(crate) fn raw_usage_row(row: &Value) -> Option<(String, String, String, f64)> {
542 let day = text(row, &["day", "date"]);
543 let workspace = text(row, &["namespace", "workspace"]);
544 let meter = text(row, &["meter", "kind", "event", "operation"]);
545 let count = number(row, &["count", "quantity", "operations", "value"])?;
546 (day.len() >= 10 && !workspace.is_empty() && !meter.is_empty())
547 .then(|| (day[..10].to_owned(), workspace.to_lowercase(), slug(meter), count))
548}
549
550/// The repos service's operation mapping, as `artifacts_usage` returns it
551/// beside the rows: for each raw meter (slugged), how many operations it
552/// is to Cloudflare (`cost_operations`) and to the customer
553/// (`billable_operations`). Repos owns this mapping
554/// (`set_operation_mapping`); billing only reads it.
555pub(crate) fn operation_mapping(body: &Value) -> BTreeMap<String, (f64, f64)> {
556 body["mapping"]
557 .as_array()
558 .map(|rows| {
559 rows.iter()
560 .filter_map(|r| {
561 let meter = r["meter"].as_str()?;
562 Some((slug(meter), (r["cost_operations"].as_f64().unwrap_or(0.0), r["billable_operations"].as_f64().unwrap_or(0.0))))
563 })
564 .collect()
565 })
566 .unwrap_or_default()
567}
568
569/// Raw counts weighted by one column of the mapping: what Cloudflare
570/// should count (`cost`), or what customers are charged for.
571pub(crate) fn weighted(raw: &[(String, f64)], mapping: &BTreeMap<String, (f64, f64)>, cost: bool) -> f64 {
572 raw.iter()
573 .map(|(meter, count)| count * mapping.get(meter).map_or(0.0, |(c, b)| if cost { *c } else { *b }))
574 .sum()
575}
576
577impl Billing {
578 /// Reads Cloudflare's bill for the days due (see `window`) into
579 /// `cost_lines`: the days read and how many lines. None without a
580 /// token. What could not be read is added to `problems`, and what AI
581 /// Gateway's analytics answered to `gateway`.
582 pub(crate) async fn read_cloudflare(
583 &self,
584 keeper: &Keeper,
585 problems: &mut Vec<String>,
586 gateway: &mut GatewayRead,
587 ) -> Result<Option<(String, String, u32)>> {
588 if !keeper.can_read_bill() {
589 return Ok(None);
590 }
591 #[derive(Deserialize)]
592 struct Last {
593 day: Option<String>,
594 }
595 let last = self
596 .db
597 .prepare("SELECT MAX(day) AS day FROM cost_lines WHERE source = ?")
598 .bind(&[SOURCE_BILLABLE.into()])?
599 .first::<Last>(None)
600 .await?
601 .and_then(|l| l.day);
602 let (since, until) = window(last.as_deref(), now_ms());
603 let fetched_at = rfc3339(now_ms());
604 let mut written = 0;
605 match keeper.billable_usage_body(&since, &until).await.map_err(|e| e.to_string()).and_then(|b| lines_from_billable(&b)) {
606 Ok(lines) => written += self.upsert_lines(&lines, &fetched_at).await?,
607 Err(error) => problems.push(format!("Cloudflare's billable usage could not be read: {error}")),
608 }
609 match keeper.graphql(artifacts_variables(keeper.account(), &since, &until)).await.map_err(|e| e.to_string()) {
610 Ok(body) => match lines_from_artifacts(&body) {
611 Ok(lines) => {
612 written += self.upsert_lines(&lines, &fetched_at).await?;
613 let owners = self.pull_owners(&pull_ids(&body)).await;
614 self.keep_cloudflare_counts(&since, &until, &artifacts_by_workspace(&body, &owners), &fetched_at).await?;
615 }
616 Err(error) => problems.push(format!("Artifacts events could not be read: {error}")),
617 },
618 Err(error) => problems.push(format!("Artifacts events could not be read: {error}")),
619 }
620 // What AI Gateway priced g1t's own provider traffic at, each day:
621 // the total the ledger's model cost is checked against (`margin`).
622 // Over its own window: until it has answered with a line, the 31
623 // days GraphQL keeps, whatever the bill's window is.
624 if !keeper.gateway().is_empty() {
625 let last = self
626 .db
627 .prepare("SELECT MAX(day) AS day FROM cost_lines WHERE source = ?")
628 .bind(&[SOURCE_GATEWAY.into()])?
629 .first::<Last>(None)
630 .await?
631 .and_then(|l| l.day);
632 let (from, to) = window(last.as_deref(), now_ms());
633 match keeper
634 .gateway_graphql(gateway_variables(keeper.account(), keeper.gateway(), &from, &to))
635 .await
636 .map_err(|e| e.to_string())
637 .and_then(|body| lines_from_gateway(&body))
638 {
639 Ok(lines) => {
640 // A day's models are read whole: one that is gone from a
641 // re-read day must not keep its old line.
642 self.db
643 .prepare("DELETE FROM cost_lines WHERE source = ?1 AND day >= ?2 AND day <= ?3")
644 .bind(&[SOURCE_GATEWAY.into(), from.as_str().into(), to.as_str().into()])?
645 .run()
646 .await?;
647 *gateway = if lines.is_empty() { GatewayRead::Empty { visible: keeper.gateway_visible().await } } else { GatewayRead::Rows };
648 written += self.upsert_lines(&lines, &fetched_at).await?;
649 }
650 Err(error) => {
651 *gateway = GatewayRead::Failed(error.clone());
652 problems.push(format!("AI Gateway's analytics could not be read: {error}"));
653 }
654 }
655 }
656 Ok(Some((since, until, written)))
657 }
658
659 /// Cloudflare's own per-workspace counts for the days, replacing what
660 /// was kept for them.
661 /// The workspace of each pull request's working copy, from repos;
662 /// nothing while repos does not answer (those events stay no one's).
663 async fn pull_owners(&self, pulls: &[String]) -> BTreeMap<String, String> {
664 let Some(repos) = &self.repos else { return BTreeMap::new() };
665 if pulls.is_empty() {
666 return BTreeMap::new();
667 }
668 match g1t_kit::call::<_, Value>(repos, "pull_owners", &json!({ "pulls": pulls })).await {
669 Ok(body) => body["owners"]
670 .as_object()
671 .map(|o| o.iter().filter_map(|(k, v)| Some((k.clone(), v.as_str()?.to_owned()))).collect())
672 .unwrap_or_default(),
673 Err(error) => {
674 worker::console_error!("pull_owners: {error}");
675 BTreeMap::new()
676 }
677 }
678 }
679
680 async fn keep_cloudflare_counts(&self, since: &str, until: &str, counts: &[(String, String, f64)], fetched_at: &str) -> Result<()> {
681 self.db
682 .prepare("DELETE FROM own_counts WHERE meter = 'cloudflare_git' AND day >= ?1 AND day <= ?2")
683 .bind(&[since.into(), until.into()])?
684 .run()
685 .await?;
686 for chunk in counts.chunks(50) {
687 let mut statements = Vec::with_capacity(chunk.len());
688 for (day, workspace, count) in chunk {
689 statements.push(
690 self.db
691 .prepare("INSERT OR REPLACE INTO own_counts (day, meter, workspace, quantity, fetched_at) VALUES (?, 'cloudflare_git', ?, ?, ?)")
692 .bind(&[day.as_str().into(), workspace.as_str().into(), (*count).into(), fetched_at.into()])?,
693 );
694 }
695 self.db.batch(statements).await?;
696 }
697 Ok(())
698 }
699
700 /// Upserts lines, a day's line replacing what was read for it before.
701 async fn upsert_lines(&self, lines: &[CostLine], fetched_at: &str) -> Result<u32> {
702 for chunk in lines.chunks(50) {
703 let mut statements = Vec::with_capacity(chunk.len());
704 for line in chunk {
705 statements.push(
706 self.db
707 .prepare(
708 "INSERT INTO cost_lines (day, source, product, meter, unit, quantity, cost_usd, raw_name, fetched_at)
709 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)
710 ON CONFLICT (day, source, product, meter) DO UPDATE SET
711 unit = ?5, quantity = ?6, cost_usd = ?7, raw_name = ?8, fetched_at = ?9",
712 )
713 .bind(&[
714 line.day.as_str().into(),
715 line.source.into(),
716 line.product.as_str().into(),
717 line.meter.as_str().into(),
718 line.unit.as_str().into(),
719 line.quantity.into(),
720 line.cost_usd.into(),
721 line.raw_name.as_str().into(),
722 fetched_at.into(),
723 ])?,
724 );
725 }
726 self.db.batch(statements).await?;
727 }
728 Ok(lines.len() as u32)
729 }
730
731 /// The mapping from Cloudflare's meters to g1t's products.
732 pub(crate) async fn rules(&self) -> Result<Vec<Rule>> {
733 #[derive(Deserialize)]
734 struct Row {
735 product: String,
736 meter: String,
737 bucket: String,
738 price_meter: Option<String>,
739 own_meter: Option<String>,
740 drift_percent: f64,
741 }
742 Ok(self
743 .db
744 .prepare("SELECT product, meter, bucket, price_meter, own_meter, drift_percent FROM cost_map")
745 .all()
746 .await?
747 .results::<Row>()?
748 .into_iter()
749 .map(|r| Rule {
750 product: r.product,
751 meter: r.meter,
752 bucket: r.bucket,
753 price_meter: r.price_meter,
754 own_meter: r.own_meter,
755 drift_percent: r.drift_percent,
756 })
757 .collect())
758 }
759
760 /// g1t's own counts for the days, all from the repos service, which
761 /// owns the mapping from raw meters to operations (`operation_mapping`):
762 ///
763 /// - `git_operations`: what customers are charged for, as repos counts
764 /// it (`git_operations`, already through its mapping).
765 /// - `cost_operations`: what g1t expects Cloudflare to bill, from the
766 /// raw meters (`artifacts_usage`) and the mapping's cost column.
767 /// - `artifacts_<meter>`: each raw meter.
768 pub(crate) async fn count_own(&self, since: &str, until: &str) -> Result<()> {
769 let Some(repos) = &self.repos else { return Ok(()) };
770 let days = days_between(since, until);
771 let fetched_at = rfc3339(now_ms());
772 let mut rows: Vec<(String, String, String, f64)> = Vec::new();
773
774 // Raw meters and repos' mapping; skipped while repos does not answer.
775 if let Ok(body) = g1t_kit::call::<_, Value>(repos, "artifacts_usage", &json!({ "from": since, "to": until })).await {
776 let mapping = operation_mapping(&body);
777 let mut by: BTreeMap<(String, String), Vec<(String, f64)>> = BTreeMap::new();
778 for (day, workspace, meter, count) in body["rows"].as_array().map(|l| l.iter().filter_map(raw_usage_row).collect::<Vec<_>>()).unwrap_or_default() {
779 by.entry((day, workspace)).or_default().push((meter, count));
780 }
781 for ((day, workspace), counts) in by {
782 // Several stores (namespaces) can give the same meter.
783 let mut merged: BTreeMap<String, f64> = BTreeMap::new();
784 for (meter, count) in &counts {
785 *merged.entry(meter.clone()).or_default() += count;
786 }
787 for (meter, count) in &merged {
788 rows.push((day.clone(), format!("artifacts_{meter}"), workspace.clone(), *count));
789 }
790 if !mapping.is_empty() {
791 rows.push((day, "cost_operations".to_owned(), workspace, weighted(&counts, &mapping, true)));
792 }
793 }
794 }
795 // What customers are charged for, as repos counts it through its mapping.
796 let mut cumulative = Vec::with_capacity(days.len());
797 for day in &days {
798 let list: Vec<WorkspaceGitOperations> = g1t_kit::call(
799 repos,
800 "git_operations",
801 &GitOperationsArgs { month: day[..7].to_owned(), since: Some(format!("{day}T00")), namespace: None },
802 )
803 .await?;
804 cumulative.push(list.into_iter().map(|w| (w.namespace.to_lowercase(), w.operations)).collect::<BTreeMap<_, _>>());
805 }
806 for (day, workspace, count) in daily_from_cumulative(&days, &cumulative) {
807 rows.push((day, "git_operations".to_owned(), workspace, count as f64));
808 }
809
810 // Each day's counts replace what was there.
811 self.db
812 .prepare("DELETE FROM own_counts WHERE day >= ?1 AND day <= ?2 AND meter NOT LIKE 'cloudflare_%'")
813 .bind(&[since.into(), until.into()])?
814 .run()
815 .await?;
816 for chunk in rows.chunks(50) {
817 let mut statements = Vec::with_capacity(chunk.len());
818 for (day, meter, workspace, quantity) in chunk {
819 statements.push(
820 self.db
821 .prepare(
822 "INSERT INTO own_counts (day, meter, workspace, quantity, fetched_at) VALUES (?1, ?2, ?3, ?4, ?5)
823 ON CONFLICT (day, meter, workspace) DO UPDATE SET quantity = ?4, fetched_at = ?5",
824 )
825 .bind(&[day.as_str().into(), meter.as_str().into(), workspace.as_str().into(), (*quantity).into(), fetched_at.as_str().into()])?,
826 );
827 }
828 self.db.batch(statements).await?;
829 }
830 Ok(())
831 }
832}
833
834#[cfg(test)]
835mod tests {
836 use super::*;
837
838 /// AI Gateway's analytics in the shape of the schema
839 /// (`aiGatewayRequestsAdaptiveGroups`: `count`, `sum`, `dimensions`).
840 fn gateway_fixture() -> Value {
841 let group = |day: &str, provider: &str, model: &str, count: u64, cost: f64, tokens: (f64, f64, f64, f64)| {
842 json!({
843 "count": count,
844 "sum": { "cost": cost, "tokensIn": tokens.0, "tokensOut": tokens.1, "cacheReadTokens": tokens.2, "cacheWriteTokens": tokens.3 },
845 "dimensions": { "date": day, "provider": provider, "model": model }
846 })
847 };
848 json!({ "data": { "viewer": { "accounts": [{
849 "own": [
850 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)),
851 group("2026-10-05", "anthropic", "claude-haiku-4-5-20251001", 300, 0.40, (400_000.0, 20_000.0, 0.0, 0.0)),
852 group("2026-10-06", "anthropic", "claude-new-1", 12, 0.0, (80_000.0, 4_000.0, 0.0, 0.0)),
853 ],
854 "wholesale": [group("2026-10-06", "openai", "gpt-x", 5, 0.10, (1_000.0, 100.0, 0.0, 0.0))],
855 }] } }, "errors": null })
856 }
857
858 #[test]
859 fn the_gateway_query_names_the_fields_its_schema_has() {
860 // As checked against Cloudflare's GraphQL schema (introspection of
861 // AccountAiGatewayRequestsAdaptiveGroups{,Sum,Dimensions,Filter}).
862 for field in ["aiGatewayRequestsAdaptiveGroups", "date_geq", "date_leq", "gateway: $gateway", "wholesale: 0", "wholesale: 1", "count", "cost", "tokensIn", "tokensOut", "cacheReadTokens", "cacheWriteTokens", "date", "provider", "model"] {
863 assert!(GATEWAY_QUERY.contains(field), "{field}");
864 }
865 // Grouped by `wholesale`, Cloudflare answers no rows and no error:
866 // it is only ever a filter.
867 for dimensions in GATEWAY_QUERY.split("dimensions {").skip(1) {
868 let dimensions = &dimensions[..dimensions.find('}').unwrap()];
869 assert!(!dimensions.contains("wholesale"), "{dimensions}");
870 }
871 let body = gateway_variables("acct", "g1t", "2026-10-01", "2026-10-07");
872 assert_eq!(body["variables"]["gateway"], "g1t");
873 }
874
875 #[test]
876 fn what_the_gateway_priced_is_kept_per_day_and_model_with_its_tokens() {
877 let lines = lines_from_gateway(&gateway_fixture()).unwrap();
878 let get = |day: &str, meter: &str| lines.iter().find(|l| l.day == day && l.meter == meter).unwrap_or_else(|| panic!("{day} {meter}"));
879 let sonnet = get("2026-10-05", "anthropic_claude_sonnet_5_5");
880 assert_eq!((sonnet.source, sonnet.product.as_str(), sonnet.quantity, sonnet.cost_usd), (SOURCE_GATEWAY, GATEWAY_PRODUCT, 120.0, 4.25));
881 assert_eq!(get("2026-10-05", "anthropic_claude_sonnet_5_5__tokens").quantity, 940_000.0);
882 assert_eq!(get("2026-10-05", "anthropic_claude_sonnet_5_5__cache_read_tokens").quantity, 3_000_000.0);
883 assert_eq!(get("2026-10-05", "anthropic_claude_sonnet_5_5__cache_write_tokens").cost_usd, 0.0);
884 assert_eq!(get("2026-10-06", "wholesale__openai_gpt_x").cost_usd, 0.10);
885 // Only the request lines carry cost: the day's total is the gateway's.
886 let day5: f64 = lines.iter().filter(|l| l.day == "2026-10-05").map(|l| l.cost_usd).sum();
887 assert!((day5 - 4.65).abs() < 1e-9);
888 // Errors are a problem for the run, not lines.
889 assert!(lines_from_gateway(&json!({ "errors": [{ "message": "not authorized for that account" }] })).is_err());
890 }
891
892 #[test]
893 fn the_gateway_lines_say_when_its_cost_cannot_be_the_providers() {
894 let rows: Vec<(String, f64, f64)> =
895 lines_from_gateway(&gateway_fixture()).unwrap().into_iter().map(|l| (l.meter, l.quantity, l.cost_usd)).collect();
896 let caveats = gateway_caveats(&rows);
897 assert_eq!(caveats.unpriced, vec!["anthropic_claude_new_1".to_owned()]);
898 assert_eq!((caveats.cache_read_tokens, caveats.cache_write_tokens), (3_000_000.0, 200_000.0));
899 assert!((caveats.wholesale_usd - 0.10).abs() < 1e-9);
900 // Every request it logged, priced or not, from both answers.
901 assert_eq!(caveats.requests, 437.0);
902 // A model priced on one day and not another is priced.
903 let mixed = vec![
904 ("m".to_owned(), 3.0, 0.5),
905 ("m__tokens".to_owned(), 10.0, 0.0),
906 ("m".to_owned(), 3.0, 0.0),
907 ("m__tokens".to_owned(), 10.0, 0.0),
908 ];
909 assert!(gateway_caveats(&mixed).unpriced.is_empty());
910 }
911
912 /// Billable usage as Cloudflare answered g1t on 2026-10-06 (two days,
913 /// trimmed), plus rows past the included amounts, which cost money.
914 fn billable_fixture() -> Value {
915 json!({
916 "success": true,
917 "errors": [],
918 "result": [
919 {
920 "ChargePeriodStart": "2026-10-03T00:00:00Z",
921 "ChargePeriodEnd": "2026-10-04T00:00:00Z",
922 "ServiceFamilyName": "Containers",
923 "ServiceName": "Container Memory, per GiB-Second (First 25 GiB-hours included)",
924 "PricingUnit": "Count",
925 "PricingQuantity": "14000",
926 "ContractedCost": 0, "BilledCost": 0, "ListCost": 0
927 },
928 {
929 "ChargePeriodStart": "2026-10-03T00:00:00Z",
930 "ServiceFamilyName": "D1",
931 "ServiceName": "D1 - Rows Read (first 25 billion included)",
932 "PricingUnit": "Count",
933 "PricingQuantity": 1200000,
934 "BilledCost": 0
935 },
936 {
937 "ChargePeriodStart": "2026-10-15T00:00:00Z",
938 "ServiceFamilyName": "Artifacts",
939 "ServiceName": "Artifacts Operations (First 10,000 included)",
940 "PricingUnit": "Count",
941 "PricingQuantity": "30000",
942 "BilledCost": "3.00",
943 "ListCost": 4.5
944 },
945 {
946 "ChargePeriodStart": "2026-10-15T00:00:00Z",
947 "ServiceFamilyName": "Artifacts",
948 "ServiceName": "Artifacts Operations (First 10,000 included)",
949 "PricingUnit": "Count",
950 "PricingQuantity": "10000",
951 "BilledCost": "1.50"
952 },
953 {
954 "ChargePeriodStart": "2026-10-15T00:00:00Z",
955 "ServiceFamilyName": "Workers",
956 "ServiceName": "Workers for Platforms Requests (First 20M are included)",
957 "PricingUnit": "Count",
958 "PricingQuantity": 25000000,
959 "ListCost": 1.5
960 },
961 { "ServiceName": "No day, no line" }
962 ]
963 })
964 }
965
966 #[test]
967 fn names_become_products_and_meters() {
968 assert_eq!(slug("Workers for Platforms CPU ms (First 60M ms are included)"), "workers_for_platforms_cpu_ms");
969 assert_eq!(product_and_meter("D1", "D1 - Rows Read (first 25 billion included)"), ("d1".into(), "d1_rows_read".into()));
970 assert_eq!(product_and_meter("Workers KV", "KV Read Operations (First 10M is included)"), ("workers_kv".into(), "kv_read_operations".into()));
971 assert_eq!(product_and_meter("", "Containers / Container vCPU"), ("containers".into(), "container_vcpu".into()));
972 assert_eq!(product_and_meter("Email", ""), ("email".into(), "usage".into()));
973 }
974
975 #[test]
976 fn billable_usage_is_parsed_and_rows_of_one_meter_are_added() {
977 let lines = lines_from_billable(&billable_fixture()).unwrap();
978 assert_eq!(lines.len(), 4, "{lines:#?}");
979 let artifacts = lines.iter().find(|l| l.product == "artifacts").unwrap();
980 // Two rows of the same day and meter: one line, added up.
981 assert_eq!(artifacts.meter, "artifacts_operations");
982 assert_eq!(artifacts.quantity, 40_000.0);
983 assert!((artifacts.cost_usd - 4.5).abs() < 1e-9, "billed, not list: {}", artifacts.cost_usd);
984 let wfp = lines.iter().find(|l| l.meter.starts_with("workers_for_platforms")).unwrap();
985 assert_eq!(wfp.product, "workers");
986 // No billed cost: the list cost.
987 assert_eq!(wfp.cost_usd, 1.5);
988 let memory = lines.iter().find(|l| l.product == "containers").unwrap();
989 assert_eq!((memory.day.as_str(), memory.quantity, memory.cost_usd), ("2026-10-03", 14_000.0, 0.0));
990 assert!(lines_from_billable(&json!({ "success": false, "errors": [{ "code": 10000 }] })).is_err());
991 }
992
993 #[test]
994 fn reading_the_same_days_twice_gives_the_same_lines() {
995 // Idempotent: the same answer aggregates to the same keys and
996 // amounts, so the upsert replaces rather than adds.
997 let once = lines_from_billable(&billable_fixture()).unwrap();
998 let twice = lines_from_billable(&billable_fixture()).unwrap();
999 assert_eq!(once, twice);
1000 let again = aggregate(once.clone());
1001 assert_eq!(again, once);
1002 }
1003
1004 #[test]
1005 fn artifacts_events_are_counted_by_type_and_day() {
1006 let body = json!({
1007 "data": { "viewer": { "accounts": [{ "artifactsEventsAdaptiveGroups": [
1008 { "count": 120, "dimensions": { "date": "2026-10-05", "eventType": "pull" } },
1009 { "count": 30, "dimensions": { "date": "2026-10-05", "eventType": "push" } },
1010 { "count": 2, "dimensions": { "date": "2026-10-05", "eventType": "rateLimited" } },
1011 { "count": 5, "dimensions": { "date": "2026-10-06", "eventType": "pull" } }
1012 ] }] } },
1013 "errors": null
1014 });
1015 let lines = lines_from_artifacts(&body).unwrap();
1016 assert_eq!(lines.len(), 4);
1017 assert_eq!(lines[0].meter, "events_pull");
1018 assert_eq!(lines[0].quantity, 120.0);
1019 assert!(lines.iter().any(|l| l.meter == "events_ratelimited"));
1020 // By workspace, from the repository's store key; forks say none.
1021 let by_repo = json!({ "data": { "viewer": { "accounts": [{ "artifactsEventsAdaptiveGroups": [
1022 { "count": 100, "dimensions": { "date": "2026-10-05", "eventType": "pull", "repositoryName": "acme--api" } },
1023 { "count": 20, "dimensions": { "date": "2026-10-05", "eventType": "push", "repositoryName": "acme--web" } },
1024 { "count": 7, "dimensions": { "date": "2026-10-05", "eventType": "fork", "repositoryName": "pulls--123" } },
1025 { "count": 3, "dimensions": { "date": "2026-10-05", "eventType": "serverError", "repositoryName": "beta--x" } }
1026 ] }] } } });
1027 assert_eq!(artifacts_by_workspace(&by_repo, &BTreeMap::new()), vec![("2026-10-05".to_string(), "acme".to_string(), 120.0)]);
1028 let none = BTreeMap::new();
1029 assert_eq!(workspace_of_store_key("Acme--api", &none), Some("acme".into()));
1030 assert_eq!(workspace_of_store_key("pulls--9", &none), None);
1031 assert_eq!(workspace_of_store_key("plain", &none), None);
1032 let owners = BTreeMap::from([("9".to_string(), "Acme".to_string())]);
1033 assert_eq!(workspace_of_store_key("pulls--9", &owners), Some("acme".into()));
1034 assert_eq!(workspace_of_store_key("pulls--10", &owners), None);
1035 let operations: f64 = lines
1036 .iter()
1037 .filter(|l| l.day == "2026-10-05" && ARTIFACTS_OPERATIONS.contains(&l.meter.as_str()))
1038 .map(|l| l.quantity)
1039 .sum();
1040 assert_eq!(operations, 150.0);
1041 let failed = json!({ "data": null, "errors": [{ "message": "unknown field" }] });
1042 assert!(lines_from_artifacts(&failed).is_err());
1043 }
1044
1045 #[test]
1046 fn the_first_run_backfills_31_days_and_later_runs_restate_a_few() {
1047 let now = g1t_contracts::time::parse_rfc3339("2026-10-06T04:17:00Z").unwrap();
1048 assert_eq!(window(None, now), ("2026-09-05".into(), "2026-10-06".into()));
1049 assert_eq!(window(Some("2026-10-06"), now), ("2026-10-02".into(), "2026-10-06".into()));
1050 // A week without a run: from the last day read.
1051 assert_eq!(window(Some("2026-09-28"), now), ("2026-09-28".into(), "2026-10-06".into()));
1052 // Never past the backfill.
1053 assert_eq!(window(Some("2026-01-01"), now), ("2026-09-05".into(), "2026-10-06".into()));
1054 assert_eq!(days_between("2026-09-29", "2026-10-02"), vec!["2026-09-29", "2026-09-30", "2026-10-01", "2026-10-02"]);
1055 assert!(days_between("2026-10-02", "2026-10-01").is_empty());
1056 }
1057
1058 #[test]
1059 fn a_days_git_operations_are_its_count_less_the_next_days() {
1060 let days: Vec<String> = ["2026-09-29", "2026-09-30", "2026-10-01", "2026-10-02"].iter().map(|d| d.to_string()).collect();
1061 let map = |pairs: &[(&str, u64)]| pairs.iter().map(|(w, n)| (w.to_string(), *n)).collect::<BTreeMap<_, _>>();
1062 // From each day to the end of its month.
1063 let cumulative = vec![map(&[("acme", 30), ("beta", 4)]), map(&[("acme", 10)]), map(&[("acme", 7)]), map(&[("acme", 2)])];
1064 let daily = daily_from_cumulative(&days, &cumulative);
1065 assert_eq!(
1066 daily,
1067 vec![
1068 ("2026-09-29".into(), "acme".into(), 20),
1069 ("2026-09-29".into(), "beta".into(), 4),
1070 // The month's last day is its own.
1071 ("2026-09-30".into(), "acme".into(), 10),
1072 ("2026-10-01".into(), "acme".into(), 5),
1073 // Today so far.
1074 ("2026-10-02".into(), "acme".into(), 2),
1075 ]
1076 );
1077 }
1078
1079 #[test]
1080 fn raw_meters_are_weighted_by_the_repos_mapping() {
1081 // As repos' artifacts_usage answers: rows, and its operation_mapping.
1082 let body = json!({
1083 "rows": [
1084 { "day": "2026-10-15", "store": "g1t", "workspace": "Acme", "meter": "git.fetch", "count": 12, "bytes_in": 0, "bytes_out": 0 },
1085 { "day": "2026-10-15", "store": "g1t", "workspace": "acme", "meter": "git.receive_pack", "count": 3, "bytes_in": 0, "bytes_out": 0 },
1086 { "day": "2026-10-15", "store": "g1t", "workspace": "acme", "meter": "binding.read_blob", "count": 400, "bytes_in": 0, "bytes_out": 0 }
1087 ],
1088 "mapping": [
1089 { "meter": "git.fetch", "cost_operations": 1, "billable_operations": 1 },
1090 { "meter": "git.receive_pack", "cost_operations": 1, "billable_operations": 1 },
1091 { "meter": "binding.read_blob", "cost_operations": 0, "billable_operations": 0 }
1092 ],
1093 "truncated": false
1094 });
1095 let row = raw_usage_row(&body["rows"][0]).unwrap();
1096 assert_eq!(row, ("2026-10-15".into(), "acme".into(), "git_fetch".into(), 12.0));
1097 assert!(raw_usage_row(&json!({ "day": "2026-10-15", "count": 1 })).is_none());
1098 let mut mapping = operation_mapping(&body);
1099 let raw: Vec<(String, f64)> = body["rows"].as_array().unwrap().iter().filter_map(raw_usage_row).map(|r| (r.2, r.3)).collect();
1100 assert_eq!(weighted(&raw, &mapping, true), 15.0);
1101 assert_eq!(weighted(&raw, &mapping, false), 15.0);
1102 // Cloudflare turns out to bill binding reads: repos changes one row
1103 // (set_operation_mapping), and the bill g1t expects follows.
1104 mapping.insert("binding_read_blob".into(), (1.0, 0.0));
1105 assert_eq!(weighted(&raw, &mapping, true), 415.0);
1106 assert_eq!(weighted(&raw, &mapping, false), 15.0);
1107 assert!(operation_mapping(&json!({})).is_empty());
1108 }
1109
1110 #[test]
1111 fn the_most_specific_rule_claims_a_line() {
1112 let rule = |product: &str, meter: &str, bucket: &str| Rule {
1113 product: product.into(),
1114 meter: meter.into(),
1115 bucket: bucket.into(),
1116 price_meter: None,
1117 own_meter: None,
1118 drift_percent: 10.0,
1119 };
1120 let rules = vec![rule("workers", "*", "platform"), rule("workers", "workers_for_platforms", "deployments"), rule("artifacts", "*", "git")];
1121 assert_eq!(classify(&rules, "workers", "workers_for_platforms_cpu_ms").unwrap().bucket, "deployments");
1122 assert_eq!(classify(&rules, "workers", "workers_cpu_ms").unwrap().bucket, "platform");
1123 assert_eq!(classify(&rules, "artifacts", "artifacts_operations").unwrap().bucket, "git");
1124 assert!(classify(&rules, "browser_rendering", "browser_hours").is_none());
1125 }
1126}