Skip to content
965 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 logged and skipped; the others still count.
37
38use std::cell::RefCell;
39use std::collections::BTreeMap;
40
41use futures_util::future::join_all;
42use g1t_contracts::billing::{
43 AdminSetPauseArgs, PauseLevel, PauseState, PlatformBreach, PlatformGuard, PlatformMetric, PlatformPause,
44};
45use g1t_contracts::time::rfc3339;
46use g1t_contracts::{FailureCode, Outcome, new_id};
47use g1t_kit::now_ms;
48use serde::Deserialize;
49use serde_json::{Value, json};
50use worker::{Env, Result};
51
52use crate::Billing;
53use crate::keeper::Keeper;
54
55const HOUR_MS: u64 = 60 * 60 * 1000;
56/// How long a pause read is kept in the isolate.
57const KEEP_MS: u64 = 30_000;
58/// A metric's breach is emailed at most this often.
59const ALERT_EVERY_MS: u64 = 6 * HOUR_MS;
60/// The spike rule needs this many hours of history first.
61const SPIKE_MIN_HOURS: usize = 24;
62
63thread_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.
69pub(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
84use PauseLevel::{Compute, Indexing, Renders, Schedules};
85
86pub(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
103pub(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)]
109pub(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
124impl 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.
156pub(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.
167pub(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)]
172struct 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)]
182struct UsageRow {
183 metric: String,
184 value: f64,
185 top_name: Option<String>,
186 top_value: Option<f64>,
187}
188
189#[derive(Deserialize)]
190struct 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
205impl 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, &note, 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.
374pub(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.
389pub(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)]
403pub(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.
410pub(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.
416pub(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`.
422pub(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`.
427pub(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.
449pub(crate) type Counted = BTreeMap<&'static str, BTreeMap<String, f64>>;
450
451/// A KV operation's metric, by its `actionType`.
452fn 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.
463pub(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)]
513pub(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.
520pub(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).
543pub(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)]
554pub(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.
566pub(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.
584pub(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.
592pub(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.
608pub(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.
629pub(crate) fn hourly_due(now: u64) -> bool {
630 let minute = now / 60_000 % 60;
631 (15..30).contains(&minute)
632}
633
634impl 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)]
812mod 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}