g1t/services/billing/src/costs.rs

819 lines38,687 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
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
One operation mapping, owned by repos; billing reads it instead of keeping its own383/// 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()
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily408}
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
One operation mapping, owned by repos; billing reads it instead of keeping its own531 /// 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.
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily539 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
One operation mapping, owned by repos; billing reads it instead of keeping its own545 // 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);
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily548 let mut by: BTreeMap<(String, String), Vec<(String, f64)>> = BTreeMap::new();
One operation mapping, owned by repos; billing reads it instead of keeping its own549 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));
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily551 }
552 for ((day, workspace), counts) in by {
One operation mapping, owned by repos; billing reads it instead of keeping its own553 // 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 }
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily564 }
565 }
One operation mapping, owned by repos; billing reads it instead of keeping its own566 // 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 }
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily580
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]
One operation mapping, owned by repos; billing reads it instead of keeping its own773 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));
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily790 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 own791 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());
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily801 }
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}