Skip to content
1,186 linesCodeBlameRaw

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

Merge platform pause and the hourly usage watcher: staff can pause compute, schedules, indexing or renders for everyone, the watcher emails on a breach and is never blind quietly, and the models proxy holds each run to its cap (billing 0051, integrations 0006)1//! Platform spend guardrails: what Cloudflare counts for all of g1t, read
2//! every hour, and a g1t-wide pause to stop it. See
The docs folder is gone, and what it held lives where people read it: how a self-hosted g1t runs and how to deploy g1t to Cloudflare are pages on docs.g1t.sh under Run g1t yourself, and speed, rate limits and operating g1t.sh are sections of CONTRIBUTING.md; code that cited a file in docs/ now points to the page or section that covers it, or says what it means itself, and applied migrations and the runner images are left as they were.3//! docs.g1t.sh/guides/deploy-to-cloudflare/#spend-guardrails.
Merge platform pause and the hourly usage watcher: staff can pause compute, schedules, indexing or renders for everyone, the watcher emails on a breach and is never blind quietly, and the models proxy holds each run to its cap (billing 0051, integrations 0006)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(
The docs folder is gone, and what it held lives where people read it: how a self-hosted g1t runs and how to deploy g1t to Cloudflare are pages on docs.g1t.sh under Run g1t yourself, and speed, rate limits and operating g1t.sh are sections of CONTRIBUTING.md; code that cited a file in docs/ now points to the page or section that covers it, or says what it means itself, and applied migrations and the runner images are left as they were.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 https://docs.g1t.sh/guides/deploy-to-cloudflare/#spend-guardrails.".to_owned(),
Merge platform pause and the hourly usage watcher: staff can pause compute, schedules, indexing or renders for everyone, the watcher emails on a breach and is never blind quietly, and the models proxy holds each run to its cap (billing 0051, integrations 0006)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(
The docs folder is gone, and what it held lives where people read it: how a self-hosted g1t runs and how to deploy g1t to Cloudflare are pages on docs.g1t.sh under Run g1t yourself, and speed, rate limits and operating g1t.sh are sections of CONTRIBUTING.md; code that cited a file in docs/ now points to the page or section that covers it, or says what it means itself, and applied migrations and the runner images are left as they were.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 https://docs.g1t.sh/guides/deploy-to-cloudflare/#spend-guardrails.".to_owned(),
Merge platform pause and the hourly usage watcher: staff can pause compute, schedules, indexing or renders for everyone, the watcher emails on a breach and is never blind quietly, and the models proxy holds each run to its cap (billing 0051, integrations 0006)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 }
The docs folder is gone, and what it held lives where people read it: how a self-hosted g1t runs and how to deploy g1t to Cloudflare are pages on docs.g1t.sh under Run g1t yourself, and speed, rate limits and operating g1t.sh are sections of CONTRIBUTING.md; code that cited a file in docs/ now points to the page or section that covers it, or says what it means itself, and applied migrations and the runner images are left as they were.1046 // The defaults when no PLATFORM_HOURLY_* variable is set.
Merge platform pause and the hourly usage watcher: staff can pause compute, schedules, indexing or renders for everyone, the watcher emails on a breach and is never blind quietly, and the models proxy holds each run to its cap (billing 0051, integrations 0006)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}

This file's history is long; its oldest lines are credited to the oldest commit read.