g1t/services/repos/src/meters.rs

868 lines35,706 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//! Sandboxes (agents, checks, builds, workflow jobs) use git like anyone
14//! else: through g1t's git endpoints with a run credential, so their
15//! clones, fetches and pushes are metered here as `git.*`, whoever runs
16//! them in the sandbox, the agent included. Only a nightly backup's clone
17//! goes to the store directly (backups.rs), and its sandbox reports it.
18//! What is asked of a pull request's working copy (`pulls--<pull id>`,
19//! where agents clone and push) is counted for the workspace of the
20//! repository it came from, looked up when the counts are written
21//! (`Pending::attributed`); before 2026-10-07 it was counted for a
22//! workspace called `pulls`, which nobody is charged as.
23//!
24//! The same place keeps how the store answered, by the minute
25//! (`store_health`), for the status page's "Git storage" part.
26//!
27//! Nothing here is on the request path. Each isolate adds up what it saw in
28//! memory, and writes it all in one batch once the answer has gone back
29//! (`flush`, from `ctx.wait_until`), every few seconds at most. A request
30//! that counts something before the next write is due plans that write in
31//! its own `wait_until`, which waits until it is (`plan_flush`,
32//! `flush_after`): what was counted is never left for a later request on
33//! the same isolate, which may never come. Before 2026-10-07 it was, and a
34//! clone's last request (its fetch, the one that is an operation) was the
35//! one most often lost when the isolate then went idle or a deploy
36//! replaced it. A failed write puts the counts back for the next one. What
37//! an isolate holds when it dies outright is lost: a few seconds' worth at
38//! most.
39
40use std::cell::RefCell;
41use std::collections::HashMap;
42
43use serde::{Deserialize, Serialize};
44use worker::wasm_bindgen::JsValue;
45use worker::{D1Database, Result};
46
47/// How an answer from the store went, for its health.
48#[derive(Clone, Copy, Debug, PartialEq, Eq)]
49pub enum Outcome {
50 Ok,
51 /// Failed, after any retries.
52 Failed,
53 /// Refused for the store's rate limit.
54 RateLimited,
55 /// Not asked: the namespace's breaker was open (resilience.rs).
56 Rejected,
57}
58
59#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
60pub struct Tally {
61 pub count: u64,
62 pub bytes_in: u64,
63 pub bytes_out: u64,
64}
65
66#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
67pub struct Health {
68 pub calls: u64,
69 pub errors: u64,
70 pub rate_limited: u64,
71 pub rejected: u64,
72 pub ms_total: u64,
73}
74
75/// One meter's place: the hour (`YYYY-MM-DDTHH`), the git store namespace,
76/// the repository's name there, its workspace, and the meter.
77#[derive(Clone, Debug, PartialEq, Eq, Hash)]
78pub struct Place {
79 pub hour: String,
80 pub store: String,
81 pub repo: String,
82 pub workspace: String,
83 pub meter: String,
84}
85
86/// What an isolate has counted and not yet written.
87#[derive(Default, Debug)]
88pub struct Pending {
89 pub usage: HashMap<Place, Tally>,
90 /// By namespace and minute (`YYYY-MM-DDTHH:MM`).
91 pub health: HashMap<(String, String), Health>,
92}
93
94impl Pending {
95 pub fn is_empty(&self) -> bool {
96 self.usage.is_empty() && self.health.is_empty()
97 }
98
99 pub fn meter(&mut self, place: Place, bytes_in: u64, bytes_out: u64) {
100 self.add(place, 1, bytes_in, bytes_out);
101 }
102
103 pub fn add(&mut self, place: Place, count: u64, bytes_in: u64, bytes_out: u64) {
104 let tally = self.usage.entry(place).or_default();
105 tally.count += count;
106 tally.bytes_in += bytes_in;
107 tally.bytes_out += bytes_out;
108 }
109
110 pub fn health(&mut self, store: &str, minute: &str, outcome: Outcome, ms: u64) {
111 let health = self.health.entry((store.to_owned(), minute.to_owned())).or_default();
112 match outcome {
113 Outcome::Rejected => health.rejected += 1,
114 _ => {
115 health.calls += 1;
116 health.ms_total += ms;
117 match outcome {
118 Outcome::Failed => health.errors += 1,
119 Outcome::RateLimited => {
120 health.errors += 1;
121 health.rate_limited += 1;
122 }
123 _ => {}
124 }
125 }
126 }
127 }
128
129 /// Puts back what a failed write took.
130 pub fn merge(&mut self, other: Pending) {
131 for (place, tally) in other.usage {
132 let kept = self.usage.entry(place).or_default();
133 kept.count += tally.count;
134 kept.bytes_in += tally.bytes_in;
135 kept.bytes_out += tally.bytes_out;
136 }
137 for (key, health) in other.health {
138 let kept = self.health.entry(key).or_default();
139 kept.calls += health.calls;
140 kept.errors += health.errors;
141 kept.rate_limited += health.rate_limited;
142 kept.rejected += health.rejected;
143 kept.ms_total += health.ms_total;
144 }
145 }
146
147 /// The pull requests whose working copies (`pulls--<pull id>`) are
148 /// counted for nobody's workspace yet: see [`Pending::attributed`].
149 pub fn unattributed_pulls(&self) -> Vec<String> {
150 let mut pulls: Vec<String> = self
151 .usage
152 .keys()
153 .filter(|place| place.workspace == crate::PULLS_NAMESPACE)
154 .filter_map(|place| working_copy(&place.repo).map(str::to_owned))
155 .collect();
156 pulls.sort();
157 pulls.dedup();
158 pulls
159 }
160
161 /// The same counts with each pull request's working copy counted for
162 /// the workspace of the repository it came from, as `owners` (pull id
163 /// to workspace) says. A working copy's path is `pulls/<pull id>`, so
164 /// what is asked of it (an agent's clone of it and its pushes to it, a
165 /// merge check, catching up, making and removing it) would otherwise be
166 /// counted for a workspace called `pulls`, which nobody is charged as.
167 /// One `owners` does not name stays as it was.
168 pub fn attributed(&self, owners: &HashMap<String, String>) -> Pending {
169 let mut out = Pending { usage: HashMap::with_capacity(self.usage.len()), health: self.health.clone() };
170 for (place, tally) in &self.usage {
171 let mut place = place.clone();
172 if place.workspace == crate::PULLS_NAMESPACE
173 && let Some(owner) = working_copy(&place.repo).and_then(|pull| owners.get(pull))
174 {
175 place.workspace = owner.clone();
176 }
177 out.add(place, tally.count, tally.bytes_in, tally.bytes_out);
178 }
179 out
180 }
181
182 /// What each workspace is counted for, by the hour, under `mapping`.
183 pub fn billable(&self, mapping: &Mapping) -> HashMap<(String, String), f64> {
184 let mut out: HashMap<(String, String), f64> = HashMap::new();
185 for (place, tally) in &self.usage {
186 let weight = mapping.billable(&place.meter);
187 if weight > 0.0 && !place.workspace.is_empty() {
188 *out.entry((place.workspace.clone(), place.hour.clone())).or_default() += weight * tally.count as f64;
189 }
190 }
191 out
192 }
193}
194
195/// Which meters are operations, and how many each is worth.
196#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
197pub struct MappingRow {
198 pub meter: String,
199 pub cost_operations: f64,
200 pub billable_operations: f64,
201 #[serde(default)]
202 pub note: Option<String>,
203 #[serde(default)]
204 pub updated_at: String,
205}
206
207#[derive(Clone, Debug, Default)]
208pub struct Mapping {
209 rows: HashMap<String, (f64, f64)>,
210}
211
212impl Mapping {
213 pub fn from_rows(rows: &[MappingRow]) -> Self {
214 Mapping {
215 rows: rows.iter().map(|row| (row.meter.clone(), (row.cost_operations, row.billable_operations))).collect(),
216 }
217 }
218
219 /// Used until the table can be read: the same as its defaults.
220 pub fn defaults() -> Self {
221 let rows = DEFAULT_OPERATIONS
222 .iter()
223 .map(|meter| MappingRow { meter: (*meter).to_owned(), cost_operations: 1.0, billable_operations: 1.0, ..MappingRow::default() })
224 .chain(COST_ONLY_OPERATIONS.iter().map(|meter| MappingRow {
225 meter: (*meter).to_owned(),
226 cost_operations: 1.0,
227 ..MappingRow::default()
228 }))
229 .collect::<Vec<_>>();
230 Mapping::from_rows(&rows)
231 }
232
233 pub fn billable(&self, meter: &str) -> f64 {
234 self.rows.get(meter).map_or(0.0, |(_, billable)| *billable)
235 }
236
237 #[cfg(test)]
238 pub fn cost(&self, meter: &str) -> f64 {
239 self.rows.get(meter).map_or(0.0, |(cost, _)| *cost)
240 }
241}
242
243/// The meters `operation_mapping` starts with as one operation each
244/// (migrations/0011): what Cloudflare most plausibly bills.
245pub const DEFAULT_OPERATIONS: [&str; 7] = [
246 "git.fetch",
247 "git.receive_pack",
248 "internal.git.fetch",
249 "internal.git.receive_pack",
250 "binding.create",
251 "binding.fork",
252 "binding.delete",
253];
254
255/// Meters that are an operation on g1t's own bill and on no workspace's:
256/// a nightly backup's clone (backups.rs, migrations/0013).
257pub const COST_ONLY_OPERATIONS: [&str; 1] = [crate::backups::FETCH_METER];
258
259/// How long a read of `operation_mapping` is used for.
260const MAPPING_TTL_MS: u64 = 5 * 60 * 1000;
261
262thread_local! {
263 static PENDING: RefCell<Pending> = RefCell::new(Pending::default());
264 static MAPPING: RefCell<Option<(Mapping, u64)>> = const { RefCell::new(None) };
265 /// The workspace each store key belongs to, as rows were read
266 /// (registry.rs `store_key`).
267 static OWNERS: RefCell<HashMap<String, String>> = RefCell::new(HashMap::new());
268}
269
270/// Notes that the repository stored under `key` is in `workspace`.
271pub fn note_owner(key: &str, workspace: &str) {
272 OWNERS.with(|owners| {
273 let mut owners = owners.borrow_mut();
274 if owners.get(key).map(String::as_str) != Some(workspace) {
275 owners.insert(key.to_owned(), workspace.to_owned());
276 }
277 });
278}
279
280/// The workspace a key belongs to: as its row said, else what its name says
281/// (`acme--rocket`), else nothing.
282fn workspace_of(key: &str, name: &str) -> String {
283 OWNERS
284 .with(|owners| owners.borrow().get(key).cloned())
285 .unwrap_or_else(|| name.split_once("--").map(|(workspace, _)| workspace.to_owned()).unwrap_or_default())
286}
287
288/// The pull id of a working copy's name in the store (`pulls--pul_7`).
289fn working_copy(name: &str) -> Option<&str> {
290 name.strip_prefix(crate::PULLS_NAMESPACE)?.strip_prefix("--").filter(|pull| !pull.is_empty())
291}
292
293/// How long the workspace a working copy is counted for is kept before it
294/// is read again (a workspace can be renamed).
295const PULL_OWNER_TTL_MS: u64 = 10 * 60 * 1000;
296
297thread_local! {
298 /// The workspace each pull request's working copy is counted for, by
299 /// pull id, as last read, and when.
300 static PULL_OWNERS: RefCell<HashMap<String, (String, u64)>> = RefCell::new(HashMap::new());
301}
302
303/// The workspace of the repository each of `pulls`' working copies came
304/// from: as kept for a while, else read in one query. Off the request
305/// path (from `flush`).
306async fn pull_owners(db: &D1Database, pulls: &[String]) -> Result<HashMap<String, String>> {
307 let now = g1t_kit::now_ms();
308 let mut owners = HashMap::new();
309 let mut missing = Vec::new();
310 PULL_OWNERS.with(|kept| {
311 let kept = kept.borrow();
312 for pull in pulls {
313 match kept.get(pull).filter(|(_, at)| now.saturating_sub(*at) < PULL_OWNER_TTL_MS) {
314 Some((owner, _)) => {
315 owners.insert(pull.clone(), owner.clone());
316 }
317 None => missing.push(pull.clone()),
318 }
319 }
320 });
321 if missing.is_empty() {
322 return Ok(owners);
323 }
324 #[derive(Deserialize)]
325 struct Row {
326 pull: String,
327 workspace: String,
328 }
329 let rows = db
330 .prepare(
331 "SELECT f.name AS pull, s.namespace AS workspace FROM repos f JOIN repos s ON s.id = f.fork_of
332 WHERE f.namespace = ?1 AND f.name IN (SELECT value FROM json_each(?2))",
333 )
334 .bind(&[crate::PULLS_NAMESPACE.into(), serde_json::to_string(&missing)?.into()])?
335 .all()
336 .await?
337 .results::<Row>()?;
338 PULL_OWNERS.with(|kept| {
339 let mut kept = kept.borrow_mut();
340 for row in &rows {
341 kept.insert(row.pull.clone(), (row.workspace.clone(), now));
342 }
343 });
344 owners.extend(rows.into_iter().map(|row| (row.pull, row.workspace)));
345 Ok(owners)
346}
347
348/// Counts one `meter` for the repository stored under `key`, with the
349/// bytes it sent and received.
350pub fn record(meter: &str, key: &str, bytes_in: u64, bytes_out: u64) {
351 let place = place(meter, key);
352 // A workspace's local count moves now, for its free limits (git_ops.rs).
353 let weight = mapping_now().billable(meter);
354 if weight > 0.0 && !place.workspace.is_empty() {
355 crate::git_ops::note_local(&place.workspace, &place.hour, weight);
356 }
357 PENDING.with(|pending| pending.borrow_mut().meter(place, bytes_in, bytes_out));
358}
359
360/// Adds bytes to a `meter` already counted for `key`.
361pub fn record_bytes(meter: &str, key: &str, bytes_in: u64, bytes_out: u64) {
362 let place = place(meter, key);
363 PENDING.with(|pending| pending.borrow_mut().add(place, 0, bytes_in, bytes_out));
364}
365
366/// Counts one `meter` for the repository behind a git `remote`.
367pub fn record_remote(meter: &str, remote: &str, bytes_in: u64, bytes_out: u64) {
368 if let Some(key) = crate::store::key_from_remote(remote) {
369 record(meter, &key, bytes_in, bytes_out);
370 }
371}
372
373fn place(meter: &str, key: &str) -> Place {
374 place_at(meter, key, hour_now())
375}
376
377fn place_at(meter: &str, key: &str, hour: String) -> Place {
378 let (store, name) = crate::store::locate(key);
379 let workspace = workspace_of(key, &name);
380 Place { hour, store, repo: name, workspace, meter: meter.to_owned() }
381}
382
383/// Records how one call to the store in namespace `store` went.
384pub fn record_health(store: &str, outcome: Outcome, ms: u64) {
385 let minute = g1t_contracts::time::rfc3339(g1t_kit::now_ms())[..16].to_owned();
386 PENDING.with(|pending| pending.borrow_mut().health(store, &minute, outcome, ms));
387}
388
389fn hour_now() -> String {
390 g1t_contracts::time::rfc3339(g1t_kit::now_ms())[..13].to_owned()
391}
392
393/// How often an isolate writes what it counted, at most, unless a lot
394/// has piled up. Each write is one D1 batch; a page view or clone does
395/// not each need one.
396const FLUSH_EVERY_MS: u64 = 5_000;
397const FLUSH_AT_PLACES: usize = 200;
398/// A planned write not begun after this long is taken as never coming (its
399/// request's `wait_until` was cut short), so another is planned.
400const PLAN_STALE_MS: u64 = 30_000;
401
402thread_local! {
403 static LAST_FLUSH: std::cell::Cell<u64> = const { std::cell::Cell::new(0) };
404 /// When the write now waiting in some request's `wait_until` was planned.
405 static PLANNED: std::cell::Cell<Option<u64>> = const { std::cell::Cell::new(None) };
406}
407
408/// How long the write a request should plan waits, in milliseconds, or
409/// `None` when it should plan none: nothing waits, or a write is planned
410/// already and has not gone stale. A pile-up is written at once, planned
411/// write or not.
412pub fn plan(waiting: usize, last: u64, planned: Option<u64>, now: u64) -> Option<u64> {
413 if waiting == 0 {
414 return None;
415 }
416 if waiting >= FLUSH_AT_PLACES {
417 return Some(0);
418 }
419 if planned.is_some_and(|at| now.saturating_sub(at) < PLAN_STALE_MS) {
420 return None;
421 }
422 Some(FLUSH_EVERY_MS.saturating_sub(now.saturating_sub(last)))
423}
424
425/// The write this request should plan, if any, as how long it waits (see
426/// [`plan`]); taken as planned. Never waits itself.
427pub fn plan_flush() -> Option<u64> {
428 let now = g1t_kit::now_ms();
429 let waiting = PENDING.with(|pending| {
430 let pending = pending.borrow();
431 pending.usage.len() + pending.health.len()
432 });
433 let wait = plan(waiting, LAST_FLUSH.with(std::cell::Cell::get), PLANNED.with(std::cell::Cell::get), now)?;
434 if wait > 0 {
435 PLANNED.with(|planned| planned.set(Some(now)));
436 }
437 Some(wait)
438}
439
440/// A planned write: waits `wait_ms`, then writes everything counted by
441/// then. For `ctx.wait_until`, after the answer has gone back.
442pub async fn flush_after(db: &D1Database, wait_ms: u64) {
443 if wait_ms > 0 {
444 worker::Delay::from(std::time::Duration::from_millis(wait_ms)).await;
445 PLANNED.with(|planned| planned.set(None));
446 }
447 LAST_FLUSH.with(|last| last.set(g1t_kit::now_ms()));
448 flush(db).await;
449}
450
451/// The mapping as this isolate last read it, else the defaults. Never
452/// waits: for deciding on the request path.
453pub fn mapping_now() -> Mapping {
454 MAPPING.with(|kept| kept.borrow().as_ref().map(|(mapping, _)| mapping.clone())).unwrap_or_else(Mapping::defaults)
455}
456
457/// The mapping, read at most every few minutes; the defaults if it cannot be.
458pub async fn mapping(db: &D1Database) -> Mapping {
459 let now = g1t_kit::now_ms();
460 if let Some(mapping) = MAPPING.with(|kept| {
461 kept.borrow().as_ref().filter(|(_, at)| now.saturating_sub(*at) < MAPPING_TTL_MS).map(|(m, _)| m.clone())
462 }) {
463 return mapping;
464 }
465 match read_mapping(db).await {
466 Ok(rows) => {
467 let mapping = Mapping::from_rows(&rows);
468 MAPPING.with(|kept| *kept.borrow_mut() = Some((mapping.clone(), now)));
469 mapping
470 }
471 Err(error) => {
472 worker::console_error!("operation_mapping not read: {error}");
473 Mapping::defaults()
474 }
475 }
476}
477
478pub async fn read_mapping(db: &D1Database) -> Result<Vec<MappingRow>> {
479 db.prepare("SELECT meter, cost_operations, billable_operations, note, updated_at FROM operation_mapping ORDER BY meter")
480 .all()
481 .await?
482 .results::<MappingRow>()
483}
484
485/// Sets how many operations a meter is worth, from now on.
486pub async fn set_mapping(db: &D1Database, row: &MappingRow, now: &str) -> Result<()> {
487 db.prepare(
488 "INSERT INTO operation_mapping (meter, cost_operations, billable_operations, note, updated_at)
489 VALUES (?1, ?2, ?3, ?4, ?5)
490 ON CONFLICT (meter) DO UPDATE SET cost_operations = ?2, billable_operations = ?3, note = ?4, updated_at = ?5",
491 )
492 .bind(&[
493 row.meter.as_str().into(),
494 row.cost_operations.into(),
495 row.billable_operations.into(),
496 row.note.as_deref().map_or(JsValue::NULL, JsValue::from),
497 now.into(),
498 ])?
499 .run()
500 .await?;
501 MAPPING.with(|kept| *kept.borrow_mut() = None);
502 Ok(())
503}
504
505/// D1 runs at most this many statements in one batch, here.
506const BATCH: usize = 50;
507
508/// Writes everything counted so far; see the module docs.
509pub async fn flush(db: &D1Database) {
510 let taken = PENDING.with(|pending| std::mem::take(&mut *pending.borrow_mut()));
511 if taken.is_empty() {
512 return;
513 }
514 let mapping = mapping(db).await;
515 // Pull requests' working copies count for their repositories' workspaces.
516 let pulls = taken.unattributed_pulls();
517 let attributed = if pulls.is_empty() {
518 None
519 } else {
520 match pull_owners(db, &pulls).await {
521 Ok(owners) => Some(taken.attributed(&owners)),
522 Err(error) => {
523 worker::console_error!("working copies' workspaces not read, counted as they are: {error}");
524 None
525 }
526 }
527 };
528 if let Err(error) = write(db, attributed.as_ref().unwrap_or(&taken), &mapping).await {
529 worker::console_error!("git store meters not written, kept for the next try: {error}");
530 PENDING.with(|pending| pending.borrow_mut().merge(taken));
531 }
532}
533
534async fn write(db: &D1Database, pending: &Pending, mapping: &Mapping) -> Result<()> {
535 let mut statements = Vec::new();
536 // By day in the table; by the hour in memory.
537 let mut days: HashMap<(String, String, String, String, String), Tally> = HashMap::new();
538 for (place, tally) in &pending.usage {
539 let day = days
540 .entry((place.hour[..10].to_owned(), place.store.clone(), place.repo.clone(), place.workspace.clone(), place.meter.clone()))
541 .or_default();
542 day.count += tally.count;
543 day.bytes_in += tally.bytes_in;
544 day.bytes_out += tally.bytes_out;
545 }
546 for ((day, store, repo, workspace, meter), tally) in days {
547 statements.push(
548 db.prepare(
549 "INSERT INTO artifacts_meters (day, store, repo, workspace, meter, count, bytes_in, bytes_out)
550 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
551 ON CONFLICT (day, store, repo, meter) DO UPDATE SET
552 count = count + ?6, bytes_in = bytes_in + ?7, bytes_out = bytes_out + ?8,
553 workspace = CASE WHEN ?4 = '' THEN workspace ELSE ?4 END",
554 )
555 .bind(&[
556 day.into(),
557 store.into(),
558 repo.into(),
559 workspace.into(),
560 meter.into(),
561 (tally.count as f64).into(),
562 (tally.bytes_in as f64).into(),
563 (tally.bytes_out as f64).into(),
564 ])?,
565 );
566 }
567 for ((store, minute), health) in &pending.health {
568 statements.push(
569 db.prepare(
570 "INSERT INTO store_health (store, minute, calls, errors, rate_limited, rejected, ms_total)
571 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
572 ON CONFLICT (store, minute) DO UPDATE SET
573 calls = calls + ?3, errors = errors + ?4, rate_limited = rate_limited + ?5,
574 rejected = rejected + ?6, ms_total = ms_total + ?7",
575 )
576 .bind(&[
577 store.as_str().into(),
578 minute.as_str().into(),
579 (health.calls as f64).into(),
580 (health.errors as f64).into(),
581 (health.rate_limited as f64).into(),
582 (health.rejected as f64).into(),
583 (health.ms_total as f64).into(),
584 ])?,
585 );
586 }
587 let billable = pending.billable(mapping);
588 for ((workspace, hour), operations) in &billable {
589 // Whole operations only; a fraction left over is dropped.
590 let operations = operations.round();
591 if operations < 1.0 {
592 continue;
593 }
594 statements.push(
595 db.prepare(
596 "INSERT INTO git_operations (namespace, hour, operations) VALUES (?1, ?2, ?3)
597 ON CONFLICT (namespace, hour) DO UPDATE SET operations = operations + ?3",
598 )
599 .bind(&[workspace.as_str().into(), hour.as_str().into(), operations.into()])?,
600 );
601 }
602 // Health and meters are kept for a while only.
603 if pending.health.keys().any(|(_, minute)| minute.ends_with(":00")) {
604 let cutoff = g1t_contracts::time::rfc3339(g1t_kit::now_ms().saturating_sub(24 * 3600 * 1000))[..16].to_owned();
605 statements.push(db.prepare("DELETE FROM store_health WHERE minute < ?1").bind(&[cutoff.into()])?);
606 }
607 while !statements.is_empty() {
608 let rest = statements.split_off(statements.len().min(BATCH));
609 db.batch(statements).await?;
610 statements = rest;
611 }
612 // What each workspace now stands at, for its free limits (git_ops.rs).
613 let workspaces: Vec<String> = billable.keys().map(|(workspace, _)| workspace.clone()).collect::<std::collections::BTreeSet<_>>().into_iter().collect();
614 for workspace in workspaces {
615 if let Err(error) = crate::git_ops::refresh(db, &workspace).await {
616 worker::console_error!("git operations of {workspace} not read back: {error}");
617 }
618 }
619 Ok(())
620}
621
622/// One meter's total, as `artifacts_usage` answers.
623#[derive(Debug, Serialize, Deserialize)]
624pub struct UsageRow {
625 pub day: String,
626 pub store: String,
627 #[serde(default)]
628 pub repo: Option<String>,
629 pub workspace: String,
630 pub meter: String,
631 pub count: f64,
632 pub bytes_in: f64,
633 pub bytes_out: f64,
634}
635
636/// `artifacts_usage`: the raw meters from `from` to `to` (days, inclusive).
637#[derive(Debug, Deserialize)]
638pub struct UsageArgs {
639 pub from: String,
640 pub to: String,
641 #[serde(default)]
642 pub workspace: Option<String>,
643 /// One row per repository; otherwise per workspace.
644 #[serde(default)]
645 pub by_repo: bool,
646}
647
648#[derive(Debug, Serialize)]
649pub struct Usage {
650 pub rows: Vec<UsageRow>,
651 /// Which meters count, and for how much, now.
652 pub mapping: Vec<MappingRow>,
653 /// Whether rows were left out (more than `MAX_USAGE_ROWS`).
654 pub truncated: bool,
655}
656
657const MAX_USAGE_ROWS: usize = 10_000;
658
659pub async fn usage(db: &D1Database, a: &UsageArgs) -> Result<Usage> {
660 let (repo, group) = if a.by_repo { ("repo", "day, store, repo, workspace, meter") } else { ("NULL AS repo", "day, store, workspace, meter") };
661 let sql = format!(
662 "SELECT day, store, {repo}, workspace, meter, SUM(count) AS count, SUM(bytes_in) AS bytes_in, SUM(bytes_out) AS bytes_out
663 FROM artifacts_meters WHERE day >= ?1 AND day <= ?2 AND (?3 IS NULL OR workspace = ?3)
664 GROUP BY {group} ORDER BY day, store, workspace, meter LIMIT {}",
665 MAX_USAGE_ROWS + 1
666 );
667 let mut rows = db
668 .prepare(sql)
669 .bind(&[a.from.as_str().into(), a.to.as_str().into(), a.workspace.as_deref().map_or(JsValue::NULL, JsValue::from)])?
670 .all()
671 .await?
672 .results::<UsageRow>()?;
673 let truncated = rows.len() > MAX_USAGE_ROWS;
674 rows.truncate(MAX_USAGE_ROWS);
675 Ok(Usage { rows, mapping: read_mapping(db).await?, truncated })
676}
677
678/// How the store has answered over the last minutes, per namespace, for
679/// the status page.
680#[derive(Debug, Default, Serialize, Deserialize, PartialEq)]
681pub struct StoreHealth {
682 pub store: String,
683 pub calls: f64,
684 pub errors: f64,
685 pub rate_limited: f64,
686 pub rejected: f64,
687 pub ms_total: f64,
688}
689
690#[derive(Debug, Serialize)]
691pub struct HealthReport {
692 pub minutes: u32,
693 pub stores: Vec<StoreHealth>,
694}
695
696#[derive(Debug, Deserialize)]
697pub struct HealthArgs {
698 #[serde(default)]
699 pub minutes: Option<u32>,
700}
701
702pub async fn health(db: &D1Database, a: &HealthArgs) -> Result<HealthReport> {
703 let minutes = a.minutes.unwrap_or(5).clamp(1, 60);
704 let since = g1t_contracts::time::rfc3339(g1t_kit::now_ms().saturating_sub(u64::from(minutes) * 60_000))[..16].to_owned();
705 let stores = db
706 .prepare(
707 "SELECT store, SUM(calls) AS calls, SUM(errors) AS errors, SUM(rate_limited) AS rate_limited,
708 SUM(rejected) AS rejected, SUM(ms_total) AS ms_total
709 FROM store_health WHERE minute >= ?1 GROUP BY store ORDER BY store",
710 )
711 .bind(&[since.into()])?
712 .all()
713 .await?
714 .results::<StoreHealth>()?;
715 Ok(HealthReport { minutes, stores })
716}
717
718#[cfg(test)]
719mod tests {
720 use super::*;
721
722 fn place(hour: &str, workspace: &str, meter: &str) -> Place {
723 Place { hour: hour.into(), store: "g1t".into(), repo: format!("{workspace}--rocket"), workspace: workspace.into(), meter: meter.into() }
724 }
725
726 #[test]
727 fn counts_are_written_now_and_then_not_on_every_request() {
728 // At once only when the last write was a while ago or much has
729 // piled up; otherwise once it is due.
730 assert_eq!(plan(1, 0, None, 100_000), Some(0));
731 assert_eq!(plan(3, 100_000, None, 100_000 + FLUSH_EVERY_MS - 1), Some(1));
732 assert_eq!(plan(3, 100_000, None, 100_000 + FLUSH_EVERY_MS), Some(0));
733 assert_eq!(plan(FLUSH_AT_PLACES, 100_000, None, 100_001), Some(0));
734 }
735
736 #[test]
737 fn what_a_request_counts_is_written_even_if_no_request_follows() {
738 // Nothing waiting: nothing planned.
739 assert_eq!(plan(0, 0, None, 100_000), None);
740 // Due: written at once.
741 assert_eq!(plan(1, 0, None, 100_000), Some(0));
742 // Counted a second after the last write: the write is planned for
743 // when it is due, not left for the next request.
744 assert_eq!(plan(3, 100_000, None, 101_000), Some(FLUSH_EVERY_MS - 1_000));
745 // One planned already: the requests after it plan none...
746 assert_eq!(plan(3, 100_000, Some(101_000), 102_000), None);
747 // ...unless much has piled up, which is written at once,
748 assert_eq!(plan(FLUSH_AT_PLACES, 100_000, Some(101_000), 102_000), Some(0));
749 // or the planned one never began (its request was cut short).
750 assert_eq!(plan(3, 100_000, Some(101_000), 101_000 + PLAN_STALE_MS), Some(0));
751 // Never planned further off than the interval.
752 assert!(plan(1, 100_000, None, 100_000).is_some_and(|wait| wait <= FLUSH_EVERY_MS));
753 }
754
755 #[test]
756 fn meters_add_up_by_place() {
757 let mut pending = Pending::default();
758 pending.meter(place("2026-10-14T09", "acme", "git.fetch"), 300, 5_000);
759 pending.meter(place("2026-10-14T09", "acme", "git.fetch"), 200, 1_000);
760 pending.meter(place("2026-10-14T09", "acme", "git.ls_refs"), 100, 50);
761 assert_eq!(pending.usage[&place("2026-10-14T09", "acme", "git.fetch")], Tally { count: 2, bytes_in: 500, bytes_out: 6_000 });
762 assert_eq!(pending.usage.len(), 2);
763 }
764
765 #[test]
766 fn what_a_workspace_is_counted_for_follows_the_mapping() {
767 let mut pending = Pending::default();
768 for _ in 0..3 {
769 pending.meter(place("2026-10-14T09", "acme", "git.fetch"), 0, 0);
770 }
771 pending.meter(place("2026-10-14T09", "acme", "git.ls_refs"), 0, 0);
772 pending.meter(place("2026-10-14T09", "acme", "binding.read_tree"), 0, 0);
773 pending.meter(place("2026-10-14T10", "acme", "git.receive_pack"), 0, 0);
774 // Unknown workspace: metered, never billed.
775 pending.meter(place("2026-10-14T10", "", "git.fetch"), 0, 0);
776 let defaults = pending.billable(&Mapping::defaults());
777 assert_eq!(defaults[&("acme".to_owned(), "2026-10-14T09".to_owned())], 3.0);
778 assert_eq!(defaults[&("acme".to_owned(), "2026-10-14T10".to_owned())], 1.0);
779 assert_eq!(defaults.len(), 2);
780 // The day Cloudflare says listing refs counts too: data, not code.
781 let mut rows: Vec<MappingRow> = DEFAULT_OPERATIONS
782 .iter()
783 .map(|meter| MappingRow { meter: (*meter).into(), cost_operations: 1.0, billable_operations: 1.0, ..MappingRow::default() })
784 .collect();
785 rows.push(MappingRow { meter: "git.ls_refs".into(), cost_operations: 1.0, billable_operations: 0.5, ..MappingRow::default() });
786 let changed = pending.billable(&Mapping::from_rows(&rows));
787 assert_eq!(changed[&("acme".to_owned(), "2026-10-14T09".to_owned())], 3.5);
788 assert_eq!(Mapping::from_rows(&rows).cost("git.ls_refs"), 1.0);
789 assert_eq!(Mapping::defaults().billable("binding.read_blob"), 0.0);
790 }
791
792 #[test]
793 fn a_working_copy_is_counted_for_its_repositorys_workspace() {
794 let copy = |meter: &str, pull: &str| Place {
795 hour: "2026-10-14T09".into(),
796 store: "g1t".into(),
797 repo: format!("pulls--{pull}"),
798 workspace: "pulls".into(),
799 meter: meter.into(),
800 };
801 let mut pending = Pending::default();
802 // A pull request's working copy is made, an agent's sandbox clones
803 // it and pushes to it, and checks clone it again; another pull
804 // request's repository is not found.
805 pending.meter(copy("binding.fork", "pul_7"), 0, 0);
806 pending.meter(copy("git.fetch", "pul_7"), 100, 9_000);
807 pending.meter(copy("git.receive_pack", "pul_7"), 4_000, 50);
808 pending.meter(copy("git.fetch", "pul_7"), 100, 9_000);
809 pending.meter(copy("git.fetch", "pul_8"), 0, 0);
810 pending.meter(place("2026-10-14T09", "acme", "git.fetch"), 0, 0);
811 pending.health("g1t", "2026-10-14T09:01", Outcome::Ok, 5);
812 assert_eq!(pending.unattributed_pulls(), ["pul_7", "pul_8"]);
813 // As they are: counted for a workspace called `pulls`.
814 let before = pending.billable(&Mapping::defaults());
815 assert_eq!(before[&("pulls".to_owned(), "2026-10-14T09".to_owned())], 5.0);
816
817 let owners = HashMap::from([("pul_7".to_owned(), "acme".to_owned())]);
818 let attributed = pending.attributed(&owners);
819 let billable = attributed.billable(&Mapping::defaults());
820 assert_eq!(billable[&("acme".to_owned(), "2026-10-14T09".to_owned())], 5.0);
821 assert_eq!(billable[&("pulls".to_owned(), "2026-10-14T09".to_owned())], 1.0);
822 // The meters keep the working copy's own name, with its workspace.
823 let mut counted = copy("git.fetch", "pul_7");
824 counted.workspace = "acme".into();
825 assert_eq!(attributed.usage[&counted], Tally { count: 2, bytes_in: 200, bytes_out: 18_000 });
826 assert_eq!(attributed.health.len(), 1);
827 assert_eq!(attributed.unattributed_pulls(), ["pul_8"]);
828 // Nothing else is a working copy.
829 assert_eq!(working_copy("pulls--pul_7"), Some("pul_7"));
830 assert_eq!(working_copy("pulls--"), None);
831 assert_eq!(working_copy("pullsx--pul_7"), None);
832 assert_eq!(working_copy("acme--pulls"), None);
833 }
834
835 #[test]
836 fn health_counts_failures_and_rejections_apart() {
837 let mut pending = Pending::default();
838 pending.health("g1t", "2026-10-14T09:01", Outcome::Ok, 40);
839 pending.health("g1t", "2026-10-14T09:01", Outcome::Failed, 900);
840 pending.health("g1t", "2026-10-14T09:01", Outcome::RateLimited, 10);
841 pending.health("g1t", "2026-10-14T09:01", Outcome::Rejected, 0);
842 let health = pending.health[&("g1t".to_owned(), "2026-10-14T09:01".to_owned())];
843 assert_eq!(health, Health { calls: 3, errors: 2, rate_limited: 1, rejected: 1, ms_total: 950 });
844 }
845
846 #[test]
847 fn a_failed_write_is_put_back() {
848 let mut kept = Pending::default();
849 kept.meter(place("2026-10-14T09", "acme", "git.fetch"), 1, 2);
850 let mut taken = Pending::default();
851 taken.meter(place("2026-10-14T09", "acme", "git.fetch"), 10, 20);
852 taken.health("g1t", "2026-10-14T09:01", Outcome::Ok, 5);
853 kept.merge(taken);
854 assert_eq!(kept.usage[&place("2026-10-14T09", "acme", "git.fetch")], Tally { count: 2, bytes_in: 11, bytes_out: 22 });
855 assert_eq!(kept.health.len(), 1);
856 }
857
858 #[test]
859 fn a_key_is_counted_for_its_workspace() {
860 assert_eq!(super::place_at("git.fetch", "g1t-us-1/acme--rocket", String::new()).store, "g1t-us-1");
861 assert_eq!(super::place_at("git.fetch", "g1t-us-1/acme--rocket", String::new()).workspace, "acme");
862 assert_eq!(super::place_at("git.fetch", "acme--rocket", String::new()).store, "g1t");
863 assert_eq!(workspace_of("zz--unknown", "zz--unknown"), "zz");
864 note_owner("rep_9", "acme");
865 assert_eq!(workspace_of("rep_9", "rep_9"), "acme");
866 assert_eq!(workspace_of("rep_8", "rep_8"), "");
867 }
868}