g1t/services/billing/src/costs.rs

1,062 lines51,745 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily1//! What Cloudflare charges g1t, day by day and product by product.
2//!
3//! Prices are what g1t pays plus 20%, so g1t has to know what it pays, as
4//! Cloudflare counts it, not as g1t assumes it. Some of what g1t runs on
5//! is in beta (Artifacts bills "operations" from 2026-10-14 without saying
6//! exactly which calls are one), so nothing here assumes a product list:
7//! every line Cloudflare bills is kept, whatever it is, and how a line
8//! maps to what g1t sells is data (`cost_map`), changed without a deploy.
9//!
10//! Once a day (`keeper::DAILY`) billing reads:
11//!
12//! - **Billable usage** (`GET /accounts/{id}/billable-usage`, FOCUS
13//! columns): one row per service per day, with its quantity and what it
14//! cost g1t. Every Cloudflare product g1t uses appears here once it is
15//! used: Workers, Workers for Platforms, D1, KV, R2, Queues, Containers,
16//! Durable Objects, Artifacts, Browser Rendering, Workers AI, Vectorize,
17//! Cloudflare for SaaS and Email.
18//! - **Artifacts events** (GraphQL `artifactsEventsAdaptiveGroups`), by
19//! event type and day: what Artifacts itself counted, to compare with
20//! what g1t counted (`margin`).
21//!
22//! Each becomes cost lines in `cost_lines`, one per (day, source,
23//! product, meter), upserted, so reading a day again replaces it. The
24//! first run reads the last 31 days; after that the last few, since
25//! Cloudflare restates recent days as usage settles.
26//!
27//! The token is `CLOUDFLARE_BILLING_TOKEN` (Account: Billing Read and
28//! Account Analytics Read), or the keeper's `CLOUDFLARE_USAGE_TOKEN`,
29//! which has both. Without either, nothing is read and nothing fails.
30
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)31use std::collections::{BTreeMap, BTreeSet};
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily32
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
Merge branch 'worktree-agent-a633ac0f7f66d419d'241/// Where a cost line came from: AI Gateway's analytics, what it priced
242/// g1t's own provider traffic at.
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
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily355/// The workspace a repository in the store belongs to: keys are
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)356/// `<workspace>--<repo>`; a pull request's working copy (`pulls--<id>`) is
357/// its repository's workspace's, from `owners` (repos' `pull_owners`), else
358/// no one's.
359pub(crate) fn workspace_of_store_key(key: &str, owners: &BTreeMap<String, String>) -> Option<String> {
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily360 let (workspace, rest) = key.split_once("--")?;
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)361 if workspace.is_empty() || rest.is_empty() {
362 return None;
363 }
364 if workspace == "pulls" {
365 return owners.get(rest).map(|w| w.to_lowercase());
366 }
367 Some(workspace.to_lowercase())
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily368}
369
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)370/// The pull request ids whose working copies Artifacts counted events for.
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
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily381/// Artifacts' billable operations per workspace and day, by repository
382/// name: how Cloudflare's own count shares out. The meter is
383/// `cloudflare_git`, which shares out the git bucket's cost (`margin`).
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)384pub(crate) fn artifacts_by_workspace(body: &Value, owners: &BTreeMap<String, String>) -> Vec<(String, String, f64)> {
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily385 let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"].as_array().cloned().unwrap_or_default();
386 let mut out: BTreeMap<(String, String), f64> = BTreeMap::new();
387 for g in &groups {
388 let (Some(day), Some(kind)) = (g["dimensions"]["date"].as_str(), g["dimensions"]["eventType"].as_str()) else { continue };
389 if day.len() < 10 || !ARTIFACTS_OPERATIONS.contains(&format!("events_{}", slug(kind)).as_str()) {
390 continue;
391 }
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)392 let Some(workspace) = g["dimensions"]["repositoryName"].as_str().and_then(|key| workspace_of_store_key(key, owners)) else { continue };
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily393 *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
One operation mapping, owned by repos; billing reads it instead of keeping its own516/// 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()
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily541}
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?;
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)573 let owners = self.pull_owners(&pull_ids(&body)).await;
574 self.keep_cloudflare_counts(&since, &until, &artifacts_by_workspace(&body, &owners), &fetched_at).await?;
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily575 }
576 Err(error) => problems.push(format!("Artifacts events could not be read: {error}")),
577 },
578 Err(error) => problems.push(format!("Artifacts events could not be read: {error}")),
579 }
Merge branch 'worktree-agent-a633ac0f7f66d419d'580 // What AI Gateway priced g1t's own provider traffic at, each day:
581 // the total the ledger's model cost is checked against (`margin`).
582 if !keeper.gateway().is_empty() {
583 match keeper
584 .graphql_either(gateway_variables(keeper.account(), keeper.gateway(), &since, &until))
585 .await
586 .map_err(|e| e.to_string())
587 .and_then(|body| lines_from_gateway(&body))
588 {
589 Ok(lines) => {
590 // A day's models are read whole: one that is gone from a
591 // re-read day must not keep its old line.
592 self.db
593 .prepare("DELETE FROM cost_lines WHERE source = ?1 AND day >= ?2 AND day <= ?3")
594 .bind(&[SOURCE_GATEWAY.into(), since.as_str().into(), until.as_str().into()])?
595 .run()
596 .await?;
597 written += self.upsert_lines(&lines, &fetched_at).await?;
598 }
599 Err(error) => problems.push(format!("AI Gateway's analytics could not be read: {error}")),
600 }
601 }
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily602 Ok(Some((since, until, written)))
603 }
604
605 /// Cloudflare's own per-workspace counts for the days, replacing what
606 /// was kept for them.
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)607 /// The workspace of each pull request's working copy, from repos;
608 /// nothing while repos does not answer (those events stay no one's).
609 async fn pull_owners(&self, pulls: &[String]) -> BTreeMap<String, String> {
610 let Some(repos) = &self.repos else { return BTreeMap::new() };
611 if pulls.is_empty() {
612 return BTreeMap::new();
613 }
614 match g1t_kit::call::<_, Value>(repos, "pull_owners", &json!({ "pulls": pulls })).await {
615 Ok(body) => body["owners"]
616 .as_object()
617 .map(|o| o.iter().filter_map(|(k, v)| Some((k.clone(), v.as_str()?.to_owned()))).collect())
618 .unwrap_or_default(),
619 Err(error) => {
620 worker::console_error!("pull_owners: {error}");
621 BTreeMap::new()
622 }
623 }
624 }
625
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily626 async fn keep_cloudflare_counts(&self, since: &str, until: &str, counts: &[(String, String, f64)], fetched_at: &str) -> Result<()> {
627 self.db
628 .prepare("DELETE FROM own_counts WHERE meter = 'cloudflare_git' AND day >= ?1 AND day <= ?2")
629 .bind(&[since.into(), until.into()])?
630 .run()
631 .await?;
632 for chunk in counts.chunks(50) {
633 let mut statements = Vec::with_capacity(chunk.len());
634 for (day, workspace, count) in chunk {
635 statements.push(
636 self.db
637 .prepare("INSERT OR REPLACE INTO own_counts (day, meter, workspace, quantity, fetched_at) VALUES (?, 'cloudflare_git', ?, ?, ?)")
638 .bind(&[day.as_str().into(), workspace.as_str().into(), (*count).into(), fetched_at.into()])?,
639 );
640 }
641 self.db.batch(statements).await?;
642 }
643 Ok(())
644 }
645
646 /// Upserts lines, a day's line replacing what was read for it before.
647 async fn upsert_lines(&self, lines: &[CostLine], fetched_at: &str) -> Result<u32> {
648 for chunk in lines.chunks(50) {
649 let mut statements = Vec::with_capacity(chunk.len());
650 for line in chunk {
651 statements.push(
652 self.db
653 .prepare(
654 "INSERT INTO cost_lines (day, source, product, meter, unit, quantity, cost_usd, raw_name, fetched_at)
655 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)
656 ON CONFLICT (day, source, product, meter) DO UPDATE SET
657 unit = ?5, quantity = ?6, cost_usd = ?7, raw_name = ?8, fetched_at = ?9",
658 )
659 .bind(&[
660 line.day.as_str().into(),
661 line.source.into(),
662 line.product.as_str().into(),
663 line.meter.as_str().into(),
664 line.unit.as_str().into(),
665 line.quantity.into(),
666 line.cost_usd.into(),
667 line.raw_name.as_str().into(),
668 fetched_at.into(),
669 ])?,
670 );
671 }
672 self.db.batch(statements).await?;
673 }
674 Ok(lines.len() as u32)
675 }
676
677 /// The mapping from Cloudflare's meters to g1t's products.
678 pub(crate) async fn rules(&self) -> Result<Vec<Rule>> {
679 #[derive(Deserialize)]
680 struct Row {
681 product: String,
682 meter: String,
683 bucket: String,
684 price_meter: Option<String>,
685 own_meter: Option<String>,
686 drift_percent: f64,
687 }
688 Ok(self
689 .db
690 .prepare("SELECT product, meter, bucket, price_meter, own_meter, drift_percent FROM cost_map")
691 .all()
692 .await?
693 .results::<Row>()?
694 .into_iter()
695 .map(|r| Rule {
696 product: r.product,
697 meter: r.meter,
698 bucket: r.bucket,
699 price_meter: r.price_meter,
700 own_meter: r.own_meter,
701 drift_percent: r.drift_percent,
702 })
703 .collect())
704 }
705
One operation mapping, owned by repos; billing reads it instead of keeping its own706 /// g1t's own counts for the days, all from the repos service, which
707 /// owns the mapping from raw meters to operations (`operation_mapping`):
708 ///
709 /// - `git_operations`: what customers are charged for, as repos counts
710 /// it (`git_operations`, already through its mapping).
711 /// - `cost_operations`: what g1t expects Cloudflare to bill, from the
712 /// raw meters (`artifacts_usage`) and the mapping's cost column.
713 /// - `artifacts_<meter>`: each raw meter.
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily714 pub(crate) async fn count_own(&self, since: &str, until: &str) -> Result<()> {
715 let Some(repos) = &self.repos else { return Ok(()) };
716 let days = days_between(since, until);
717 let fetched_at = rfc3339(now_ms());
718 let mut rows: Vec<(String, String, String, f64)> = Vec::new();
719
One operation mapping, owned by repos; billing reads it instead of keeping its own720 // Raw meters and repos' mapping; skipped while repos does not answer.
721 if let Ok(body) = g1t_kit::call::<_, Value>(repos, "artifacts_usage", &json!({ "from": since, "to": until })).await {
722 let mapping = operation_mapping(&body);
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily723 let mut by: BTreeMap<(String, String), Vec<(String, f64)>> = BTreeMap::new();
One operation mapping, owned by repos; billing reads it instead of keeping its own724 for (day, workspace, meter, count) in body["rows"].as_array().map(|l| l.iter().filter_map(raw_usage_row).collect::<Vec<_>>()).unwrap_or_default() {
725 by.entry((day, workspace)).or_default().push((meter, count));
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily726 }
727 for ((day, workspace), counts) in by {
One operation mapping, owned by repos; billing reads it instead of keeping its own728 // Several stores (namespaces) can give the same meter.
729 let mut merged: BTreeMap<String, f64> = BTreeMap::new();
730 for (meter, count) in &counts {
731 *merged.entry(meter.clone()).or_default() += count;
732 }
733 for (meter, count) in &merged {
734 rows.push((day.clone(), format!("artifacts_{meter}"), workspace.clone(), *count));
735 }
736 if !mapping.is_empty() {
737 rows.push((day, "cost_operations".to_owned(), workspace, weighted(&counts, &mapping, true)));
738 }
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily739 }
740 }
One operation mapping, owned by repos; billing reads it instead of keeping its own741 // What customers are charged for, as repos counts it through its mapping.
742 let mut cumulative = Vec::with_capacity(days.len());
743 for day in &days {
744 let list: Vec<WorkspaceGitOperations> = g1t_kit::call(
745 repos,
746 "git_operations",
747 &GitOperationsArgs { month: day[..7].to_owned(), since: Some(format!("{day}T00")), namespace: None },
748 )
749 .await?;
750 cumulative.push(list.into_iter().map(|w| (w.namespace.to_lowercase(), w.operations)).collect::<BTreeMap<_, _>>());
751 }
752 for (day, workspace, count) in daily_from_cumulative(&days, &cumulative) {
753 rows.push((day, "git_operations".to_owned(), workspace, count as f64));
754 }
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily755
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
Merge branch 'worktree-agent-a633ac0f7f66d419d'784 /// AI Gateway's analytics in the shape of the schema
785 /// (`aiGatewayRequestsAdaptiveGroups`: `count`, `sum`, `dimensions`).
786 fn gateway_fixture() -> Value {
787 let group = |day: &str, provider: &str, model: &str, wholesale: u8, count: u64, cost: f64, tokens: (f64, f64, f64, f64)| {
788 json!({
789 "count": count,
790 "sum": { "cost": cost, "tokensIn": tokens.0, "tokensOut": tokens.1, "cacheReadTokens": tokens.2, "cacheWriteTokens": tokens.3 },
791 "dimensions": { "date": day, "provider": provider, "model": model, "wholesale": wholesale }
792 })
793 };
794 json!({ "data": { "viewer": { "accounts": [{ "aiGatewayRequestsAdaptiveGroups": [
795 group("2026-10-05", "anthropic", "claude-sonnet-5-5", 0, 120, 4.25, (900_000.0, 40_000.0, 3_000_000.0, 200_000.0)),
796 group("2026-10-05", "anthropic", "claude-haiku-4-5-20251001", 0, 300, 0.40, (400_000.0, 20_000.0, 0.0, 0.0)),
797 group("2026-10-06", "anthropic", "claude-new-1", 0, 12, 0.0, (80_000.0, 4_000.0, 0.0, 0.0)),
798 group("2026-10-06", "openai", "gpt-x", 1, 5, 0.10, (1_000.0, 100.0, 0.0, 0.0)),
799 ] }] } }, "errors": null })
800 }
801
802 #[test]
803 fn the_gateway_query_names_the_fields_its_schema_has() {
804 // As checked against Cloudflare's GraphQL schema (introspection of
805 // AccountAiGatewayRequestsAdaptiveGroups{,Sum,Dimensions,Filter}).
806 for field in ["aiGatewayRequestsAdaptiveGroups", "date_geq", "date_leq", "gateway: $gateway", "count", "cost", "tokensIn", "tokensOut", "cacheReadTokens", "cacheWriteTokens", "date", "provider", "model", "wholesale"] {
807 assert!(GATEWAY_QUERY.contains(field), "{field}");
808 }
809 let body = gateway_variables("acct", "g1t", "2026-10-01", "2026-10-07");
810 assert_eq!(body["variables"]["gateway"], "g1t");
811 }
812
813 #[test]
814 fn what_the_gateway_priced_is_kept_per_day_and_model_with_its_tokens() {
815 let lines = lines_from_gateway(&gateway_fixture()).unwrap();
816 let get = |day: &str, meter: &str| lines.iter().find(|l| l.day == day && l.meter == meter).unwrap_or_else(|| panic!("{day} {meter}"));
817 let sonnet = get("2026-10-05", "anthropic_claude_sonnet_5_5");
818 assert_eq!((sonnet.source, sonnet.product.as_str(), sonnet.quantity, sonnet.cost_usd), (SOURCE_GATEWAY, GATEWAY_PRODUCT, 120.0, 4.25));
819 assert_eq!(get("2026-10-05", "anthropic_claude_sonnet_5_5__tokens").quantity, 940_000.0);
820 assert_eq!(get("2026-10-05", "anthropic_claude_sonnet_5_5__cache_read_tokens").quantity, 3_000_000.0);
821 assert_eq!(get("2026-10-05", "anthropic_claude_sonnet_5_5__cache_write_tokens").cost_usd, 0.0);
822 assert_eq!(get("2026-10-06", "wholesale__openai_gpt_x").cost_usd, 0.10);
823 // Only the request lines carry cost: the day's total is the gateway's.
824 let day5: f64 = lines.iter().filter(|l| l.day == "2026-10-05").map(|l| l.cost_usd).sum();
825 assert!((day5 - 4.65).abs() < 1e-9);
826 // Errors are a problem for the run, not lines.
827 assert!(lines_from_gateway(&json!({ "errors": [{ "message": "not authorized for that account" }] })).is_err());
828 }
829
830 #[test]
831 fn the_gateway_lines_say_when_its_cost_cannot_be_the_providers() {
832 let rows: Vec<(String, f64, f64)> =
833 lines_from_gateway(&gateway_fixture()).unwrap().into_iter().map(|l| (l.meter, l.quantity, l.cost_usd)).collect();
834 let caveats = gateway_caveats(&rows);
835 assert_eq!(caveats.unpriced, vec!["anthropic_claude_new_1".to_owned()]);
836 assert_eq!((caveats.cache_read_tokens, caveats.cache_write_tokens), (3_000_000.0, 200_000.0));
837 assert!((caveats.wholesale_usd - 0.10).abs() < 1e-9);
838 // A model priced on one day and not another is priced.
839 let mixed = vec![
840 ("m".to_owned(), 3.0, 0.5),
841 ("m__tokens".to_owned(), 10.0, 0.0),
842 ("m".to_owned(), 3.0, 0.0),
843 ("m__tokens".to_owned(), 10.0, 0.0),
844 ];
845 assert!(gateway_caveats(&mixed).unpriced.is_empty());
846 }
847
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily848 /// Billable usage as Cloudflare answered g1t on 2026-10-06 (two days,
849 /// trimmed), plus rows past the included amounts, which cost money.
850 fn billable_fixture() -> Value {
851 json!({
852 "success": true,
853 "errors": [],
854 "result": [
855 {
856 "ChargePeriodStart": "2026-10-03T00:00:00Z",
857 "ChargePeriodEnd": "2026-10-04T00:00:00Z",
858 "ServiceFamilyName": "Containers",
859 "ServiceName": "Container Memory, per GiB-Second (First 25 GiB-hours included)",
860 "PricingUnit": "Count",
861 "PricingQuantity": "14000",
862 "ContractedCost": 0, "BilledCost": 0, "ListCost": 0
863 },
864 {
865 "ChargePeriodStart": "2026-10-03T00:00:00Z",
866 "ServiceFamilyName": "D1",
867 "ServiceName": "D1 - Rows Read (first 25 billion included)",
868 "PricingUnit": "Count",
869 "PricingQuantity": 1200000,
870 "BilledCost": 0
871 },
872 {
873 "ChargePeriodStart": "2026-10-15T00:00:00Z",
874 "ServiceFamilyName": "Artifacts",
875 "ServiceName": "Artifacts Operations (First 10,000 included)",
876 "PricingUnit": "Count",
877 "PricingQuantity": "30000",
878 "BilledCost": "3.00",
879 "ListCost": 4.5
880 },
881 {
882 "ChargePeriodStart": "2026-10-15T00:00:00Z",
883 "ServiceFamilyName": "Artifacts",
884 "ServiceName": "Artifacts Operations (First 10,000 included)",
885 "PricingUnit": "Count",
886 "PricingQuantity": "10000",
887 "BilledCost": "1.50"
888 },
889 {
890 "ChargePeriodStart": "2026-10-15T00:00:00Z",
891 "ServiceFamilyName": "Workers",
892 "ServiceName": "Workers for Platforms Requests (First 20M are included)",
893 "PricingUnit": "Count",
894 "PricingQuantity": 25000000,
895 "ListCost": 1.5
896 },
897 { "ServiceName": "No day, no line" }
898 ]
899 })
900 }
901
902 #[test]
903 fn names_become_products_and_meters() {
904 assert_eq!(slug("Workers for Platforms CPU ms (First 60M ms are included)"), "workers_for_platforms_cpu_ms");
905 assert_eq!(product_and_meter("D1", "D1 - Rows Read (first 25 billion included)"), ("d1".into(), "d1_rows_read".into()));
906 assert_eq!(product_and_meter("Workers KV", "KV Read Operations (First 10M is included)"), ("workers_kv".into(), "kv_read_operations".into()));
907 assert_eq!(product_and_meter("", "Containers / Container vCPU"), ("containers".into(), "container_vcpu".into()));
908 assert_eq!(product_and_meter("Email", ""), ("email".into(), "usage".into()));
909 }
910
911 #[test]
912 fn billable_usage_is_parsed_and_rows_of_one_meter_are_added() {
913 let lines = lines_from_billable(&billable_fixture()).unwrap();
914 assert_eq!(lines.len(), 4, "{lines:#?}");
915 let artifacts = lines.iter().find(|l| l.product == "artifacts").unwrap();
916 // Two rows of the same day and meter: one line, added up.
917 assert_eq!(artifacts.meter, "artifacts_operations");
918 assert_eq!(artifacts.quantity, 40_000.0);
919 assert!((artifacts.cost_usd - 4.5).abs() < 1e-9, "billed, not list: {}", artifacts.cost_usd);
920 let wfp = lines.iter().find(|l| l.meter.starts_with("workers_for_platforms")).unwrap();
921 assert_eq!(wfp.product, "workers");
922 // No billed cost: the list cost.
923 assert_eq!(wfp.cost_usd, 1.5);
924 let memory = lines.iter().find(|l| l.product == "containers").unwrap();
925 assert_eq!((memory.day.as_str(), memory.quantity, memory.cost_usd), ("2026-10-03", 14_000.0, 0.0));
926 assert!(lines_from_billable(&json!({ "success": false, "errors": [{ "code": 10000 }] })).is_err());
927 }
928
929 #[test]
930 fn reading_the_same_days_twice_gives_the_same_lines() {
931 // Idempotent: the same answer aggregates to the same keys and
932 // amounts, so the upsert replaces rather than adds.
933 let once = lines_from_billable(&billable_fixture()).unwrap();
934 let twice = lines_from_billable(&billable_fixture()).unwrap();
935 assert_eq!(once, twice);
936 let again = aggregate(once.clone());
937 assert_eq!(again, once);
938 }
939
940 #[test]
941 fn artifacts_events_are_counted_by_type_and_day() {
942 let body = json!({
943 "data": { "viewer": { "accounts": [{ "artifactsEventsAdaptiveGroups": [
944 { "count": 120, "dimensions": { "date": "2026-10-05", "eventType": "pull" } },
945 { "count": 30, "dimensions": { "date": "2026-10-05", "eventType": "push" } },
946 { "count": 2, "dimensions": { "date": "2026-10-05", "eventType": "rateLimited" } },
947 { "count": 5, "dimensions": { "date": "2026-10-06", "eventType": "pull" } }
948 ] }] } },
949 "errors": null
950 });
951 let lines = lines_from_artifacts(&body).unwrap();
952 assert_eq!(lines.len(), 4);
953 assert_eq!(lines[0].meter, "events_pull");
954 assert_eq!(lines[0].quantity, 120.0);
955 assert!(lines.iter().any(|l| l.meter == "events_ratelimited"));
956 // By workspace, from the repository's store key; forks say none.
957 let by_repo = json!({ "data": { "viewer": { "accounts": [{ "artifactsEventsAdaptiveGroups": [
958 { "count": 100, "dimensions": { "date": "2026-10-05", "eventType": "pull", "repositoryName": "acme--api" } },
959 { "count": 20, "dimensions": { "date": "2026-10-05", "eventType": "push", "repositoryName": "acme--web" } },
960 { "count": 7, "dimensions": { "date": "2026-10-05", "eventType": "fork", "repositoryName": "pulls--123" } },
961 { "count": 3, "dimensions": { "date": "2026-10-05", "eventType": "serverError", "repositoryName": "beta--x" } }
962 ] }] } } });
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)963 assert_eq!(artifacts_by_workspace(&by_repo, &BTreeMap::new()), vec![("2026-10-05".to_string(), "acme".to_string(), 120.0)]);
964 let none = BTreeMap::new();
965 assert_eq!(workspace_of_store_key("Acme--api", &none), Some("acme".into()));
966 assert_eq!(workspace_of_store_key("pulls--9", &none), None);
967 assert_eq!(workspace_of_store_key("plain", &none), None);
968 let owners = BTreeMap::from([("9".to_string(), "Acme".to_string())]);
969 assert_eq!(workspace_of_store_key("pulls--9", &owners), Some("acme".into()));
970 assert_eq!(workspace_of_store_key("pulls--10", &owners), None);
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily971 let operations: f64 = lines
972 .iter()
973 .filter(|l| l.day == "2026-10-05" && ARTIFACTS_OPERATIONS.contains(&l.meter.as_str()))
974 .map(|l| l.quantity)
975 .sum();
976 assert_eq!(operations, 150.0);
977 let failed = json!({ "data": null, "errors": [{ "message": "unknown field" }] });
978 assert!(lines_from_artifacts(&failed).is_err());
979 }
980
981 #[test]
982 fn the_first_run_backfills_31_days_and_later_runs_restate_a_few() {
983 let now = g1t_contracts::time::parse_rfc3339("2026-10-06T04:17:00Z").unwrap();
984 assert_eq!(window(None, now), ("2026-09-05".into(), "2026-10-06".into()));
985 assert_eq!(window(Some("2026-10-06"), now), ("2026-10-02".into(), "2026-10-06".into()));
986 // A week without a run: from the last day read.
987 assert_eq!(window(Some("2026-09-28"), now), ("2026-09-28".into(), "2026-10-06".into()));
988 // Never past the backfill.
989 assert_eq!(window(Some("2026-01-01"), now), ("2026-09-05".into(), "2026-10-06".into()));
990 assert_eq!(days_between("2026-09-29", "2026-10-02"), vec!["2026-09-29", "2026-09-30", "2026-10-01", "2026-10-02"]);
991 assert!(days_between("2026-10-02", "2026-10-01").is_empty());
992 }
993
994 #[test]
995 fn a_days_git_operations_are_its_count_less_the_next_days() {
996 let days: Vec<String> = ["2026-09-29", "2026-09-30", "2026-10-01", "2026-10-02"].iter().map(|d| d.to_string()).collect();
997 let map = |pairs: &[(&str, u64)]| pairs.iter().map(|(w, n)| (w.to_string(), *n)).collect::<BTreeMap<_, _>>();
998 // From each day to the end of its month.
999 let cumulative = vec![map(&[("acme", 30), ("beta", 4)]), map(&[("acme", 10)]), map(&[("acme", 7)]), map(&[("acme", 2)])];
1000 let daily = daily_from_cumulative(&days, &cumulative);
1001 assert_eq!(
1002 daily,
1003 vec![
1004 ("2026-09-29".into(), "acme".into(), 20),
1005 ("2026-09-29".into(), "beta".into(), 4),
1006 // The month's last day is its own.
1007 ("2026-09-30".into(), "acme".into(), 10),
1008 ("2026-10-01".into(), "acme".into(), 5),
1009 // Today so far.
1010 ("2026-10-02".into(), "acme".into(), 2),
1011 ]
1012 );
1013 }
1014
1015 #[test]
One operation mapping, owned by repos; billing reads it instead of keeping its own1016 fn raw_meters_are_weighted_by_the_repos_mapping() {
1017 // As repos' artifacts_usage answers: rows, and its operation_mapping.
1018 let body = json!({
1019 "rows": [
1020 { "day": "2026-10-15", "store": "g1t", "workspace": "Acme", "meter": "git.fetch", "count": 12, "bytes_in": 0, "bytes_out": 0 },
1021 { "day": "2026-10-15", "store": "g1t", "workspace": "acme", "meter": "git.receive_pack", "count": 3, "bytes_in": 0, "bytes_out": 0 },
1022 { "day": "2026-10-15", "store": "g1t", "workspace": "acme", "meter": "binding.read_blob", "count": 400, "bytes_in": 0, "bytes_out": 0 }
1023 ],
1024 "mapping": [
1025 { "meter": "git.fetch", "cost_operations": 1, "billable_operations": 1 },
1026 { "meter": "git.receive_pack", "cost_operations": 1, "billable_operations": 1 },
1027 { "meter": "binding.read_blob", "cost_operations": 0, "billable_operations": 0 }
1028 ],
1029 "truncated": false
1030 });
1031 let row = raw_usage_row(&body["rows"][0]).unwrap();
1032 assert_eq!(row, ("2026-10-15".into(), "acme".into(), "git_fetch".into(), 12.0));
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily1033 assert!(raw_usage_row(&json!({ "day": "2026-10-15", "count": 1 })).is_none());
One operation mapping, owned by repos; billing reads it instead of keeping its own1034 let mut mapping = operation_mapping(&body);
1035 let raw: Vec<(String, f64)> = body["rows"].as_array().unwrap().iter().filter_map(raw_usage_row).map(|r| (r.2, r.3)).collect();
1036 assert_eq!(weighted(&raw, &mapping, true), 15.0);
1037 assert_eq!(weighted(&raw, &mapping, false), 15.0);
1038 // Cloudflare turns out to bill binding reads: repos changes one row
1039 // (set_operation_mapping), and the bill g1t expects follows.
1040 mapping.insert("binding_read_blob".into(), (1.0, 0.0));
1041 assert_eq!(weighted(&raw, &mapping, true), 415.0);
1042 assert_eq!(weighted(&raw, &mapping, false), 15.0);
1043 assert!(operation_mapping(&json!({})).is_empty());
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily1044 }
1045
1046 #[test]
1047 fn the_most_specific_rule_claims_a_line() {
1048 let rule = |product: &str, meter: &str, bucket: &str| Rule {
1049 product: product.into(),
1050 meter: meter.into(),
1051 bucket: bucket.into(),
1052 price_meter: None,
1053 own_meter: None,
1054 drift_percent: 10.0,
1055 };
1056 let rules = vec![rule("workers", "*", "platform"), rule("workers", "workers_for_platforms", "deployments"), rule("artifacts", "*", "git")];
1057 assert_eq!(classify(&rules, "workers", "workers_for_platforms_cpu_ms").unwrap().bucket, "deployments");
1058 assert_eq!(classify(&rules, "workers", "workers_cpu_ms").unwrap().bucket, "platform");
1059 assert_eq!(classify(&rules, "artifacts", "artifacts_operations").unwrap().bucket, "git");
1060 assert!(classify(&rules, "browser_rendering", "browser_hours").is_none());
1061 }
1062}

This file's history is long; its oldest lines are credited to the oldest commit read.