| 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 | .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. |
| 187 | pub 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. |
| 198 | const MAPPING_TTL_MS: u64 = 5 * 60 * 1000; |
| 199 | |
| 200 | thread_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`. |
| 209 | pub 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. |
| 220 | fn 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. |
| 228 | pub 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`. |
| 239 | pub 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`. |
| 245 | pub 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 | |
| 251 | fn place(meter: &str, key: &str) -> Place { |
| 252 | place_at(meter, key, hour_now()) |
| 253 | } |
| 254 | |
| 255 | fn 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. |
| 262 | pub 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 | |
| 267 | fn 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. |
| 274 | const FLUSH_EVERY_MS: u64 = 5_000; |
| 275 | const FLUSH_AT_PLACES: usize = 200; |
| 276 | |
| 277 | thread_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. |
| 283 | pub 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. |
| 289 | pub 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. |
| 304 | pub 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. |
| 309 | pub 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 | |
| 329 | pub 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. |
| 337 | pub 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. |
| 357 | const BATCH: usize = 50; |
| 358 | |
| 359 | /// Writes everything counted so far; see the module docs. |
| 360 | pub 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 | |
| 372 | async 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)] |
| 462 | pub 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)] |
| 476 | pub 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)] |
| 487 | pub 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 | |
| 495 | const MAX_USAGE_ROWS: usize = 10_000; |
| 496 | |
| 497 | pub 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)] |
| 519 | pub 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)] |
| 529 | pub struct HealthReport { |
| 530 | pub minutes: u32, |
| 531 | pub stores: Vec<StoreHealth>, |
| 532 | } |
| 533 | |
| 534 | #[derive(Debug, Deserialize)] |
| 535 | pub struct HealthArgs { |
| 536 | #[serde(default)] |
| 537 | pub minutes: Option<u32>, |
| 538 | } |
| 539 | |
| 540 | pub 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)] |
| 557 | mod 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 | } |