g1t/services/billing/src/costs.rs

862 lines40,729 bytesCodeBlame

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

Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily1//! What Cloudflare charges g1t, day by day and product by product.
2//!
3//! Prices are what g1t pays plus 20%, so g1t has to know what it pays, as
4//! Cloudflare counts it, not as g1t assumes it. Some of what g1t runs on
5//! is in beta (Artifacts bills "operations" from 2026-10-14 without saying
6//! exactly which calls are one), so nothing here assumes a product list:
7//! every line Cloudflare bills is kept, whatever it is, and how a line
8//! maps to what g1t sells is data (`cost_map`), changed without a deploy.
9//!
10//! Once a day (`keeper::DAILY`) billing reads:
11//!
12//! - **Billable usage** (`GET /accounts/{id}/billable-usage`, FOCUS
13//! columns): one row per service per day, with its quantity and what it
14//! cost g1t. Every Cloudflare product g1t uses appears here once it is
15//! used: Workers, Workers for Platforms, D1, KV, R2, Queues, Containers,
16//! Durable Objects, Artifacts, Browser Rendering, Workers AI, Vectorize,
17//! Cloudflare for SaaS and Email.
18//! - **Artifacts events** (GraphQL `artifactsEventsAdaptiveGroups`), by
19//! event type and day: what Artifacts itself counted, to compare with
20//! what g1t counted (`margin`).
21//!
22//! Each becomes cost lines in `cost_lines`, one per (day, source,
23//! product, meter), upserted, so reading a day again replaces it. The
24//! first run reads the last 31 days; after that the last few, since
25//! Cloudflare restates recent days as usage settles.
26//!
27//! The token is `CLOUDFLARE_BILLING_TOKEN` (Account: Billing Read and
28//! Account Analytics Read), or the keeper's `CLOUDFLARE_USAGE_TOKEN`,
29//! which has both. Without either, nothing is read and nothing fails.
30
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)31use std::collections::{BTreeMap, BTreeSet};
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily32
33use g1t_contracts::repos::{GitOperationsArgs, WorkspaceGitOperations};
34use g1t_contracts::time::rfc3339;
35use g1t_kit::now_ms;
36use serde::Deserialize;
37use serde_json::{Value, json};
38use worker::Result;
39
40use crate::Billing;
41use crate::keeper::Keeper;
42
43/// Where a cost line came from.
44pub(crate) const SOURCE_BILLABLE: &str = "billable_usage";
45pub(crate) const SOURCE_ARTIFACTS: &str = "artifacts_events";
46
47/// How far back the first run reads: the 31 days GraphQL keeps.
48pub(crate) const BACKFILL_DAYS: u64 = 31;
49/// How many recent days every later run reads again.
50pub(crate) const RESTATE_DAYS: u64 = 4;
51
52pub(crate) const DAY_MS: u64 = 24 * 60 * 60 * 1000;
53
54/// One day of one meter of one Cloudflare product.
55#[derive(Clone, Debug, PartialEq)]
56pub(crate) struct CostLine {
57 /// YYYY-MM-DD, UTC.
58 pub day: String,
59 pub source: &'static str,
60 /// `containers`, `workers`, `workers_kv`, `artifacts`, … from
61 /// Cloudflare's own family name.
62 pub product: String,
63 /// The service within it, such as `container_memory_per_gib_second`
64 /// or `events_push`.
65 pub meter: String,
66 pub unit: String,
67 pub quantity: f64,
68 /// What g1t pays, in dollars: contracted, billed, or list.
69 pub cost_usd: f64,
70 /// The name as Cloudflare gave it, for people.
71 pub raw_name: String,
72}
73
74/// `Workers for Platforms CPU ms (First 60M ms are included)` →
75/// `workers_for_platforms_cpu_ms`: lower case, words joined by `_`, and
76/// what is in parentheses (the included amount, which changes) left out.
77pub(crate) fn slug(text: &str) -> String {
78 let mut out = String::new();
79 let mut depth = 0u32;
80 let mut gap = false;
81 for c in text.chars() {
82 match c {
83 '(' => depth += 1,
84 ')' => depth = depth.saturating_sub(1),
85 _ if depth > 0 => {}
86 c if c.is_ascii_alphanumeric() => {
87 if gap && !out.is_empty() {
88 out.push('_');
89 }
90 gap = false;
91 out.push(c.to_ascii_lowercase());
92 }
93 _ => gap = true,
94 }
95 }
96 out
97}
98
99/// The product and meter for a family and service name, such as
100/// (`D1`, `D1 - Rows Read (first 25 billion included)`) → (`d1`,
101/// `d1_rows_read`). Without a family, the service's first word.
102pub(crate) fn product_and_meter(family: &str, service: &str) -> (String, String) {
103 let (family, service) = if family.trim().is_empty() {
104 match service.split_once(" / ") {
105 Some((family, service)) => (family, service),
106 None => (service.split_whitespace().next().unwrap_or("other"), service),
107 }
108 } else {
109 (family, service)
110 };
111 let product = slug(family);
112 let product = if product.is_empty() { "other".to_owned() } else { product };
113 // Kept whole: `Workers for Platforms Requests` under `Workers` must stay
114 // `workers_for_platforms_requests`.
115 let meter = slug(service);
116 let meter = if meter.is_empty() { "usage".to_owned() } else { meter };
117 (product, meter)
118}
119
120fn text<'a>(row: &'a Value, keys: &[&str]) -> &'a str {
121 keys.iter().find_map(|k| row[*k].as_str()).unwrap_or_default()
122}
123
124fn number(row: &Value, keys: &[&str]) -> Option<f64> {
125 keys.iter()
126 .find_map(|k| row[*k].as_f64().or_else(|| row[*k].as_str().and_then(|s| s.trim().parse().ok())))
127}
128
129/// One row of billable usage, read leniently: the API is new, and its
130/// field names are FOCUS's (in either case). None without a service or a
131/// day.
132pub(crate) fn line_from_focus(row: &Value) -> Option<CostLine> {
133 let service = text(row, &["ServiceName", "service_name", "service"]);
134 let day = text(row, &["ChargePeriodStart", "charge_period_start", "UsageDate", "date"]);
135 if service.is_empty() || day.len() < 10 {
136 return None;
137 }
138 let family = text(row, &["ServiceFamilyName", "service_family_name", "ServiceCategory"]);
139 let (product, meter) = product_and_meter(family, service);
140 // What g1t pays: contracted, else billed, else list.
141 let cost = [
142 &["ContractedCost", "contracted_cost"][..],
143 &["BilledCost", "billed_cost"][..],
144 &["EffectiveCost", "effective_cost"][..],
145 &["ListCost", "list_cost"][..],
146 ]
147 .iter()
148 .find_map(|keys| number(row, keys).filter(|c| *c > 0.0))
149 .unwrap_or(0.0);
150 Some(CostLine {
151 day: day[..10].to_owned(),
152 source: SOURCE_BILLABLE,
153 product,
154 meter,
155 unit: text(row, &["PricingUnit", "pricing_unit", "ConsumedUnit", "consumed_unit"]).to_owned(),
156 quantity: number(row, &["PricingQuantity", "pricing_quantity", "ConsumedQuantity", "consumed_quantity"]).unwrap_or(0.0),
157 cost_usd: cost,
158 raw_name: if family.is_empty() { service.to_owned() } else { format!("{family} / {service}") },
159 })
160}
161
162/// Lines with the same key added together, in key order. Cloudflare can
163/// give one service several rows a day (regions, tiers); an upsert of
164/// each would keep only the last.
165pub(crate) fn aggregate(lines: Vec<CostLine>) -> Vec<CostLine> {
166 let mut out: Vec<CostLine> = Vec::new();
167 for line in lines {
168 match out
169 .iter_mut()
170 .find(|l| l.day == line.day && l.source == line.source && l.product == line.product && l.meter == line.meter)
171 {
172 Some(existing) => {
173 existing.quantity += line.quantity;
174 existing.cost_usd += line.cost_usd;
175 }
176 None => out.push(line),
177 }
178 }
179 out.sort_by(|a, b| (&a.day, a.source, &a.product, &a.meter).cmp(&(&b.day, b.source, &b.product, &b.meter)));
180 out
181}
182
183/// Every line of a billable-usage answer, aggregated. Errors when
184/// Cloudflare says it failed.
185pub(crate) fn lines_from_billable(body: &Value) -> std::result::Result<Vec<CostLine>, String> {
186 if body["success"] == Value::Bool(false) {
187 return Err(format!("billable usage failed: {}", body["errors"]));
188 }
189 let rows = body["result"].as_array().cloned().unwrap_or_default();
190 Ok(aggregate(rows.iter().filter_map(line_from_focus).collect()))
191}
192
193/// The GraphQL query for Artifacts' events by type and day.
194pub(crate) const ARTIFACTS_QUERY: &str = "query ($account: String!, $since: Date!, $until: Date!) {
195 viewer { accounts(filter: { accountTag: $account }) {
196 artifactsEventsAdaptiveGroups(limit: 10000, filter: { date_geq: $since, date_leq: $until }) {
197 count
198 dimensions { date eventType repositoryName }
199 }
200 } }
201}";
202
203pub(crate) fn artifacts_variables(account: &str, since: &str, until: &str) -> Value {
204 json!({ "query": ARTIFACTS_QUERY, "variables": { "account": account, "since": since, "until": until } })
205}
206
207/// Artifacts' events as lines (`events_push`, `events_pull`, …), with no
208/// cost: billable usage carries the cost. Errors when GraphQL does.
209pub(crate) fn lines_from_artifacts(body: &Value) -> std::result::Result<Vec<CostLine>, String> {
210 if let Some(errors) = body["errors"].as_array().filter(|e| !e.is_empty()) {
211 return Err(format!("Artifacts events failed: {}", Value::Array(errors.clone())));
212 }
213 let groups = body["data"]["viewer"]["accounts"][0]["artifactsEventsAdaptiveGroups"]
214 .as_array()
215 .cloned()
216 .unwrap_or_default();
217 Ok(aggregate(
218 groups
219 .iter()
220 .filter_map(|g| {
221 let day = g["dimensions"]["date"].as_str()?;
222 let kind = g["dimensions"]["eventType"].as_str()?;
223 if day.len() < 10 {
224 return None;
225 }
226 Some(CostLine {
227 day: day[..10].to_owned(),
228 source: SOURCE_ARTIFACTS,
229 product: "artifacts".to_owned(),
230 meter: format!("events_{}", slug(kind)),
231 unit: "events".to_owned(),
232 quantity: g["count"].as_f64().unwrap_or(0.0),
233 cost_usd: 0.0,
234 raw_name: format!("Artifacts / {kind}"),
235 })
236 })
237 .collect(),
238 ))
239}
240
241/// The workspace a repository in the store belongs to: keys are
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)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> {
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily246 let (workspace, rest) = key.split_once("--")?;
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)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())
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily254}
255
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)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
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily267/// 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`).
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)270pub(crate) fn artifacts_by_workspace(body: &Value, owners: &BTreeMap<String, String>) -> Vec<(String, String, f64)> {
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily271 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 }
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)278 let Some(workspace) = g["dimensions"]["repositoryName"].as_str().and_then(|key| workspace_of_store_key(key, owners)) else { continue };
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily279 *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
One operation mapping, owned by repos; billing reads it instead of keeping its own402/// 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()
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily427}
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?;
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)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?;
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily461 }
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.
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)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
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily490 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
One operation mapping, owned by repos; billing reads it instead of keeping its own570 /// 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.
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily578 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
One operation mapping, owned by repos; billing reads it instead of keeping its own584 // 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);
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily587 let mut by: BTreeMap<(String, String), Vec<(String, f64)>> = BTreeMap::new();
One operation mapping, owned by repos; billing reads it instead of keeping its own588 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));
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily590 }
591 for ((day, workspace), counts) in by {
One operation mapping, owned by repos; billing reads it instead of keeping its own592 // 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 }
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily603 }
604 }
One operation mapping, owned by repos; billing reads it instead of keeping its own605 // 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 }
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily619
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 ] }] } } });
Costs: Cloudflare's count for a pull request's working copy is shared out to its repository's workspace (repos pull_owners)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);
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily771 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]
One operation mapping, owned by repos; billing reads it instead of keeping its own816 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));
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily833 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 own834 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());
Git storage hardened, pages in tens of milliseconds, honest security alerts, and costs reconciled daily844 }
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}