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