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