g1t/services/billing/src/costs.rs

862 lines40,729 bytesCodeBlame
1//! What Cloudflare charges g1t, day by day and product by product.
2//!
3//! Prices are what g1t pays plus 20%, so g1t has to know what it pays, as
4//! Cloudflare counts it, not as g1t assumes it. Some of what g1t runs on
5//! is in beta (Artifacts bills "operations" from 2026-10-14 without saying
6//! exactly which calls are one), so nothing here assumes a product list:
7//! every line Cloudflare bills is kept, whatever it is, and how a line
8//! maps to what g1t sells is data (`cost_map`), changed without a deploy.
9//!
10//! Once a day (`keeper::DAILY`) billing reads:
11//!
12//! - **Billable usage** (`GET /accounts/{id}/billable-usage`, FOCUS
13//! columns): one row per service per day, with its quantity and what it
14//! cost g1t. Every Cloudflare product g1t uses appears here once it is
15//! used: Workers, Workers for Platforms, D1, KV, R2, Queues, Containers,
16//! Durable Objects, Artifacts, Browser Rendering, Workers AI, Vectorize,
17//! Cloudflare for SaaS and Email.
18//! - **Artifacts events** (GraphQL `artifactsEventsAdaptiveGroups`), by
19//! event type and day: what Artifacts itself counted, to compare with
20//! what g1t counted (`margin`).
21//!
22//! Each becomes cost lines in `cost_lines`, one per (day, source,
23//! product, meter), upserted, so reading a day again replaces it. The
24//! first run reads the last 31 days; after that the last few, since
25//! Cloudflare restates recent days as usage settles.
26//!
27//! The token is `CLOUDFLARE_BILLING_TOKEN` (Account: Billing Read and
28//! Account Analytics Read), or the keeper's `CLOUDFLARE_USAGE_TOKEN`,
29//! which has both. Without either, nothing is read and nothing fails.
30
31use std::collections::{BTreeMap, BTreeSet};
32
33use g1t_contracts::repos::{GitOperationsArgs, WorkspaceGitOperations};
34use g1t_contracts::time::rfc3339;
35use g1t_kit::now_ms;
36use serde::Deserialize;
37use serde_json::{Value, json};
38use worker::Result;
39
40use crate::Billing;
41use crate::keeper::Keeper;
42
43/// Where a cost line came from.
44pub(crate) const SOURCE_BILLABLE: &str = "billable_usage";
45pub(crate) const SOURCE_ARTIFACTS: &str = "artifacts_events";
46
47/// How far back the first run reads: the 31 days GraphQL keeps.
48pub(crate) const BACKFILL_DAYS: u64 = 31;
49/// How many recent days every later run reads again.
50pub(crate) const RESTATE_DAYS: u64 = 4;
51
52pub(crate) const DAY_MS: u64 = 24 * 60 * 60 * 1000;
53
54/// One day of one meter of one Cloudflare product.
55#[derive(Clone, Debug, PartialEq)]
56pub(crate) struct CostLine {
57 /// YYYY-MM-DD, UTC.
58 pub day: String,
59 pub source: &'static str,
60 /// `containers`, `workers`, `workers_kv`, `artifacts`, … from
61 /// Cloudflare's own family name.
62 pub product: String,
63 /// The service within it, such as `container_memory_per_gib_second`
64 /// or `events_push`.
65 pub meter: String,
66 pub unit: String,
67 pub quantity: f64,
68 /// What g1t pays, in dollars: contracted, billed, or list.
69 pub cost_usd: f64,
70 /// The name as Cloudflare gave it, for people.
71 pub raw_name: String,
72}
73
74/// `Workers for Platforms CPU ms (First 60M ms are included)` →
75/// `workers_for_platforms_cpu_ms`: lower case, words joined by `_`, and
76/// what is in parentheses (the included amount, which changes) left out.
77pub(crate) fn slug(text: &str) -> String {
78 let mut out = String::new();
79 let mut depth = 0u32;
80 let mut gap = false;
81 for c in text.chars() {
82 match c {
83 '(' => depth += 1,
84 ')' => depth = depth.saturating_sub(1),
85 _ if depth > 0 => {}
86 c if c.is_ascii_alphanumeric() => {
87 if gap && !out.is_empty() {
88 out.push('_');
89 }
90 gap = false;
91 out.push(c.to_ascii_lowercase());
92 }
93 _ => gap = true,
94 }
95 }
96 out
97}
98
99/// The product and meter for a family and service name, such as
100/// (`D1`, `D1 - Rows Read (first 25 billion included)`) → (`d1`,
101/// `d1_rows_read`). Without a family, the service's first word.
102pub(crate) fn product_and_meter(family: &str, service: &str) -> (String, String) {
103 let (family, service) = if family.trim().is_empty() {
104 match service.split_once(" / ") {
105 Some((family, service)) => (family, service),
106 None => (service.split_whitespace().next().unwrap_or("other"), service),
107 }
108 } else {
109 (family, service)
110 };
111 let product = slug(family);
112 let product = if product.is_empty() { "other".to_owned() } else { product };
113 // Kept whole: `Workers for Platforms Requests` under `Workers` must stay
114 // `workers_for_platforms_requests`.
115 let meter = slug(service);
116 let meter = if meter.is_empty() { "usage".to_owned() } else { meter };
117 (product, meter)
118}
119
120fn text<'a>(row: &'a Value, keys: &[&str]) -> &'a str {
121 keys.iter().find_map(|k| row[*k].as_str()).unwrap_or_default()
122}
123
124fn number(row: &Value, keys: &[&str]) -> Option<f64> {
125 keys.iter()
126 .find_map(|k| row[*k].as_f64().or_else(|| row[*k].as_str().and_then(|s| s.trim().parse().ok())))
127}
128
129/// One row of billable usage, read leniently: the API is new, and its
130/// field names are FOCUS's (in either case). None without a service or a
131/// day.
132pub(crate) fn line_from_focus(row: &Value) -> Option<CostLine> {
133 let service = text(row, &["ServiceName", "service_name", "service"]);
134 let day = text(row, &["ChargePeriodStart", "charge_period_start", "UsageDate", "date"]);
135 if service.is_empty() || day.len() < 10 {
136 return None;
137 }
138 let family = text(row, &["ServiceFamilyName", "service_family_name", "ServiceCategory"]);
139 let (product, meter) = product_and_meter(family, service);
140 // What g1t pays: contracted, else billed, else list.
141 let cost = [
142 &["ContractedCost", "contracted_cost"][..],
143 &["BilledCost", "billed_cost"][..],
144 &["EffectiveCost", "effective_cost"][..],
145 &["ListCost", "list_cost"][..],
146 ]
147 .iter()
148 .find_map(|keys| number(row, keys).filter(|c| *c > 0.0))
149 .unwrap_or(0.0);
150 Some(CostLine {
151 day: day[..10].to_owned(),
152 source: SOURCE_BILLABLE,
153 product,
154 meter,
155 unit: text(row, &["PricingUnit", "pricing_unit", "ConsumedUnit", "consumed_unit"]).to_owned(),
156 quantity: number(row, &["PricingQuantity", "pricing_quantity", "ConsumedQuantity", "consumed_quantity"]).unwrap_or(0.0),
157 cost_usd: cost,
158 raw_name: if family.is_empty() { service.to_owned() } else { format!("{family} / {service}") },
159 })
160}
161
162/// Lines with the same key added together, in key order. Cloudflare can
163/// give one service several rows a day (regions, tiers); an upsert of
164/// each would keep only the last.
165pub(crate) fn aggregate(lines: Vec<CostLine>) -> Vec<CostLine> {
166 let mut out: Vec<CostLine> = Vec::new();
167 for line in lines {
168 match out
169 .iter_mut()
170 .find(|l| l.day == line.day && l.source == line.source && l.product == line.product && l.meter == line.meter)
171 {
172 Some(existing) => {
173 existing.quantity += line.quantity;
174 existing.cost_usd += line.cost_usd;
175 }
176 None => out.push(line),
177 }
178 }
179 out.sort_by(|a, b| (&a.day, a.source, &a.product, &a.meter).cmp(&(&b.day, b.source, &b.product, &b.meter)));
180 out
181}
182
183/// Every line of a billable-usage answer, aggregated. Errors when
184/// Cloudflare says it failed.
185pub(crate) fn lines_from_billable(body: &Value) -> std::result::Result<Vec<CostLine>, String> {
186 if body["success"] == Value::Bool(false) {
187 return Err(format!("billable usage failed: {}", body["errors"]));
188 }
189 let rows = body["result"].as_array().cloned().unwrap_or_default();
190 Ok(aggregate(rows.iter().filter_map(line_from_focus).collect()))
191}
192
193/// The GraphQL query for Artifacts' events by type and day.
194pub(crate) const ARTIFACTS_QUERY: &str = "query ($account: String!, $since: Date!, $until: Date!) {
195 viewer { accounts(filter: { accountTag: $account }) {
196 artifactsEventsAdaptiveGroups(limit: 10000, filter: { date_geq: $since, date_leq: $until }) {
197 count
198 dimensions { date eventType repositoryName }
199 }
200 } }
201}";
202
203pub(crate) fn artifacts_variables(account: &str, since: &str, until: &str) -> Value {
204 json!({ "query": ARTIFACTS_QUERY, "variables": { "account": account, "since": since, "until": until } })
205}
206
207/// Artifacts' events as lines (`events_push`, `events_pull`, …), with no
208/// cost: billable usage carries the cost. Errors when GraphQL does.
209pub(crate) fn lines_from_artifacts(body: &Value) -> std::result::Result<Vec<CostLine>, String> {
210 if let Some(errors) = body["errors"].as_array().filter(|e| !e.is_empty()) {
211 return Err(format!("Artifacts events failed: {}", Value::Array(errors.clone())));
212 }
213 let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"]
214 .as_array()
215 .cloned()
216 .unwrap_or_default();
217 Ok(aggregate(
218 groups
219 .iter()
220 .filter_map(|g| {
221 let day = g["dimensions"]["date"].as_str()?;
222 let kind = g["dimensions"]["eventType"].as_str()?;
223 if day.len() < 10 {
224 return None;
225 }
226 Some(CostLine {
227 day: day[..10].to_owned(),
228 source: SOURCE_ARTIFACTS,
229 product: "artifacts".to_owned(),
230 meter: format!("events_{}", slug(kind)),
231 unit: "events".to_owned(),
232 quantity: g["count"].as_f64().unwrap_or(0.0),
233 cost_usd: 0.0,
234 raw_name: format!("Artifacts / {kind}"),
235 })
236 })
237 .collect(),
238 ))
239}
240
241/// The workspace a repository in the store belongs to: keys are
242/// `<workspace>--<repo>`; a pull request's working copy (`pulls--<id>`) is
243/// its repository's workspace's, from `owners` (repos' `pull_owners`), else
244/// no one's.
245pub(crate) fn workspace_of_store_key(key: &str, owners: &BTreeMap<String, String>) -> Option<String> {
246 let (workspace, rest) = key.split_once("--")?;
247 if workspace.is_empty() || rest.is_empty() {
248 return None;
249 }
250 if workspace == "pulls" {
251 return owners.get(rest).map(|w| w.to_lowercase());
252 }
253 Some(workspace.to_lowercase())
254}
255
256/// The pull request ids whose working copies Artifacts counted events for.
257pub(crate) fn pull_ids(body: &Value) -> Vec<String> {
258 let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"].as_array().cloned().unwrap_or_default();
259 let ids: BTreeSet<String> = groups
260 .iter()
261 .filter_map(|g| g["dimensions"]["repositoryName"].as_str()?.strip_prefix("pulls--").map(str::to_owned))
262 .filter(|id| !id.is_empty())
263 .collect();
264 ids.into_iter().collect()
265}
266
267/// Artifacts' billable operations per workspace and day, by repository
268/// name: how Cloudflare's own count shares out. The meter is
269/// `cloudflare_git`, which shares out the git bucket's cost (`margin`).
270pub(crate) fn artifacts_by_workspace(body: &Value, owners: &BTreeMap<String, String>) -> Vec<(String, String, f64)> {
271 let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"].as_array().cloned().unwrap_or_default();
272 let mut out: BTreeMap<(String, String), f64> = BTreeMap::new();
273 for g in &groups {
274 let (Some(day), Some(kind)) = (g["dimensions"]["date"].as_str(), g["dimensions"]["eventType"].as_str()) else { continue };
275 if day.len() < 10 || !ARTIFACTS_OPERATIONS.contains(&format!("events_{}", slug(kind)).as_str()) {
276 continue;
277 }
278 let Some(workspace) = g["dimensions"]["repositoryName"].as_str().and_then(|key| workspace_of_store_key(key, owners)) else { continue };
279 *out.entry((day[..10].to_owned(), workspace)).or_default() += g["count"].as_f64().unwrap_or(0.0);
280 }
281 out.into_iter().map(|((day, workspace), count)| (day, workspace, count)).collect()
282}
283
284/// Artifacts' event types that are operations it bills: what its
285/// pricing names (create, push, pull, clone) and the metrics list.
286/// Errors such as `rateLimited` are not.
287pub(crate) const ARTIFACTS_OPERATIONS: [&str; 5] = ["events_create", "events_fork", "events_push", "events_pull", "events_delete"];
288
289/// The days to read: the 31 before today on the first run (`last` is
290/// None), else the last few, never further back than the backfill.
291pub(crate) fn window(last_fetched_day: Option<&str>, now_ms: u64) -> (String, String) {
292 let day = |ms: u64| g1t_contracts::time::rfc3339(ms)[..10].to_owned();
293 let today = day(now_ms);
294 let since = match last_fetched_day {
295 None => day(now_ms.saturating_sub(BACKFILL_DAYS * DAY_MS)),
296 Some(last) => {
297 let restate = day(now_ms.saturating_sub(RESTATE_DAYS * DAY_MS));
298 let floor = day(now_ms.saturating_sub(BACKFILL_DAYS * DAY_MS));
299 // A gap since the last run is read too, up to the backfill.
300 let from = if last < restate.as_str() { last.to_owned() } else { restate };
301 if from < floor { floor } else { from }
302 }
303 };
304 (since, today)
305}
306
307/// Every day from `since` to `until`, inclusive.
308pub(crate) fn days_between(since: &str, until: &str) -> Vec<String> {
309 let start = g1t_contracts::time::parse_rfc3339(&format!("{since}T00:00:00Z"));
310 let end = g1t_contracts::time::parse_rfc3339(&format!("{until}T00:00:00Z"));
311 match (start, end) {
312 (Some(start), Some(end)) if start <= end => (0..=((end - start) / DAY_MS))
313 .map(|n| g1t_contracts::time::rfc3339(start + n * DAY_MS)[..10].to_owned())
314 .collect(),
315 _ => Vec::new(),
316 }
317}
318
319/// A rule from `cost_map`: which Cloudflare lines feed which of g1t's
320/// products.
321#[derive(Clone, Debug, PartialEq, serde::Deserialize)]
322pub(crate) struct Rule {
323 /// Cloudflare's product, as slugged here.
324 pub product: String,
325 /// A meter prefix, or `*` for any meter of the product.
326 pub meter: String,
327 /// g1t's product it is a cost of: `sandboxes`, `git`, `platform`, …
328 pub bucket: String,
329 /// The price book meter whose cost it measures, if any.
330 pub price_meter: Option<String>,
331 /// g1t's own count of the same units, to compare quantities.
332 pub own_meter: Option<String>,
333 /// How far g1t's count may be from Cloudflare's before it is drift.
334 pub drift_percent: f64,
335}
336
337/// The rule for a line: its product's rule with the longest matching
338/// meter prefix, `*` last. None means no one decided what pays for it:
339/// a leak until someone does.
340pub(crate) fn classify<'a>(rules: &'a [Rule], product: &str, meter: &str) -> Option<&'a Rule> {
341 rules
342 .iter()
343 .filter(|r| r.product == product && (r.meter == "*" || meter.starts_with(r.meter.as_str())))
344 .max_by_key(|r| if r.meter == "*" { 0 } else { r.meter.len() + 1 })
345}
346
347/// What a bucket is called in sudo.
348pub(crate) fn bucket_title(bucket: &str) -> String {
349 match bucket {
350 "sandboxes" => "Sandboxes and builds".into(),
351 "deployments" => "Deployments".into(),
352 "git" => "Git operations".into(),
353 "repo_storage" => "Repository storage".into(),
354 "actions_cache" => "Actions cache".into(),
355 "embeddings" => "Search embeddings".into(),
356 "security" => "Security scans".into(),
357 "domains" => "Custom domains".into(),
358 "models" => "Models".into(),
359 "platform" => "Running g1t (paid by the plan)".into(),
360 UNMAPPED => "Not mapped".into(),
361 other => other.replace('_', " "),
362 }
363}
364
365/// The bucket of a line no rule claims.
366pub(crate) const UNMAPPED: &str = "unmapped";
367
368/// A day's count per workspace from counts "from this day to the end of
369/// its month" (what the repos service's `git_operations` answers with a
370/// `since`): each day's is its own less the next day's, within a month.
371/// The last day (today, so far) is its own.
372pub(crate) fn daily_from_cumulative(days: &[String], cumulative: &[BTreeMap<String, u64>]) -> Vec<(String, String, u64)> {
373 let mut out = Vec::new();
374 for (index, day) in days.iter().enumerate() {
375 let Some(today) = cumulative.get(index) else { break };
376 let next = days
377 .get(index + 1)
378 .filter(|next| next[..7] == day[..7])
379 .and_then(|_| cumulative.get(index + 1));
380 for (workspace, count) in today {
381 let later = next.and_then(|n| n.get(workspace)).copied().unwrap_or(0);
382 let own = count.saturating_sub(later);
383 if own > 0 {
384 out.push((day.clone(), workspace.clone(), own));
385 }
386 }
387 }
388 out
389}
390
391/// One row of the repos service's `artifacts_usage`, read leniently while
392/// its shape settles: a day, a workspace, a raw meter and a count.
393pub(crate) fn raw_usage_row(row: &Value) -> Option<(String, String, String, f64)> {
394 let day = text(row, &["day", "date"]);
395 let workspace = text(row, &["namespace", "workspace"]);
396 let meter = text(row, &["meter", "kind", "event", "operation"]);
397 let count = number(row, &["count", "quantity", "operations", "value"])?;
398 (day.len() >= 10 && !workspace.is_empty() && !meter.is_empty())
399 .then(|| (day[..10].to_owned(), workspace.to_lowercase(), slug(meter), count))
400}
401
402/// The repos service's operation mapping, as `artifacts_usage` returns it
403/// beside the rows: for each raw meter (slugged), how many operations it
404/// is to Cloudflare (`cost_operations`) and to the customer
405/// (`billable_operations`). Repos owns this mapping
406/// (`set_operation_mapping`); billing only reads it.
407pub(crate) fn operation_mapping(body: &Value) -> BTreeMap<String, (f64, f64)> {
408 body["mapping"]
409 .as_array()
410 .map(|rows| {
411 rows.iter()
412 .filter_map(|r| {
413 let meter = r["meter"].as_str()?;
414 Some((slug(meter), (r["cost_operations"].as_f64().unwrap_or(0.0), r["billable_operations"].as_f64().unwrap_or(0.0))))
415 })
416 .collect()
417 })
418 .unwrap_or_default()
419}
420
421/// Raw counts weighted by one column of the mapping: what Cloudflare
422/// should count (`cost`), or what customers are charged for.
423pub(crate) fn weighted(raw: &[(String, f64)], mapping: &BTreeMap<String, (f64, f64)>, cost: bool) -> f64 {
424 raw.iter()
425 .map(|(meter, count)| count * mapping.get(meter).map_or(0.0, |(c, b)| if cost { *c } else { *b }))
426 .sum()
427}
428
429impl Billing {
430 /// Reads Cloudflare's bill for the days due (see `window`) into
431 /// `cost_lines`: the days read and how many lines. None without a
432 /// token. What could not be read is added to `problems`.
433 pub(crate) async fn read_cloudflare(&self, keeper: &Keeper, problems: &mut Vec<String>) -> Result<Option<(String, String, u32)>> {
434 if !keeper.can_read_bill() {
435 return Ok(None);
436 }
437 #[derive(Deserialize)]
438 struct Last {
439 day: Option<String>,
440 }
441 let last = self
442 .db
443 .prepare("SELECT MAX(day) AS day FROM cost_lines WHERE source = ?")
444 .bind(&[SOURCE_BILLABLE.into()])?
445 .first::<Last>(None)
446 .await?
447 .and_then(|l| l.day);
448 let (since, until) = window(last.as_deref(), now_ms());
449 let fetched_at = rfc3339(now_ms());
450 let mut written = 0;
451 match keeper.billable_usage_body(&since, &until).await.map_err(|e| e.to_string()).and_then(|b| lines_from_billable(&b)) {
452 Ok(lines) => written += self.upsert_lines(&lines, &fetched_at).await?,
453 Err(error) => problems.push(format!("Cloudflare's billable usage could not be read: {error}")),
454 }
455 match keeper.graphql(artifacts_variables(keeper.account(), &since, &until)).await.map_err(|e| e.to_string()) {
456 Ok(body) => match lines_from_artifacts(&body) {
457 Ok(lines) => {
458 written += self.upsert_lines(&lines, &fetched_at).await?;
459 let owners = self.pull_owners(&pull_ids(&body)).await;
460 self.keep_cloudflare_counts(&since, &until, &artifacts_by_workspace(&body, &owners), &fetched_at).await?;
461 }
462 Err(error) => problems.push(format!("Artifacts events could not be read: {error}")),
463 },
464 Err(error) => problems.push(format!("Artifacts events could not be read: {error}")),
465 }
466 Ok(Some((since, until, written)))
467 }
468
469 /// Cloudflare's own per-workspace counts for the days, replacing what
470 /// was kept for them.
471 /// The workspace of each pull request's working copy, from repos;
472 /// nothing while repos does not answer (those events stay no one's).
473 async fn pull_owners(&self, pulls: &[String]) -> BTreeMap<String, String> {
474 let Some(repos) = &self.repos else { return BTreeMap::new() };
475 if pulls.is_empty() {
476 return BTreeMap::new();
477 }
478 match g1t_kit::call::<_, Value>(repos, "pull_owners", &json!({ "pulls": pulls })).await {
479 Ok(body) => body["owners"]
480 .as_object()
481 .map(|o| o.iter().filter_map(|(k, v)| Some((k.clone(), v.as_str()?.to_owned()))).collect())
482 .unwrap_or_default(),
483 Err(error) => {
484 worker::console_error!("pull_owners: {error}");
485 BTreeMap::new()
486 }
487 }
488 }
489
490 async fn keep_cloudflare_counts(&self, since: &str, until: &str, counts: &[(String, String, f64)], fetched_at: &str) -> Result<()> {
491 self.db
492 .prepare("DELETE FROM own_counts WHERE meter = 'cloudflare_git' AND day >= ?1 AND day <= ?2")
493 .bind(&[since.into(), until.into()])?
494 .run()
495 .await?;
496 for chunk in counts.chunks(50) {
497 let mut statements = Vec::with_capacity(chunk.len());
498 for (day, workspace, count) in chunk {
499 statements.push(
500 self.db
501 .prepare("INSERT OR REPLACE INTO own_counts (day, meter, workspace, quantity, fetched_at) VALUES (?, 'cloudflare_git', ?, ?, ?)")
502 .bind(&[day.as_str().into(), workspace.as_str().into(), (*count).into(), fetched_at.into()])?,
503 );
504 }
505 self.db.batch(statements).await?;
506 }
507 Ok(())
508 }
509
510 /// Upserts lines, a day's line replacing what was read for it before.
511 async fn upsert_lines(&self, lines: &[CostLine], fetched_at: &str) -> Result<u32> {
512 for chunk in lines.chunks(50) {
513 let mut statements = Vec::with_capacity(chunk.len());
514 for line in chunk {
515 statements.push(
516 self.db
517 .prepare(
518 "INSERT INTO cost_lines (day, source, product, meter, unit, quantity, cost_usd, raw_name, fetched_at)
519 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)
520 ON CONFLICT (day, source, product, meter) DO UPDATE SET
521 unit = ?5, quantity = ?6, cost_usd = ?7, raw_name = ?8, fetched_at = ?9",
522 )
523 .bind(&[
524 line.day.as_str().into(),
525 line.source.into(),
526 line.product.as_str().into(),
527 line.meter.as_str().into(),
528 line.unit.as_str().into(),
529 line.quantity.into(),
530 line.cost_usd.into(),
531 line.raw_name.as_str().into(),
532 fetched_at.into(),
533 ])?,
534 );
535 }
536 self.db.batch(statements).await?;
537 }
538 Ok(lines.len() as u32)
539 }
540
541 /// The mapping from Cloudflare's meters to g1t's products.
542 pub(crate) async fn rules(&self) -> Result<Vec<Rule>> {
543 #[derive(Deserialize)]
544 struct Row {
545 product: String,
546 meter: String,
547 bucket: String,
548 price_meter: Option<String>,
549 own_meter: Option<String>,
550 drift_percent: f64,
551 }
552 Ok(self
553 .db
554 .prepare("SELECT product, meter, bucket, price_meter, own_meter, drift_percent FROM cost_map")
555 .all()
556 .await?
557 .results::<Row>()?
558 .into_iter()
559 .map(|r| Rule {
560 product: r.product,
561 meter: r.meter,
562 bucket: r.bucket,
563 price_meter: r.price_meter,
564 own_meter: r.own_meter,
565 drift_percent: r.drift_percent,
566 })
567 .collect())
568 }
569
570 /// g1t's own counts for the days, all from the repos service, which
571 /// owns the mapping from raw meters to operations (`operation_mapping`):
572 ///
573 /// - `git_operations`: what customers are charged for, as repos counts
574 /// it (`git_operations`, already through its mapping).
575 /// - `cost_operations`: what g1t expects Cloudflare to bill, from the
576 /// raw meters (`artifacts_usage`) and the mapping's cost column.
577 /// - `artifacts_<meter>`: each raw meter.
578 pub(crate) async fn count_own(&self, since: &str, until: &str) -> Result<()> {
579 let Some(repos) = &self.repos else { return Ok(()) };
580 let days = days_between(since, until);
581 let fetched_at = rfc3339(now_ms());
582 let mut rows: Vec<(String, String, String, f64)> = Vec::new();
583
584 // Raw meters and repos' mapping; skipped while repos does not answer.
585 if let Ok(body) = g1t_kit::call::<_, Value>(repos, "artifacts_usage", &json!({ "from": since, "to": until })).await {
586 let mapping = operation_mapping(&body);
587 let mut by: BTreeMap<(String, String), Vec<(String, f64)>> = BTreeMap::new();
588 for (day, workspace, meter, count) in body["rows"].as_array().map(|l| l.iter().filter_map(raw_usage_row).collect::<Vec<_>>()).unwrap_or_default() {
589 by.entry((day, workspace)).or_default().push((meter, count));
590 }
591 for ((day, workspace), counts) in by {
592 // Several stores (namespaces) can give the same meter.
593 let mut merged: BTreeMap<String, f64> = BTreeMap::new();
594 for (meter, count) in &counts {
595 *merged.entry(meter.clone()).or_default() += count;
596 }
597 for (meter, count) in &merged {
598 rows.push((day.clone(), format!("artifacts_{meter}"), workspace.clone(), *count));
599 }
600 if !mapping.is_empty() {
601 rows.push((day, "cost_operations".to_owned(), workspace, weighted(&counts, &mapping, true)));
602 }
603 }
604 }
605 // What customers are charged for, as repos counts it through its mapping.
606 let mut cumulative = Vec::with_capacity(days.len());
607 for day in &days {
608 let list: Vec<WorkspaceGitOperations> = g1t_kit::call(
609 repos,
610 "git_operations",
611 &GitOperationsArgs { month: day[..7].to_owned(), since: Some(format!("{day}T00")), namespace: None },
612 )
613 .await?;
614 cumulative.push(list.into_iter().map(|w| (w.namespace.to_lowercase(), w.operations)).collect::<BTreeMap<_, _>>());
615 }
616 for (day, workspace, count) in daily_from_cumulative(&days, &cumulative) {
617 rows.push((day, "git_operations".to_owned(), workspace, count as f64));
618 }
619
620 // Each day's counts replace what was there.
621 self.db
622 .prepare("DELETE FROM own_counts WHERE day >= ?1 AND day <= ?2 AND meter NOT LIKE 'cloudflare_%'")
623 .bind(&[since.into(), until.into()])?
624 .run()
625 .await?;
626 for chunk in rows.chunks(50) {
627 let mut statements = Vec::with_capacity(chunk.len());
628 for (day, meter, workspace, quantity) in chunk {
629 statements.push(
630 self.db
631 .prepare(
632 "INSERT INTO own_counts (day, meter, workspace, quantity, fetched_at) VALUES (?1, ?2, ?3, ?4, ?5)
633 ON CONFLICT (day, meter, workspace) DO UPDATE SET quantity = ?4, fetched_at = ?5",
634 )
635 .bind(&[day.as_str().into(), meter.as_str().into(), workspace.as_str().into(), (*quantity).into(), fetched_at.as_str().into()])?,
636 );
637 }
638 self.db.batch(statements).await?;
639 }
640 Ok(())
641 }
642}
643
644#[cfg(test)]
645mod tests {
646 use super::*;
647
648 /// Billable usage as Cloudflare answered g1t on 2026-10-06 (two days,
649 /// trimmed), plus rows past the included amounts, which cost money.
650 fn billable_fixture() -> Value {
651 json!({
652 "success": true,
653 "errors": [],
654 "result": [
655 {
656 "ChargePeriodStart": "2026-10-03T00:00:00Z",
657 "ChargePeriodEnd": "2026-10-04T00:00:00Z",
658 "ServiceFamilyName": "Containers",
659 "ServiceName": "Container Memory, per GiB-Second (First 25 GiB-hours included)",
660 "PricingUnit": "Count",
661 "PricingQuantity": "14000",
662 "ContractedCost": 0, "BilledCost": 0, "ListCost": 0
663 },
664 {
665 "ChargePeriodStart": "2026-10-03T00:00:00Z",
666 "ServiceFamilyName": "D1",
667 "ServiceName": "D1 - Rows Read (first 25 billion included)",
668 "PricingUnit": "Count",
669 "PricingQuantity": 1200000,
670 "BilledCost": 0
671 },
672 {
673 "ChargePeriodStart": "2026-10-15T00:00:00Z",
674 "ServiceFamilyName": "Artifacts",
675 "ServiceName": "Artifacts Operations (First 10,000 included)",
676 "PricingUnit": "Count",
677 "PricingQuantity": "30000",
678 "BilledCost": "3.00",
679 "ListCost": 4.5
680 },
681 {
682 "ChargePeriodStart": "2026-10-15T00:00:00Z",
683 "ServiceFamilyName": "Artifacts",
684 "ServiceName": "Artifacts Operations (First 10,000 included)",
685 "PricingUnit": "Count",
686 "PricingQuantity": "10000",
687 "BilledCost": "1.50"
688 },
689 {
690 "ChargePeriodStart": "2026-10-15T00:00:00Z",
691 "ServiceFamilyName": "Workers",
692 "ServiceName": "Workers for Platforms Requests (First 20M are included)",
693 "PricingUnit": "Count",
694 "PricingQuantity": 25000000,
695 "ListCost": 1.5
696 },
697 { "ServiceName": "No day, no line" }
698 ]
699 })
700 }
701
702 #[test]
703 fn names_become_products_and_meters() {
704 assert_eq!(slug("Workers for Platforms CPU ms (First 60M ms are included)"), "workers_for_platforms_cpu_ms");
705 assert_eq!(product_and_meter("D1", "D1 - Rows Read (first 25 billion included)"), ("d1".into(), "d1_rows_read".into()));
706 assert_eq!(product_and_meter("Workers KV", "KV Read Operations (First 10M is included)"), ("workers_kv".into(), "kv_read_operations".into()));
707 assert_eq!(product_and_meter("", "Containers / Container vCPU"), ("containers".into(), "container_vcpu".into()));
708 assert_eq!(product_and_meter("Email", ""), ("email".into(), "usage".into()));
709 }
710
711 #[test]
712 fn billable_usage_is_parsed_and_rows_of_one_meter_are_added() {
713 let lines = lines_from_billable(&billable_fixture()).unwrap();
714 assert_eq!(lines.len(), 4, "{lines:#?}");
715 let artifacts = lines.iter().find(|l| l.product == "artifacts").unwrap();
716 // Two rows of the same day and meter: one line, added up.
717 assert_eq!(artifacts.meter, "artifacts_operations");
718 assert_eq!(artifacts.quantity, 40_000.0);
719 assert!((artifacts.cost_usd - 4.5).abs() < 1e-9, "billed, not list: {}", artifacts.cost_usd);
720 let wfp = lines.iter().find(|l| l.meter.starts_with("workers_for_platforms")).unwrap();
721 assert_eq!(wfp.product, "workers");
722 // No billed cost: the list cost.
723 assert_eq!(wfp.cost_usd, 1.5);
724 let memory = lines.iter().find(|l| l.product == "containers").unwrap();
725 assert_eq!((memory.day.as_str(), memory.quantity, memory.cost_usd), ("2026-10-03", 14_000.0, 0.0));
726 assert!(lines_from_billable(&json!({ "success": false, "errors": [{ "code": 10000 }] })).is_err());
727 }
728
729 #[test]
730 fn reading_the_same_days_twice_gives_the_same_lines() {
731 // Idempotent: the same answer aggregates to the same keys and
732 // amounts, so the upsert replaces rather than adds.
733 let once = lines_from_billable(&billable_fixture()).unwrap();
734 let twice = lines_from_billable(&billable_fixture()).unwrap();
735 assert_eq!(once, twice);
736 let again = aggregate(once.clone());
737 assert_eq!(again, once);
738 }
739
740 #[test]
741 fn artifacts_events_are_counted_by_type_and_day() {
742 let body = json!({
743 "data": { "viewer": { "accounts": [{ "artifactsEventsAdaptiveGroups": [
744 { "count": 120, "dimensions": { "date": "2026-10-05", "eventType": "pull" } },
745 { "count": 30, "dimensions": { "date": "2026-10-05", "eventType": "push" } },
746 { "count": 2, "dimensions": { "date": "2026-10-05", "eventType": "rateLimited" } },
747 { "count": 5, "dimensions": { "date": "2026-10-06", "eventType": "pull" } }
748 ] }] } },
749 "errors": null
750 });
751 let lines = lines_from_artifacts(&body).unwrap();
752 assert_eq!(lines.len(), 4);
753 assert_eq!(lines[0].meter, "events_pull");
754 assert_eq!(lines[0].quantity, 120.0);
755 assert!(lines.iter().any(|l| l.meter == "events_ratelimited"));
756 // By workspace, from the repository's store key; forks say none.
757 let by_repo = json!({ "data": { "viewer": { "accounts": [{ "artifactsEventsAdaptiveGroups": [
758 { "count": 100, "dimensions": { "date": "2026-10-05", "eventType": "pull", "repositoryName": "acme--api" } },
759 { "count": 20, "dimensions": { "date": "2026-10-05", "eventType": "push", "repositoryName": "acme--web" } },
760 { "count": 7, "dimensions": { "date": "2026-10-05", "eventType": "fork", "repositoryName": "pulls--123" } },
761 { "count": 3, "dimensions": { "date": "2026-10-05", "eventType": "serverError", "repositoryName": "beta--x" } }
762 ] }] } } });
763 assert_eq!(artifacts_by_workspace(&by_repo, &BTreeMap::new()), vec![("2026-10-05".to_string(), "acme".to_string(), 120.0)]);
764 let none = BTreeMap::new();
765 assert_eq!(workspace_of_store_key("Acme--api", &none), Some("acme".into()));
766 assert_eq!(workspace_of_store_key("pulls--9", &none), None);
767 assert_eq!(workspace_of_store_key("plain", &none), None);
768 let owners = BTreeMap::from([("9".to_string(), "Acme".to_string())]);
769 assert_eq!(workspace_of_store_key("pulls--9", &owners), Some("acme".into()));
770 assert_eq!(workspace_of_store_key("pulls--10", &owners), None);
771 let operations: f64 = lines
772 .iter()
773 .filter(|l| l.day == "2026-10-05" && ARTIFACTS_OPERATIONS.contains(&l.meter.as_str()))
774 .map(|l| l.quantity)
775 .sum();
776 assert_eq!(operations, 150.0);
777 let failed = json!({ "data": null, "errors": [{ "message": "unknown field" }] });
778 assert!(lines_from_artifacts(&failed).is_err());
779 }
780
781 #[test]
782 fn the_first_run_backfills_31_days_and_later_runs_restate_a_few() {
783 let now = g1t_contracts::time::parse_rfc3339("2026-10-06T04:17:00Z").unwrap();
784 assert_eq!(window(None, now), ("2026-09-05".into(), "2026-10-06".into()));
785 assert_eq!(window(Some("2026-10-06"), now), ("2026-10-02".into(), "2026-10-06".into()));
786 // A week without a run: from the last day read.
787 assert_eq!(window(Some("2026-09-28"), now), ("2026-09-28".into(), "2026-10-06".into()));
788 // Never past the backfill.
789 assert_eq!(window(Some("2026-01-01"), now), ("2026-09-05".into(), "2026-10-06".into()));
790 assert_eq!(days_between("2026-09-29", "2026-10-02"), vec!["2026-09-29", "2026-09-30", "2026-10-01", "2026-10-02"]);
791 assert!(days_between("2026-10-02", "2026-10-01").is_empty());
792 }
793
794 #[test]
795 fn a_days_git_operations_are_its_count_less_the_next_days() {
796 let days: Vec<String> = ["2026-09-29", "2026-09-30", "2026-10-01", "2026-10-02"].iter().map(|d| d.to_string()).collect();
797 let map = |pairs: &[(&str, u64)]| pairs.iter().map(|(w, n)| (w.to_string(), *n)).collect::<BTreeMap<_, _>>();
798 // From each day to the end of its month.
799 let cumulative = vec![map(&[("acme", 30), ("beta", 4)]), map(&[("acme", 10)]), map(&[("acme", 7)]), map(&[("acme", 2)])];
800 let daily = daily_from_cumulative(&days, &cumulative);
801 assert_eq!(
802 daily,
803 vec![
804 ("2026-09-29".into(), "acme".into(), 20),
805 ("2026-09-29".into(), "beta".into(), 4),
806 // The month's last day is its own.
807 ("2026-09-30".into(), "acme".into(), 10),
808 ("2026-10-01".into(), "acme".into(), 5),
809 // Today so far.
810 ("2026-10-02".into(), "acme".into(), 2),
811 ]
812 );
813 }
814
815 #[test]
816 fn raw_meters_are_weighted_by_the_repos_mapping() {
817 // As repos' artifacts_usage answers: rows, and its operation_mapping.
818 let body = json!({
819 "rows": [
820 { "day": "2026-10-15", "store": "g1t", "workspace": "Acme", "meter": "git.fetch", "count": 12, "bytes_in": 0, "bytes_out": 0 },
821 { "day": "2026-10-15", "store": "g1t", "workspace": "acme", "meter": "git.receive_pack", "count": 3, "bytes_in": 0, "bytes_out": 0 },
822 { "day": "2026-10-15", "store": "g1t", "workspace": "acme", "meter": "binding.read_blob", "count": 400, "bytes_in": 0, "bytes_out": 0 }
823 ],
824 "mapping": [
825 { "meter": "git.fetch", "cost_operations": 1, "billable_operations": 1 },
826 { "meter": "git.receive_pack", "cost_operations": 1, "billable_operations": 1 },
827 { "meter": "binding.read_blob", "cost_operations": 0, "billable_operations": 0 }
828 ],
829 "truncated": false
830 });
831 let row = raw_usage_row(&body["rows"][0]).unwrap();
832 assert_eq!(row, ("2026-10-15".into(), "acme".into(), "git_fetch".into(), 12.0));
833 assert!(raw_usage_row(&json!({ "day": "2026-10-15", "count": 1 })).is_none());
834 let mut mapping = operation_mapping(&body);
835 let raw: Vec<(String, f64)> = body["rows"].as_array().unwrap().iter().filter_map(raw_usage_row).map(|r| (r.2, r.3)).collect();
836 assert_eq!(weighted(&raw, &mapping, true), 15.0);
837 assert_eq!(weighted(&raw, &mapping, false), 15.0);
838 // Cloudflare turns out to bill binding reads: repos changes one row
839 // (set_operation_mapping), and the bill g1t expects follows.
840 mapping.insert("binding_read_blob".into(), (1.0, 0.0));
841 assert_eq!(weighted(&raw, &mapping, true), 415.0);
842 assert_eq!(weighted(&raw, &mapping, false), 15.0);
843 assert!(operation_mapping(&json!({})).is_empty());
844 }
845
846 #[test]
847 fn the_most_specific_rule_claims_a_line() {
848 let rule = |product: &str, meter: &str, bucket: &str| Rule {
849 product: product.into(),
850 meter: meter.into(),
851 bucket: bucket.into(),
852 price_meter: None,
853 own_meter: None,
854 drift_percent: 10.0,
855 };
856 let rules = vec![rule("workers", "*", "platform"), rule("workers", "workers_for_platforms", "deployments"), rule("artifacts", "*", "git")];
857 assert_eq!(classify(&rules, "workers", "workers_for_platforms_cpu_ms").unwrap().bucket, "deployments");
858 assert_eq!(classify(&rules, "workers", "workers_cpu_ms").unwrap().bucket, "platform");
859 assert_eq!(classify(&rules, "artifacts", "artifacts_operations").unwrap().bucket, "git");
860 assert!(classify(&rules, "browser_rendering", "browser_hours").is_none());
861 }
862}