g1t/services/billing/src/costs.rs

790 lines36,980 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
383/// A price meter's units from raw counts and `billable_units` weights:
384/// what g1t charges for, from what was counted.
385pub(crate) fn billable(raw: &[(String, f64)], weights: &BTreeMap<String, f64>) -> f64 {
386 raw.iter().map(|(meter, count)| count * weights.get(meter).copied().unwrap_or(0.0)).sum()
387}
388
389impl Billing {
390 /// Reads Cloudflare's bill for the days due (see `window`) into
391 /// `cost_lines`: the days read and how many lines. None without a
392 /// token. What could not be read is added to `problems`.
393 pub(crate) async fn read_cloudflare(&self, keeper: &Keeper, problems: &mut Vec<String>) -> Result<Option<(String, String, u32)>> {
394 if !keeper.can_read_bill() {
395 return Ok(None);
396 }
397 #[derive(Deserialize)]
398 struct Last {
399 day: Option<String>,
400 }
401 let last = self
402 .db
403 .prepare("SELECT MAX(day) AS day FROM cost_lines WHERE source = ?")
404 .bind(&[SOURCE_BILLABLE.into()])?
405 .first::<Last>(None)
406 .await?
407 .and_then(|l| l.day);
408 let (since, until) = window(last.as_deref(), now_ms());
409 let fetched_at = rfc3339(now_ms());
410 let mut written = 0;
411 match keeper.billable_usage_body(&since, &until).await.map_err(|e| e.to_string()).and_then(|b| lines_from_billable(&b)) {
412 Ok(lines) => written += self.upsert_lines(&lines, &fetched_at).await?,
413 Err(error) => problems.push(format!("Cloudflare's billable usage could not be read: {error}")),
414 }
415 match keeper.graphql(artifacts_variables(keeper.account(), &since, &until)).await.map_err(|e| e.to_string()) {
416 Ok(body) => match lines_from_artifacts(&body) {
417 Ok(lines) => {
418 written += self.upsert_lines(&lines, &fetched_at).await?;
419 self.keep_cloudflare_counts(&since, &until, &artifacts_by_workspace(&body), &fetched_at).await?;
420 }
421 Err(error) => problems.push(format!("Artifacts events could not be read: {error}")),
422 },
423 Err(error) => problems.push(format!("Artifacts events could not be read: {error}")),
424 }
425 Ok(Some((since, until, written)))
426 }
427
428 /// Cloudflare's own per-workspace counts for the days, replacing what
429 /// was kept for them.
430 async fn keep_cloudflare_counts(&self, since: &str, until: &str, counts: &[(String, String, f64)], fetched_at: &str) -> Result<()> {
431 self.db
432 .prepare("DELETE FROM own_counts WHERE meter = 'cloudflare_git' AND day >= ?1 AND day <= ?2")
433 .bind(&[since.into(), until.into()])?
434 .run()
435 .await?;
436 for chunk in counts.chunks(50) {
437 let mut statements = Vec::with_capacity(chunk.len());
438 for (day, workspace, count) in chunk {
439 statements.push(
440 self.db
441 .prepare("INSERT OR REPLACE INTO own_counts (day, meter, workspace, quantity, fetched_at) VALUES (?, 'cloudflare_git', ?, ?, ?)")
442 .bind(&[day.as_str().into(), workspace.as_str().into(), (*count).into(), fetched_at.into()])?,
443 );
444 }
445 self.db.batch(statements).await?;
446 }
447 Ok(())
448 }
449
450 /// Upserts lines, a day's line replacing what was read for it before.
451 async fn upsert_lines(&self, lines: &[CostLine], fetched_at: &str) -> Result<u32> {
452 for chunk in lines.chunks(50) {
453 let mut statements = Vec::with_capacity(chunk.len());
454 for line in chunk {
455 statements.push(
456 self.db
457 .prepare(
458 "INSERT INTO cost_lines (day, source, product, meter, unit, quantity, cost_usd, raw_name, fetched_at)
459 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)
460 ON CONFLICT (day, source, product, meter) DO UPDATE SET
461 unit = ?5, quantity = ?6, cost_usd = ?7, raw_name = ?8, fetched_at = ?9",
462 )
463 .bind(&[
464 line.day.as_str().into(),
465 line.source.into(),
466 line.product.as_str().into(),
467 line.meter.as_str().into(),
468 line.unit.as_str().into(),
469 line.quantity.into(),
470 line.cost_usd.into(),
471 line.raw_name.as_str().into(),
472 fetched_at.into(),
473 ])?,
474 );
475 }
476 self.db.batch(statements).await?;
477 }
478 Ok(lines.len() as u32)
479 }
480
481 /// The mapping from Cloudflare's meters to g1t's products.
482 pub(crate) async fn rules(&self) -> Result<Vec<Rule>> {
483 #[derive(Deserialize)]
484 struct Row {
485 product: String,
486 meter: String,
487 bucket: String,
488 price_meter: Option<String>,
489 own_meter: Option<String>,
490 drift_percent: f64,
491 }
492 Ok(self
493 .db
494 .prepare("SELECT product, meter, bucket, price_meter, own_meter, drift_percent FROM cost_map")
495 .all()
496 .await?
497 .results::<Row>()?
498 .into_iter()
499 .map(|r| Rule {
500 product: r.product,
501 meter: r.meter,
502 bucket: r.bucket,
503 price_meter: r.price_meter,
504 own_meter: r.own_meter,
505 drift_percent: r.drift_percent,
506 })
507 .collect())
508 }
509
510 /// g1t's own counts for the days: raw Artifacts meters from the repos
511 /// service's `artifacts_usage` when it answers, and git operations,
512 /// from those raw meters and `billable_units` weights when both exist,
513 /// else from its `git_operations`.
514 pub(crate) async fn count_own(&self, since: &str, until: &str) -> Result<()> {
515 let Some(repos) = &self.repos else { return Ok(()) };
516 let days = days_between(since, until);
517 let fetched_at = rfc3339(now_ms());
518 let mut rows: Vec<(String, String, String, f64)> = Vec::new();
519
520 // Raw meters, while the repos service may not have them yet.
521 let raw: Vec<(String, String, String, f64)> =
522 match g1t_kit::call::<_, Value>(repos, "artifacts_usage", &json!({ "from": since, "to": until })).await {
523 Ok(Value::Array(list)) => list.iter().filter_map(raw_usage_row).collect(),
524 Ok(other) => other["rows"].as_array().map(|l| l.iter().filter_map(raw_usage_row).collect()).unwrap_or_default(),
525 Err(_) => Vec::new(),
526 };
527 #[derive(Deserialize)]
528 struct Weight {
529 raw_meter: String,
530 weight: f64,
531 }
532 let weights: BTreeMap<String, f64> = self
533 .db
534 .prepare("SELECT raw_meter, weight FROM billable_units WHERE price_meter = 'git_operations'")
535 .all()
536 .await?
537 .results::<Weight>()?
538 .into_iter()
539 .map(|w| (slug(&w.raw_meter), w.weight))
540 .collect();
541 for (day, workspace, meter, count) in &raw {
542 rows.push((day.clone(), format!("artifacts_{meter}"), workspace.clone(), *count));
543 }
544 if !raw.is_empty() && !weights.is_empty() {
545 let mut by: BTreeMap<(String, String), Vec<(String, f64)>> = BTreeMap::new();
546 for (day, workspace, meter, count) in &raw {
547 by.entry((day.clone(), workspace.clone())).or_default().push((meter.clone(), *count));
548 }
549 for ((day, workspace), counts) in by {
550 rows.push((day, "git_operations".to_owned(), workspace, billable(&counts, &weights)));
551 }
552 } else {
553 let mut cumulative = Vec::with_capacity(days.len());
554 for day in &days {
555 let list: Vec<WorkspaceGitOperations> = g1t_kit::call(
556 repos,
557 "git_operations",
558 &GitOperationsArgs { month: day[..7].to_owned(), since: Some(format!("{day}T00")), namespace: None },
559 )
560 .await?;
561 cumulative.push(list.into_iter().map(|w| (w.namespace.to_lowercase(), w.operations)).collect::<BTreeMap<_, _>>());
562 }
563 for (day, workspace, count) in daily_from_cumulative(&days, &cumulative) {
564 rows.push((day, "git_operations".to_owned(), workspace, count as f64));
565 }
566 }
567
568 // Each day's counts replace what was there.
569 self.db
570 .prepare("DELETE FROM own_counts WHERE day >= ?1 AND day <= ?2 AND meter NOT LIKE 'cloudflare_%'")
571 .bind(&[since.into(), until.into()])?
572 .run()
573 .await?;
574 for chunk in rows.chunks(50) {
575 let mut statements = Vec::with_capacity(chunk.len());
576 for (day, meter, workspace, quantity) in chunk {
577 statements.push(
578 self.db
579 .prepare(
580 "INSERT INTO own_counts (day, meter, workspace, quantity, fetched_at) VALUES (?1, ?2, ?3, ?4, ?5)
581 ON CONFLICT (day, meter, workspace) DO UPDATE SET quantity = ?4, fetched_at = ?5",
582 )
583 .bind(&[day.as_str().into(), meter.as_str().into(), workspace.as_str().into(), (*quantity).into(), fetched_at.as_str().into()])?,
584 );
585 }
586 self.db.batch(statements).await?;
587 }
588 Ok(())
589 }
590}
591
592#[cfg(test)]
593mod tests {
594 use super::*;
595
596 /// Billable usage as Cloudflare answered g1t on 2026-10-06 (two days,
597 /// trimmed), plus rows past the included amounts, which cost money.
598 fn billable_fixture() -> Value {
599 json!({
600 "success": true,
601 "errors": [],
602 "result": [
603 {
604 "ChargePeriodStart": "2026-10-03T00:00:00Z",
605 "ChargePeriodEnd": "2026-10-04T00:00:00Z",
606 "ServiceFamilyName": "Containers",
607 "ServiceName": "Container Memory, per GiB-Second (First 25 GiB-hours included)",
608 "PricingUnit": "Count",
609 "PricingQuantity": "14000",
610 "ContractedCost": 0, "BilledCost": 0, "ListCost": 0
611 },
612 {
613 "ChargePeriodStart": "2026-10-03T00:00:00Z",
614 "ServiceFamilyName": "D1",
615 "ServiceName": "D1 - Rows Read (first 25 billion included)",
616 "PricingUnit": "Count",
617 "PricingQuantity": 1200000,
618 "BilledCost": 0
619 },
620 {
621 "ChargePeriodStart": "2026-10-15T00:00:00Z",
622 "ServiceFamilyName": "Artifacts",
623 "ServiceName": "Artifacts Operations (First 10,000 included)",
624 "PricingUnit": "Count",
625 "PricingQuantity": "30000",
626 "BilledCost": "3.00",
627 "ListCost": 4.5
628 },
629 {
630 "ChargePeriodStart": "2026-10-15T00:00:00Z",
631 "ServiceFamilyName": "Artifacts",
632 "ServiceName": "Artifacts Operations (First 10,000 included)",
633 "PricingUnit": "Count",
634 "PricingQuantity": "10000",
635 "BilledCost": "1.50"
636 },
637 {
638 "ChargePeriodStart": "2026-10-15T00:00:00Z",
639 "ServiceFamilyName": "Workers",
640 "ServiceName": "Workers for Platforms Requests (First 20M are included)",
641 "PricingUnit": "Count",
642 "PricingQuantity": 25000000,
643 "ListCost": 1.5
644 },
645 { "ServiceName": "No day, no line" }
646 ]
647 })
648 }
649
650 #[test]
651 fn names_become_products_and_meters() {
652 assert_eq!(slug("Workers for Platforms CPU ms (First 60M ms are included)"), "workers_for_platforms_cpu_ms");
653 assert_eq!(product_and_meter("D1", "D1 - Rows Read (first 25 billion included)"), ("d1".into(), "d1_rows_read".into()));
654 assert_eq!(product_and_meter("Workers KV", "KV Read Operations (First 10M is included)"), ("workers_kv".into(), "kv_read_operations".into()));
655 assert_eq!(product_and_meter("", "Containers / Container vCPU"), ("containers".into(), "container_vcpu".into()));
656 assert_eq!(product_and_meter("Email", ""), ("email".into(), "usage".into()));
657 }
658
659 #[test]
660 fn billable_usage_is_parsed_and_rows_of_one_meter_are_added() {
661 let lines = lines_from_billable(&billable_fixture()).unwrap();
662 assert_eq!(lines.len(), 4, "{lines:#?}");
663 let artifacts = lines.iter().find(|l| l.product == "artifacts").unwrap();
664 // Two rows of the same day and meter: one line, added up.
665 assert_eq!(artifacts.meter, "artifacts_operations");
666 assert_eq!(artifacts.quantity, 40_000.0);
667 assert!((artifacts.cost_usd - 4.5).abs() < 1e-9, "billed, not list: {}", artifacts.cost_usd);
668 let wfp = lines.iter().find(|l| l.meter.starts_with("workers_for_platforms")).unwrap();
669 assert_eq!(wfp.product, "workers");
670 // No billed cost: the list cost.
671 assert_eq!(wfp.cost_usd, 1.5);
672 let memory = lines.iter().find(|l| l.product == "containers").unwrap();
673 assert_eq!((memory.day.as_str(), memory.quantity, memory.cost_usd), ("2026-10-03", 14_000.0, 0.0));
674 assert!(lines_from_billable(&json!({ "success": false, "errors": [{ "code": 10000 }] })).is_err());
675 }
676
677 #[test]
678 fn reading_the_same_days_twice_gives_the_same_lines() {
679 // Idempotent: the same answer aggregates to the same keys and
680 // amounts, so the upsert replaces rather than adds.
681 let once = lines_from_billable(&billable_fixture()).unwrap();
682 let twice = lines_from_billable(&billable_fixture()).unwrap();
683 assert_eq!(once, twice);
684 let again = aggregate(once.clone());
685 assert_eq!(again, once);
686 }
687
688 #[test]
689 fn artifacts_events_are_counted_by_type_and_day() {
690 let body = json!({
691 "data": { "viewer": { "accounts": [{ "artifactsEventsAdaptiveGroups": [
692 { "count": 120, "dimensions": { "date": "2026-10-05", "eventType": "pull" } },
693 { "count": 30, "dimensions": { "date": "2026-10-05", "eventType": "push" } },
694 { "count": 2, "dimensions": { "date": "2026-10-05", "eventType": "rateLimited" } },
695 { "count": 5, "dimensions": { "date": "2026-10-06", "eventType": "pull" } }
696 ] }] } },
697 "errors": null
698 });
699 let lines = lines_from_artifacts(&body).unwrap();
700 assert_eq!(lines.len(), 4);
701 assert_eq!(lines[0].meter, "events_pull");
702 assert_eq!(lines[0].quantity, 120.0);
703 assert!(lines.iter().any(|l| l.meter == "events_ratelimited"));
704 // By workspace, from the repository's store key; forks say none.
705 let by_repo = json!({ "data": { "viewer": { "accounts": [{ "artifactsEventsAdaptiveGroups": [
706 { "count": 100, "dimensions": { "date": "2026-10-05", "eventType": "pull", "repositoryName": "acme--api" } },
707 { "count": 20, "dimensions": { "date": "2026-10-05", "eventType": "push", "repositoryName": "acme--web" } },
708 { "count": 7, "dimensions": { "date": "2026-10-05", "eventType": "fork", "repositoryName": "pulls--123" } },
709 { "count": 3, "dimensions": { "date": "2026-10-05", "eventType": "serverError", "repositoryName": "beta--x" } }
710 ] }] } } });
711 assert_eq!(artifacts_by_workspace(&by_repo), vec![("2026-10-05".to_string(), "acme".to_string(), 120.0)]);
712 assert_eq!(workspace_of_store_key("Acme--api"), Some("acme".into()));
713 assert_eq!(workspace_of_store_key("pulls--9"), None);
714 assert_eq!(workspace_of_store_key("plain"), None);
715 let operations: f64 = lines
716 .iter()
717 .filter(|l| l.day == "2026-10-05" && ARTIFACTS_OPERATIONS.contains(&l.meter.as_str()))
718 .map(|l| l.quantity)
719 .sum();
720 assert_eq!(operations, 150.0);
721 let failed = json!({ "data": null, "errors": [{ "message": "unknown field" }] });
722 assert!(lines_from_artifacts(&failed).is_err());
723 }
724
725 #[test]
726 fn the_first_run_backfills_31_days_and_later_runs_restate_a_few() {
727 let now = g1t_contracts::time::parse_rfc3339("2026-10-06T04:17:00Z").unwrap();
728 assert_eq!(window(None, now), ("2026-09-05".into(), "2026-10-06".into()));
729 assert_eq!(window(Some("2026-10-06"), now), ("2026-10-02".into(), "2026-10-06".into()));
730 // A week without a run: from the last day read.
731 assert_eq!(window(Some("2026-09-28"), now), ("2026-09-28".into(), "2026-10-06".into()));
732 // Never past the backfill.
733 assert_eq!(window(Some("2026-01-01"), now), ("2026-09-05".into(), "2026-10-06".into()));
734 assert_eq!(days_between("2026-09-29", "2026-10-02"), vec!["2026-09-29", "2026-09-30", "2026-10-01", "2026-10-02"]);
735 assert!(days_between("2026-10-02", "2026-10-01").is_empty());
736 }
737
738 #[test]
739 fn a_days_git_operations_are_its_count_less_the_next_days() {
740 let days: Vec<String> = ["2026-09-29", "2026-09-30", "2026-10-01", "2026-10-02"].iter().map(|d| d.to_string()).collect();
741 let map = |pairs: &[(&str, u64)]| pairs.iter().map(|(w, n)| (w.to_string(), *n)).collect::<BTreeMap<_, _>>();
742 // From each day to the end of its month.
743 let cumulative = vec![map(&[("acme", 30), ("beta", 4)]), map(&[("acme", 10)]), map(&[("acme", 7)]), map(&[("acme", 2)])];
744 let daily = daily_from_cumulative(&days, &cumulative);
745 assert_eq!(
746 daily,
747 vec![
748 ("2026-09-29".into(), "acme".into(), 20),
749 ("2026-09-29".into(), "beta".into(), 4),
750 // The month's last day is its own.
751 ("2026-09-30".into(), "acme".into(), 10),
752 ("2026-10-01".into(), "acme".into(), 5),
753 // Today so far.
754 ("2026-10-02".into(), "acme".into(), 2),
755 ]
756 );
757 }
758
759 #[test]
760 fn raw_meters_become_billable_units_by_their_weights() {
761 let row = raw_usage_row(&json!({ "day": "2026-10-15", "namespace": "Acme", "meter": "upload_pack", "count": 12 })).unwrap();
762 assert_eq!(row, ("2026-10-15".into(), "acme".into(), "upload_pack".into(), 12.0));
763 assert!(raw_usage_row(&json!({ "day": "2026-10-15", "count": 1 })).is_none());
764 let weights: BTreeMap<String, f64> = [("upload_pack".to_string(), 1.0), ("receive_pack".to_string(), 1.0), ("binding_read".to_string(), 0.0)].into();
765 let raw = vec![("upload_pack".to_string(), 12.0), ("receive_pack".to_string(), 3.0), ("binding_read".to_string(), 400.0), ("ls_refs".to_string(), 9.0)];
766 assert_eq!(billable(&raw, &weights), 15.0);
767 // Cloudflare turns out to count binding reads: one row changes, and
768 // so does what is counted from then on.
769 let mut weights = weights;
770 weights.insert("binding_read".into(), 1.0);
771 assert_eq!(billable(&raw, &weights), 415.0);
772 }
773
774 #[test]
775 fn the_most_specific_rule_claims_a_line() {
776 let rule = |product: &str, meter: &str, bucket: &str| Rule {
777 product: product.into(),
778 meter: meter.into(),
779 bucket: bucket.into(),
780 price_meter: None,
781 own_meter: None,
782 drift_percent: 10.0,
783 };
784 let rules = vec![rule("workers", "*", "platform"), rule("workers", "workers_for_platforms", "deployments"), rule("artifacts", "*", "git")];
785 assert_eq!(classify(&rules, "workers", "workers_for_platforms_cpu_ms").unwrap().bucket, "deployments");
786 assert_eq!(classify(&rules, "workers", "workers_cpu_ms").unwrap().bucket, "platform");
787 assert_eq!(classify(&rules, "artifacts", "artifacts_operations").unwrap().bucket, "git");
788 assert!(classify(&rules, "browser_rendering", "browser_hours").is_none());
789 }
790}