g1t/services/billing/src/costs.rs

819 lines38,687 bytesCodeBlame
1//! What Cloudflare charges g1t, day by day and product by product.
2//!
3//! Prices are what g1t pays plus 20%, so g1t has to know what it pays, as
4//! Cloudflare counts it, not as g1t assumes it. Some of what g1t runs on
5//! is in beta (Artifacts bills "operations" from 2026-10-14 without saying
6//! exactly which calls are one), so nothing here assumes a product list:
7//! every line Cloudflare bills is kept, whatever it is, and how a line
8//! maps to what g1t sells is data (`cost_map`), changed without a deploy.
9//!
10//! Once a day (`keeper::DAILY`) billing reads:
11//!
12//! - **Billable usage** (`GET /accounts/{id}/billable-usage`, FOCUS
13//! columns): one row per service per day, with its quantity and what it
14//! cost g1t. Every Cloudflare product g1t uses appears here once it is
15//! used: Workers, Workers for Platforms, D1, KV, R2, Queues, Containers,
16//! Durable Objects, Artifacts, Browser Rendering, Workers AI, Vectorize,
17//! Cloudflare for SaaS and Email.
18//! - **Artifacts events** (GraphQL `artifactsEventsAdaptiveGroups`), by
19//! event type and day: what Artifacts itself counted, to compare with
20//! what g1t counted (`margin`).
21//!
22//! Each becomes cost lines in `cost_lines`, one per (day, source,
23//! product, meter), upserted, so reading a day again replaces it. The
24//! first run reads the last 31 days; after that the last few, since
25//! Cloudflare restates recent days as usage settles.
26//!
27//! The token is `CLOUDFLARE_BILLING_TOKEN` (Account: Billing Read and
28//! Account Analytics Read), or the keeper's `CLOUDFLARE_USAGE_TOKEN`,
29//! which has both. Without either, nothing is read and nothing fails.
30
31use std::collections::BTreeMap;
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/// The workspace a repository in the store belongs to: keys are
242/// `<workspace>--<repo>`; a pull request's fork (`pulls--<id>`) says none.
243pub(crate) fn workspace_of_store_key(key: &str) -> Option<String> {
244 let (workspace, rest) = key.split_once("--")?;
245 (!workspace.is_empty() && !rest.is_empty() && workspace != "pulls").then(|| workspace.to_lowercase())
246}
247
248/// Artifacts' billable operations per workspace and day, by repository
249/// name: how Cloudflare's own count shares out. The meter is
250/// `cloudflare_git`, which shares out the git bucket's cost (`margin`).
251pub(crate) fn artifacts_by_workspace(body: &Value) -> Vec<(String, String, f64)> {
252 let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"].as_array().cloned().unwrap_or_default();
253 let mut out: BTreeMap<(String, String), f64> = BTreeMap::new();
254 for g in &groups {
255 let (Some(day), Some(kind)) = (g["dimensions"]["date"].as_str(), g["dimensions"]["eventType"].as_str()) else { continue };
256 if day.len() < 10 || !ARTIFACTS_OPERATIONS.contains(&format!("events_{}", slug(kind)).as_str()) {
257 continue;
258 }
259 let Some(workspace) = g["dimensions"]["repositoryName"].as_str().and_then(workspace_of_store_key) else { continue };
260 *out.entry((day[..10].to_owned(), workspace)).or_default() += g["count"].as_f64().unwrap_or(0.0);
261 }
262 out.into_iter().map(|((day, workspace), count)| (day, workspace, count)).collect()
263}
264
265/// Artifacts' event types that are operations it bills: what its
266/// pricing names (create, push, pull, clone) and the metrics list.
267/// Errors such as `rateLimited` are not.
268pub(crate) const ARTIFACTS_OPERATIONS: [&str; 5] = ["events_create", "events_fork", "events_push", "events_pull", "events_delete"];
269
270/// The days to read: the 31 before today on the first run (`last` is
271/// None), else the last few, never further back than the backfill.
272pub(crate) fn window(last_fetched_day: Option<&str>, now_ms: u64) -> (String, String) {
273 let day = |ms: u64| g1t_contracts::time::rfc3339(ms)[..10].to_owned();
274 let today = day(now_ms);
275 let since = match last_fetched_day {
276 None => day(now_ms.saturating_sub(BACKFILL_DAYS * DAY_MS)),
277 Some(last) => {
278 let restate = day(now_ms.saturating_sub(RESTATE_DAYS * DAY_MS));
279 let floor = day(now_ms.saturating_sub(BACKFILL_DAYS * DAY_MS));
280 // A gap since the last run is read too, up to the backfill.
281 let from = if last < restate.as_str() { last.to_owned() } else { restate };
282 if from < floor { floor } else { from }
283 }
284 };
285 (since, today)
286}
287
288/// Every day from `since` to `until`, inclusive.
289pub(crate) fn days_between(since: &str, until: &str) -> Vec<String> {
290 let start = g1t_contracts::time::parse_rfc3339(&format!("{since}T00:00:00Z"));
291 let end = g1t_contracts::time::parse_rfc3339(&format!("{until}T00:00:00Z"));
292 match (start, end) {
293 (Some(start), Some(end)) if start <= end => (0..=((end - start) / DAY_MS))
294 .map(|n| g1t_contracts::time::rfc3339(start + n * DAY_MS)[..10].to_owned())
295 .collect(),
296 _ => Vec::new(),
297 }
298}
299
300/// A rule from `cost_map`: which Cloudflare lines feed which of g1t's
301/// products.
302#[derive(Clone, Debug, PartialEq, serde::Deserialize)]
303pub(crate) struct Rule {
304 /// Cloudflare's product, as slugged here.
305 pub product: String,
306 /// A meter prefix, or `*` for any meter of the product.
307 pub meter: String,
308 /// g1t's product it is a cost of: `sandboxes`, `git`, `platform`, …
309 pub bucket: String,
310 /// The price book meter whose cost it measures, if any.
311 pub price_meter: Option<String>,
312 /// g1t's own count of the same units, to compare quantities.
313 pub own_meter: Option<String>,
314 /// How far g1t's count may be from Cloudflare's before it is drift.
315 pub drift_percent: f64,
316}
317
318/// The rule for a line: its product's rule with the longest matching
319/// meter prefix, `*` last. None means no one decided what pays for it:
320/// a leak until someone does.
321pub(crate) fn classify<'a>(rules: &'a [Rule], product: &str, meter: &str) -> Option<&'a Rule> {
322 rules
323 .iter()
324 .filter(|r| r.product == product && (r.meter == "*" || meter.starts_with(r.meter.as_str())))
325 .max_by_key(|r| if r.meter == "*" { 0 } else { r.meter.len() + 1 })
326}
327
328/// What a bucket is called in sudo.
329pub(crate) fn bucket_title(bucket: &str) -> String {
330 match bucket {
331 "sandboxes" => "Sandboxes and builds".into(),
332 "deployments" => "Deployments".into(),
333 "git" => "Git operations".into(),
334 "repo_storage" => "Repository storage".into(),
335 "actions_cache" => "Actions cache".into(),
336 "embeddings" => "Search embeddings".into(),
337 "security" => "Security scans".into(),
338 "domains" => "Custom domains".into(),
339 "models" => "Models".into(),
340 "platform" => "Running g1t (paid by the plan)".into(),
341 UNMAPPED => "Not mapped".into(),
342 other => other.replace('_', " "),
343 }
344}
345
346/// The bucket of a line no rule claims.
347pub(crate) const UNMAPPED: &str = "unmapped";
348
349/// A day's count per workspace from counts "from this day to the end of
350/// its month" (what the repos service's `git_operations` answers with a
351/// `since`): each day's is its own less the next day's, within a month.
352/// The last day (today, so far) is its own.
353pub(crate) fn daily_from_cumulative(days: &[String], cumulative: &[BTreeMap<String, u64>]) -> Vec<(String, String, u64)> {
354 let mut out = Vec::new();
355 for (index, day) in days.iter().enumerate() {
356 let Some(today) = cumulative.get(index) else { break };
357 let next = days
358 .get(index + 1)
359 .filter(|next| next[..7] == day[..7])
360 .and_then(|_| cumulative.get(index + 1));
361 for (workspace, count) in today {
362 let later = next.and_then(|n| n.get(workspace)).copied().unwrap_or(0);
363 let own = count.saturating_sub(later);
364 if own > 0 {
365 out.push((day.clone(), workspace.clone(), own));
366 }
367 }
368 }
369 out
370}
371
372/// One row of the repos service's `artifacts_usage`, read leniently while
373/// its shape settles: a day, a workspace, a raw meter and a count.
374pub(crate) fn raw_usage_row(row: &Value) -> Option<(String, String, String, f64)> {
375 let day = text(row, &["day", "date"]);
376 let workspace = text(row, &["namespace", "workspace"]);
377 let meter = text(row, &["meter", "kind", "event", "operation"]);
378 let count = number(row, &["count", "quantity", "operations", "value"])?;
379 (day.len() >= 10 && !workspace.is_empty() && !meter.is_empty())
380 .then(|| (day[..10].to_owned(), workspace.to_lowercase(), slug(meter), count))
381}
382
383/// The repos service's operation mapping, as `artifacts_usage` returns it
384/// beside the rows: for each raw meter (slugged), how many operations it
385/// is to Cloudflare (`cost_operations`) and to the customer
386/// (`billable_operations`). Repos owns this mapping
387/// (`set_operation_mapping`); billing only reads it.
388pub(crate) fn operation_mapping(body: &Value) -> BTreeMap<String, (f64, f64)> {
389 body["mapping"]
390 .as_array()
391 .map(|rows| {
392 rows.iter()
393 .filter_map(|r| {
394 let meter = r["meter"].as_str()?;
395 Some((slug(meter), (r["cost_operations"].as_f64().unwrap_or(0.0), r["billable_operations"].as_f64().unwrap_or(0.0))))
396 })
397 .collect()
398 })
399 .unwrap_or_default()
400}
401
402/// Raw counts weighted by one column of the mapping: what Cloudflare
403/// should count (`cost`), or what customers are charged for.
404pub(crate) fn weighted(raw: &[(String, f64)], mapping: &BTreeMap<String, (f64, f64)>, cost: bool) -> f64 {
405 raw.iter()
406 .map(|(meter, count)| count * mapping.get(meter).map_or(0.0, |(c, b)| if cost { *c } else { *b }))
407 .sum()
408}
409
410impl Billing {
411 /// Reads Cloudflare's bill for the days due (see `window`) into
412 /// `cost_lines`: the days read and how many lines. None without a
413 /// token. What could not be read is added to `problems`.
414 pub(crate) async fn read_cloudflare(&self, keeper: &Keeper, problems: &mut Vec<String>) -> Result<Option<(String, String, u32)>> {
415 if !keeper.can_read_bill() {
416 return Ok(None);
417 }
418 #[derive(Deserialize)]
419 struct Last {
420 day: Option<String>,
421 }
422 let last = self
423 .db
424 .prepare("SELECT MAX(day) AS day FROM cost_lines WHERE source = ?")
425 .bind(&[SOURCE_BILLABLE.into()])?
426 .first::<Last>(None)
427 .await?
428 .and_then(|l| l.day);
429 let (since, until) = window(last.as_deref(), now_ms());
430 let fetched_at = rfc3339(now_ms());
431 let mut written = 0;
432 match keeper.billable_usage_body(&since, &until).await.map_err(|e| e.to_string()).and_then(|b| lines_from_billable(&b)) {
433 Ok(lines) => written += self.upsert_lines(&lines, &fetched_at).await?,
434 Err(error) => problems.push(format!("Cloudflare's billable usage could not be read: {error}")),
435 }
436 match keeper.graphql(artifacts_variables(keeper.account(), &since, &until)).await.map_err(|e| e.to_string()) {
437 Ok(body) => match lines_from_artifacts(&body) {
438 Ok(lines) => {
439 written += self.upsert_lines(&lines, &fetched_at).await?;
440 self.keep_cloudflare_counts(&since, &until, &artifacts_by_workspace(&body), &fetched_at).await?;
441 }
442 Err(error) => problems.push(format!("Artifacts events could not be read: {error}")),
443 },
444 Err(error) => problems.push(format!("Artifacts events could not be read: {error}")),
445 }
446 Ok(Some((since, until, written)))
447 }
448
449 /// Cloudflare's own per-workspace counts for the days, replacing what
450 /// was kept for them.
451 async fn keep_cloudflare_counts(&self, since: &str, until: &str, counts: &[(String, String, f64)], fetched_at: &str) -> Result<()> {
452 self.db
453 .prepare("DELETE FROM own_counts WHERE meter = 'cloudflare_git' AND day >= ?1 AND day <= ?2")
454 .bind(&[since.into(), until.into()])?
455 .run()
456 .await?;
457 for chunk in counts.chunks(50) {
458 let mut statements = Vec::with_capacity(chunk.len());
459 for (day, workspace, count) in chunk {
460 statements.push(
461 self.db
462 .prepare("INSERT OR REPLACE INTO own_counts (day, meter, workspace, quantity, fetched_at) VALUES (?, 'cloudflare_git', ?, ?, ?)")
463 .bind(&[day.as_str().into(), workspace.as_str().into(), (*count).into(), fetched_at.into()])?,
464 );
465 }
466 self.db.batch(statements).await?;
467 }
468 Ok(())
469 }
470
471 /// Upserts lines, a day's line replacing what was read for it before.
472 async fn upsert_lines(&self, lines: &[CostLine], fetched_at: &str) -> Result<u32> {
473 for chunk in lines.chunks(50) {
474 let mut statements = Vec::with_capacity(chunk.len());
475 for line in chunk {
476 statements.push(
477 self.db
478 .prepare(
479 "INSERT INTO cost_lines (day, source, product, meter, unit, quantity, cost_usd, raw_name, fetched_at)
480 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)
481 ON CONFLICT (day, source, product, meter) DO UPDATE SET
482 unit = ?5, quantity = ?6, cost_usd = ?7, raw_name = ?8, fetched_at = ?9",
483 )
484 .bind(&[
485 line.day.as_str().into(),
486 line.source.into(),
487 line.product.as_str().into(),
488 line.meter.as_str().into(),
489 line.unit.as_str().into(),
490 line.quantity.into(),
491 line.cost_usd.into(),
492 line.raw_name.as_str().into(),
493 fetched_at.into(),
494 ])?,
495 );
496 }
497 self.db.batch(statements).await?;
498 }
499 Ok(lines.len() as u32)
500 }
501
502 /// The mapping from Cloudflare's meters to g1t's products.
503 pub(crate) async fn rules(&self) -> Result<Vec<Rule>> {
504 #[derive(Deserialize)]
505 struct Row {
506 product: String,
507 meter: String,
508 bucket: String,
509 price_meter: Option<String>,
510 own_meter: Option<String>,
511 drift_percent: f64,
512 }
513 Ok(self
514 .db
515 .prepare("SELECT product, meter, bucket, price_meter, own_meter, drift_percent FROM cost_map")
516 .all()
517 .await?
518 .results::<Row>()?
519 .into_iter()
520 .map(|r| Rule {
521 product: r.product,
522 meter: r.meter,
523 bucket: r.bucket,
524 price_meter: r.price_meter,
525 own_meter: r.own_meter,
526 drift_percent: r.drift_percent,
527 })
528 .collect())
529 }
530
531 /// g1t's own counts for the days, all from the repos service, which
532 /// owns the mapping from raw meters to operations (`operation_mapping`):
533 ///
534 /// - `git_operations`: what customers are charged for, as repos counts
535 /// it (`git_operations`, already through its mapping).
536 /// - `cost_operations`: what g1t expects Cloudflare to bill, from the
537 /// raw meters (`artifacts_usage`) and the mapping's cost column.
538 /// - `artifacts_<meter>`: each raw meter.
539 pub(crate) async fn count_own(&self, since: &str, until: &str) -> Result<()> {
540 let Some(repos) = &self.repos else { return Ok(()) };
541 let days = days_between(since, until);
542 let fetched_at = rfc3339(now_ms());
543 let mut rows: Vec<(String, String, String, f64)> = Vec::new();
544
545 // Raw meters and repos' mapping; skipped while repos does not answer.
546 if let Ok(body) = g1t_kit::call::<_, Value>(repos, "artifacts_usage", &json!({ "from": since, "to": until })).await {
547 let mapping = operation_mapping(&body);
548 let mut by: BTreeMap<(String, String), Vec<(String, f64)>> = BTreeMap::new();
549 for (day, workspace, meter, count) in body["rows"].as_array().map(|l| l.iter().filter_map(raw_usage_row).collect::<Vec<_>>()).unwrap_or_default() {
550 by.entry((day, workspace)).or_default().push((meter, count));
551 }
552 for ((day, workspace), counts) in by {
553 // Several stores (namespaces) can give the same meter.
554 let mut merged: BTreeMap<String, f64> = BTreeMap::new();
555 for (meter, count) in &counts {
556 *merged.entry(meter.clone()).or_default() += count;
557 }
558 for (meter, count) in &merged {
559 rows.push((day.clone(), format!("artifacts_{meter}"), workspace.clone(), *count));
560 }
561 if !mapping.is_empty() {
562 rows.push((day, "cost_operations".to_owned(), workspace, weighted(&counts, &mapping, true)));
563 }
564 }
565 }
566 // What customers are charged for, as repos counts it through its mapping.
567 let mut cumulative = Vec::with_capacity(days.len());
568 for day in &days {
569 let list: Vec<WorkspaceGitOperations> = g1t_kit::call(
570 repos,
571 "git_operations",
572 &GitOperationsArgs { month: day[..7].to_owned(), since: Some(format!("{day}T00")), namespace: None },
573 )
574 .await?;
575 cumulative.push(list.into_iter().map(|w| (w.namespace.to_lowercase(), w.operations)).collect::<BTreeMap<_, _>>());
576 }
577 for (day, workspace, count) in daily_from_cumulative(&days, &cumulative) {
578 rows.push((day, "git_operations".to_owned(), workspace, count as f64));
579 }
580
581 // Each day's counts replace what was there.
582 self.db
583 .prepare("DELETE FROM own_counts WHERE day >= ?1 AND day <= ?2 AND meter NOT LIKE 'cloudflare_%'")
584 .bind(&[since.into(), until.into()])?
585 .run()
586 .await?;
587 for chunk in rows.chunks(50) {
588 let mut statements = Vec::with_capacity(chunk.len());
589 for (day, meter, workspace, quantity) in chunk {
590 statements.push(
591 self.db
592 .prepare(
593 "INSERT INTO own_counts (day, meter, workspace, quantity, fetched_at) VALUES (?1, ?2, ?3, ?4, ?5)
594 ON CONFLICT (day, meter, workspace) DO UPDATE SET quantity = ?4, fetched_at = ?5",
595 )
596 .bind(&[day.as_str().into(), meter.as_str().into(), workspace.as_str().into(), (*quantity).into(), fetched_at.as_str().into()])?,
597 );
598 }
599 self.db.batch(statements).await?;
600 }
601 Ok(())
602 }
603}
604
605#[cfg(test)]
606mod tests {
607 use super::*;
608
609 /// Billable usage as Cloudflare answered g1t on 2026-10-06 (two days,
610 /// trimmed), plus rows past the included amounts, which cost money.
611 fn billable_fixture() -> Value {
612 json!({
613 "success": true,
614 "errors": [],
615 "result": [
616 {
617 "ChargePeriodStart": "2026-10-03T00:00:00Z",
618 "ChargePeriodEnd": "2026-10-04T00:00:00Z",
619 "ServiceFamilyName": "Containers",
620 "ServiceName": "Container Memory, per GiB-Second (First 25 GiB-hours included)",
621 "PricingUnit": "Count",
622 "PricingQuantity": "14000",
623 "ContractedCost": 0, "BilledCost": 0, "ListCost": 0
624 },
625 {
626 "ChargePeriodStart": "2026-10-03T00:00:00Z",
627 "ServiceFamilyName": "D1",
628 "ServiceName": "D1 - Rows Read (first 25 billion included)",
629 "PricingUnit": "Count",
630 "PricingQuantity": 1200000,
631 "BilledCost": 0
632 },
633 {
634 "ChargePeriodStart": "2026-10-15T00:00:00Z",
635 "ServiceFamilyName": "Artifacts",
636 "ServiceName": "Artifacts Operations (First 10,000 included)",
637 "PricingUnit": "Count",
638 "PricingQuantity": "30000",
639 "BilledCost": "3.00",
640 "ListCost": 4.5
641 },
642 {
643 "ChargePeriodStart": "2026-10-15T00:00:00Z",
644 "ServiceFamilyName": "Artifacts",
645 "ServiceName": "Artifacts Operations (First 10,000 included)",
646 "PricingUnit": "Count",
647 "PricingQuantity": "10000",
648 "BilledCost": "1.50"
649 },
650 {
651 "ChargePeriodStart": "2026-10-15T00:00:00Z",
652 "ServiceFamilyName": "Workers",
653 "ServiceName": "Workers for Platforms Requests (First 20M are included)",
654 "PricingUnit": "Count",
655 "PricingQuantity": 25000000,
656 "ListCost": 1.5
657 },
658 { "ServiceName": "No day, no line" }
659 ]
660 })
661 }
662
663 #[test]
664 fn names_become_products_and_meters() {
665 assert_eq!(slug("Workers for Platforms CPU ms (First 60M ms are included)"), "workers_for_platforms_cpu_ms");
666 assert_eq!(product_and_meter("D1", "D1 - Rows Read (first 25 billion included)"), ("d1".into(), "d1_rows_read".into()));
667 assert_eq!(product_and_meter("Workers KV", "KV Read Operations (First 10M is included)"), ("workers_kv".into(), "kv_read_operations".into()));
668 assert_eq!(product_and_meter("", "Containers / Container vCPU"), ("containers".into(), "container_vcpu".into()));
669 assert_eq!(product_and_meter("Email", ""), ("email".into(), "usage".into()));
670 }
671
672 #[test]
673 fn billable_usage_is_parsed_and_rows_of_one_meter_are_added() {
674 let lines = lines_from_billable(&billable_fixture()).unwrap();
675 assert_eq!(lines.len(), 4, "{lines:#?}");
676 let artifacts = lines.iter().find(|l| l.product == "artifacts").unwrap();
677 // Two rows of the same day and meter: one line, added up.
678 assert_eq!(artifacts.meter, "artifacts_operations");
679 assert_eq!(artifacts.quantity, 40_000.0);
680 assert!((artifacts.cost_usd - 4.5).abs() < 1e-9, "billed, not list: {}", artifacts.cost_usd);
681 let wfp = lines.iter().find(|l| l.meter.starts_with("workers_for_platforms")).unwrap();
682 assert_eq!(wfp.product, "workers");
683 // No billed cost: the list cost.
684 assert_eq!(wfp.cost_usd, 1.5);
685 let memory = lines.iter().find(|l| l.product == "containers").unwrap();
686 assert_eq!((memory.day.as_str(), memory.quantity, memory.cost_usd), ("2026-10-03", 14_000.0, 0.0));
687 assert!(lines_from_billable(&json!({ "success": false, "errors": [{ "code": 10000 }] })).is_err());
688 }
689
690 #[test]
691 fn reading_the_same_days_twice_gives_the_same_lines() {
692 // Idempotent: the same answer aggregates to the same keys and
693 // amounts, so the upsert replaces rather than adds.
694 let once = lines_from_billable(&billable_fixture()).unwrap();
695 let twice = lines_from_billable(&billable_fixture()).unwrap();
696 assert_eq!(once, twice);
697 let again = aggregate(once.clone());
698 assert_eq!(again, once);
699 }
700
701 #[test]
702 fn artifacts_events_are_counted_by_type_and_day() {
703 let body = json!({
704 "data": { "viewer": { "accounts": [{ "artifactsEventsAdaptiveGroups": [
705 { "count": 120, "dimensions": { "date": "2026-10-05", "eventType": "pull" } },
706 { "count": 30, "dimensions": { "date": "2026-10-05", "eventType": "push" } },
707 { "count": 2, "dimensions": { "date": "2026-10-05", "eventType": "rateLimited" } },
708 { "count": 5, "dimensions": { "date": "2026-10-06", "eventType": "pull" } }
709 ] }] } },
710 "errors": null
711 });
712 let lines = lines_from_artifacts(&body).unwrap();
713 assert_eq!(lines.len(), 4);
714 assert_eq!(lines[0].meter, "events_pull");
715 assert_eq!(lines[0].quantity, 120.0);
716 assert!(lines.iter().any(|l| l.meter == "events_ratelimited"));
717 // By workspace, from the repository's store key; forks say none.
718 let by_repo = json!({ "data": { "viewer": { "accounts": [{ "artifactsEventsAdaptiveGroups": [
719 { "count": 100, "dimensions": { "date": "2026-10-05", "eventType": "pull", "repositoryName": "acme--api" } },
720 { "count": 20, "dimensions": { "date": "2026-10-05", "eventType": "push", "repositoryName": "acme--web" } },
721 { "count": 7, "dimensions": { "date": "2026-10-05", "eventType": "fork", "repositoryName": "pulls--123" } },
722 { "count": 3, "dimensions": { "date": "2026-10-05", "eventType": "serverError", "repositoryName": "beta--x" } }
723 ] }] } } });
724 assert_eq!(artifacts_by_workspace(&by_repo), vec![("2026-10-05".to_string(), "acme".to_string(), 120.0)]);
725 assert_eq!(workspace_of_store_key("Acme--api"), Some("acme".into()));
726 assert_eq!(workspace_of_store_key("pulls--9"), None);
727 assert_eq!(workspace_of_store_key("plain"), None);
728 let operations: f64 = lines
729 .iter()
730 .filter(|l| l.day == "2026-10-05" && ARTIFACTS_OPERATIONS.contains(&l.meter.as_str()))
731 .map(|l| l.quantity)
732 .sum();
733 assert_eq!(operations, 150.0);
734 let failed = json!({ "data": null, "errors": [{ "message": "unknown field" }] });
735 assert!(lines_from_artifacts(&failed).is_err());
736 }
737
738 #[test]
739 fn the_first_run_backfills_31_days_and_later_runs_restate_a_few() {
740 let now = g1t_contracts::time::parse_rfc3339("2026-10-06T04:17:00Z").unwrap();
741 assert_eq!(window(None, now), ("2026-09-05".into(), "2026-10-06".into()));
742 assert_eq!(window(Some("2026-10-06"), now), ("2026-10-02".into(), "2026-10-06".into()));
743 // A week without a run: from the last day read.
744 assert_eq!(window(Some("2026-09-28"), now), ("2026-09-28".into(), "2026-10-06".into()));
745 // Never past the backfill.
746 assert_eq!(window(Some("2026-01-01"), now), ("2026-09-05".into(), "2026-10-06".into()));
747 assert_eq!(days_between("2026-09-29", "2026-10-02"), vec!["2026-09-29", "2026-09-30", "2026-10-01", "2026-10-02"]);
748 assert!(days_between("2026-10-02", "2026-10-01").is_empty());
749 }
750
751 #[test]
752 fn a_days_git_operations_are_its_count_less_the_next_days() {
753 let days: Vec<String> = ["2026-09-29", "2026-09-30", "2026-10-01", "2026-10-02"].iter().map(|d| d.to_string()).collect();
754 let map = |pairs: &[(&str, u64)]| pairs.iter().map(|(w, n)| (w.to_string(), *n)).collect::<BTreeMap<_, _>>();
755 // From each day to the end of its month.
756 let cumulative = vec![map(&[("acme", 30), ("beta", 4)]), map(&[("acme", 10)]), map(&[("acme", 7)]), map(&[("acme", 2)])];
757 let daily = daily_from_cumulative(&days, &cumulative);
758 assert_eq!(
759 daily,
760 vec![
761 ("2026-09-29".into(), "acme".into(), 20),
762 ("2026-09-29".into(), "beta".into(), 4),
763 // The month's last day is its own.
764 ("2026-09-30".into(), "acme".into(), 10),
765 ("2026-10-01".into(), "acme".into(), 5),
766 // Today so far.
767 ("2026-10-02".into(), "acme".into(), 2),
768 ]
769 );
770 }
771
772 #[test]
773 fn raw_meters_are_weighted_by_the_repos_mapping() {
774 // As repos' artifacts_usage answers: rows, and its operation_mapping.
775 let body = json!({
776 "rows": [
777 { "day": "2026-10-15", "store": "g1t", "workspace": "Acme", "meter": "git.fetch", "count": 12, "bytes_in": 0, "bytes_out": 0 },
778 { "day": "2026-10-15", "store": "g1t", "workspace": "acme", "meter": "git.receive_pack", "count": 3, "bytes_in": 0, "bytes_out": 0 },
779 { "day": "2026-10-15", "store": "g1t", "workspace": "acme", "meter": "binding.read_blob", "count": 400, "bytes_in": 0, "bytes_out": 0 }
780 ],
781 "mapping": [
782 { "meter": "git.fetch", "cost_operations": 1, "billable_operations": 1 },
783 { "meter": "git.receive_pack", "cost_operations": 1, "billable_operations": 1 },
784 { "meter": "binding.read_blob", "cost_operations": 0, "billable_operations": 0 }
785 ],
786 "truncated": false
787 });
788 let row = raw_usage_row(&body["rows"][0]).unwrap();
789 assert_eq!(row, ("2026-10-15".into(), "acme".into(), "git_fetch".into(), 12.0));
790 assert!(raw_usage_row(&json!({ "day": "2026-10-15", "count": 1 })).is_none());
791 let mut mapping = operation_mapping(&body);
792 let raw: Vec<(String, f64)> = body["rows"].as_array().unwrap().iter().filter_map(raw_usage_row).map(|r| (r.2, r.3)).collect();
793 assert_eq!(weighted(&raw, &mapping, true), 15.0);
794 assert_eq!(weighted(&raw, &mapping, false), 15.0);
795 // Cloudflare turns out to bill binding reads: repos changes one row
796 // (set_operation_mapping), and the bill g1t expects follows.
797 mapping.insert("binding_read_blob".into(), (1.0, 0.0));
798 assert_eq!(weighted(&raw, &mapping, true), 415.0);
799 assert_eq!(weighted(&raw, &mapping, false), 15.0);
800 assert!(operation_mapping(&json!({})).is_empty());
801 }
802
803 #[test]
804 fn the_most_specific_rule_claims_a_line() {
805 let rule = |product: &str, meter: &str, bucket: &str| Rule {
806 product: product.into(),
807 meter: meter.into(),
808 bucket: bucket.into(),
809 price_meter: None,
810 own_meter: None,
811 drift_percent: 10.0,
812 };
813 let rules = vec![rule("workers", "*", "platform"), rule("workers", "workers_for_platforms", "deployments"), rule("artifacts", "*", "git")];
814 assert_eq!(classify(&rules, "workers", "workers_for_platforms_cpu_ms").unwrap().bucket, "deployments");
815 assert_eq!(classify(&rules, "workers", "workers_cpu_ms").unwrap().bucket, "platform");
816 assert_eq!(classify(&rules, "artifacts", "artifacts_operations").unwrap().bucket, "git");
817 assert!(classify(&rules, "browser_rendering", "browser_hours").is_none());
818 }
819}