| 1 | //! Platform spend guardrails: what Cloudflare counts for all of g1t, read |
| 2 | //! every hour, and a g1t-wide pause to stop it. See |
| 3 | //! docs/SPEND-GUARDRAILS.md. |
| 4 | //! |
| 5 | //! The caps in `budget` hold what g1t pays for agents and sandboxes; this |
| 6 | //! is about the platform underneath, which no workspace's limit covers and |
| 7 | //! Cloudflare never caps: Workers requests and CPU, D1 rows, Queue |
| 8 | //! operations, Durable Objects, KV and Artifacts. A loop in any of them is |
| 9 | //! otherwise found on the bill. |
| 10 | //! |
| 11 | //! - **The pause** (`platform_pause`): four independent levels, set by |
| 12 | //! staff in sudo (Costs & margin, Platform pause) or by the watcher on a |
| 13 | //! severe breach. `compute` refuses every reservation through `reserve` |
| 14 | //! but embeddings; `indexing` refuses embeddings, and context's and |
| 15 | //! search's backfills check it themselves; `schedules` stops Actions' |
| 16 | //! cron runs and the runner's sweep; `renders` makes social cards a |
| 17 | //! static image. Callers keep the answer 30 seconds in their isolate |
| 18 | //! (`g1t_kit::pause`, `@g1t/contracts` `platformPaused`), and so does |
| 19 | //! billing (`pause_now`): never a database read per request. When the |
| 20 | //! flag cannot be read, nothing is paused: a pause is pulled on purpose, |
| 21 | //! and a billing outage must not stop the platform with it. |
| 22 | //! - **The watcher** (`watch_platform`): at the quarter past each hour, the |
| 23 | //! hour before and the month so far, from Cloudflare's GraphQL Analytics |
| 24 | //! with the costs reconciliation's token (`CLOUDFLARE_BILLING_TOKEN`, else |
| 25 | //! `CLOUDFLARE_USAGE_TOKEN`; Account Analytics Read). Each hour goes in |
| 26 | //! `platform_usage` with the script, queue, database or namespace that |
| 27 | //! counted most. A metric over its hourly threshold (`PLATFORM_HOURLY_*`), |
| 28 | //! or over `PLATFORM_SPIKE_FACTOR` times the last week's median hour (and |
| 29 | //! at least `PLATFORM_SPIKE_FLOOR_PERCENT` of its threshold), is a breach: |
| 30 | //! staff are emailed (`COSTS_ALERT_EMAIL`, once per metric every 6 |
| 31 | //! hours) and sudo shows it. Over `PLATFORM_SEVERE_FACTOR` times its |
| 32 | //! threshold, the levels that metric feeds are paused, of those |
| 33 | //! `AUTO_PAUSE` allows (`schedules,indexing` unless set). |
| 34 | //! |
| 35 | //! A dataset GraphQL refuses (a field renamed, a dataset the token cannot |
| 36 | //! read) is skipped and the others still count, but never quietly: each |
| 37 | //! run records what it could not see (`platform_watch_runs`), sudo shows |
| 38 | //! it, and a query failing `BLIND_RUNS` runs in a row, or every dataset |
| 39 | //! answering empty that long (the wrong account or token), emails staff |
| 40 | //! once a day (`watch_health`). |
| 41 | |
| 42 | use std::cell::RefCell; |
| 43 | use std::collections::BTreeMap; |
| 44 | |
| 45 | use futures_util::future::join_all; |
| 46 | use g1t_contracts::billing::{ |
| 47 | AdminSetPauseArgs, BlindQuery, PauseLevel, PauseState, PlatformBreach, PlatformGuard, PlatformMetric, PlatformPause, |
| 48 | }; |
| 49 | use g1t_contracts::time::rfc3339; |
| 50 | use g1t_contracts::{FailureCode, Outcome, new_id}; |
| 51 | use g1t_kit::now_ms; |
| 52 | use serde::Deserialize; |
| 53 | use serde_json::{Value, json}; |
| 54 | use worker::{Env, Result}; |
| 55 | |
| 56 | use crate::Billing; |
| 57 | use crate::keeper::Keeper; |
| 58 | |
| 59 | const HOUR_MS: u64 = 60 * 60 * 1000; |
| 60 | /// How long a pause read is kept in the isolate. |
| 61 | const KEEP_MS: u64 = 30_000; |
| 62 | /// A metric's breach is emailed at most this often. |
| 63 | const ALERT_EVERY_MS: u64 = 6 * HOUR_MS; |
| 64 | /// The spike rule needs this many hours of history first. |
| 65 | const SPIKE_MIN_HOURS: usize = 24; |
| 66 | |
| 67 | thread_local! { |
| 68 | static KEPT: RefCell<Option<(PlatformPause, u64)>> = const { RefCell::new(None) }; |
| 69 | } |
| 70 | |
| 71 | /// One metric the watcher reads, with its threshold's variable and the |
| 72 | /// pause levels it feeds. |
| 73 | pub(crate) struct Metric { |
| 74 | pub key: &'static str, |
| 75 | pub title: &'static str, |
| 76 | /// `PLATFORM_HOURLY_<VAR>`. |
| 77 | pub var: &'static str, |
| 78 | /// The hourly threshold when the variable is not set: about a dollar |
| 79 | /// to a few dollars an hour at Cloudflare's list prices, far above a |
| 80 | /// small alpha's normal hour. |
| 81 | pub default: f64, |
| 82 | /// What kind of thing `top_name` is, for the alert. |
| 83 | pub of: &'static str, |
| 84 | /// The levels a severe breach may pause (of those `AUTO_PAUSE` allows). |
| 85 | pub levels: &'static [PauseLevel], |
| 86 | } |
| 87 | |
| 88 | use PauseLevel::{Compute, Indexing, Renders, Schedules}; |
| 89 | |
| 90 | pub(crate) const METRICS: &[Metric] = &[ |
| 91 | Metric { key: "workers_requests", title: "Workers requests", var: "WORKERS_REQUESTS", default: 20_000_000.0, of: "script", levels: &[Schedules, Indexing, Renders] }, |
| 92 | Metric { key: "workers_cpu_ms", title: "Workers CPU (ms)", var: "WORKERS_CPU_MS", default: 100_000_000.0, of: "script", levels: &[Schedules, Indexing, Renders] }, |
| 93 | Metric { key: "d1_rows_read", title: "D1 rows read", var: "D1_ROWS_READ", default: 2_000_000_000.0, of: "D1 database id", levels: &[Schedules, Indexing] }, |
| 94 | Metric { key: "d1_rows_written", title: "D1 rows written", var: "D1_ROWS_WRITTEN", default: 5_000_000.0, of: "D1 database id", levels: &[Schedules, Indexing] }, |
| 95 | Metric { key: "queue_operations", title: "Queue operations", var: "QUEUE_OPERATIONS", default: 5_000_000.0, of: "queue id", levels: &[Schedules, Indexing] }, |
| 96 | Metric { key: "do_requests", title: "Durable Object requests", var: "DO_REQUESTS", default: 20_000_000.0, of: "script", levels: &[Compute, Schedules] }, |
| 97 | Metric { key: "do_rows_written", title: "Durable Object rows written", var: "DO_ROWS_WRITTEN", default: 5_000_000.0, of: "Durable Object namespace id", levels: &[Compute, Schedules] }, |
| 98 | Metric { key: "do_storage_write_units", title: "Durable Object storage writes", var: "DO_STORAGE_WRITE_UNITS", default: 5_000_000.0, of: "Durable Object namespace id", levels: &[Compute, Schedules] }, |
| 99 | Metric { key: "do_active_seconds", title: "Durable Object active seconds", var: "DO_ACTIVE_SECONDS", default: 3_000_000.0, of: "Durable Object namespace id", levels: &[Compute, Schedules] }, |
| 100 | Metric { key: "kv_reads", title: "KV reads", var: "KV_READS", default: 10_000_000.0, of: "KV namespace id", levels: &[Indexing, Renders] }, |
| 101 | Metric { key: "kv_writes", title: "KV writes", var: "KV_WRITES", default: 200_000.0, of: "KV namespace id", levels: &[Indexing, Renders] }, |
| 102 | Metric { key: "kv_deletes", title: "KV deletes", var: "KV_DELETES", default: 200_000.0, of: "KV namespace id", levels: &[Indexing, Renders] }, |
| 103 | Metric { key: "kv_lists", title: "KV lists", var: "KV_LISTS", default: 200_000.0, of: "KV namespace id", levels: &[Indexing, Renders] }, |
| 104 | Metric { key: "artifacts_events", title: "Artifacts events", var: "ARTIFACTS_EVENTS", default: 1_000_000.0, of: "repository", levels: &[Compute, Schedules] }, |
| 105 | ]; |
| 106 | |
| 107 | pub(crate) fn metric(key: &str) -> Option<&'static Metric> { |
| 108 | METRICS.iter().find(|m| m.key == key) |
| 109 | } |
| 110 | |
| 111 | /// The watcher's numbers, from the billing service's variables. |
| 112 | #[derive(Clone, Debug)] |
| 113 | pub(crate) struct Watch { |
| 114 | /// Each metric's hourly threshold; zero turns its threshold rule off. |
| 115 | pub thresholds: BTreeMap<&'static str, f64>, |
| 116 | /// `PLATFORM_SPIKE_FACTOR`: an hour above this many times the week's |
| 117 | /// median hour is a spike. Zero: no spike rule. |
| 118 | pub spike_factor: f64, |
| 119 | /// `PLATFORM_SPIKE_FLOOR_PERCENT`: a spike must also be at least this |
| 120 | /// share of the metric's threshold, so a quiet metric doubling is not one. |
| 121 | pub spike_floor_percent: f64, |
| 122 | /// `PLATFORM_SEVERE_FACTOR`: this many times the threshold pauses. |
| 123 | pub severe_factor: f64, |
| 124 | /// `AUTO_PAUSE`: the levels a severe breach may pause. |
| 125 | pub auto_pause: Vec<PauseLevel>, |
| 126 | } |
| 127 | |
| 128 | impl Watch { |
| 129 | pub(crate) fn from_env(env: &Env) -> Self { |
| 130 | let var = |name: &str| env.var(name).ok().map(|v| v.to_string()); |
| 131 | Watch::from_vars(var) |
| 132 | } |
| 133 | |
| 134 | pub(crate) fn from_vars(var: impl Fn(&str) -> Option<String>) -> Self { |
| 135 | let number = |name: &str, default: f64| { |
| 136 | var(name).and_then(|v| v.trim().replace('_', "").parse::<f64>().ok()).filter(|n| n.is_finite() && *n >= 0.0).unwrap_or(default) |
| 137 | }; |
| 138 | let thresholds = METRICS.iter().map(|m| (m.key, number(&format!("PLATFORM_HOURLY_{}", m.var), m.default))).collect(); |
| 139 | let auto_pause = match var("AUTO_PAUSE") { |
| 140 | Some(list) => list.split(',').filter_map(PauseLevel::parse).collect(), |
| 141 | None => vec![Schedules, Indexing], |
| 142 | }; |
| 143 | Watch { |
| 144 | thresholds, |
| 145 | spike_factor: number("PLATFORM_SPIKE_FACTOR", 10.0), |
| 146 | spike_floor_percent: number("PLATFORM_SPIKE_FLOOR_PERCENT", 10.0), |
| 147 | severe_factor: number("PLATFORM_SEVERE_FACTOR", 5.0).max(1.0), |
| 148 | auto_pause, |
| 149 | } |
| 150 | } |
| 151 | |
| 152 | pub(crate) fn threshold(&self, key: &str) -> f64 { |
| 153 | self.thresholds.get(key).copied().unwrap_or(0.0) |
| 154 | } |
| 155 | } |
| 156 | |
| 157 | // ---- The pause ------------------------------------------------------------ |
| 158 | |
| 159 | /// What a refused reservation is told while `level` is paused. |
| 160 | pub(crate) fn pause_refusal(level: PauseLevel) -> String { |
| 161 | let what = match level { |
| 162 | Compute => "new agent runs, checks, workflow jobs and builds", |
| 163 | Indexing => "new indexing and semantic search embeddings", |
| 164 | Schedules => "scheduled workflow runs and queued agents", |
| 165 | Renders => "new social card images", |
| 166 | }; |
| 167 | format!("g1t has paused {what} across the platform while it looks into unusual usage. Work already running finishes; try again later.") |
| 168 | } |
| 169 | |
| 170 | /// The level a reservation of `kind` waits on. |
| 171 | pub(crate) fn level_for(kind: g1t_contracts::billing::ComputeKind) -> PauseLevel { |
| 172 | if kind == g1t_contracts::billing::ComputeKind::Embedding { Indexing } else { Compute } |
| 173 | } |
| 174 | |
| 175 | #[derive(Deserialize)] |
| 176 | struct PauseRow { |
| 177 | level: String, |
| 178 | paused: i64, |
| 179 | note: Option<String>, |
| 180 | set_by: Option<String>, |
| 181 | set_at: Option<String>, |
| 182 | auto: i64, |
| 183 | } |
| 184 | |
| 185 | #[derive(Deserialize)] |
| 186 | struct UsageRow { |
| 187 | metric: String, |
| 188 | value: f64, |
| 189 | top_name: Option<String>, |
| 190 | top_value: Option<f64>, |
| 191 | } |
| 192 | |
| 193 | #[derive(Deserialize)] |
| 194 | struct AlertRow { |
| 195 | id: String, |
| 196 | metric: String, |
| 197 | hour: String, |
| 198 | rule: String, |
| 199 | value: f64, |
| 200 | threshold: f64, |
| 201 | severe: i64, |
| 202 | top_name: Option<String>, |
| 203 | detail: String, |
| 204 | paused: Option<String>, |
| 205 | opened_at: String, |
| 206 | emailed_at: Option<String>, |
| 207 | } |
| 208 | |
| 209 | impl Billing { |
| 210 | async fn pause_rows(&self) -> Result<Vec<PauseRow>> { |
| 211 | self.db.prepare("SELECT level, paused, note, set_by, set_at, auto FROM platform_pause").all().await?.results::<PauseRow>() |
| 212 | } |
| 213 | |
| 214 | /// Every level, kept in the isolate for 30 seconds. Nothing paused |
| 215 | /// when the table cannot be read. |
| 216 | pub(crate) async fn pause_now(&self) -> PlatformPause { |
| 217 | let now = now_ms(); |
| 218 | if let Some(pause) = KEPT.with(|kept| kept.borrow().filter(|(_, until)| *until > now).map(|(p, _)| p)) { |
| 219 | return pause; |
| 220 | } |
| 221 | let pause = match self.pause_rows().await { |
| 222 | Ok(rows) => { |
| 223 | let mut pause = PlatformPause::default(); |
| 224 | for row in rows { |
| 225 | if let Some(level) = PauseLevel::parse(&row.level) { |
| 226 | pause.set(level, row.paused != 0); |
| 227 | } |
| 228 | } |
| 229 | pause |
| 230 | } |
| 231 | Err(error) => { |
| 232 | worker::console_error!("platform pause unreadable, so nothing is paused: {error}"); |
| 233 | PlatformPause::default() |
| 234 | } |
| 235 | }; |
| 236 | KEPT.with(|kept| *kept.borrow_mut() = Some((pause, now + KEEP_MS))); |
| 237 | pause |
| 238 | } |
| 239 | |
| 240 | /// Why a reservation of `kind` is refused by a platform pause, if it is. |
| 241 | pub(crate) async fn platform_refuses(&self, kind: g1t_contracts::billing::ComputeKind) -> Option<String> { |
| 242 | let level = level_for(kind); |
| 243 | self.pause_now().await.is(level).then(|| pause_refusal(level)) |
| 244 | } |
| 245 | |
| 246 | /// Pauses or resumes a level, and records who and why in the audit log. |
| 247 | async fn set_pause(&self, level: PauseLevel, paused: bool, note: &str, by: &str, auto: bool) -> Result<()> { |
| 248 | let now = rfc3339(now_ms()); |
| 249 | self.db |
| 250 | .prepare( |
| 251 | "INSERT INTO platform_pause (level, paused, note, set_by, set_at, auto) VALUES (?1, ?2, ?3, ?4, ?5, ?6) |
| 252 | ON CONFLICT (level) DO UPDATE SET paused = ?2, note = ?3, set_by = ?4, set_at = ?5, auto = ?6", |
| 253 | ) |
| 254 | .bind(&[level.as_str().into(), i32::from(paused).into(), note.into(), by.into(), now.as_str().into(), i32::from(auto).into()])? |
| 255 | .run() |
| 256 | .await?; |
| 257 | KEPT.with(|kept| *kept.borrow_mut() = None); |
| 258 | let action = if paused { "platform_paused" } else { "platform_resumed" }; |
| 259 | self.audit("costs", action, &format!("{}: {note}", level.as_str()), by).await |
| 260 | } |
| 261 | |
| 262 | /// `admin_set_pause`: staff pause or resume one level, with why. |
| 263 | pub(crate) async fn admin_set_pause(&self, a: AdminSetPauseArgs, keeper: &Keeper) -> Result<Outcome<PlatformGuard>> { |
| 264 | let (by, note) = (a.by.trim(), a.note.trim()); |
| 265 | let Some(level) = PauseLevel::parse(&a.level) else { |
| 266 | return Ok(Outcome::fail(FailureCode::Invalid, "Pick compute, schedules, indexing or renders.")); |
| 267 | }; |
| 268 | if by.is_empty() || note.chars().count() < 5 { |
| 269 | return Ok(Outcome::fail(FailureCode::Invalid, "Say who is changing it, and why, in the note.")); |
| 270 | } |
| 271 | let note: String = note.chars().take(500).collect(); |
| 272 | self.set_pause(level, a.paused, ¬e, by, false).await?; |
| 273 | Ok(Outcome::Ok(self.platform_guard(keeper).await?)) |
| 274 | } |
| 275 | |
| 276 | /// `admin_platform_guard`: the pauses, the last hour, the month so far |
| 277 | /// and the last day's breaches. |
| 278 | pub(crate) async fn platform_guard(&self, keeper: &Keeper) -> Result<PlatformGuard> { |
| 279 | let watch = Watch::from_env(&self.env); |
| 280 | let rows = self.pause_rows().await?; |
| 281 | let levels = PauseLevel::ALL |
| 282 | .iter() |
| 283 | .map(|level| match rows.iter().find(|r| r.level == level.as_str()) { |
| 284 | Some(r) => PauseState { |
| 285 | level: r.level.clone(), |
| 286 | paused: r.paused != 0, |
| 287 | note: r.note.clone(), |
| 288 | set_by: r.set_by.clone(), |
| 289 | set_at: r.set_at.clone(), |
| 290 | auto: r.auto != 0, |
| 291 | }, |
| 292 | None => PauseState { level: level.as_str().to_owned(), ..PauseState::default() }, |
| 293 | }) |
| 294 | .collect(); |
| 295 | #[derive(Deserialize)] |
| 296 | struct Hour { |
| 297 | hour: Option<String>, |
| 298 | } |
| 299 | let hour = self.db.prepare("SELECT MAX(hour) AS hour FROM platform_usage").first::<Hour>(None).await?.and_then(|h| h.hour); |
| 300 | let to_metrics = |rows: Vec<UsageRow>| -> Vec<PlatformMetric> { |
| 301 | METRICS |
| 302 | .iter() |
| 303 | .filter_map(|m| { |
| 304 | let row = rows.iter().find(|r| r.metric == m.key)?; |
| 305 | Some(PlatformMetric { |
| 306 | metric: m.key.to_owned(), |
| 307 | title: m.title.to_owned(), |
| 308 | value: row.value, |
| 309 | threshold: watch.threshold(m.key), |
| 310 | top_name: row.top_name.clone(), |
| 311 | top_value: row.top_value, |
| 312 | }) |
| 313 | }) |
| 314 | .collect() |
| 315 | }; |
| 316 | let last_hour = match &hour { |
| 317 | Some(hour) => to_metrics( |
| 318 | self.db |
| 319 | .prepare("SELECT metric, value, top_name, top_value FROM platform_usage WHERE hour = ?") |
| 320 | .bind(&[hour.as_str().into()])? |
| 321 | .all() |
| 322 | .await? |
| 323 | .results::<UsageRow>()?, |
| 324 | ), |
| 325 | None => vec![], |
| 326 | }; |
| 327 | let month = rfc3339(now_ms())[..7].to_owned(); |
| 328 | let month_to_date = to_metrics( |
| 329 | self.db |
| 330 | .prepare("SELECT metric, value, top_name, top_value FROM platform_usage_month WHERE month = ?") |
| 331 | .bind(&[month.as_str().into()])? |
| 332 | .all() |
| 333 | .await? |
| 334 | .results::<UsageRow>()?, |
| 335 | ); |
| 336 | let breaches = self |
| 337 | .db |
| 338 | .prepare( |
| 339 | "SELECT id, metric, hour, rule, value, threshold, severe, top_name, detail, paused, opened_at, emailed_at |
| 340 | FROM platform_alerts WHERE opened_at >= ? ORDER BY opened_at DESC LIMIT 50", |
| 341 | ) |
| 342 | .bind(&[rfc3339(now_ms().saturating_sub(24 * HOUR_MS)).into()])? |
| 343 | .all() |
| 344 | .await? |
| 345 | .results::<AlertRow>()? |
| 346 | .into_iter() |
| 347 | .map(|r| PlatformBreach { |
| 348 | id: r.id, |
| 349 | metric: r.metric, |
| 350 | hour: r.hour, |
| 351 | rule: r.rule, |
| 352 | value: r.value, |
| 353 | threshold: r.threshold, |
| 354 | severe: r.severe != 0, |
| 355 | top_name: r.top_name, |
| 356 | detail: r.detail, |
| 357 | paused: r.paused.map(|p| p.split(',').filter(|s| !s.is_empty()).map(str::to_owned).collect()).unwrap_or_default(), |
| 358 | opened_at: r.opened_at, |
| 359 | emailed_at: r.emailed_at, |
| 360 | }) |
| 361 | .collect(); |
| 362 | // What the latest run could not see. |
| 363 | let latest = self.watch_runs(1).await?.into_iter().next(); |
| 364 | let blind = latest.as_ref().map_or(vec![], |(_, failed, _)| { |
| 365 | failed.iter().map(|(key, error)| BlindQuery { key: key.clone(), dataset: dataset_of(key).to_owned(), error: error.clone() }).collect() |
| 366 | }); |
| 367 | Ok(PlatformGuard { |
| 368 | blind, |
| 369 | empty: latest.as_ref().is_some_and(|(_, _, empty)| *empty), |
| 370 | last_run: latest.map(|(hour, _, _)| hour), |
| 371 | levels, |
| 372 | hour, |
| 373 | last_hour, |
| 374 | month, |
| 375 | month_to_date, |
| 376 | breaches, |
| 377 | can_read: keeper.can_read_bill(), |
| 378 | auto_pause: watch.auto_pause.iter().map(|l| l.as_str().to_owned()).collect(), |
| 379 | }) |
| 380 | } |
| 381 | } |
| 382 | |
| 383 | // ---- Reading Cloudflare's analytics --------------------------------------- |
| 384 | |
| 385 | /// One GraphQL query: a dataset, what to sum, and what to group by. |
| 386 | pub(crate) struct Query { |
| 387 | pub key: &'static str, |
| 388 | pub dataset: &'static str, |
| 389 | /// What the dataset is selected with: `sum { … }`, or `count`. |
| 390 | pub select: &'static str, |
| 391 | /// The dimension that names what counted (and `actionType` for KV). |
| 392 | pub dimensions: &'static str, |
| 393 | /// The filter fields for an hour: `datetime` (Time) on most datasets, |
| 394 | /// `datetimeHour` on D1's. |
| 395 | pub hour_filter: &'static str, |
| 396 | } |
| 397 | |
| 398 | /// Each metric from its own query where a field is less certain, so one |
| 399 | /// GraphQL refuses does not take the others with it. Field names as |
| 400 | /// Cloudflare documents them; a refused one is logged and skipped. |
| 401 | pub(crate) const QUERIES: &[Query] = &[ |
| 402 | Query { key: "workers", dataset: "workersInvocationsAdaptive", select: "sum { requests }", dimensions: "scriptName", hour_filter: "datetime" }, |
| 403 | Query { key: "workers_cpu", dataset: "workersInvocationsAdaptive", select: "sum { cpuTimeUs }", dimensions: "scriptName", hour_filter: "datetime" }, |
| 404 | Query { key: "d1", dataset: "d1AnalyticsAdaptiveGroups", select: "sum { rowsRead rowsWritten }", dimensions: "databaseId", hour_filter: "datetimeHour" }, |
| 405 | Query { key: "queues", dataset: "queueMessageOperationsAdaptiveGroups", select: "sum { billableOperations }", dimensions: "queueId", hour_filter: "datetime" }, |
| 406 | Query { key: "do_invocations", dataset: "durableObjectsInvocationsAdaptiveGroups", select: "sum { requests }", dimensions: "scriptName", hour_filter: "datetime" }, |
| 407 | Query { key: "do_periodic", dataset: "durableObjectsPeriodicGroups", select: "sum { activeTime storageWriteUnits }", dimensions: "namespaceId", hour_filter: "datetime" }, |
| 408 | Query { key: "do_sql", dataset: "durableObjectsPeriodicGroups", select: "sum { rowsWritten }", dimensions: "namespaceId", hour_filter: "datetime" }, |
| 409 | Query { key: "kv", dataset: "kvOperationsAdaptiveGroups", select: "sum { requests }", dimensions: "namespaceId actionType", hour_filter: "datetime" }, |
| 410 | Query { key: "artifacts", dataset: "artifactsEventsAdaptiveGroups", select: "count", dimensions: "repositoryName", hour_filter: "datetime" }, |
| 411 | ]; |
| 412 | |
| 413 | /// A window to read: one hour (`Time` bounds) or days (`Date` bounds). |
| 414 | #[derive(Clone, Debug, PartialEq, Eq)] |
| 415 | pub(crate) enum Window { |
| 416 | Hour { since: String, until: String }, |
| 417 | Days { since: String, until: String }, |
| 418 | } |
| 419 | |
| 420 | /// The hour before the one `now` is in, `[since, until)`, as Cloudflare's |
| 421 | /// `Time` takes it. |
| 422 | pub(crate) fn last_hour(now: u64) -> Window { |
| 423 | let until = now / HOUR_MS * HOUR_MS; |
| 424 | Window::Hour { since: hour_label(until - HOUR_MS), until: hour_label(until) } |
| 425 | } |
| 426 | |
| 427 | /// The month so far (UTC), today included. |
| 428 | pub(crate) fn month_so_far(now: u64) -> Window { |
| 429 | let today = rfc3339(now)[..10].to_owned(); |
| 430 | Window::Days { since: format!("{}-01", &today[..7]), until: today } |
| 431 | } |
| 432 | |
| 433 | /// `2026-10-08T13:00:00Z`. |
| 434 | pub(crate) fn hour_label(ms: u64) -> String { |
| 435 | format!("{}:00:00Z", &rfc3339(ms)[..13]) |
| 436 | } |
| 437 | |
| 438 | /// The request body for `query` over `window`, its dataset aliased `rows`. |
| 439 | pub(crate) fn query_body(account: &str, query: &Query, window: &Window) -> Value { |
| 440 | let (filter, kind, since, until) = match window { |
| 441 | Window::Hour { since, until } => (format!("{f}_geq: $since, {f}_lt: $until", f = query.hour_filter), "Time", since, until), |
| 442 | Window::Days { since, until } => ("date_geq: $since, date_leq: $until".to_owned(), "Date", since, until), |
| 443 | }; |
| 444 | let text = format!( |
| 445 | "query ($account: String!, $since: {kind}!, $until: {kind}!) {{ |
| 446 | viewer {{ accounts(filter: {{ accountTag: $account }}) {{ |
| 447 | rows: {dataset}(limit: 10000, filter: {{ {filter} }}) {{ |
| 448 | {select} |
| 449 | dimensions {{ {dimensions} }} |
| 450 | }} |
| 451 | }} }} |
| 452 | }}", |
| 453 | dataset = query.dataset, |
| 454 | select = query.select, |
| 455 | dimensions = query.dimensions, |
| 456 | ); |
| 457 | json!({ "query": text, "variables": { "account": account, "since": since, "until": until } }) |
| 458 | } |
| 459 | |
| 460 | /// What one query counted, per metric: the total and each name's part. |
| 461 | pub(crate) type Counted = BTreeMap<&'static str, BTreeMap<String, f64>>; |
| 462 | |
| 463 | /// A KV operation's metric, by its `actionType`. |
| 464 | fn kv_metric(action: &str) -> Option<&'static str> { |
| 465 | match action.to_ascii_lowercase().as_str() { |
| 466 | "read" | "get" => Some("kv_reads"), |
| 467 | "write" | "put" => Some("kv_writes"), |
| 468 | "delete" => Some("kv_deletes"), |
| 469 | "list" => Some("kv_lists"), |
| 470 | _ => None, |
| 471 | } |
| 472 | } |
| 473 | |
| 474 | /// Reads one query's answer into metrics. Errors when GraphQL does. |
| 475 | pub(crate) fn counted(query: &Query, body: &Value) -> std::result::Result<Counted, String> { |
| 476 | if let Some(errors) = body["errors"].as_array().filter(|e| !e.is_empty()) { |
| 477 | let messages: Vec<&str> = errors.iter().filter_map(|e| e["message"].as_str()).collect(); |
| 478 | return Err(format!("{} ({}): {}", query.dataset, query.key, messages.join("; "))); |
| 479 | } |
| 480 | let Some(groups) = body["data"]["viewer"]["accounts"][0]["rows"].as_array() else { |
| 481 | return Err(format!("{} ({}): no rows in the answer", query.dataset, query.key)); |
| 482 | }; |
| 483 | let mut out: Counted = BTreeMap::new(); |
| 484 | let mut add = |metric: &'static str, name: &str, value: f64| { |
| 485 | if value.is_finite() && value > 0.0 { |
| 486 | *out.entry(metric).or_default().entry(name.to_owned()).or_default() += value; |
| 487 | } |
| 488 | }; |
| 489 | for g in groups { |
| 490 | let sum = &g["sum"]; |
| 491 | let number = |field: &str| sum[field].as_f64().unwrap_or(0.0); |
| 492 | let dims = &g["dimensions"]; |
| 493 | let name_of = |field: &str| dims[field].as_str().filter(|s| !s.is_empty()).unwrap_or("(unnamed)").to_owned(); |
| 494 | match query.key { |
| 495 | "workers" => add("workers_requests", &name_of("scriptName"), number("requests")), |
| 496 | "workers_cpu" => add("workers_cpu_ms", &name_of("scriptName"), number("cpuTimeUs") / 1000.0), |
| 497 | "d1" => { |
| 498 | let name = name_of("databaseId"); |
| 499 | add("d1_rows_read", &name, number("rowsRead")); |
| 500 | add("d1_rows_written", &name, number("rowsWritten")); |
| 501 | } |
| 502 | "queues" => add("queue_operations", &name_of("queueId"), number("billableOperations")), |
| 503 | "do_invocations" => add("do_requests", &name_of("scriptName"), number("requests")), |
| 504 | "do_periodic" => { |
| 505 | let name = name_of("namespaceId"); |
| 506 | // activeTime is in microseconds. |
| 507 | add("do_active_seconds", &name, number("activeTime") / 1_000_000.0); |
| 508 | add("do_storage_write_units", &name, number("storageWriteUnits")); |
| 509 | } |
| 510 | "do_sql" => add("do_rows_written", &name_of("namespaceId"), number("rowsWritten")), |
| 511 | "kv" => { |
| 512 | if let Some(metric) = dims["actionType"].as_str().and_then(kv_metric) { |
| 513 | add(metric, &name_of("namespaceId"), number("requests")); |
| 514 | } |
| 515 | } |
| 516 | "artifacts" => add("artifacts_events", &name_of("repositoryName"), g["count"].as_f64().unwrap_or(0.0)), |
| 517 | _ => {} |
| 518 | } |
| 519 | } |
| 520 | Ok(out) |
| 521 | } |
| 522 | |
| 523 | /// A metric's total, and the name that counted most. |
| 524 | #[derive(Clone, Debug, PartialEq)] |
| 525 | pub(crate) struct Total { |
| 526 | pub value: f64, |
| 527 | pub top_name: Option<String>, |
| 528 | pub top_value: Option<f64>, |
| 529 | } |
| 530 | |
| 531 | /// Every query's answers merged into one total per metric. |
| 532 | pub(crate) fn totals(answers: &[Counted]) -> BTreeMap<&'static str, Total> { |
| 533 | let mut merged: Counted = BTreeMap::new(); |
| 534 | for answer in answers { |
| 535 | for (metric, names) in answer { |
| 536 | let entry = merged.entry(metric).or_default(); |
| 537 | for (name, value) in names { |
| 538 | *entry.entry(name.clone()).or_default() += value; |
| 539 | } |
| 540 | } |
| 541 | } |
| 542 | merged |
| 543 | .into_iter() |
| 544 | .map(|(metric, names)| { |
| 545 | let value = names.values().sum(); |
| 546 | let top = names.into_iter().max_by(|a, b| a.1.total_cmp(&b.1)); |
| 547 | (metric, Total { value, top_name: top.as_ref().map(|t| t.0.clone()), top_value: top.map(|t| t.1) }) |
| 548 | }) |
| 549 | .collect() |
| 550 | } |
| 551 | |
| 552 | // ---- Deciding what is a breach -------------------------------------------- |
| 553 | |
| 554 | /// The middle of the week's hours (the mean of the middle two when even). |
| 555 | pub(crate) fn median(values: &mut [f64]) -> Option<f64> { |
| 556 | if values.is_empty() { |
| 557 | return None; |
| 558 | } |
| 559 | values.sort_by(f64::total_cmp); |
| 560 | let mid = values.len() / 2; |
| 561 | Some(if values.len() % 2 == 0 { (values[mid - 1] + values[mid]) / 2.0 } else { values[mid] }) |
| 562 | } |
| 563 | |
| 564 | /// A breach of one metric, before it is recorded. |
| 565 | #[derive(Clone, Debug, PartialEq)] |
| 566 | pub(crate) struct Breach { |
| 567 | pub metric: &'static str, |
| 568 | /// `threshold` or `spike`. |
| 569 | pub rule: &'static str, |
| 570 | pub value: f64, |
| 571 | /// What it was held to: the threshold, or the spike's line. |
| 572 | pub limit: f64, |
| 573 | pub severe: bool, |
| 574 | } |
| 575 | |
| 576 | /// Whether an hour's `value` of `metric` is a breach, given the week's |
| 577 | /// earlier hours (`history`). The threshold rule wins over the spike rule. |
| 578 | pub(crate) fn judge(watch: &Watch, metric: &'static str, value: f64, history: &[f64]) -> Option<Breach> { |
| 579 | let threshold = watch.threshold(metric); |
| 580 | if threshold > 0.0 && value > threshold { |
| 581 | return Some(Breach { metric, rule: "threshold", value, limit: threshold, severe: value >= threshold * watch.severe_factor }); |
| 582 | } |
| 583 | if watch.spike_factor <= 0.0 || history.len() < SPIKE_MIN_HOURS { |
| 584 | return None; |
| 585 | } |
| 586 | let mut week = history.to_vec(); |
| 587 | let usual = median(&mut week)?; |
| 588 | let line = usual * watch.spike_factor; |
| 589 | let floor = threshold * watch.spike_floor_percent / 100.0; |
| 590 | // A metric with no threshold never spikes: there is no floor to hold it to. |
| 591 | (threshold > 0.0 && value > line && value >= floor).then_some(Breach { metric, rule: "spike", value, limit: line, severe: false }) |
| 592 | } |
| 593 | |
| 594 | /// The levels a severe breach of `metric` pauses: those it feeds that |
| 595 | /// `AUTO_PAUSE` allows. |
| 596 | pub(crate) fn levels_to_pause(watch: &Watch, breach: &Breach) -> Vec<PauseLevel> { |
| 597 | if !breach.severe { |
| 598 | return vec![]; |
| 599 | } |
| 600 | metric(breach.metric).map_or(vec![], |m| m.levels.iter().copied().filter(|l| watch.auto_pause.contains(l)).collect()) |
| 601 | } |
| 602 | |
| 603 | /// `20,000,000`, or `1.5` for small fractions. |
| 604 | pub(crate) fn amount(n: f64) -> String { |
| 605 | if n.fract().abs() > 0.0 && n.abs() < 100.0 { |
| 606 | return format!("{n:.1}"); |
| 607 | } |
| 608 | let digits = format!("{:.0}", n.abs()); |
| 609 | let mut out = String::new(); |
| 610 | for (i, c) in digits.chars().enumerate() { |
| 611 | if i > 0 && (digits.len() - i).is_multiple_of(3) { |
| 612 | out.push(','); |
| 613 | } |
| 614 | out.push(c); |
| 615 | } |
| 616 | if n < 0.0 { format!("-{out}") } else { out } |
| 617 | } |
| 618 | |
| 619 | /// The sentence an alert and sudo carry for a breach. |
| 620 | pub(crate) fn describe(breach: &Breach, hour: &str, top: Option<(&str, f64)>) -> String { |
| 621 | let m = metric(breach.metric); |
| 622 | let title = m.map_or(breach.metric, |m| m.title); |
| 623 | let of = m.map_or("source", |m| m.of); |
| 624 | let what = match breach.rule { |
| 625 | "spike" => format!("more than its spike line of {} (the last week's usual hour times PLATFORM_SPIKE_FACTOR)", amount(breach.limit)), |
| 626 | _ => format!( |
| 627 | "over its hourly threshold of {} (PLATFORM_HOURLY_{})", |
| 628 | amount(breach.limit), |
| 629 | m.map_or("?", |m| m.var) |
| 630 | ), |
| 631 | }; |
| 632 | let mut text = format!("{title}: {} in the hour from {hour}, {what}.", amount(breach.value)); |
| 633 | if let Some((name, value)) = top { |
| 634 | text.push_str(&format!(" Most of it from the {of} {name} ({}).", amount(value))); |
| 635 | } |
| 636 | text |
| 637 | } |
| 638 | |
| 639 | /// Whether the quarter-hour tick at `now` is the hour's watch: the one at |
| 640 | /// a quarter past, when the hour before is in Cloudflare's analytics. |
| 641 | pub(crate) fn hourly_due(now: u64) -> bool { |
| 642 | let minute = now / 60_000 % 60; |
| 643 | (15..30).contains(&minute) |
| 644 | } |
| 645 | |
| 646 | impl Billing { |
| 647 | /// Each hour: the hour before and the month so far, from Cloudflare, |
| 648 | /// kept and judged; staff emailed and levels paused on a breach. |
| 649 | /// Returns how many metrics were read and how many breached. |
| 650 | pub(crate) async fn watch_platform(&self, keeper: &Keeper) -> Result<(usize, usize)> { |
| 651 | if !keeper.can_read_bill() { |
| 652 | worker::console_log!("platform watch: no Cloudflare token with Account Analytics Read, so nothing is read"); |
| 653 | return Ok((0, 0)); |
| 654 | } |
| 655 | let now = now_ms(); |
| 656 | let watch = Watch::from_env(&self.env); |
| 657 | let hour_window = last_hour(now); |
| 658 | let month_window = month_so_far(now); |
| 659 | let Window::Hour { since: hour, .. } = &hour_window else { unreachable!() }; |
| 660 | let hour = hour.clone(); |
| 661 | let month = rfc3339(now)[..7].to_owned(); |
| 662 | let read = |window: &Window| { |
| 663 | join_all(QUERIES.iter().map(|query| { |
| 664 | let body = query_body(keeper.account(), query, window); |
| 665 | async move { |
| 666 | let answer = match keeper.graphql(body).await { |
| 667 | Ok(answer) => counted(query, &answer).map(|c| (c, rows_in(&answer))), |
| 668 | Err(error) => Err(format!("{} ({}): {error}", query.dataset, query.key)), |
| 669 | }; |
| 670 | (query.key, answer) |
| 671 | } |
| 672 | })) |
| 673 | }; |
| 674 | let (hour_answers, month_answers) = futures_util::future::join(read(&hour_window), read(&month_window)).await; |
| 675 | // What the watcher could not see this run, by query: loud in sudo, |
| 676 | // and emailed when it lasts (`watch_health`). |
| 677 | let mut sight = Sight::default(); |
| 678 | let mut keep = |answers: Vec<(&'static str, std::result::Result<(Counted, usize), String>)>| { |
| 679 | let mut kept = vec![]; |
| 680 | for (key, answer) in answers { |
| 681 | match answer { |
| 682 | Ok((counted, rows)) => { |
| 683 | sight.answered += 1; |
| 684 | sight.rows += rows; |
| 685 | kept.push(counted); |
| 686 | } |
| 687 | Err(error) => { |
| 688 | worker::console_error!("platform watch: skipped {error}"); |
| 689 | sight.fail(key, &error); |
| 690 | } |
| 691 | } |
| 692 | } |
| 693 | kept |
| 694 | }; |
| 695 | let hourly = totals(&keep(hour_answers)); |
| 696 | let monthly = totals(&keep(month_answers)); |
| 697 | let read_at = rfc3339(now); |
| 698 | if let Err(error) = self.watch_health(&hour, &read_at, &sight).await { |
| 699 | worker::console_error!("platform watch: could not record what it could not see: {error}"); |
| 700 | } |
| 701 | let mut writes = vec![]; |
| 702 | for (table, period, values) in [("platform_usage", "hour", &hourly), ("platform_usage_month", "month", &monthly)] { |
| 703 | let key = if period == "hour" { hour.as_str() } else { month.as_str() }; |
| 704 | for (metric, total) in values.iter() { |
| 705 | writes.push( |
| 706 | self.db |
| 707 | .prepare(format!( |
| 708 | "INSERT INTO {table} ({period}, metric, value, top_name, top_value, read_at) VALUES (?1, ?2, ?3, ?4, ?5, ?6) |
| 709 | ON CONFLICT ({period}, metric) DO UPDATE SET value = ?3, top_name = ?4, top_value = ?5, read_at = ?6" |
| 710 | )) |
| 711 | .bind(&[ |
| 712 | key.into(), |
| 713 | (*metric).into(), |
| 714 | total.value.into(), |
| 715 | crate::optional(total.top_name.as_deref()), |
| 716 | total.top_value.map_or(worker::wasm_bindgen::JsValue::NULL, Into::into), |
| 717 | read_at.as_str().into(), |
| 718 | ])?, |
| 719 | ); |
| 720 | } |
| 721 | } |
| 722 | if !writes.is_empty() { |
| 723 | self.db.batch(writes).await?; |
| 724 | } |
| 725 | // The week before this hour, for the spike rule. |
| 726 | #[derive(Deserialize)] |
| 727 | struct Past { |
| 728 | metric: String, |
| 729 | value: f64, |
| 730 | } |
| 731 | let past = self |
| 732 | .db |
| 733 | .prepare("SELECT metric, value FROM platform_usage WHERE hour >= ? AND hour < ?") |
| 734 | .bind(&[hour_label(now / HOUR_MS * HOUR_MS - 8 * 24 * HOUR_MS).into(), hour.as_str().into()])? |
| 735 | .all() |
| 736 | .await? |
| 737 | .results::<Past>()?; |
| 738 | let mut breaches = vec![]; |
| 739 | for (metric, total) in &hourly { |
| 740 | let history: Vec<f64> = past.iter().filter(|p| p.metric == *metric).map(|p| p.value).collect(); |
| 741 | if let Some(breach) = judge(&watch, metric, total.value, &history) { |
| 742 | breaches.push((breach, total.clone())); |
| 743 | } |
| 744 | } |
| 745 | let found = breaches.len(); |
| 746 | if found > 0 { |
| 747 | self.on_breaches(&watch, &hour, breaches).await?; |
| 748 | } |
| 749 | Ok((hourly.len(), found)) |
| 750 | } |
| 751 | |
| 752 | /// Records each breach, pauses what a severe one should, and emails |
| 753 | /// staff about the metrics not emailed in the last 6 hours. |
| 754 | async fn on_breaches(&self, watch: &Watch, hour: &str, breaches: Vec<(Breach, Total)>) -> Result<()> { |
| 755 | let now = now_ms(); |
| 756 | let opened_at = rfc3339(now); |
| 757 | let current = self.pause_now().await; |
| 758 | let mut lines = vec![]; |
| 759 | let mut ids = vec![]; |
| 760 | let mut paused_now: Vec<PauseLevel> = vec![]; |
| 761 | for (breach, total) in breaches { |
| 762 | let top = total.top_name.as_deref().zip(total.top_value); |
| 763 | let mut detail = describe(&breach, hour, top); |
| 764 | let pause: Vec<PauseLevel> = |
| 765 | levels_to_pause(watch, &breach).into_iter().filter(|l| !current.is(*l) && !paused_now.contains(l)).collect(); |
| 766 | for level in &pause { |
| 767 | self.set_pause(*level, true, &format!("Automatic: {detail}"), "g1t-billing's usage watcher", true).await?; |
| 768 | paused_now.push(*level); |
| 769 | } |
| 770 | if !pause.is_empty() { |
| 771 | let names: Vec<&str> = pause.iter().map(|l| l.as_str()).collect(); |
| 772 | detail.push_str(&format!(" Severe (over {}× the threshold): paused {}.", amount(watch.severe_factor), names.join(" and "))); |
| 773 | } |
| 774 | let id = new_id("pal", now); |
| 775 | let paused_text = pause.iter().map(|l| l.as_str()).collect::<Vec<_>>().join(","); |
| 776 | self.db |
| 777 | .prepare( |
| 778 | "INSERT INTO platform_alerts (id, metric, hour, rule, value, threshold, severe, top_name, top_value, detail, paused, opened_at) |
| 779 | VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", |
| 780 | ) |
| 781 | .bind(&[ |
| 782 | id.as_str().into(), |
| 783 | breach.metric.into(), |
| 784 | hour.into(), |
| 785 | breach.rule.into(), |
| 786 | breach.value.into(), |
| 787 | breach.limit.into(), |
| 788 | i32::from(breach.severe).into(), |
| 789 | crate::optional(total.top_name.as_deref()), |
| 790 | total.top_value.map_or(worker::wasm_bindgen::JsValue::NULL, Into::into), |
| 791 | detail.as_str().into(), |
| 792 | crate::optional(Some(paused_text.as_str()).filter(|t| !t.is_empty())), |
| 793 | opened_at.as_str().into(), |
| 794 | ])? |
| 795 | .run() |
| 796 | .await?; |
| 797 | // Emailed at most once per metric every 6 hours; a new pause is |
| 798 | // always said. |
| 799 | #[derive(Deserialize)] |
| 800 | struct Last { |
| 801 | at: Option<String>, |
| 802 | } |
| 803 | let last = self |
| 804 | .db |
| 805 | .prepare("SELECT MAX(emailed_at) AS at FROM platform_alerts WHERE metric = ? AND emailed_at >= ?") |
| 806 | .bind(&[breach.metric.into(), rfc3339(now.saturating_sub(ALERT_EVERY_MS)).into()])? |
| 807 | .first::<Last>(None) |
| 808 | .await? |
| 809 | .and_then(|l| l.at); |
| 810 | if last.is_none() || !pause.is_empty() { |
| 811 | lines.push(detail); |
| 812 | ids.push(id); |
| 813 | } |
| 814 | } |
| 815 | let alert_to = &self.caps.alert_to; |
| 816 | if lines.is_empty() || alert_to.is_empty() { |
| 817 | return Ok(()); |
| 818 | } |
| 819 | let subject = if paused_now.is_empty() { |
| 820 | format!("g1t: platform usage breach in the hour from {hour}") |
| 821 | } else { |
| 822 | let names: Vec<&str> = paused_now.iter().map(|l| l.as_str()).collect(); |
| 823 | format!("g1t: platform usage breach, {} paused", names.join(" and ")) |
| 824 | }; |
| 825 | lines.push( |
| 826 | "Ids are Cloudflare's: `node scripts/ops/platform-usage.mjs` names them and shows the last 24 hours. To pause or resume a level: sudo, Costs & margin, Platform pause. Thresholds: PLATFORM_HOURLY_* in services/billing/wrangler.jsonc. See docs/SPEND-GUARDRAILS.md.".to_owned(), |
| 827 | ); |
| 828 | match crate::margin::email_staff_page(&self.env, alert_to, &subject, &lines, ("Platform pause", "https://sudo.g1t.sh/costs#platform"), "g1t-billing's usage watcher").await { |
| 829 | Ok(()) => { |
| 830 | let at = rfc3339(now_ms()); |
| 831 | for id in ids { |
| 832 | self.db.prepare("UPDATE platform_alerts SET emailed_at = ? WHERE id = ?").bind(&[at.as_str().into(), id.as_str().into()])?.run().await?; |
| 833 | } |
| 834 | } |
| 835 | Err(error) => worker::console_error!("could not email the platform usage breach: {error}"), |
| 836 | } |
| 837 | Ok(()) |
| 838 | } |
| 839 | } |
| 840 | |
| 841 | // ---- What the watcher cannot see ------------------------------------------ |
| 842 | |
| 843 | /// Runs in a row a query must fail (or every dataset answer empty) before |
| 844 | /// staff are emailed: one bad hour is Cloudflare's, three are ours. |
| 845 | pub(crate) const BLIND_RUNS: usize = 3; |
| 846 | /// A blind query is emailed at most this often. |
| 847 | const BLIND_EVERY_MS: u64 = 24 * HOUR_MS; |
| 848 | /// The key `platform_watch_alerts` uses for "every dataset answered empty". |
| 849 | pub(crate) const ALL_EMPTY: &str = "all_empty"; |
| 850 | |
| 851 | /// What one run could and could not see. |
| 852 | #[derive(Clone, Debug, Default, PartialEq)] |
| 853 | pub(crate) struct Sight { |
| 854 | /// Each query that failed, in either window, with its first error. |
| 855 | pub failed: BTreeMap<String, String>, |
| 856 | /// Queries that answered, and the rows they answered with in all. |
| 857 | pub answered: usize, |
| 858 | pub rows: usize, |
| 859 | } |
| 860 | |
| 861 | impl Sight { |
| 862 | pub(crate) fn fail(&mut self, key: &str, error: &str) { |
| 863 | self.failed.entry(key.to_owned()).or_insert_with(|| error.chars().take(500).collect()); |
| 864 | } |
| 865 | |
| 866 | /// Every dataset that answered answered with nothing: no Worker ran |
| 867 | /// all month, or (likelier) the wrong account or a token that cannot |
| 868 | /// see its analytics. |
| 869 | pub(crate) fn empty(&self) -> bool { |
| 870 | self.answered > 0 && self.rows == 0 |
| 871 | } |
| 872 | } |
| 873 | |
| 874 | /// How many groups a GraphQL answer has under `rows`. |
| 875 | pub(crate) fn rows_in(body: &Value) -> usize { |
| 876 | body["data"]["viewer"]["accounts"][0]["rows"].as_array().map_or(0, Vec::len) |
| 877 | } |
| 878 | |
| 879 | /// What has been blind for the last `BLIND_RUNS` runs (newest first): the |
| 880 | /// queries failing in every one, and `all_empty` when every one was empty. |
| 881 | pub(crate) fn blind_for_long(runs: &[(BTreeMap<String, String>, bool)]) -> Vec<String> { |
| 882 | if runs.len() < BLIND_RUNS { |
| 883 | return vec![]; |
| 884 | } |
| 885 | let recent = &runs[..BLIND_RUNS]; |
| 886 | let mut out: Vec<String> = recent[0].0.keys().filter(|key| recent.iter().all(|(failed, _)| failed.contains_key(*key))).cloned().collect(); |
| 887 | if recent.iter().all(|(_, empty)| *empty) { |
| 888 | out.push(ALL_EMPTY.to_owned()); |
| 889 | } |
| 890 | out |
| 891 | } |
| 892 | |
| 893 | /// The dataset behind a query key, for people. |
| 894 | pub(crate) fn dataset_of(key: &str) -> &str { |
| 895 | QUERIES.iter().find(|q| q.key == key).map_or(key, |q| q.dataset) |
| 896 | } |
| 897 | |
| 898 | #[derive(Deserialize)] |
| 899 | struct RunRow { |
| 900 | hour: String, |
| 901 | failed: Option<String>, |
| 902 | empty: i64, |
| 903 | } |
| 904 | |
| 905 | impl Billing { |
| 906 | /// Records what this run could not see, and emails staff about what |
| 907 | /// has stayed unseen for `BLIND_RUNS` runs, once a day per query. |
| 908 | async fn watch_health(&self, hour: &str, read_at: &str, sight: &Sight) -> Result<()> { |
| 909 | let failed = serde_json::to_string(&sight.failed)?; |
| 910 | self.db |
| 911 | .batch(vec![ |
| 912 | self.db |
| 913 | .prepare( |
| 914 | "INSERT INTO platform_watch_runs (hour, read_at, failed, empty) VALUES (?1, ?2, ?3, ?4) |
| 915 | ON CONFLICT (hour) DO UPDATE SET read_at = ?2, failed = ?3, empty = ?4", |
| 916 | ) |
| 917 | .bind(&[hour.into(), read_at.into(), failed.as_str().into(), i32::from(sight.empty()).into()])?, |
| 918 | self.db.prepare("DELETE FROM platform_watch_runs WHERE hour < ?").bind(&[rfc3339(now_ms().saturating_sub(14 * 24 * HOUR_MS)).into()])?, |
| 919 | ]) |
| 920 | .await?; |
| 921 | let runs = self.watch_runs(BLIND_RUNS as u32).await?; |
| 922 | let blind = blind_for_long(&runs.iter().map(|(_, failed, empty)| (failed.clone(), *empty)).collect::<Vec<_>>()); |
| 923 | if blind.is_empty() || self.caps.alert_to.is_empty() { |
| 924 | return Ok(()); |
| 925 | } |
| 926 | #[derive(Deserialize)] |
| 927 | struct Sent { |
| 928 | key: String, |
| 929 | } |
| 930 | let since = rfc3339(now_ms().saturating_sub(BLIND_EVERY_MS)); |
| 931 | let sent: Vec<String> = self |
| 932 | .db |
| 933 | .prepare("SELECT key FROM platform_watch_alerts WHERE emailed_at >= ?") |
| 934 | .bind(&[since.into()])? |
| 935 | .all() |
| 936 | .await? |
| 937 | .results::<Sent>()? |
| 938 | .into_iter() |
| 939 | .map(|s| s.key) |
| 940 | .collect(); |
| 941 | let due: Vec<&String> = blind.iter().filter(|key| !sent.contains(key)).collect(); |
| 942 | if due.is_empty() { |
| 943 | return Ok(()); |
| 944 | } |
| 945 | let mut lines = vec![]; |
| 946 | for key in &due { |
| 947 | if key.as_str() == ALL_EMPTY { |
| 948 | lines.push(format!( |
| 949 | "Every dataset answered with no rows for the last {BLIND_RUNS} hourly runs, the month so far included. That is almost certainly the wrong account (CLOUDFLARE_ACCOUNT_ID) or a token without Account Analytics Read on it, so the watcher sees nothing at all." |
| 950 | )); |
| 951 | } else { |
| 952 | let error = sight.failed.get(key.as_str()).map_or("(no error this run)", String::as_str); |
| 953 | lines.push(format!( |
| 954 | "The watcher can't see {} ({key}): it failed for the last {BLIND_RUNS} hourly runs. Its metrics are not watched until it is fixed. Last error: {error}", |
| 955 | dataset_of(key) |
| 956 | )); |
| 957 | } |
| 958 | } |
| 959 | lines.push( |
| 960 | "Check with `node scripts/ops/platform-usage.mjs` and a token: every dataset should report rows, and it exits non-zero naming any that errored. A renamed field is fixed in QUERIES in services/billing/src/platform.rs. See docs/SPEND-GUARDRAILS.md.".to_owned(), |
| 961 | ); |
| 962 | let subject = format!("g1t: the platform usage watcher is blind on {}", due.iter().map(|k| k.as_str()).collect::<Vec<_>>().join(", ")); |
| 963 | match crate::margin::email_staff_page(&self.env, &self.caps.alert_to, &subject, &lines, ("Platform pause", "https://sudo.g1t.sh/costs#platform"), "g1t-billing's usage watcher").await { |
| 964 | Ok(()) => { |
| 965 | let at = rfc3339(now_ms()); |
| 966 | let mut writes = vec![]; |
| 967 | for key in due { |
| 968 | writes.push( |
| 969 | self.db |
| 970 | .prepare("INSERT INTO platform_watch_alerts (key, emailed_at) VALUES (?1, ?2) ON CONFLICT (key) DO UPDATE SET emailed_at = ?2") |
| 971 | .bind(&[key.as_str().into(), at.as_str().into()])?, |
| 972 | ); |
| 973 | } |
| 974 | self.db.batch(writes).await?; |
| 975 | } |
| 976 | Err(error) => worker::console_error!("could not email that the watcher is blind: {error}"), |
| 977 | } |
| 978 | Ok(()) |
| 979 | } |
| 980 | |
| 981 | /// The last `n` runs, newest first: their hour, what failed, and |
| 982 | /// whether every dataset answered empty. |
| 983 | async fn watch_runs(&self, n: u32) -> Result<Vec<(String, BTreeMap<String, String>, bool)>> { |
| 984 | Ok(self |
| 985 | .db |
| 986 | .prepare("SELECT hour, failed, empty FROM platform_watch_runs ORDER BY hour DESC LIMIT ?") |
| 987 | .bind(&[n.into()])? |
| 988 | .all() |
| 989 | .await? |
| 990 | .results::<RunRow>()? |
| 991 | .into_iter() |
| 992 | .map(|r| (r.hour, r.failed.and_then(|f| serde_json::from_str(&f).ok()).unwrap_or_default(), r.empty != 0)) |
| 993 | .collect()) |
| 994 | } |
| 995 | } |
| 996 | |
| 997 | #[cfg(test)] |
| 998 | mod tests { |
| 999 | use super::*; |
| 1000 | |
| 1001 | #[test] |
| 1002 | fn a_query_failing_three_runs_in_a_row_is_blind_and_so_is_every_dataset_empty() { |
| 1003 | let failed = |keys: &[&str]| keys.iter().map(|k| ((*k).to_owned(), "unknown field".to_owned())).collect::<BTreeMap<_, _>>(); |
| 1004 | let runs = vec![ |
| 1005 | (failed(&["workers_cpu", "do_sql"]), false), |
| 1006 | (failed(&["workers_cpu", "do_sql"]), false), |
| 1007 | (failed(&["workers_cpu"]), false), |
| 1008 | (failed(&[]), false), |
| 1009 | ]; |
| 1010 | assert_eq!(blind_for_long(&runs), vec!["workers_cpu".to_owned()]); |
| 1011 | // Two runs are not yet three. |
| 1012 | assert!(blind_for_long(&runs[..2]).is_empty()); |
| 1013 | let empty = vec![(failed(&[]), true); 3]; |
| 1014 | assert_eq!(blind_for_long(&empty), vec![ALL_EMPTY.to_owned()]); |
| 1015 | // One full hour among them clears it. |
| 1016 | assert!(blind_for_long(&[(failed(&[]), true), (failed(&[]), false), (failed(&[]), true)]).is_empty()); |
| 1017 | assert_eq!(dataset_of("do_sql"), "durableObjectsPeriodicGroups"); |
| 1018 | } |
| 1019 | |
| 1020 | #[test] |
| 1021 | fn a_run_is_empty_only_when_something_answered_with_nothing() { |
| 1022 | let mut sight = Sight::default(); |
| 1023 | // Nothing answered at all: failures, not emptiness. |
| 1024 | sight.fail("kv", "403"); |
| 1025 | sight.fail("kv", "a second error is not kept"); |
| 1026 | assert!(!sight.empty()); |
| 1027 | assert_eq!(sight.failed["kv"], "403"); |
| 1028 | sight.answered = 8; |
| 1029 | assert!(sight.empty()); |
| 1030 | sight.rows = 1; |
| 1031 | assert!(!sight.empty()); |
| 1032 | assert_eq!(rows_in(&json!({ "data": { "viewer": { "accounts": [{ "rows": [{}, {}] }] } } })), 2); |
| 1033 | assert_eq!(rows_in(&json!({ "errors": [{}] })), 0); |
| 1034 | } |
| 1035 | |
| 1036 | fn watch() -> Watch { |
| 1037 | Watch::from_vars(|_| None) |
| 1038 | } |
| 1039 | |
| 1040 | #[test] |
| 1041 | fn every_metric_has_a_threshold_and_a_query_that_counts_it() { |
| 1042 | let w = watch(); |
| 1043 | for m in METRICS { |
| 1044 | assert!(w.threshold(m.key) > 0.0, "{}", m.key); |
| 1045 | } |
| 1046 | // The thresholds named in docs/SPEND-GUARDRAILS.md. |
| 1047 | assert_eq!(w.threshold("kv_lists"), 200_000.0); |
| 1048 | assert_eq!(w.threshold("queue_operations"), 5_000_000.0); |
| 1049 | assert_eq!(w.threshold("d1_rows_read"), 2_000_000_000.0); |
| 1050 | assert_eq!(w.threshold("do_rows_written"), 5_000_000.0); |
| 1051 | assert_eq!(w.threshold("workers_requests"), 20_000_000.0); |
| 1052 | assert_eq!(w.auto_pause, vec![Schedules, Indexing]); |
| 1053 | } |
| 1054 | |
| 1055 | #[test] |
| 1056 | fn variables_change_thresholds_and_auto_pause() { |
| 1057 | let w = Watch::from_vars(|name| match name { |
| 1058 | "PLATFORM_HOURLY_KV_LISTS" => Some("50_000".into()), |
| 1059 | "PLATFORM_HOURLY_D1_ROWS_READ" => Some("0".into()), |
| 1060 | "AUTO_PAUSE" => Some("compute, renders,nonsense".into()), |
| 1061 | "PLATFORM_SEVERE_FACTOR" => Some("0.5".into()), |
| 1062 | _ => None, |
| 1063 | }); |
| 1064 | assert_eq!(w.threshold("kv_lists"), 50_000.0); |
| 1065 | assert_eq!(w.threshold("d1_rows_read"), 0.0); |
| 1066 | assert_eq!(w.auto_pause, vec![Compute, Renders]); |
| 1067 | // Never below the threshold itself. |
| 1068 | assert_eq!(w.severe_factor, 1.0); |
| 1069 | // Empty: nothing pauses by itself. |
| 1070 | assert!(Watch::from_vars(|n| (n == "AUTO_PAUSE").then(String::new)).auto_pause.is_empty()); |
| 1071 | } |
| 1072 | |
| 1073 | #[test] |
| 1074 | fn the_hour_read_is_the_one_before_at_a_quarter_past() { |
| 1075 | // 2026-10-08 13:17:00 UTC. |
| 1076 | let now = 1_791_465_420_000; |
| 1077 | assert_eq!(rfc3339(now), "2026-10-08T13:17:00.000Z"); |
| 1078 | assert!(hourly_due(now)); |
| 1079 | assert!(!hourly_due(now - 15 * 60_000)); |
| 1080 | assert!(!hourly_due(now + 15 * 60_000)); |
| 1081 | assert_eq!(last_hour(now), Window::Hour { since: "2026-10-08T12:00:00Z".into(), until: "2026-10-08T13:00:00Z".into() }); |
| 1082 | assert_eq!(month_so_far(now), Window::Days { since: "2026-10-01".into(), until: "2026-10-08".into() }); |
| 1083 | } |
| 1084 | |
| 1085 | #[test] |
| 1086 | fn queries_alias_their_dataset_and_filter_by_the_window() { |
| 1087 | let d1 = QUERIES.iter().find(|q| q.key == "d1").unwrap(); |
| 1088 | let hour = query_body("acct", d1, &last_hour(1_791_465_420_000)); |
| 1089 | let text = hour["query"].as_str().unwrap(); |
| 1090 | assert!(text.contains("rows: d1AnalyticsAdaptiveGroups(limit: 10000, filter: { datetimeHour_geq: $since, datetimeHour_lt: $until })"), "{text}"); |
| 1091 | assert!(text.contains("$since: Time!") && text.contains("sum { rowsRead rowsWritten }") && text.contains("dimensions { databaseId }")); |
| 1092 | assert_eq!(hour["variables"]["since"], "2026-10-08T12:00:00Z"); |
| 1093 | let month = query_body("acct", d1, &month_so_far(1_791_465_420_000)); |
| 1094 | let text = month["query"].as_str().unwrap(); |
| 1095 | assert!(text.contains("date_geq: $since, date_leq: $until") && text.contains("$since: Date!"), "{text}"); |
| 1096 | } |
| 1097 | |
| 1098 | fn answer(rows: Value) -> Value { |
| 1099 | json!({ "data": { "viewer": { "accounts": [{ "rows": rows }] } }, "errors": null }) |
| 1100 | } |
| 1101 | |
| 1102 | fn query(key: &str) -> &'static Query { |
| 1103 | QUERIES.iter().find(|q| q.key == key).unwrap() |
| 1104 | } |
| 1105 | |
| 1106 | #[test] |
| 1107 | fn answers_become_metrics_by_name_and_a_refused_dataset_is_skipped() { |
| 1108 | let kv = counted( |
| 1109 | query("kv"), |
| 1110 | &answer(json!([ |
| 1111 | { "sum": { "requests": 150000 }, "dimensions": { "namespaceId": "ns_a", "actionType": "list" } }, |
| 1112 | { "sum": { "requests": 90000 }, "dimensions": { "namespaceId": "ns_b", "actionType": "list" } }, |
| 1113 | { "sum": { "requests": 4000 }, "dimensions": { "namespaceId": "ns_a", "actionType": "read" } }, |
| 1114 | { "sum": { "requests": 7 }, "dimensions": { "namespaceId": "ns_a", "actionType": "mystery" } }, |
| 1115 | ])), |
| 1116 | ) |
| 1117 | .unwrap(); |
| 1118 | let d1 = counted(query("d1"), &answer(json!([{ "sum": { "rowsRead": 10.0, "rowsWritten": 2.0 }, "dimensions": { "databaseId": "db1" } }]))).unwrap(); |
| 1119 | let cpu = counted(query("workers_cpu"), &answer(json!([{ "sum": { "cpuTimeUs": 5_000_000 }, "dimensions": { "scriptName": "g1t-web" } }]))).unwrap(); |
| 1120 | let refused = counted(query("do_sql"), &json!({ "data": null, "errors": [{ "message": "unknown field rowsWritten" }] })); |
| 1121 | assert!(refused.unwrap_err().contains("unknown field rowsWritten")); |
| 1122 | assert!(counted(query("queues"), &json!({ "data": { "viewer": { "accounts": [] } } })).is_err()); |
| 1123 | let all = totals(&[kv, d1, cpu]); |
| 1124 | assert_eq!(all["kv_lists"], Total { value: 240_000.0, top_name: Some("ns_a".into()), top_value: Some(150_000.0) }); |
| 1125 | assert_eq!(all["kv_reads"].value, 4_000.0); |
| 1126 | assert_eq!(all["d1_rows_written"].value, 2.0); |
| 1127 | assert_eq!(all["workers_cpu_ms"].value, 5_000.0); |
| 1128 | assert!(!all.contains_key("do_rows_written")); |
| 1129 | } |
| 1130 | |
| 1131 | #[test] |
| 1132 | fn over_the_threshold_is_a_breach_and_five_times_it_is_severe() { |
| 1133 | let w = watch(); |
| 1134 | assert_eq!(judge(&w, "kv_lists", 200_000.0, &[]), None); |
| 1135 | let breach = judge(&w, "kv_lists", 240_000.0, &[]).unwrap(); |
| 1136 | assert_eq!((breach.rule, breach.severe), ("threshold", false)); |
| 1137 | assert!(levels_to_pause(&w, &breach).is_empty()); |
| 1138 | let severe = judge(&w, "kv_lists", 1_000_000.0, &[]).unwrap(); |
| 1139 | assert!(severe.severe); |
| 1140 | // KV feeds indexing and renders; AUTO_PAUSE allows indexing only. |
| 1141 | assert_eq!(levels_to_pause(&w, &severe), vec![Indexing]); |
| 1142 | // Durable Objects feed compute, which is never paused by itself by default. |
| 1143 | let dos = judge(&w, "do_rows_written", 30_000_000.0, &[]).unwrap(); |
| 1144 | assert_eq!(levels_to_pause(&w, &dos), vec![Schedules]); |
| 1145 | } |
| 1146 | |
| 1147 | #[test] |
| 1148 | fn a_spike_is_ten_times_the_weeks_usual_hour_above_a_floor() { |
| 1149 | let w = watch(); |
| 1150 | let quiet = vec![1_000.0; 168]; |
| 1151 | // Ten times the usual hour, but under 10% of the 5M threshold: not one. |
| 1152 | assert_eq!(judge(&w, "queue_operations", 20_000.0, &quiet), None); |
| 1153 | let usual = vec![60_000.0; 168]; |
| 1154 | let spike = judge(&w, "queue_operations", 700_000.0, &usual).unwrap(); |
| 1155 | assert_eq!((spike.rule, spike.limit, spike.severe), ("spike", 600_000.0, false)); |
| 1156 | assert_eq!(judge(&w, "queue_operations", 590_000.0, &usual), None); |
| 1157 | // Under a day of history: no spike rule yet. |
| 1158 | assert_eq!(judge(&w, "queue_operations", 700_000.0, &usual[..23]), None); |
| 1159 | assert_eq!(median(&mut [3.0, 1.0, 2.0, 10.0]), Some(2.5)); |
| 1160 | assert_eq!(median(&mut []), None); |
| 1161 | } |
| 1162 | |
| 1163 | #[test] |
| 1164 | fn an_alert_names_where_to_look() { |
| 1165 | let breach = Breach { metric: "kv_lists", rule: "threshold", value: 1_250_000.0, limit: 200_000.0, severe: true }; |
| 1166 | let text = describe(&breach, "2026-10-08T12:00:00Z", Some(("e627b571", 1_200_000.0))); |
| 1167 | assert_eq!( |
| 1168 | text, |
| 1169 | "KV lists: 1,250,000 in the hour from 2026-10-08T12:00:00Z, over its hourly threshold of 200,000 (PLATFORM_HOURLY_KV_LISTS). Most of it from the KV namespace id e627b571 (1,200,000)." |
| 1170 | ); |
| 1171 | assert_eq!(amount(1.5), "1.5"); |
| 1172 | assert_eq!(amount(5.0), "5"); |
| 1173 | } |
| 1174 | |
| 1175 | #[test] |
| 1176 | fn a_pause_refuses_the_level_a_reservation_waits_on() { |
| 1177 | use g1t_contracts::billing::ComputeKind; |
| 1178 | assert_eq!(level_for(ComputeKind::Embedding), Indexing); |
| 1179 | for kind in [ComputeKind::Agent, ComputeKind::Check, ComputeKind::Workflow, ComputeKind::Queue, ComputeKind::Deploy] { |
| 1180 | assert_eq!(level_for(kind), Compute); |
| 1181 | } |
| 1182 | assert!(pause_refusal(Compute).contains("Work already running finishes")); |
| 1183 | assert_eq!(PauseLevel::parse(" indexing "), Some(Indexing)); |
| 1184 | assert_eq!(PauseLevel::parse("everything"), None); |
| 1185 | } |
| 1186 | } |