g1t/services/repos/src/meters.rs

643 lines25,008 bytesCodeBlame
1//! Raw meters of everything g1t asks of the git store, and how it answered.
2//!
3//! Cloudflare bills Artifacts per "operation" from 2026-10-14 without
4//! having said exactly which calls are operations, so every interaction is
5//! counted by kind, per day, git store namespace, repository and workspace
6//! (`artifacts_meters`): git over HTTPS from people and tools, the same by
7//! g1t itself, and every call on the binding, with the bytes each moved
8//! where known. Which meters are operations, for g1t's own bill and for
9//! what workspaces are charged, is data (`operation_mapping`), read here
10//! and changed without a deploy. What a workspace is counted for goes to
11//! `git_operations` by the hour, as before, which billing reads.
12//!
13//! The same place keeps how the store answered, by the minute
14//! (`store_health`), for the status page's "Git storage" part.
15//!
16//! Nothing here is on the request path. Each isolate adds up what it saw in
17//! memory, and writes it all in one batch once the answer has gone back
18//! (`flush`, from `ctx.wait_until`). A failed write puts the counts back
19//! for the next one. What an isolate holds when it is evicted is lost: a
20//! few seconds' worth at most, since every request ends with a flush.
21
22use std::cell::RefCell;
23use std::collections::HashMap;
24
25use serde::{Deserialize, Serialize};
26use worker::wasm_bindgen::JsValue;
27use worker::{D1Database, Result};
28
29/// How an answer from the store went, for its health.
30#[derive(Clone, Copy, Debug, PartialEq, Eq)]
31pub enum Outcome {
32 Ok,
33 /// Failed, after any retries.
34 Failed,
35 /// Refused for the store's rate limit.
36 RateLimited,
37 /// Not asked: the namespace's breaker was open (resilience.rs).
38 Rejected,
39}
40
41#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
42pub struct Tally {
43 pub count: u64,
44 pub bytes_in: u64,
45 pub bytes_out: u64,
46}
47
48#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
49pub struct Health {
50 pub calls: u64,
51 pub errors: u64,
52 pub rate_limited: u64,
53 pub rejected: u64,
54 pub ms_total: u64,
55}
56
57/// One meter's place: the hour (`YYYY-MM-DDTHH`), the git store namespace,
58/// the repository's name there, its workspace, and the meter.
59#[derive(Clone, Debug, PartialEq, Eq, Hash)]
60pub struct Place {
61 pub hour: String,
62 pub store: String,
63 pub repo: String,
64 pub workspace: String,
65 pub meter: String,
66}
67
68/// What an isolate has counted and not yet written.
69#[derive(Default, Debug)]
70pub struct Pending {
71 pub usage: HashMap<Place, Tally>,
72 /// By namespace and minute (`YYYY-MM-DDTHH:MM`).
73 pub health: HashMap<(String, String), Health>,
74}
75
76impl Pending {
77 pub fn is_empty(&self) -> bool {
78 self.usage.is_empty() && self.health.is_empty()
79 }
80
81 pub fn meter(&mut self, place: Place, bytes_in: u64, bytes_out: u64) {
82 self.add(place, 1, bytes_in, bytes_out);
83 }
84
85 pub fn add(&mut self, place: Place, count: u64, bytes_in: u64, bytes_out: u64) {
86 let tally = self.usage.entry(place).or_default();
87 tally.count += count;
88 tally.bytes_in += bytes_in;
89 tally.bytes_out += bytes_out;
90 }
91
92 pub fn health(&mut self, store: &str, minute: &str, outcome: Outcome, ms: u64) {
93 let health = self.health.entry((store.to_owned(), minute.to_owned())).or_default();
94 match outcome {
95 Outcome::Rejected => health.rejected += 1,
96 _ => {
97 health.calls += 1;
98 health.ms_total += ms;
99 match outcome {
100 Outcome::Failed => health.errors += 1,
101 Outcome::RateLimited => {
102 health.errors += 1;
103 health.rate_limited += 1;
104 }
105 _ => {}
106 }
107 }
108 }
109 }
110
111 /// Puts back what a failed write took.
112 pub fn merge(&mut self, other: Pending) {
113 for (place, tally) in other.usage {
114 let kept = self.usage.entry(place).or_default();
115 kept.count += tally.count;
116 kept.bytes_in += tally.bytes_in;
117 kept.bytes_out += tally.bytes_out;
118 }
119 for (key, health) in other.health {
120 let kept = self.health.entry(key).or_default();
121 kept.calls += health.calls;
122 kept.errors += health.errors;
123 kept.rate_limited += health.rate_limited;
124 kept.rejected += health.rejected;
125 kept.ms_total += health.ms_total;
126 }
127 }
128
129 /// What each workspace is counted for, by the hour, under `mapping`.
130 pub fn billable(&self, mapping: &Mapping) -> HashMap<(String, String), f64> {
131 let mut out: HashMap<(String, String), f64> = HashMap::new();
132 for (place, tally) in &self.usage {
133 let weight = mapping.billable(&place.meter);
134 if weight > 0.0 && !place.workspace.is_empty() {
135 *out.entry((place.workspace.clone(), place.hour.clone())).or_default() += weight * tally.count as f64;
136 }
137 }
138 out
139 }
140}
141
142/// Which meters are operations, and how many each is worth.
143#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
144pub struct MappingRow {
145 pub meter: String,
146 pub cost_operations: f64,
147 pub billable_operations: f64,
148 #[serde(default)]
149 pub note: Option<String>,
150 #[serde(default)]
151 pub updated_at: String,
152}
153
154#[derive(Clone, Debug, Default)]
155pub struct Mapping {
156 rows: HashMap<String, (f64, f64)>,
157}
158
159impl Mapping {
160 pub fn from_rows(rows: &[MappingRow]) -> Self {
161 Mapping {
162 rows: rows.iter().map(|row| (row.meter.clone(), (row.cost_operations, row.billable_operations))).collect(),
163 }
164 }
165
166 /// Used until the table can be read: the same as its defaults.
167 pub fn defaults() -> Self {
168 let rows = DEFAULT_OPERATIONS
169 .iter()
170 .map(|meter| MappingRow { meter: (*meter).to_owned(), cost_operations: 1.0, billable_operations: 1.0, ..MappingRow::default() })
171 .collect::<Vec<_>>();
172 Mapping::from_rows(&rows)
173 }
174
175 pub fn billable(&self, meter: &str) -> f64 {
176 self.rows.get(meter).map_or(0.0, |(_, billable)| *billable)
177 }
178
179 #[cfg(test)]
180 pub fn cost(&self, meter: &str) -> f64 {
181 self.rows.get(meter).map_or(0.0, |(cost, _)| *cost)
182 }
183}
184
185/// The meters `operation_mapping` starts with as one operation each
186/// (migrations/0011): what Cloudflare most plausibly bills.
187pub const DEFAULT_OPERATIONS: [&str; 7] = [
188 "git.fetch",
189 "git.receive_pack",
190 "internal.git.fetch",
191 "internal.git.receive_pack",
192 "binding.create",
193 "binding.fork",
194 "binding.delete",
195];
196
197/// How long a read of `operation_mapping` is used for.
198const MAPPING_TTL_MS: u64 = 5 * 60 * 1000;
199
200thread_local! {
201 static PENDING: RefCell<Pending> = RefCell::new(Pending::default());
202 static MAPPING: RefCell<Option<(Mapping, u64)>> = const { RefCell::new(None) };
203 /// The workspace each store key belongs to, as rows were read
204 /// (registry.rs `store_key`).
205 static OWNERS: RefCell<HashMap<String, String>> = RefCell::new(HashMap::new());
206}
207
208/// Notes that the repository stored under `key` is in `workspace`.
209pub fn note_owner(key: &str, workspace: &str) {
210 OWNERS.with(|owners| {
211 let mut owners = owners.borrow_mut();
212 if owners.get(key).map(String::as_str) != Some(workspace) {
213 owners.insert(key.to_owned(), workspace.to_owned());
214 }
215 });
216}
217
218/// The workspace a key belongs to: as its row said, else what its name says
219/// (`acme--rocket`), else nothing.
220fn workspace_of(key: &str, name: &str) -> String {
221 OWNERS
222 .with(|owners| owners.borrow().get(key).cloned())
223 .unwrap_or_else(|| name.split_once("--").map(|(workspace, _)| workspace.to_owned()).unwrap_or_default())
224}
225
226/// Counts one `meter` for the repository stored under `key`, with the
227/// bytes it sent and received.
228pub fn record(meter: &str, key: &str, bytes_in: u64, bytes_out: u64) {
229 let place = place(meter, key);
230 // A workspace's local count moves now, for its free limits (git_ops.rs).
231 let weight = mapping_now().billable(meter);
232 if weight > 0.0 && !place.workspace.is_empty() {
233 crate::git_ops::note_local(&place.workspace, &place.hour, weight);
234 }
235 PENDING.with(|pending| pending.borrow_mut().meter(place, bytes_in, bytes_out));
236}
237
238/// Adds bytes to a `meter` already counted for `key`.
239pub fn record_bytes(meter: &str, key: &str, bytes_in: u64, bytes_out: u64) {
240 let place = place(meter, key);
241 PENDING.with(|pending| pending.borrow_mut().add(place, 0, bytes_in, bytes_out));
242}
243
244/// Counts one `meter` for the repository behind a git `remote`.
245pub fn record_remote(meter: &str, remote: &str, bytes_in: u64, bytes_out: u64) {
246 if let Some(key) = crate::store::key_from_remote(remote) {
247 record(meter, &key, bytes_in, bytes_out);
248 }
249}
250
251fn place(meter: &str, key: &str) -> Place {
252 place_at(meter, key, hour_now())
253}
254
255fn place_at(meter: &str, key: &str, hour: String) -> Place {
256 let (store, name) = crate::store::locate(key);
257 let workspace = workspace_of(key, &name);
258 Place { hour, store, repo: name, workspace, meter: meter.to_owned() }
259}
260
261/// Records how one call to the store in namespace `store` went.
262pub fn record_health(store: &str, outcome: Outcome, ms: u64) {
263 let minute = g1t_contracts::time::rfc3339(g1t_kit::now_ms())[..16].to_owned();
264 PENDING.with(|pending| pending.borrow_mut().health(store, &minute, outcome, ms));
265}
266
267fn hour_now() -> String {
268 g1t_contracts::time::rfc3339(g1t_kit::now_ms())[..13].to_owned()
269}
270
271/// How often an isolate writes what it counted, at most, unless a lot
272/// has piled up. Each write is one D1 batch; a page view or clone does
273/// not each need one.
274const FLUSH_EVERY_MS: u64 = 5_000;
275const FLUSH_AT_PLACES: usize = 200;
276
277thread_local! {
278 static LAST_FLUSH: std::cell::Cell<u64> = const { std::cell::Cell::new(0) };
279}
280
281/// Whether a write is due: something waits, and the last write was a
282/// while ago or much has piled up.
283pub fn flush_due(waiting: usize, last: u64, now: u64) -> bool {
284 waiting > 0 && (now.saturating_sub(last) >= FLUSH_EVERY_MS || waiting >= FLUSH_AT_PLACES)
285}
286
287/// Whether this isolate should write what it counted now; if so, the
288/// write is taken as begun.
289pub fn take_due() -> bool {
290 let now = g1t_kit::now_ms();
291 let waiting = PENDING.with(|pending| {
292 let pending = pending.borrow();
293 pending.usage.len() + pending.health.len()
294 });
295 let due = flush_due(waiting, LAST_FLUSH.with(std::cell::Cell::get), now);
296 if due {
297 LAST_FLUSH.with(|last| last.set(now));
298 }
299 due
300}
301
302/// The mapping as this isolate last read it, else the defaults. Never
303/// waits: for deciding on the request path.
304pub fn mapping_now() -> Mapping {
305 MAPPING.with(|kept| kept.borrow().as_ref().map(|(mapping, _)| mapping.clone())).unwrap_or_else(Mapping::defaults)
306}
307
308/// The mapping, read at most every few minutes; the defaults if it cannot be.
309pub async fn mapping(db: &D1Database) -> Mapping {
310 let now = g1t_kit::now_ms();
311 if let Some(mapping) = MAPPING.with(|kept| {
312 kept.borrow().as_ref().filter(|(_, at)| now.saturating_sub(*at) < MAPPING_TTL_MS).map(|(m, _)| m.clone())
313 }) {
314 return mapping;
315 }
316 match read_mapping(db).await {
317 Ok(rows) => {
318 let mapping = Mapping::from_rows(&rows);
319 MAPPING.with(|kept| *kept.borrow_mut() = Some((mapping.clone(), now)));
320 mapping
321 }
322 Err(error) => {
323 worker::console_error!("operation_mapping not read: {error}");
324 Mapping::defaults()
325 }
326 }
327}
328
329pub async fn read_mapping(db: &D1Database) -> Result<Vec<MappingRow>> {
330 db.prepare("SELECT meter, cost_operations, billable_operations, note, updated_at FROM operation_mapping ORDER BY meter")
331 .all()
332 .await?
333 .results::<MappingRow>()
334}
335
336/// Sets how many operations a meter is worth, from now on.
337pub async fn set_mapping(db: &D1Database, row: &MappingRow, now: &str) -> Result<()> {
338 db.prepare(
339 "INSERT INTO operation_mapping (meter, cost_operations, billable_operations, note, updated_at)
340 VALUES (?1, ?2, ?3, ?4, ?5)
341 ON CONFLICT (meter) DO UPDATE SET cost_operations = ?2, billable_operations = ?3, note = ?4, updated_at = ?5",
342 )
343 .bind(&[
344 row.meter.as_str().into(),
345 row.cost_operations.into(),
346 row.billable_operations.into(),
347 row.note.as_deref().map_or(JsValue::NULL, JsValue::from),
348 now.into(),
349 ])?
350 .run()
351 .await?;
352 MAPPING.with(|kept| *kept.borrow_mut() = None);
353 Ok(())
354}
355
356/// D1 runs at most this many statements in one batch, here.
357const BATCH: usize = 50;
358
359/// Writes everything counted so far; see the module docs.
360pub async fn flush(db: &D1Database) {
361 let taken = PENDING.with(|pending| std::mem::take(&mut *pending.borrow_mut()));
362 if taken.is_empty() {
363 return;
364 }
365 let mapping = mapping(db).await;
366 if let Err(error) = write(db, &taken, &mapping).await {
367 worker::console_error!("git store meters not written, kept for the next try: {error}");
368 PENDING.with(|pending| pending.borrow_mut().merge(taken));
369 }
370}
371
372async fn write(db: &D1Database, pending: &Pending, mapping: &Mapping) -> Result<()> {
373 let mut statements = Vec::new();
374 // By day in the table; by the hour in memory.
375 let mut days: HashMap<(String, String, String, String, String), Tally> = HashMap::new();
376 for (place, tally) in &pending.usage {
377 let day = days
378 .entry((place.hour[..10].to_owned(), place.store.clone(), place.repo.clone(), place.workspace.clone(), place.meter.clone()))
379 .or_default();
380 day.count += tally.count;
381 day.bytes_in += tally.bytes_in;
382 day.bytes_out += tally.bytes_out;
383 }
384 for ((day, store, repo, workspace, meter), tally) in days {
385 statements.push(
386 db.prepare(
387 "INSERT INTO artifacts_meters (day, store, repo, workspace, meter, count, bytes_in, bytes_out)
388 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
389 ON CONFLICT (day, store, repo, meter) DO UPDATE SET
390 count = count + ?6, bytes_in = bytes_in + ?7, bytes_out = bytes_out + ?8,
391 workspace = CASE WHEN ?4 = '' THEN workspace ELSE ?4 END",
392 )
393 .bind(&[
394 day.into(),
395 store.into(),
396 repo.into(),
397 workspace.into(),
398 meter.into(),
399 (tally.count as f64).into(),
400 (tally.bytes_in as f64).into(),
401 (tally.bytes_out as f64).into(),
402 ])?,
403 );
404 }
405 for ((store, minute), health) in &pending.health {
406 statements.push(
407 db.prepare(
408 "INSERT INTO store_health (store, minute, calls, errors, rate_limited, rejected, ms_total)
409 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
410 ON CONFLICT (store, minute) DO UPDATE SET
411 calls = calls + ?3, errors = errors + ?4, rate_limited = rate_limited + ?5,
412 rejected = rejected + ?6, ms_total = ms_total + ?7",
413 )
414 .bind(&[
415 store.as_str().into(),
416 minute.as_str().into(),
417 (health.calls as f64).into(),
418 (health.errors as f64).into(),
419 (health.rate_limited as f64).into(),
420 (health.rejected as f64).into(),
421 (health.ms_total as f64).into(),
422 ])?,
423 );
424 }
425 let billable = pending.billable(mapping);
426 for ((workspace, hour), operations) in &billable {
427 // Whole operations only; a fraction left over is dropped.
428 let operations = operations.round();
429 if operations < 1.0 {
430 continue;
431 }
432 statements.push(
433 db.prepare(
434 "INSERT INTO git_operations (namespace, hour, operations) VALUES (?1, ?2, ?3)
435 ON CONFLICT (namespace, hour) DO UPDATE SET operations = operations + ?3",
436 )
437 .bind(&[workspace.as_str().into(), hour.as_str().into(), operations.into()])?,
438 );
439 }
440 // Health and meters are kept for a while only.
441 if pending.health.keys().any(|(_, minute)| minute.ends_with(":00")) {
442 let cutoff = g1t_contracts::time::rfc3339(g1t_kit::now_ms().saturating_sub(24 * 3600 * 1000))[..16].to_owned();
443 statements.push(db.prepare("DELETE FROM store_health WHERE minute < ?1").bind(&[cutoff.into()])?);
444 }
445 while !statements.is_empty() {
446 let rest = statements.split_off(statements.len().min(BATCH));
447 db.batch(statements).await?;
448 statements = rest;
449 }
450 // What each workspace now stands at, for its free limits (git_ops.rs).
451 let workspaces: Vec<String> = billable.keys().map(|(workspace, _)| workspace.clone()).collect::<std::collections::BTreeSet<_>>().into_iter().collect();
452 for workspace in workspaces {
453 if let Err(error) = crate::git_ops::refresh(db, &workspace).await {
454 worker::console_error!("git operations of {workspace} not read back: {error}");
455 }
456 }
457 Ok(())
458}
459
460/// One meter's total, as `artifacts_usage` answers.
461#[derive(Debug, Serialize, Deserialize)]
462pub struct UsageRow {
463 pub day: String,
464 pub store: String,
465 #[serde(default)]
466 pub repo: Option<String>,
467 pub workspace: String,
468 pub meter: String,
469 pub count: f64,
470 pub bytes_in: f64,
471 pub bytes_out: f64,
472}
473
474/// `artifacts_usage`: the raw meters from `from` to `to` (days, inclusive).
475#[derive(Debug, Deserialize)]
476pub struct UsageArgs {
477 pub from: String,
478 pub to: String,
479 #[serde(default)]
480 pub workspace: Option<String>,
481 /// One row per repository; otherwise per workspace.
482 #[serde(default)]
483 pub by_repo: bool,
484}
485
486#[derive(Debug, Serialize)]
487pub struct Usage {
488 pub rows: Vec<UsageRow>,
489 /// Which meters count, and for how much, now.
490 pub mapping: Vec<MappingRow>,
491 /// Whether rows were left out (more than `MAX_USAGE_ROWS`).
492 pub truncated: bool,
493}
494
495const MAX_USAGE_ROWS: usize = 10_000;
496
497pub async fn usage(db: &D1Database, a: &UsageArgs) -> Result<Usage> {
498 let (repo, group) = if a.by_repo { ("repo", "day, store, repo, workspace, meter") } else { ("NULL AS repo", "day, store, workspace, meter") };
499 let sql = format!(
500 "SELECT day, store, {repo}, workspace, meter, SUM(count) AS count, SUM(bytes_in) AS bytes_in, SUM(bytes_out) AS bytes_out
501 FROM artifacts_meters WHERE day >= ?1 AND day <= ?2 AND (?3 IS NULL OR workspace = ?3)
502 GROUP BY {group} ORDER BY day, store, workspace, meter LIMIT {}",
503 MAX_USAGE_ROWS + 1
504 );
505 let mut rows = db
506 .prepare(sql)
507 .bind(&[a.from.as_str().into(), a.to.as_str().into(), a.workspace.as_deref().map_or(JsValue::NULL, JsValue::from)])?
508 .all()
509 .await?
510 .results::<UsageRow>()?;
511 let truncated = rows.len() > MAX_USAGE_ROWS;
512 rows.truncate(MAX_USAGE_ROWS);
513 Ok(Usage { rows, mapping: read_mapping(db).await?, truncated })
514}
515
516/// How the store has answered over the last minutes, per namespace, for
517/// the status page.
518#[derive(Debug, Default, Serialize, Deserialize, PartialEq)]
519pub struct StoreHealth {
520 pub store: String,
521 pub calls: f64,
522 pub errors: f64,
523 pub rate_limited: f64,
524 pub rejected: f64,
525 pub ms_total: f64,
526}
527
528#[derive(Debug, Serialize)]
529pub struct HealthReport {
530 pub minutes: u32,
531 pub stores: Vec<StoreHealth>,
532}
533
534#[derive(Debug, Deserialize)]
535pub struct HealthArgs {
536 #[serde(default)]
537 pub minutes: Option<u32>,
538}
539
540pub async fn health(db: &D1Database, a: &HealthArgs) -> Result<HealthReport> {
541 let minutes = a.minutes.unwrap_or(5).clamp(1, 60);
542 let since = g1t_contracts::time::rfc3339(g1t_kit::now_ms().saturating_sub(u64::from(minutes) * 60_000))[..16].to_owned();
543 let stores = db
544 .prepare(
545 "SELECT store, SUM(calls) AS calls, SUM(errors) AS errors, SUM(rate_limited) AS rate_limited,
546 SUM(rejected) AS rejected, SUM(ms_total) AS ms_total
547 FROM store_health WHERE minute >= ?1 GROUP BY store ORDER BY store",
548 )
549 .bind(&[since.into()])?
550 .all()
551 .await?
552 .results::<StoreHealth>()?;
553 Ok(HealthReport { minutes, stores })
554}
555
556#[cfg(test)]
557mod tests {
558 use super::*;
559
560 fn place(hour: &str, workspace: &str, meter: &str) -> Place {
561 Place { hour: hour.into(), store: "g1t".into(), repo: format!("{workspace}--rocket"), workspace: workspace.into(), meter: meter.into() }
562 }
563
564 #[test]
565 fn counts_are_written_now_and_then_not_on_every_request() {
566 assert!(!flush_due(0, 0, 100_000));
567 assert!(flush_due(1, 0, 100_000));
568 assert!(!flush_due(3, 100_000, 100_000 + FLUSH_EVERY_MS - 1));
569 assert!(flush_due(3, 100_000, 100_000 + FLUSH_EVERY_MS));
570 assert!(flush_due(FLUSH_AT_PLACES, 100_000, 100_001));
571 }
572
573 #[test]
574 fn meters_add_up_by_place() {
575 let mut pending = Pending::default();
576 pending.meter(place("2026-10-14T09", "acme", "git.fetch"), 300, 5_000);
577 pending.meter(place("2026-10-14T09", "acme", "git.fetch"), 200, 1_000);
578 pending.meter(place("2026-10-14T09", "acme", "git.ls_refs"), 100, 50);
579 assert_eq!(pending.usage[&place("2026-10-14T09", "acme", "git.fetch")], Tally { count: 2, bytes_in: 500, bytes_out: 6_000 });
580 assert_eq!(pending.usage.len(), 2);
581 }
582
583 #[test]
584 fn what_a_workspace_is_counted_for_follows_the_mapping() {
585 let mut pending = Pending::default();
586 for _ in 0..3 {
587 pending.meter(place("2026-10-14T09", "acme", "git.fetch"), 0, 0);
588 }
589 pending.meter(place("2026-10-14T09", "acme", "git.ls_refs"), 0, 0);
590 pending.meter(place("2026-10-14T09", "acme", "binding.read_tree"), 0, 0);
591 pending.meter(place("2026-10-14T10", "acme", "git.receive_pack"), 0, 0);
592 // Unknown workspace: metered, never billed.
593 pending.meter(place("2026-10-14T10", "", "git.fetch"), 0, 0);
594 let defaults = pending.billable(&Mapping::defaults());
595 assert_eq!(defaults[&("acme".to_owned(), "2026-10-14T09".to_owned())], 3.0);
596 assert_eq!(defaults[&("acme".to_owned(), "2026-10-14T10".to_owned())], 1.0);
597 assert_eq!(defaults.len(), 2);
598 // The day Cloudflare says listing refs counts too: data, not code.
599 let mut rows: Vec<MappingRow> = DEFAULT_OPERATIONS
600 .iter()
601 .map(|meter| MappingRow { meter: (*meter).into(), cost_operations: 1.0, billable_operations: 1.0, ..MappingRow::default() })
602 .collect();
603 rows.push(MappingRow { meter: "git.ls_refs".into(), cost_operations: 1.0, billable_operations: 0.5, ..MappingRow::default() });
604 let changed = pending.billable(&Mapping::from_rows(&rows));
605 assert_eq!(changed[&("acme".to_owned(), "2026-10-14T09".to_owned())], 3.5);
606 assert_eq!(Mapping::from_rows(&rows).cost("git.ls_refs"), 1.0);
607 assert_eq!(Mapping::defaults().billable("binding.read_blob"), 0.0);
608 }
609
610 #[test]
611 fn health_counts_failures_and_rejections_apart() {
612 let mut pending = Pending::default();
613 pending.health("g1t", "2026-10-14T09:01", Outcome::Ok, 40);
614 pending.health("g1t", "2026-10-14T09:01", Outcome::Failed, 900);
615 pending.health("g1t", "2026-10-14T09:01", Outcome::RateLimited, 10);
616 pending.health("g1t", "2026-10-14T09:01", Outcome::Rejected, 0);
617 let health = pending.health[&("g1t".to_owned(), "2026-10-14T09:01".to_owned())];
618 assert_eq!(health, Health { calls: 3, errors: 2, rate_limited: 1, rejected: 1, ms_total: 950 });
619 }
620
621 #[test]
622 fn a_failed_write_is_put_back() {
623 let mut kept = Pending::default();
624 kept.meter(place("2026-10-14T09", "acme", "git.fetch"), 1, 2);
625 let mut taken = Pending::default();
626 taken.meter(place("2026-10-14T09", "acme", "git.fetch"), 10, 20);
627 taken.health("g1t", "2026-10-14T09:01", Outcome::Ok, 5);
628 kept.merge(taken);
629 assert_eq!(kept.usage[&place("2026-10-14T09", "acme", "git.fetch")], Tally { count: 2, bytes_in: 11, bytes_out: 22 });
630 assert_eq!(kept.health.len(), 1);
631 }
632
633 #[test]
634 fn a_key_is_counted_for_its_workspace() {
635 assert_eq!(super::place_at("git.fetch", "g1t-us-1/acme--rocket", String::new()).store, "g1t-us-1");
636 assert_eq!(super::place_at("git.fetch", "g1t-us-1/acme--rocket", String::new()).workspace, "acme");
637 assert_eq!(super::place_at("git.fetch", "acme--rocket", String::new()).store, "g1t");
638 assert_eq!(workspace_of("zz--unknown", "zz--unknown"), "zz");
639 note_owner("rep_9", "acme");
640 assert_eq!(workspace_of("rep_9", "rep_9"), "acme");
641 assert_eq!(workspace_of("rep_8", "rep_8"), "");
642 }
643}