Skip to content
1,186 linesCodeBlameRaw
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
42use std::cell::RefCell;
43use std::collections::BTreeMap;
44
45use futures_util::future::join_all;
46use g1t_contracts::billing::{
47 AdminSetPauseArgs, BlindQuery, PauseLevel, PauseState, PlatformBreach, PlatformGuard, PlatformMetric, PlatformPause,
48};
49use g1t_contracts::time::rfc3339;
50use g1t_contracts::{FailureCode, Outcome, new_id};
51use g1t_kit::now_ms;
52use serde::Deserialize;
53use serde_json::{Value, json};
54use worker::{Env, Result};
55
56use crate::Billing;
57use crate::keeper::Keeper;
58
59const HOUR_MS: u64 = 60 * 60 * 1000;
60/// How long a pause read is kept in the isolate.
61const KEEP_MS: u64 = 30_000;
62/// A metric's breach is emailed at most this often.
63const ALERT_EVERY_MS: u64 = 6 * HOUR_MS;
64/// The spike rule needs this many hours of history first.
65const SPIKE_MIN_HOURS: usize = 24;
66
67thread_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.
73pub(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
88use PauseLevel::{Compute, Indexing, Renders, Schedules};
89
90pub(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
107pub(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)]
113pub(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
128impl 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.
160pub(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.
171pub(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)]
176struct 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)]
186struct UsageRow {
187 metric: String,
188 value: f64,
189 top_name: Option<String>,
190 top_value: Option<f64>,
191}
192
193#[derive(Deserialize)]
194struct 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
209impl 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, &note, 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.
386pub(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.
401pub(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)]
415pub(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.
422pub(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.
428pub(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`.
434pub(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`.
439pub(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.
461pub(crate) type Counted = BTreeMap<&'static str, BTreeMap<String, f64>>;
462
463/// A KV operation's metric, by its `actionType`.
464fn 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.
475pub(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)]
525pub(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.
532pub(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).
555pub(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)]
566pub(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.
578pub(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.
596pub(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.
604pub(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.
620pub(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.
641pub(crate) fn hourly_due(now: u64) -> bool {
642 let minute = now / 60_000 % 60;
643 (15..30).contains(&minute)
644}
645
646impl 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.
845pub(crate) const BLIND_RUNS: usize = 3;
846/// A blind query is emailed at most this often.
847const BLIND_EVERY_MS: u64 = 24 * HOUR_MS;
848/// The key `platform_watch_alerts` uses for "every dataset answered empty".
849pub(crate) const ALL_EMPTY: &str = "all_empty";
850
851/// What one run could and could not see.
852#[derive(Clone, Debug, Default, PartialEq)]
853pub(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
861impl 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`.
875pub(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.
881pub(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.
894pub(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)]
899struct RunRow {
900 hour: String,
901 failed: Option<String>,
902 empty: i64,
903}
904
905impl 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)]
998mod 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}