Skip to content

Commit

Staff can pause compute, schedules, indexing or renders across g1t, and billing reads the platform's Cloudflare usage every hour: thresholds, spikes, staff emails and automatic pauses (billing 0050)

syntaqxcommitted Parent67b02efBrowse files
29 files+1939−180/29 viewed
+16−0
852852 workflow job is recorded as failed with "Not started:" and the reason.
853853 Runs already under way finish, so usage can go slightly past a limit.
854854
855+### When g1t pauses work for everyone
856+
857+If usage across all of g1t climbs far past normal, g1t can pause some
858+kinds of work for every workspace while it looks into it. Each pause is
859+separate, and runs already under way finish:
860+
861+| Paused | What you see |
862+| --- | --- |
863+| Compute | New agent runs, checks, workflow jobs and deploy builds are refused with `paused` and a message that g1t has paused them across the platform. Try again later. |
864+| Schedules | Workflows on `schedule:` skip the minutes while it lasts; they are not run late. Issues waiting for an agent stay in the queue. |
865+| Indexing | **Rebuild** in the [context hub](/guides/context-hub/) is refused, and new semantic search embeddings wait. Text search still answers, and search keeps up with new pushes. |
866+| Renders | A link to g1t shows g1t's logo instead of the page's own card. |
867+
868+Nothing in your workspace changes, and nothing is charged while work
869+waits.
870+
855871 ## Security on every plan
856872
857873 Every workspace, free or on the plan, has the [audit log](/guides/audit-log/),
+17−2
44 */
55 import { data, redirect } from "react-router";
66
7−import { parseCostSettings, parseMapping, parseRange } from "./costs";
7+import { parseCostSettings, parseMapping, parsePauseLevel, parseRange } from "./costs";
88 import { admin } from "./services.server";
99 import { settle } from "./settle";
1010 import { requireStaff } from "./staff";
1717 mapping: "Mapping saved. It applies from the next run; read the bill now to see it.",
1818 removed: "Mapping removed.",
1919 lifted: "Breaker lifted for the rest of today (UTC). Hosted-model runs start again; it is recorded in the audit log.",
20+ paused: "Paused across g1t. Every service sees it within 30 seconds; it is recorded in the audit log.",
21+ resumed: "Resumed across g1t. Every service sees it within 30 seconds; it is recorded in the audit log.",
2022 };
2123
2224 export type CostsActionResult = { error: string; section: string; values?: Record<string, string> };
2527 requireStaff(context as Parameters<typeof requireStaff>[0]);
2628 const url = new URL(request.url);
2729 const range = parseRange(url.searchParams.get("days"));
28− const report = await settle(admin.costs(range));
30+ const [report, guard] = await Promise.all([settle(admin.costs(range)), settle(admin.platformGuard())]);
2931 const done = url.searchParams.get("done");
3032 return {
3133 range,
3234 bucket: url.searchParams.get("product"),
3335 report: report.ok ? report.value : null,
3436 error: report.ok ? null : report.error,
37+ // The platform pause and usage watch (billing's platform.rs).
38+ guard: guard.ok ? guard.value : null,
39+ guardError: guard.ok ? null : guard.error,
3540 done: done ? (DONE[done] ?? null) : null,
3641 };
3742 }
6267 if (!result.value.ok) return fail("lift", result.value.error.message);
6368 throw back("lifted", "#spend");
6469 }
70+ if (intent === "pause" || intent === "resume") {
71+ const level = parsePauseLevel(form.get("level"));
72+ if (!level) return fail("platform", "Pick compute, schedules, indexing or renders.");
73+ const note = String(form.get("note") ?? "").trim().slice(0, 500);
74+ if (note.length < 5) return fail(`pause-${level}`, "Say why, for whoever looks next.");
75+ const result = await settle(admin.setPause(level, intent === "pause", note, staff.email));
76+ if (!result.ok) return fail(`pause-${level}`, `Billing did not answer: ${result.error}`);
77+ if (!result.value.ok) return fail(`pause-${level}`, result.value.error.message);
78+ throw back(intent === "pause" ? "paused" : "resumed", "#platform");
79+ }
6580 if (intent === "decide") {
6681 const id = String(form.get("id") ?? "");
6782 const decision = form.get("decision") === "approve" ? "approve" : "reject";
+38−1
11 import assert from "node:assert/strict";
22 import { test } from "node:test";
33
4−import type { CostDay, SpendCaps } from "@g1t/contracts";
4+import type { CostDay, PlatformGuard, SpendCaps } from "@g1t/contracts";
55
66 import {
7+ count,
78 daySeries,
89 daysBetween,
910 marginOnPrice,
1213 parseBucket,
1314 parseCostSettings,
1415 parseMapping,
16+ parsePauseLevel,
17+ pauseBanner,
1518 parseRange,
1619 percentLabel,
1720 proposalOutcome,
1821 spendBanner,
1922 spendRows,
2023 subscriptionsOver,
24+ thresholdShare,
2125 unitDollars,
2226 versionCells,
2327 whoPaid,
200204 assert.equal(proposalOutcome({ status: "superseded", decidedBy: null, effectiveAt: null }), "replaced by a later measurement");
201205 assert.equal(proposalOutcome({ status: "open", decidedBy: null, effectiveAt: null }), null);
202206 });
207+
208+test("the pause banner names every paused level, and who paused it when it was not a person", () => {
209+ const level = (name: "compute" | "schedules" | "indexing" | "renders", paused: boolean, auto = false) => ({
210+ level: name,
211+ paused,
212+ note: null,
213+ set_by: null,
214+ set_at: null,
215+ auto,
216+ });
217+ const guard = (levels: ReturnType<typeof level>[]): PlatformGuard => ({
218+ levels,
219+ hour: null,
220+ last_hour: [],
221+ month: "2026-10",
222+ month_to_date: [],
223+ breaches: [],
224+ can_read: true,
225+ auto_pause: ["schedules", "indexing"],
226+ });
227+ assert.equal(pauseBanner(guard([level("compute", false), level("renders", false)])), null);
228+ assert.equal(pauseBanner(guard([level("schedules", true, true)])), "Paused across g1t: schedules (by the usage watcher).");
229+ assert.equal(
230+ pauseBanner(guard([level("compute", true), level("schedules", true), level("indexing", true, true)])),
231+ "Paused across g1t: compute, schedules and indexing (by the usage watcher).",
232+ );
233+ assert.equal(parsePauseLevel("renders"), "renders");
234+ assert.equal(parsePauseLevel("everything"), null);
235+ assert.equal(count(240_000), "240k");
236+ assert.equal(count(2_000_000_000), "2.0B");
237+ assert.equal(thresholdShare(240_000, 200_000), 120);
238+ assert.equal(thresholdShare(5, 0), null);
239+});
+40−1
22 * Costs & margin: the arithmetic behind the page, apart from the SVG and
33 * the Workers runtime so it can be tested under Node. Money is in micros.
44 */
5−import type { CostDay, CostMappingInput, CostSettings, SpendCaps } from "@g1t/contracts";
5+import type { CostDay, CostMappingInput, CostSettings, PauseLevel, PlatformGuard, SpendCaps } from "@g1t/contracts";
66
77 import { parseDollars, usd } from "./money.ts";
88
266266 if (capMicros <= 0) return 0;
267267 return Math.max(0, Math.min(100, (usedMicros / capMicros) * 100));
268268 }
269+
270+// --- Platform pause and usage watch (billing's platform.rs) -----------------
271+
272+/** What each level of the platform pause stops, for the page and the banner. */
273+export const PAUSE_LEVELS: { level: PauseLevel; title: string; stops: string }[] = [
274+ { level: "compute", title: "Compute", stops: "New agent runs, checks, workflow jobs and deploy builds, for every workspace. Runs already going finish." },
275+ { level: "schedules", title: "Schedules", stops: "Actions' cron-triggered runs, and the runner's sweep that starts queued agents." },
276+ { level: "indexing", title: "Indexing", stops: "Context embeddings and backfills, and search's backfills (they go on from where they were when resumed)." },
277+ { level: "renders", title: "Renders", stops: "Social card images: a cache miss gets the brand card or the static logo." },
278+];
279+
280+/** Whether a form's level is one of the four. */
281+export function parsePauseLevel(value: unknown): PauseLevel | null {
282+ const level = String(value ?? "");
283+ return PAUSE_LEVELS.some((l) => l.level === level) ? (level as PauseLevel) : null;
284+}
285+
286+/** The red bar on every sudo page while any level is paused. Null when none is. */
287+export function pauseBanner(guard: PlatformGuard): string | null {
288+ const paused = guard.levels.filter((l) => l.paused);
289+ if (paused.length === 0) return null;
290+ const names = paused.map((l) => (l.auto ? `${l.level} (by the usage watcher)` : l.level));
291+ const list = names.length === 1 ? names[0] : `${names.slice(0, -1).join(", ")} and ${names[names.length - 1]}`;
292+ return `Paused across g1t: ${list}.`;
293+}
294+
295+/** `1.2M`, `240k`, `2.0B`: a count in a few characters. */
296+export function count(value: number): string {
297+ const n = Math.abs(value);
298+ if (n >= 1e9) return `${(value / 1e9).toFixed(1)}B`;
299+ if (n >= 1e6) return `${(value / 1e6).toFixed(1)}M`;
300+ if (n >= 1e3) return `${Math.round(value / 1e3)}k`;
301+ return `${Math.round(value)}`;
302+}
303+
304+/** An hour's value as a share of its threshold, rounded; null with no threshold. */
305+export function thresholdShare(value: number, threshold: number): number | null {
306+ return threshold > 0 ? Math.round((value / threshold) * 100) : null;
307+}
+15−3
77 import { MobileBar, Sidebar } from "./components/shell";
88 import { ButtonLink } from "./components/ui";
99 import type { NavCounts } from "./lib/nav";
10−import { spendBanner } from "./lib/costs";
10+import { pauseBanner, spendBanner } from "./lib/costs";
1111 import { admin, identity, statusAdmin } from "./lib/services.server";
1212 import { settle } from "./lib/settle";
1313 import { requireStaff, zoneContext } from "./lib/staff";
2727 export async function loader({ context }: Route.LoaderArgs) {
2828 const { email } = requireStaff(context);
2929 // The sidebar's counts: a service that does not answer shows none.
30− const [waitlist, incidents, alerts, caps] = await Promise.all([
30+ const [waitlist, incidents, alerts, caps, guard] = await Promise.all([
3131 settle(identity.waitlistPending()),
3232 settle(statusAdmin.openCount()),
3333 settle(admin.costAlerts()),
3434 settle(admin.spendCaps()),
35+ settle(admin.platformGuard()),
3536 ]);
3637 const counts: NavCounts = { waitlist: waitlist.ok ? waitlist.value : 0, incidents: incidents.ok ? incidents.value : 0 };
3738 // Every page says times in this zone (components/ui.tsx `When`).
4445 // g1t's own spend (billing's budget): the daily breaker open, or a comped
4546 // account's monthly budget used up. Red until it clears or staff act.
4647 const spend = caps.ok ? spendBanner(caps.value) : null;
47− return { email, counts, zone, zoneChosen: chosen, margin, spend };
48+ // A platform pause (billing's platform.rs): red on every page while any
49+ // level is paused, by staff or by the usage watcher.
50+ const paused = guard.ok ? pauseBanner(guard.value) : null;
51+ return { email, counts, zone, zoneChosen: chosen, margin, spend, paused };
4852 }
4953
5054 export function Layout({ children }: { children: React.ReactNode }) {
7983 </a>
8084 </div>
8185 )}
86+ {root?.paused && (
87+ <div role="alert" className="border-b border-danger/40 bg-danger/12 px-4 py-2 text-sm text-danger">
88+ <span className="font-medium">Platform pause:</span> {root.paused}{" "}
89+ <a href="/costs#platform" className="underline underline-offset-2">
90+ Platform pause
91+ </a>
92+ </div>
93+ )}
8294 {children}
8395 </div>
8496 {/* No <Scripts />: sudo ships no JavaScript, and its policy allows none. */}
+157−2
11 import type { ReactNode } from "react";
22 import { Link } from "react-router";
33
4−import type { CostsReport } from "@g1t/contracts";
4+import type { CostsReport, PlatformGuard, PlatformMetric } from "@g1t/contracts";
55
66 import type { Route } from "./+types/costs";
77 import { DaysChart } from "~/components/costs";
88 import { CostsHeader, chip, costsHref } from "~/components/costs-header";
99 import { Badge, Button, Field, Input, Notice, Section, Stat, When } from "~/components/ui";
10−import { capPercent, daySeries, marginOnPrice, marginTone, parseBucket, percentLabel, spendRows, subscriptionsOver, whoPaid } from "~/lib/costs";
10+import { PAUSE_LEVELS, capPercent, count, daySeries, marginOnPrice, marginTone, parseBucket, percentLabel, spendRows, subscriptionsOver, thresholdShare, whoPaid } from "~/lib/costs";
1111 import { type CostsActionResult, costsAction, costsLoader } from "~/lib/costs-route.server";
1212 import { usd } from "~/lib/money";
1313
3434 <div className="mt-5">
3535 <Notice tone="warn">Billing did not answer for costs: {error}</Notice>
3636 </div>
37+ <PlatformSection guard={loaderData.guard} unavailable={loaderData.guardError} failed={failed} />
3738 </main>
3839 );
3940 }
5051
5152 <SpendSection caps={report.caps} range={range} error={failed?.section === "lift" ? failed.error : null} />
5253
54+ <PlatformSection guard={loaderData.guard} unavailable={loaderData.guardError} failed={failed} />
55+
5356 <Section
5457 className="mt-6"
5558 title={product ? `${product.title}, by day` : "By day"}
504507 );
505508 }
506509
510+
511+/**
512+ * The platform pause and the hourly usage watch (billing's platform.rs,
513+ * docs/SPEND-GUARDRAILS.md): four levels staff can pause across g1t, what
514+ * Cloudflare counted in the last hour against each threshold, and the
515+ * last day's breaches.
516+ */
517+function PlatformSection({
518+ guard,
519+ unavailable,
520+ failed,
521+}: {
522+ guard: PlatformGuard | null;
523+ unavailable: string | null;
524+ failed: CostsActionResult | null;
525+}) {
526+ return (
527+ <Section
528+ className="mt-6"
529+ id="platform"
530+ title="Platform pause"
531+ description="What Cloudflare counted for all of g1t, read at a quarter past each hour: Workers, D1, Queues, Durable Objects, KV and Artifacts, each against an hourly threshold (PLATFORM_HOURLY_* in billing's wrangler.jsonc). A breach emails staff once per metric every 6 hours. Five times a threshold pauses what that metric feeds, of the levels AUTO_PAUSE names. Any level can be paused here by hand; every service sees a change within 30 seconds."
532+ >
533+ {!guard ? (
534+ <Notice tone="warn">Billing did not answer for the platform pause: {unavailable}</Notice>
535+ ) : (
536+ <>
537+ {!guard.can_read && (
538+ <div className="mb-4">
539+ <Notice tone="warn">Billing has no Cloudflare token with Account Analytics Read, so the hourly watch reads nothing. The pause still works.</Notice>
540+ </div>
541+ )}
542+ {failed?.section === "platform" && (
543+ <div className="mb-4">
544+ <Notice tone="error">{failed.error}</Notice>
545+ </div>
546+ )}
547+ <div className="grid gap-3 lg:grid-cols-2">
548+ {PAUSE_LEVELS.map(({ level, title, stops }) => {
549+ const state = guard.levels.find((l) => l.level === level);
550+ const paused = state?.paused ?? false;
551+ const error = failed?.section === `pause-${level}` ? failed.error : null;
552+ return (
553+ <div key={level} className={`rounded-lg border p-4 ${paused ? "border-danger/50 bg-danger/8" : "border-line"}`}>
554+ <div className="flex items-center justify-between gap-2">
555+ <span className="font-medium">{title}</span>
556+ {paused ? <Badge tone="danger">Paused</Badge> : <Badge tone="mint">Running</Badge>}
557+ </div>
558+ <p className="mt-1 text-xs text-muted">{stops}</p>
559+ {state?.set_at && (
560+ <p className="mt-2 text-xs text-faint">
561+ {paused ? "Paused" : "Resumed"} <When at={state.set_at} time /> by {state.auto ? "the usage watcher" : state.set_by}
562+ {state.note ? `: “${state.note}”` : ""}
563+ </p>
564+ )}
565+ <form method="post" action="#platform" className="mt-3 flex flex-col gap-2 sm:flex-row sm:items-end">
566+ <input type="hidden" name="intent" value={paused ? "resume" : "pause"} />
567+ <input type="hidden" name="level" value={level} />
568+ <Field label={paused ? "Why resume it" : "Why pause it"} hint="Recorded in the audit log.">
569+ <Input name="note" required minLength={5} maxLength={500} placeholder={paused ? "e.g. Fixed the loop in search" : "e.g. KV lists runaway"} />
570+ </Field>
571+ <Button type="submit" variant={paused ? "primary" : "danger"}>
572+ {paused ? "Resume" : "Pause"}
573+ </Button>
574+ </form>
575+ {error && (
576+ <div className="mt-2">
577+ <Notice tone="error">{error}</Notice>
578+ </div>
579+ )}
580+ </div>
581+ );
582+ })}
583+ </div>
584+ <p className="mt-3 text-xs text-muted">
585+ A severe breach may pause: {guard.auto_pause.length ? guard.auto_pause.join(", ") : "nothing (AUTO_PAUSE is empty)"}.
586+ </p>
587+
588+ <h3 className="mt-6 text-sm font-medium">Breaches in the last 24 hours</h3>
589+ {guard.breaches.length === 0 ? (
590+ <p className="mt-2 text-sm text-muted">None.</p>
591+ ) : (
592+ <ul className="mt-2 space-y-2 text-sm">
593+ {guard.breaches.map((b) => (
594+ <li key={b.id} className="rounded-md border border-line px-3 py-2">
595+ <div className="flex flex-wrap items-center gap-2">
596+ <Badge tone={b.severe ? "danger" : "warn"}>{b.severe ? "Severe" : b.rule === "spike" ? "Spike" : "Over threshold"}</Badge>
597+ <span className="text-xs text-faint">
598+ <When at={b.opened_at} time />
599+ {b.emailed_at ? ", emailed" : ""}
600+ </span>
601+ </div>
602+ <p className="mt-1 text-fg-soft">{b.detail}</p>
603+ </li>
604+ ))}
605+ </ul>
606+ )}
607+
608+ <UsageTable title={guard.hour ? `The hour from ${guard.hour}` : "The last hour"} metrics={guard.last_hour} hourly />
609+ <UsageTable title={`${guard.month}, so far`} metrics={guard.month_to_date} hourly={false} />
610+ <p className="mt-3 text-xs text-muted">
611+ Ids are Cloudflare&apos;s. <code>node scripts/ops/platform-usage.mjs</code> names them and shows the last 24 hours per script, queue, database and
612+ namespace.
613+ </p>
614+ </>
615+ )}
616+ </Section>
617+ );
618+}
619+
620+function UsageTable({ title, metrics, hourly }: { title: string; metrics: PlatformMetric[]; hourly: boolean }) {
621+ return (
622+ <>
623+ <h3 className="mt-6 text-sm font-medium">{title}</h3>
624+ {metrics.length === 0 ? (
625+ <p className="mt-2 text-sm text-muted">Nothing read yet.</p>
626+ ) : (
627+ <div className="-mx-4 mt-2 overflow-x-auto sm:-mx-5">
628+ <table className="w-full min-w-[36rem] text-sm">
629+ <thead>
630+ <tr className="border-b border-line text-left text-xs text-muted">
631+ <th className="px-4 py-2 font-medium sm:px-5">Metric</th>
632+ <th className="px-4 py-2 text-right font-medium">Count</th>
633+ {hourly && <th className="px-4 py-2 text-right font-medium">Of threshold</th>}
634+ <th className="px-4 py-2 font-medium sm:pr-5">Most from</th>
635+ </tr>
636+ </thead>
637+ <tbody>
638+ {metrics.map((m) => {
639+ const share = thresholdShare(m.value, m.threshold);
640+ return (
641+ <tr key={m.metric} className="border-b border-line last:border-0">
642+ <td className="px-4 py-2.5 sm:px-5">{m.title}</td>
643+ <td className="tabular px-4 py-2.5 text-right">{count(m.value)}</td>
644+ {hourly && (
645+ <td className={`tabular px-4 py-2.5 text-right ${share != null && share > 100 ? "text-danger" : "text-muted"}`}>
646+ {share == null ? "—" : `${share}% of ${count(m.threshold)}`}
647+ </td>
648+ )}
649+ <td className="max-w-[16rem] truncate px-4 py-2.5 text-xs text-faint sm:pr-5">
650+ {m.top_name ? `${m.top_name} (${count(m.top_value ?? 0)})` : "—"}
651+ </td>
652+ </tr>
653+ );
654+ })}
655+ </tbody>
656+ </table>
657+ </div>
658+ )}
659+ </>
660+ );
661+}
+151−0
36443644 pub by: String,
36453645 }
36463646
3647+// ---- Platform pauses and the usage watcher (docs/SPEND-GUARDRAILS.md) ----
3648+
3649+/// A g1t-wide pause, set by staff in sudo or by billing's hourly usage
3650+/// watcher on a severe breach. Each level is independent.
3651+#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
3652+#[serde(rename_all = "snake_case")]
3653+pub enum PauseLevel {
3654+ /// Agents, sandboxes, Actions hosted jobs, deploy builds: every
3655+ /// reservation through billing's `reserve` but embeddings.
3656+ Compute,
3657+ /// Actions' cron-triggered runs, and the runner's sweep that starts
3658+ /// queued agents.
3659+ Schedules,
3660+ /// Context embeddings and backfills, and search's backfills.
3661+ Indexing,
3662+ /// Social card rendering, which falls back to a static image.
3663+ Renders,
3664+}
3665+
3666+impl PauseLevel {
3667+ pub const ALL: [PauseLevel; 4] = [PauseLevel::Compute, PauseLevel::Schedules, PauseLevel::Indexing, PauseLevel::Renders];
3668+
3669+ pub fn as_str(self) -> &'static str {
3670+ match self {
3671+ PauseLevel::Compute => "compute",
3672+ PauseLevel::Schedules => "schedules",
3673+ PauseLevel::Indexing => "indexing",
3674+ PauseLevel::Renders => "renders",
3675+ }
3676+ }
3677+
3678+ pub fn parse(text: &str) -> Option<PauseLevel> {
3679+ PauseLevel::ALL.into_iter().find(|level| level.as_str() == text.trim())
3680+ }
3681+}
3682+
3683+/// `platform_pause`: which levels are paused now. Takes nothing. Callers
3684+/// keep the answer about 30 seconds; one that cannot read it runs (fails
3685+/// open).
3686+#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize)]
3687+pub struct PlatformPause {
3688+ #[serde(default)]
3689+ pub compute: bool,
3690+ #[serde(default)]
3691+ pub schedules: bool,
3692+ #[serde(default)]
3693+ pub indexing: bool,
3694+ #[serde(default)]
3695+ pub renders: bool,
3696+}
3697+
3698+impl PlatformPause {
3699+ pub fn is(&self, level: PauseLevel) -> bool {
3700+ match level {
3701+ PauseLevel::Compute => self.compute,
3702+ PauseLevel::Schedules => self.schedules,
3703+ PauseLevel::Indexing => self.indexing,
3704+ PauseLevel::Renders => self.renders,
3705+ }
3706+ }
3707+
3708+ pub fn set(&mut self, level: PauseLevel, paused: bool) {
3709+ match level {
3710+ PauseLevel::Compute => self.compute = paused,
3711+ PauseLevel::Schedules => self.schedules = paused,
3712+ PauseLevel::Indexing => self.indexing = paused,
3713+ PauseLevel::Renders => self.renders = paused,
3714+ }
3715+ }
3716+
3717+ pub fn any(&self) -> bool {
3718+ PauseLevel::ALL.into_iter().any(|level| self.is(level))
3719+ }
3720+}
3721+
3722+/// One level as sudo shows it: who paused it, when and why.
3723+#[derive(Clone, Debug, Default, Serialize, Deserialize)]
3724+pub struct PauseState {
3725+ pub level: String,
3726+ pub paused: bool,
3727+ pub note: Option<String>,
3728+ pub set_by: Option<String>,
3729+ pub set_at: Option<String>,
3730+ /// Set by the usage watcher, not a person.
3731+ pub auto: bool,
3732+}
3733+
3734+/// One metric's usage over an hour or the month so far.
3735+#[derive(Clone, Debug, Default, Serialize, Deserialize)]
3736+pub struct PlatformMetric {
3737+ pub metric: String,
3738+ pub title: String,
3739+ pub value: f64,
3740+ /// Its hourly threshold (`PLATFORM_HOURLY_*`); zero: none.
3741+ pub threshold: f64,
3742+ /// The script, queue, database or namespace that counted most.
3743+ pub top_name: Option<String>,
3744+ pub top_value: Option<f64>,
3745+}
3746+
3747+/// A breach the watcher found.
3748+#[derive(Clone, Debug, Default, Serialize, Deserialize)]
3749+pub struct PlatformBreach {
3750+ pub id: String,
3751+ pub metric: String,
3752+ pub hour: String,
3753+ /// `threshold` or `spike`.
3754+ pub rule: String,
3755+ pub value: f64,
3756+ pub threshold: f64,
3757+ pub severe: bool,
3758+ pub top_name: Option<String>,
3759+ pub detail: String,
3760+ /// Levels it paused.
3761+ pub paused: Vec<String>,
3762+ pub opened_at: String,
3763+ pub emailed_at: Option<String>,
3764+}
3765+
3766+/// `admin_platform_guard`: the pauses, the last hour read, the month so
3767+/// far and the last day's breaches, for sudo's Costs page and its banner.
3768+#[derive(Clone, Debug, Default, Serialize, Deserialize)]
3769+pub struct PlatformGuard {
3770+ pub levels: Vec<PauseState>,
3771+ /// The last hour the watcher read (`YYYY-MM-DDTHH:00:00Z`), if any.
3772+ pub hour: Option<String>,
3773+ pub last_hour: Vec<PlatformMetric>,
3774+ pub month: String,
3775+ pub month_to_date: Vec<PlatformMetric>,
3776+ /// The last 24 hours' breaches, newest first.
3777+ pub breaches: Vec<PlatformBreach>,
3778+ /// Whether the watcher can read Cloudflare's analytics at all.
3779+ pub can_read: bool,
3780+ /// `AUTO_PAUSE`: the levels a severe breach may pause.
3781+ pub auto_pause: Vec<String>,
3782+}
3783+
3784+/// `admin_platform_guard`. Returns `PlatformGuard`.
3785+#[derive(Debug, Default, Serialize, Deserialize)]
3786+pub struct AdminPlatformGuardArgs {}
3787+
3788+/// `admin_set_pause`: pauses or resumes one level, with why. Recorded in
3789+/// the audit log. Returns `Outcome<PlatformGuard>`.
3790+#[derive(Debug, Serialize, Deserialize)]
3791+pub struct AdminSetPauseArgs {
3792+ pub level: String,
3793+ pub paused: bool,
3794+ pub note: String,
3795+ pub by: String,
3796+}
3797+
36473798 /// `admin_cost_alerts`: the open margin alerts, for sudo's banner.
36483799 /// Returns `Vec<MarginAlert>`.
36493800 #[derive(Debug, Default, Serialize, Deserialize)]
+1−0
6161 }
6262
6363 pub mod d1;
64+pub mod pause;
6465 pub mod wire;
6566
6667 /// Helpers for bindings that workers-rs has no typed wrapper for, such as
+67−0
1+//! g1t-wide pauses (billing's `platform_pause`; docs/SPEND-GUARDRAILS.md),
2+//! read cheaply: one call to billing per isolate every 30 seconds at most,
3+//! never a database read per request.
4+//!
5+//! When billing cannot say, nothing is paused: a pause is a brake staff
6+//! (or billing's usage watcher) pull on purpose, and a billing outage must
7+//! not stop schedules and indexing everywhere. The failure is kept for the
8+//! same 30 seconds, so an outage is not asked about on every request.
9+
10+use std::cell::RefCell;
11+
12+use g1t_contracts::billing::{PauseLevel, PlatformPause};
13+use worker::Fetcher;
14+
15+/// How long an answer is kept in the isolate.
16+pub const KEEP_MS: u64 = 30_000;
17+
18+thread_local! {
19+ static KEPT: RefCell<Option<(PlatformPause, u64)>> = const { RefCell::new(None) };
20+}
21+
22+/// A kept answer, while it is fresh.
23+fn fresh(kept: Option<(PlatformPause, u64)>, now: u64) -> Option<PlatformPause> {
24+ kept.filter(|(_, until)| *until > now).map(|(pause, _)| pause)
25+}
26+
27+/// Whether `level` is paused, through the billing service's binding.
28+pub async fn paused(billing: &Fetcher, level: PauseLevel) -> bool {
29+ current(billing).await.is(level)
30+}
31+
32+/// Every level, kept for `KEEP_MS`. Nothing paused when billing cannot say.
33+pub async fn current(billing: &Fetcher) -> PlatformPause {
34+ let now = crate::now_ms();
35+ if let Some(pause) = KEPT.with(|kept| fresh(*kept.borrow(), now)) {
36+ return pause;
37+ }
38+ let pause = match crate::call::<_, PlatformPause>(billing, "platform_pause", &serde_json::json!({})).await {
39+ Ok(pause) => pause,
40+ Err(error) => {
41+ worker::console_error!("platform pause unreadable, so nothing is paused: {error}");
42+ PlatformPause::default()
43+ }
44+ };
45+ KEPT.with(|kept| *kept.borrow_mut() = Some((pause, now + KEEP_MS)));
46+ pause
47+}
48+
49+#[cfg(test)]
50+mod tests {
51+ use super::*;
52+
53+ #[test]
54+ fn an_answer_is_kept_thirty_seconds() {
55+ let paused = PlatformPause { schedules: true, ..PlatformPause::default() };
56+ assert_eq!(fresh(Some((paused, 1_000 + KEEP_MS)), 1_000), Some(paused));
57+ assert_eq!(fresh(Some((paused, 1_000 + KEEP_MS)), 1_000 + KEEP_MS), None);
58+ assert_eq!(fresh(None, 0), None);
59+ }
60+
61+ #[test]
62+ fn an_unread_pause_pauses_nothing() {
63+ let none = PlatformPause::default();
64+ assert!(PauseLevel::ALL.into_iter().all(|level| !none.is(level)));
65+ assert!(!none.any());
66+ }
67+}
+60−0
11 import type { ComputeKind, Reservation } from "./compute";
2+import type { PauseLevel } from "./platform";
23 import type { User, Viewer } from "./identity";
34 import type { RepoPath } from "./repos";
45 import type { Result } from "./result";
675676 spendCaps(): Promise<SpendCaps>;
676677 /** Lets hosted-model runs start again for the rest of today (UTC); needs a note. */
677678 liftBreaker(note: string, by: string): Promise<Result<SpendCaps>>;
679+ /** Platform pauses, the last hour of platform usage, the month so far and the last day's breaches (billing's platform.rs). */
680+ platformGuard(): Promise<PlatformGuard>;
681+ /** Pauses or resumes one level across g1t; needs a note, recorded in the audit log. */
682+ setPause(level: PauseLevel, paused: boolean, note: string, by: string): Promise<Result<PlatformGuard>>;
678683 /** Approve or reject a price proposal; a rejection needs a note. An approved rise waits out the notice period. */
679684 decideProposal(id: string, decision: "approve" | "reject", note: string, by: string): Promise<Result<PriceProposal>>;
680685 setCostSettings(settings: CostSettings, by: string): Promise<Result<CostSettings>>;
16771682 caps: SpendCaps;
16781683 };
16791684
1685+/** One level of the platform pause, as sudo shows it. Snake case, as billing sends it. */
1686+export type PauseState = {
1687+ level: PauseLevel;
1688+ paused: boolean;
1689+ note: string | null;
1690+ set_by: string | null;
1691+ set_at: string | null;
1692+ /** Set by billing's usage watcher, not a person. */
1693+ auto: boolean;
1694+};
1695+
1696+/** One platform metric over an hour or the month so far. */
1697+export type PlatformMetric = {
1698+ metric: string;
1699+ title: string;
1700+ value: number;
1701+ /** Its hourly threshold (`PLATFORM_HOURLY_*`); 0: none. */
1702+ threshold: number;
1703+ /** The script, queue, database or namespace that counted most. */
1704+ top_name: string | null;
1705+ top_value: number | null;
1706+};
1707+
1708+/** A breach billing's usage watcher found. */
1709+export type PlatformBreach = {
1710+ id: string;
1711+ metric: string;
1712+ hour: string;
1713+ rule: "threshold" | "spike";
1714+ value: number;
1715+ threshold: number;
1716+ severe: boolean;
1717+ top_name: string | null;
1718+ detail: string;
1719+ /** Levels it paused. */
1720+ paused: PauseLevel[];
1721+ opened_at: string;
1722+ emailed_at: string | null;
1723+};
1724+
1725+/** Billing's `admin_platform_guard`: the platform pause and usage watcher (docs/SPEND-GUARDRAILS.md). */
1726+export type PlatformGuard = {
1727+ levels: PauseState[];
1728+ /** The last hour read, `YYYY-MM-DDTHH:00:00Z`. */
1729+ hour: string | null;
1730+ last_hour: PlatformMetric[];
1731+ month: string;
1732+ month_to_date: PlatformMetric[];
1733+ breaches: PlatformBreach[];
1734+ /** Whether billing can read Cloudflare's analytics. */
1735+ can_read: boolean;
1736+ /** `AUTO_PAUSE`: the levels a severe breach may pause. */
1737+ auto_pause: PauseLevel[];
1738+};
1739+
16801740 /** What g1t pays for itself, at cost, against its caps (billing's `budget`). */
16811741 export type SpendCaps = {
16821742 /** Today (UTC), YYYY-MM-DD, and this month, YYYY-MM. */
+3−1
673673 costAlerts: () => call("admin_cost_alerts", {}),
674674 spendCaps: () => call("admin_spend_caps", {}),
675675 liftBreaker: (note, by) => call("admin_lift_breaker", { note, by }),
676+ platformGuard: () => call("admin_platform_guard", {}),
677+ setPause: (level, paused, note, by) => call("admin_set_pause", { level, paused, note, by }),
676678 decideProposal: (id, decision, note, by) => call("admin_decide_proposal", { id, decision, note, by }),
677679 setCostSettings: (settings, by) => call("admin_set_cost_settings", { settings, by }),
678680 setCostMapping: (mapping, by) =>
753755 gatewayUpstream: (workspace) => call("gateway_upstream", { workspace }),
754756 gatewayProviders: (workspace) => call("gateway_providers", { workspace }),
755757 closeModelSessions: (tokenHashes) => call("close_model_sessions", { token_hashes: tokenHashes }),
758+ capModelSessions: (tokenHashes, capMicros) => call("cap_model_sessions", { token_hashes: tokenHashes, cap_micros: capMicros }),
756759 routes: (workspace, viewer) => call("routes", { workspace, viewer }),
757760 setRoutes: (actor, workspace, routes) => call("set_routes", { actor, workspace, routes }),
758− capModelSessions: (tokenHashes, capMicros) => call("cap_model_sessions", { token_hashes: tokenHashes, cap_micros: capMicros }),
759761 };
760762 }
761763
+1−0
2929 export * from "./oauth";
3030 export * from "./og";
3131 export * from "./packages";
32+export * from "./platform";
3233 export * from "./projects";
3334 export * from "./repos";
3435 export * from "./result";
+78−0
1+/**
2+ * g1t-wide pauses: staff (or billing's hourly usage watcher) can stop
3+ * whole kinds of work across the platform while unusual usage is looked
4+ * into. See docs/SPEND-GUARDRAILS.md and the billing service's
5+ * `platform.rs`.
6+ *
7+ * - `compute`: agents, sandboxes, Actions hosted jobs and builds. Billing's
8+ * `reserve` refuses them itself, so every `ComputeGate.admit` caller gets
9+ * it without asking here.
10+ * - `schedules`: Actions' cron-triggered runs, and the runner's sweep that
11+ * starts queued agents.
12+ * - `indexing`: context embeddings and backfills, search backfills.
13+ * - `renders`: social card rendering, which falls back to a static image.
14+ *
15+ * Read through billing's `platform_pause`, kept for 30 seconds in the
16+ * isolate: never a database read per request. When billing cannot say,
17+ * nothing is paused (fails open): a pause is pulled on purpose, and a
18+ * billing outage must not stop the platform with it. The failure is kept
19+ * for the same 30 seconds.
20+ *
21+ * Only type imports, so services' unit tests can load it on its own.
22+ */
23+import type { ServiceBinding } from "./clients";
24+
25+export type PauseLevel = "compute" | "schedules" | "indexing" | "renders";
26+
27+export const PAUSE_LEVELS: readonly PauseLevel[] = ["compute", "schedules", "indexing", "renders"];
28+
29+/** Billing's `platform_pause`: which levels are paused now. */
30+export type PlatformPause = Record<PauseLevel, boolean>;
31+
32+/** How long an answer is kept in the isolate. */
33+export const PAUSE_KEPT_MS = 30_000;
34+
35+const NOTHING_PAUSED: PlatformPause = { compute: false, schedules: false, indexing: false, renders: false };
36+
37+let kept: { value: PlatformPause; until: number } | null = null;
38+
39+/** Billing's answer as a pause, anything it does not say as not paused. */
40+export function readPause(answer: unknown): PlatformPause {
41+ const raw = answer && typeof answer === "object" ? (answer as Record<string, unknown>) : {};
42+ return {
43+ compute: raw.compute === true,
44+ schedules: raw.schedules === true,
45+ indexing: raw.indexing === true,
46+ renders: raw.renders === true,
47+ };
48+}
49+
50+/** Every level, kept for 30 seconds. Nothing paused when billing cannot say. Never throws. */
51+export async function platformPause(billing: ServiceBinding | undefined, now = Date.now()): Promise<PlatformPause> {
52+ if (!billing) return NOTHING_PAUSED;
53+ if (kept && kept.until > now) return kept.value;
54+ let value = NOTHING_PAUSED;
55+ try {
56+ const response = await billing.fetch("https://service/rpc/platform_pause", {
57+ method: "POST",
58+ headers: { "content-type": "application/json" },
59+ body: "{}",
60+ });
61+ if (!response.ok) throw new Error(`platform_pause failed with status ${response.status}`);
62+ value = readPause(await response.json());
63+ } catch (error) {
64+ console.error("platform pause unreadable, so nothing is paused", String(error));
65+ }
66+ kept = { value, until: now + PAUSE_KEPT_MS };
67+ return value;
68+}
69+
70+/** Whether `level` is paused across g1t. */
71+export async function platformPaused(billing: ServiceBinding | undefined, level: PauseLevel): Promise<boolean> {
72+ return (await platformPause(billing))[level];
73+}
74+
75+/** For tests: forget the kept answer. */
76+export function forgetPlatformPause(): void {
77+ kept = null;
78+}
+6−1
22802280
22812281 pub async fn on_minute(&self, now_ms: u64) -> Result<()> {
22822282 let minute = now_ms / 60_000 * 60_000;
2283− if let Err(error) = self.run_schedules(minute).await {
2283+ // Staff (or billing's usage watcher) paused scheduled runs across
2284+ // g1t: this minute's schedules are skipped, not queued for later.
2285+ // Kept 30 seconds in the isolate (g1t_kit::pause).
2286+ if g1t_kit::pause::paused(&self.billing, g1t_contracts::billing::PauseLevel::Schedules).await {
2287+ worker::console_log!("actions: schedules are paused across g1t; skipped this minute's");
2288+ } else if let Err(error) = self.run_schedules(minute).await {
22842289 worker::console_error!("actions: schedules failed: {error}");
22852290 }
22862291 // Jobs whose sandbox went quiet or ran past their time.
+61−0
1+-- Platform spend guardrails (src/platform.rs, docs/SPEND-GUARDRAILS.md).
2+--
3+-- platform_pause: g1t-wide pauses, one row per level (compute, schedules,
4+-- indexing, renders), set by staff in sudo or by the hourly usage watcher
5+-- on a severe breach. No row, or paused = 0: running.
6+CREATE TABLE IF NOT EXISTS platform_pause (
7+ level TEXT PRIMARY KEY,
8+ paused INTEGER NOT NULL DEFAULT 0,
9+ note TEXT,
10+ set_by TEXT,
11+ set_at TEXT,
12+ -- 1 when the usage watcher set it, not a person.
13+ auto INTEGER NOT NULL DEFAULT 0
14+);
15+
16+-- platform_usage: what Cloudflare counted each hour (UTC) for each metric
17+-- (workers_requests, d1_rows_read, kv_lists, …), with the script, queue,
18+-- database or namespace that counted most. The spike rule reads the last
19+-- week of it.
20+CREATE TABLE IF NOT EXISTS platform_usage (
21+ hour TEXT NOT NULL,
22+ metric TEXT NOT NULL,
23+ value REAL NOT NULL DEFAULT 0,
24+ top_name TEXT,
25+ top_value REAL,
26+ read_at TEXT NOT NULL,
27+ PRIMARY KEY (hour, metric)
28+);
29+CREATE INDEX IF NOT EXISTS platform_usage_by_metric ON platform_usage (metric, hour);
30+
31+-- platform_usage_month: the month so far for each metric, as last read.
32+CREATE TABLE IF NOT EXISTS platform_usage_month (
33+ month TEXT NOT NULL,
34+ metric TEXT NOT NULL,
35+ value REAL NOT NULL DEFAULT 0,
36+ top_name TEXT,
37+ top_value REAL,
38+ read_at TEXT NOT NULL,
39+ PRIMARY KEY (month, metric)
40+);
41+
42+-- platform_alerts: each breach the watcher found: over its hourly
43+-- threshold, or a spike over the week's usual hour. Emailed at most once
44+-- per metric every 6 hours.
45+CREATE TABLE IF NOT EXISTS platform_alerts (
46+ id TEXT PRIMARY KEY,
47+ metric TEXT NOT NULL,
48+ hour TEXT NOT NULL,
49+ rule TEXT NOT NULL,
50+ value REAL NOT NULL,
51+ threshold REAL NOT NULL,
52+ severe INTEGER NOT NULL DEFAULT 0,
53+ top_name TEXT,
54+ top_value REAL,
55+ detail TEXT NOT NULL,
56+ -- The levels it paused, comma-separated; NULL when none.
57+ paused TEXT,
58+ opened_at TEXT NOT NULL,
59+ emailed_at TEXT
60+);
61+CREATE INDEX IF NOT EXISTS platform_alerts_by_metric ON platform_alerts (metric, opened_at);
+5−0
532532 let now = now_ms();
533533 let expires_at = rfc3339(now + RESERVATION_HOURS * 60 * 60 * 1000);
534534 let repo = format!("{}/{}", a.repo.namespace, a.repo.name).to_lowercase();
535+ // A platform pause holds for everyone, a g1t that does not charge
536+ // included (platform.rs): kept 30 seconds in the isolate.
537+ if let Some(why) = self.platform_refuses(a.kind).await {
538+ return Ok(Outcome::fail(FailureCode::Paused, why));
539+ }
535540 // A g1t that does not charge holds nothing.
536541 if self.stripe.is_none() {
537542 return Ok(Outcome::Ok(Reservation { id: new_id("rsv", now), paid_by: PaidBy::OnDemand, held_micros: 0, expires_at }));
+2−0
3434
3535 /// The cron that also checks costs against Cloudflare's bill.
3636 pub(crate) const DAILY: &str = "17 4 * * *";
37+/// The quarter-hourly tick: settling, and once an hour the platform watch.
38+pub(crate) const QUARTER_HOURLY: &str = "*/15 * * * *";
3739
3840 /// A run is settled once its logs have had time to land.
3941 const SETTLE_AFTER_MS: u64 = 5 * 60 * 1000;
+17−0
2626 mod compute;
2727 mod costs;
2828 mod margin;
29+mod platform;
2930 mod pricing;
3031 mod report;
3132 mod details;
11951196 if let Err(error) = billing.watch_spend().await {
11961197 worker::console_error!("watching g1t's own spend failed: {error}");
11971198 }
1199+ // Once an hour, at the quarter past: what Cloudflare counted for the
1200+ // whole platform in the hour before, against its thresholds
1201+ // (platform.rs).
1202+ if event.cron() == keeper::QUARTER_HOURLY && platform::hourly_due(now_ms()) {
1203+ match billing.watch_platform(&keeper).await {
1204+ Ok((read, breached)) => worker::console_log!("platform watch: {read} metrics, {breached} breaches"),
1205+ Err(error) => worker::console_error!("watching platform usage failed: {error}"),
1206+ }
1207+ }
11981208 if let Ok(identity) = env.service("IDENTITY")
11991209 && let Err(error) = billing.warn_limits(&identity).await {
12001210 worker::console_error!("warning owners failed: {error}");
14011411 "admin_cost_alerts" => reply(&billing.admin_cost_alerts(args(body)?).await?),
14021412 "admin_spend_caps" => reply(&billing.spend_caps().await?),
14031413 "admin_lift_breaker" => reply(&billing.admin_lift_breaker(args(body)?).await?),
1414+ // Platform pauses (src/platform.rs): read by every service that
1415+ // honours one, kept 30 seconds in each isolate.
1416+ "platform_pause" => reply(&billing.pause_now().await),
1417+ "admin_platform_guard" => reply(&billing.platform_guard(&keeper::Keeper::from_env(&env)).await?),
1418+ "admin_set_pause" => reply(&billing.admin_set_pause(args(body)?, &keeper::Keeper::from_env(&env)).await?),
14041419 "admin_decide_proposal" => reply(&billing.admin_decide_proposal(args(body)?).await?),
14051420 "admin_set_cost_settings" => reply(&billing.admin_set_cost_settings(args(body)?).await?),
14061421 "admin_set_cost_mapping" => reply(&billing.admin_set_cost_mapping(args(body)?).await?),
15111526 include_str!("../migrations/0046_reset_costs.sql"),
15121527 include_str!("../migrations/0047_gateway_formats.sql"),
15131528 include_str!("../migrations/0048_model_catalogue.sql"),
1529+ include_str!("../migrations/0049_superseded_proposals.sql"),
1530+ include_str!("../migrations/0050_platform_guardrails.sql"),
15141531 ];
15151532
15161533 /// The columns of `table` after the migrations: each with whether an
+965−0
1+//! Platform spend guardrails: what Cloudflare counts for all of g1t, read
2+//! every hour, and a g1t-wide pause to stop it. See
3+//! docs/SPEND-GUARDRAILS.md.
4+//!
5+//! The caps in `budget` hold what g1t pays for agents and sandboxes; this
6+//! is about the platform underneath, which no workspace's limit covers and
7+//! Cloudflare never caps: Workers requests and CPU, D1 rows, Queue
8+//! operations, Durable Objects, KV and Artifacts. A loop in any of them is
9+//! otherwise found on the bill.
10+//!
11+//! - **The pause** (`platform_pause`): four independent levels, set by
12+//! staff in sudo (Costs & margin, Platform pause) or by the watcher on a
13+//! severe breach. `compute` refuses every reservation through `reserve`
14+//! but embeddings; `indexing` refuses embeddings, and context's and
15+//! search's backfills check it themselves; `schedules` stops Actions'
16+//! cron runs and the runner's sweep; `renders` makes social cards a
17+//! static image. Callers keep the answer 30 seconds in their isolate
18+//! (`g1t_kit::pause`, `@g1t/contracts` `platformPaused`), and so does
19+//! billing (`pause_now`): never a database read per request. When the
20+//! flag cannot be read, nothing is paused: a pause is pulled on purpose,
21+//! and a billing outage must not stop the platform with it.
22+//! - **The watcher** (`watch_platform`): at the quarter past each hour, the
23+//! hour before and the month so far, from Cloudflare's GraphQL Analytics
24+//! with the costs reconciliation's token (`CLOUDFLARE_BILLING_TOKEN`, else
25+//! `CLOUDFLARE_USAGE_TOKEN`; Account Analytics Read). Each hour goes in
26+//! `platform_usage` with the script, queue, database or namespace that
27+//! counted most. A metric over its hourly threshold (`PLATFORM_HOURLY_*`),
28+//! or over `PLATFORM_SPIKE_FACTOR` times the last week's median hour (and
29+//! at least `PLATFORM_SPIKE_FLOOR_PERCENT` of its threshold), is a breach:
30+//! staff are emailed (`COSTS_ALERT_EMAIL`, once per metric every 6
31+//! hours) and sudo shows it. Over `PLATFORM_SEVERE_FACTOR` times its
32+//! threshold, the levels that metric feeds are paused, of those
33+//! `AUTO_PAUSE` allows (`schedules,indexing` unless set).
34+//!
35+//! A dataset GraphQL refuses (a field renamed, a dataset the token cannot
36+//! read) is logged and skipped; the others still count.
37+
38+use std::cell::RefCell;
39+use std::collections::BTreeMap;
40+
41+use futures_util::future::join_all;
42+use g1t_contracts::billing::{
43+ AdminSetPauseArgs, PauseLevel, PauseState, PlatformBreach, PlatformGuard, PlatformMetric, PlatformPause,
44+};
45+use g1t_contracts::time::rfc3339;
46+use g1t_contracts::{FailureCode, Outcome, new_id};
47+use g1t_kit::now_ms;
48+use serde::Deserialize;
49+use serde_json::{Value, json};
50+use worker::{Env, Result};
51+
52+use crate::Billing;
53+use crate::keeper::Keeper;
54+
55+const HOUR_MS: u64 = 60 * 60 * 1000;
56+/// How long a pause read is kept in the isolate.
57+const KEEP_MS: u64 = 30_000;
58+/// A metric's breach is emailed at most this often.
59+const ALERT_EVERY_MS: u64 = 6 * HOUR_MS;
60+/// The spike rule needs this many hours of history first.
61+const SPIKE_MIN_HOURS: usize = 24;
62+
63+thread_local! {
64+ static KEPT: RefCell<Option<(PlatformPause, u64)>> = const { RefCell::new(None) };
65+}
66+
67+/// One metric the watcher reads, with its threshold's variable and the
68+/// pause levels it feeds.
69+pub(crate) struct Metric {
70+ pub key: &'static str,
71+ pub title: &'static str,
72+ /// `PLATFORM_HOURLY_<VAR>`.
73+ pub var: &'static str,
74+ /// The hourly threshold when the variable is not set: about a dollar
75+ /// to a few dollars an hour at Cloudflare's list prices, far above a
76+ /// small alpha's normal hour.
77+ pub default: f64,
78+ /// What kind of thing `top_name` is, for the alert.
79+ pub of: &'static str,
80+ /// The levels a severe breach may pause (of those `AUTO_PAUSE` allows).
81+ pub levels: &'static [PauseLevel],
82+}
83+
84+use PauseLevel::{Compute, Indexing, Renders, Schedules};
85+
86+pub(crate) const METRICS: &[Metric] = &[
87+ Metric { key: "workers_requests", title: "Workers requests", var: "WORKERS_REQUESTS", default: 20_000_000.0, of: "script", levels: &[Schedules, Indexing, Renders] },
88+ Metric { key: "workers_cpu_ms", title: "Workers CPU (ms)", var: "WORKERS_CPU_MS", default: 100_000_000.0, of: "script", levels: &[Schedules, Indexing, Renders] },
89+ Metric { key: "d1_rows_read", title: "D1 rows read", var: "D1_ROWS_READ", default: 2_000_000_000.0, of: "D1 database id", levels: &[Schedules, Indexing] },
90+ Metric { key: "d1_rows_written", title: "D1 rows written", var: "D1_ROWS_WRITTEN", default: 5_000_000.0, of: "D1 database id", levels: &[Schedules, Indexing] },
91+ Metric { key: "queue_operations", title: "Queue operations", var: "QUEUE_OPERATIONS", default: 5_000_000.0, of: "queue id", levels: &[Schedules, Indexing] },
92+ Metric { key: "do_requests", title: "Durable Object requests", var: "DO_REQUESTS", default: 20_000_000.0, of: "script", levels: &[Compute, Schedules] },
93+ Metric { key: "do_rows_written", title: "Durable Object rows written", var: "DO_ROWS_WRITTEN", default: 5_000_000.0, of: "Durable Object namespace id", levels: &[Compute, Schedules] },
94+ Metric { key: "do_storage_write_units", title: "Durable Object storage writes", var: "DO_STORAGE_WRITE_UNITS", default: 5_000_000.0, of: "Durable Object namespace id", levels: &[Compute, Schedules] },
95+ Metric { key: "do_active_seconds", title: "Durable Object active seconds", var: "DO_ACTIVE_SECONDS", default: 3_000_000.0, of: "Durable Object namespace id", levels: &[Compute, Schedules] },
96+ Metric { key: "kv_reads", title: "KV reads", var: "KV_READS", default: 10_000_000.0, of: "KV namespace id", levels: &[Indexing, Renders] },
97+ Metric { key: "kv_writes", title: "KV writes", var: "KV_WRITES", default: 200_000.0, of: "KV namespace id", levels: &[Indexing, Renders] },
98+ Metric { key: "kv_deletes", title: "KV deletes", var: "KV_DELETES", default: 200_000.0, of: "KV namespace id", levels: &[Indexing, Renders] },
99+ Metric { key: "kv_lists", title: "KV lists", var: "KV_LISTS", default: 200_000.0, of: "KV namespace id", levels: &[Indexing, Renders] },
100+ Metric { key: "artifacts_events", title: "Artifacts events", var: "ARTIFACTS_EVENTS", default: 1_000_000.0, of: "repository", levels: &[Compute, Schedules] },
101+];
102+
103+pub(crate) fn metric(key: &str) -> Option<&'static Metric> {
104+ METRICS.iter().find(|m| m.key == key)
105+}
106+
107+/// The watcher's numbers, from the billing service's variables.
108+#[derive(Clone, Debug)]
109+pub(crate) struct Watch {
110+ /// Each metric's hourly threshold; zero turns its threshold rule off.
111+ pub thresholds: BTreeMap<&'static str, f64>,
112+ /// `PLATFORM_SPIKE_FACTOR`: an hour above this many times the week's
113+ /// median hour is a spike. Zero: no spike rule.
114+ pub spike_factor: f64,
115+ /// `PLATFORM_SPIKE_FLOOR_PERCENT`: a spike must also be at least this
116+ /// share of the metric's threshold, so a quiet metric doubling is not one.
117+ pub spike_floor_percent: f64,
118+ /// `PLATFORM_SEVERE_FACTOR`: this many times the threshold pauses.
119+ pub severe_factor: f64,
120+ /// `AUTO_PAUSE`: the levels a severe breach may pause.
121+ pub auto_pause: Vec<PauseLevel>,
122+}
123+
124+impl Watch {
125+ pub(crate) fn from_env(env: &Env) -> Self {
126+ let var = |name: &str| env.var(name).ok().map(|v| v.to_string());
127+ Watch::from_vars(var)
128+ }
129+
130+ pub(crate) fn from_vars(var: impl Fn(&str) -> Option<String>) -> Self {
131+ let number = |name: &str, default: f64| {
132+ var(name).and_then(|v| v.trim().replace('_', "").parse::<f64>().ok()).filter(|n| n.is_finite() && *n >= 0.0).unwrap_or(default)
133+ };
134+ let thresholds = METRICS.iter().map(|m| (m.key, number(&format!("PLATFORM_HOURLY_{}", m.var), m.default))).collect();
135+ let auto_pause = match var("AUTO_PAUSE") {
136+ Some(list) => list.split(',').filter_map(PauseLevel::parse).collect(),
137+ None => vec![Schedules, Indexing],
138+ };
139+ Watch {
140+ thresholds,
141+ spike_factor: number("PLATFORM_SPIKE_FACTOR", 10.0),
142+ spike_floor_percent: number("PLATFORM_SPIKE_FLOOR_PERCENT", 10.0),
143+ severe_factor: number("PLATFORM_SEVERE_FACTOR", 5.0).max(1.0),
144+ auto_pause,
145+ }
146+ }
147+
148+ pub(crate) fn threshold(&self, key: &str) -> f64 {
149+ self.thresholds.get(key).copied().unwrap_or(0.0)
150+ }
151+}
152+
153+// ---- The pause ------------------------------------------------------------
154+
155+/// What a refused reservation is told while `level` is paused.
156+pub(crate) fn pause_refusal(level: PauseLevel) -> String {
157+ let what = match level {
158+ Compute => "new agent runs, checks, workflow jobs and builds",
159+ Indexing => "new indexing and semantic search embeddings",
160+ Schedules => "scheduled workflow runs and queued agents",
161+ Renders => "new social card images",
162+ };
163+ format!("g1t has paused {what} across the platform while it looks into unusual usage. Work already running finishes; try again later.")
164+}
165+
166+/// The level a reservation of `kind` waits on.
167+pub(crate) fn level_for(kind: g1t_contracts::billing::ComputeKind) -> PauseLevel {
168+ if kind == g1t_contracts::billing::ComputeKind::Embedding { Indexing } else { Compute }
169+}
170+
171+#[derive(Deserialize)]
172+struct PauseRow {
173+ level: String,
174+ paused: i64,
175+ note: Option<String>,
176+ set_by: Option<String>,
177+ set_at: Option<String>,
178+ auto: i64,
179+}
180+
181+#[derive(Deserialize)]
182+struct UsageRow {
183+ metric: String,
184+ value: f64,
185+ top_name: Option<String>,
186+ top_value: Option<f64>,
187+}
188+
189+#[derive(Deserialize)]
190+struct AlertRow {
191+ id: String,
192+ metric: String,
193+ hour: String,
194+ rule: String,
195+ value: f64,
196+ threshold: f64,
197+ severe: i64,
198+ top_name: Option<String>,
199+ detail: String,
200+ paused: Option<String>,
201+ opened_at: String,
202+ emailed_at: Option<String>,
203+}
204+
205+impl Billing {
206+ async fn pause_rows(&self) -> Result<Vec<PauseRow>> {
207+ self.db.prepare("SELECT level, paused, note, set_by, set_at, auto FROM platform_pause").all().await?.results::<PauseRow>()
208+ }
209+
210+ /// Every level, kept in the isolate for 30 seconds. Nothing paused
211+ /// when the table cannot be read.
212+ pub(crate) async fn pause_now(&self) -> PlatformPause {
213+ let now = now_ms();
214+ if let Some(pause) = KEPT.with(|kept| kept.borrow().filter(|(_, until)| *until > now).map(|(p, _)| p)) {
215+ return pause;
216+ }
217+ let pause = match self.pause_rows().await {
218+ Ok(rows) => {
219+ let mut pause = PlatformPause::default();
220+ for row in rows {
221+ if let Some(level) = PauseLevel::parse(&row.level) {
222+ pause.set(level, row.paused != 0);
223+ }
224+ }
225+ pause
226+ }
227+ Err(error) => {
228+ worker::console_error!("platform pause unreadable, so nothing is paused: {error}");
229+ PlatformPause::default()
230+ }
231+ };
232+ KEPT.with(|kept| *kept.borrow_mut() = Some((pause, now + KEEP_MS)));
233+ pause
234+ }
235+
236+ /// Why a reservation of `kind` is refused by a platform pause, if it is.
237+ pub(crate) async fn platform_refuses(&self, kind: g1t_contracts::billing::ComputeKind) -> Option<String> {
238+ let level = level_for(kind);
239+ self.pause_now().await.is(level).then(|| pause_refusal(level))
240+ }
241+
242+ /// Pauses or resumes a level, and records who and why in the audit log.
243+ async fn set_pause(&self, level: PauseLevel, paused: bool, note: &str, by: &str, auto: bool) -> Result<()> {
244+ let now = rfc3339(now_ms());
245+ self.db
246+ .prepare(
247+ "INSERT INTO platform_pause (level, paused, note, set_by, set_at, auto) VALUES (?1, ?2, ?3, ?4, ?5, ?6)
248+ ON CONFLICT (level) DO UPDATE SET paused = ?2, note = ?3, set_by = ?4, set_at = ?5, auto = ?6",
249+ )
250+ .bind(&[level.as_str().into(), i32::from(paused).into(), note.into(), by.into(), now.as_str().into(), i32::from(auto).into()])?
251+ .run()
252+ .await?;
253+ KEPT.with(|kept| *kept.borrow_mut() = None);
254+ let action = if paused { "platform_paused" } else { "platform_resumed" };
255+ self.audit("costs", action, &format!("{}: {note}", level.as_str()), by).await
256+ }
257+
258+ /// `admin_set_pause`: staff pause or resume one level, with why.
259+ pub(crate) async fn admin_set_pause(&self, a: AdminSetPauseArgs, keeper: &Keeper) -> Result<Outcome<PlatformGuard>> {
260+ let (by, note) = (a.by.trim(), a.note.trim());
261+ let Some(level) = PauseLevel::parse(&a.level) else {
262+ return Ok(Outcome::fail(FailureCode::Invalid, "Pick compute, schedules, indexing or renders."));
263+ };
264+ if by.is_empty() || note.chars().count() < 5 {
265+ return Ok(Outcome::fail(FailureCode::Invalid, "Say who is changing it, and why, in the note."));
266+ }
267+ let note: String = note.chars().take(500).collect();
268+ self.set_pause(level, a.paused, &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.
374+pub(crate) struct Query {
375+ pub key: &'static str,
376+ pub dataset: &'static str,
377+ /// What the dataset is selected with: `sum { … }`, or `count`.
378+ pub select: &'static str,
379+ /// The dimension that names what counted (and `actionType` for KV).
380+ pub dimensions: &'static str,
381+ /// The filter fields for an hour: `datetime` (Time) on most datasets,
382+ /// `datetimeHour` on D1's.
383+ pub hour_filter: &'static str,
384+}
385+
386+/// Each metric from its own query where a field is less certain, so one
387+/// GraphQL refuses does not take the others with it. Field names as
388+/// Cloudflare documents them; a refused one is logged and skipped.
389+pub(crate) const QUERIES: &[Query] = &[
390+ Query { key: "workers", dataset: "workersInvocationsAdaptive", select: "sum { requests }", dimensions: "scriptName", hour_filter: "datetime" },
391+ Query { key: "workers_cpu", dataset: "workersInvocationsAdaptive", select: "sum { cpuTimeUs }", dimensions: "scriptName", hour_filter: "datetime" },
392+ Query { key: "d1", dataset: "d1AnalyticsAdaptiveGroups", select: "sum { rowsRead rowsWritten }", dimensions: "databaseId", hour_filter: "datetimeHour" },
393+ Query { key: "queues", dataset: "queueMessageOperationsAdaptiveGroups", select: "sum { billableOperations }", dimensions: "queueId", hour_filter: "datetime" },
394+ Query { key: "do_invocations", dataset: "durableObjectsInvocationsAdaptiveGroups", select: "sum { requests }", dimensions: "scriptName", hour_filter: "datetime" },
395+ Query { key: "do_periodic", dataset: "durableObjectsPeriodicGroups", select: "sum { activeTime storageWriteUnits }", dimensions: "namespaceId", hour_filter: "datetime" },
396+ Query { key: "do_sql", dataset: "durableObjectsPeriodicGroups", select: "sum { rowsWritten }", dimensions: "namespaceId", hour_filter: "datetime" },
397+ Query { key: "kv", dataset: "kvOperationsAdaptiveGroups", select: "sum { requests }", dimensions: "namespaceId actionType", hour_filter: "datetime" },
398+ Query { key: "artifacts", dataset: "artifactsEventsAdaptiveGroups", select: "count", dimensions: "repositoryName", hour_filter: "datetime" },
399+];
400+
401+/// A window to read: one hour (`Time` bounds) or days (`Date` bounds).
402+#[derive(Clone, Debug, PartialEq, Eq)]
403+pub(crate) enum Window {
404+ Hour { since: String, until: String },
405+ Days { since: String, until: String },
406+}
407+
408+/// The hour before the one `now` is in, `[since, until)`, as Cloudflare's
409+/// `Time` takes it.
410+pub(crate) fn last_hour(now: u64) -> Window {
411+ let until = now / HOUR_MS * HOUR_MS;
412+ Window::Hour { since: hour_label(until - HOUR_MS), until: hour_label(until) }
413+}
414+
415+/// The month so far (UTC), today included.
416+pub(crate) fn month_so_far(now: u64) -> Window {
417+ let today = rfc3339(now)[..10].to_owned();
418+ Window::Days { since: format!("{}-01", &today[..7]), until: today }
419+}
420+
421+/// `2026-10-08T13:00:00Z`.
422+pub(crate) fn hour_label(ms: u64) -> String {
423+ format!("{}:00:00Z", &rfc3339(ms)[..13])
424+}
425+
426+/// The request body for `query` over `window`, its dataset aliased `rows`.
427+pub(crate) fn query_body(account: &str, query: &Query, window: &Window) -> Value {
428+ let (filter, kind, since, until) = match window {
429+ Window::Hour { since, until } => (format!("{f}_geq: $since, {f}_lt: $until", f = query.hour_filter), "Time", since, until),
430+ Window::Days { since, until } => ("date_geq: $since, date_leq: $until".to_owned(), "Date", since, until),
431+ };
432+ let text = format!(
433+ "query ($account: String!, $since: {kind}!, $until: {kind}!) {{
434+ viewer {{ accounts(filter: {{ accountTag: $account }}) {{
435+ rows: {dataset}(limit: 10000, filter: {{ {filter} }}) {{
436+ {select}
437+ dimensions {{ {dimensions} }}
438+ }}
439+ }} }}
440+}}",
441+ dataset = query.dataset,
442+ select = query.select,
443+ dimensions = query.dimensions,
444+ );
445+ json!({ "query": text, "variables": { "account": account, "since": since, "until": until } })
446+}
447+
448+/// What one query counted, per metric: the total and each name's part.
449+pub(crate) type Counted = BTreeMap<&'static str, BTreeMap<String, f64>>;
450+
451+/// A KV operation's metric, by its `actionType`.
452+fn kv_metric(action: &str) -> Option<&'static str> {
453+ match action.to_ascii_lowercase().as_str() {
454+ "read" | "get" => Some("kv_reads"),
455+ "write" | "put" => Some("kv_writes"),
456+ "delete" => Some("kv_deletes"),
457+ "list" => Some("kv_lists"),
458+ _ => None,
459+ }
460+}
461+
462+/// Reads one query's answer into metrics. Errors when GraphQL does.
463+pub(crate) fn counted(query: &Query, body: &Value) -> std::result::Result<Counted, String> {
464+ if let Some(errors) = body["errors"].as_array().filter(|e| !e.is_empty()) {
465+ let messages: Vec<&str> = errors.iter().filter_map(|e| e["message"].as_str()).collect();
466+ return Err(format!("{} ({}): {}", query.dataset, query.key, messages.join("; ")));
467+ }
468+ let Some(groups) = body["data"]["viewer"]["accounts"][0]["rows"].as_array() else {
469+ return Err(format!("{} ({}): no rows in the answer", query.dataset, query.key));
470+ };
471+ let mut out: Counted = BTreeMap::new();
472+ let mut add = |metric: &'static str, name: &str, value: f64| {
473+ if value.is_finite() && value > 0.0 {
474+ *out.entry(metric).or_default().entry(name.to_owned()).or_default() += value;
475+ }
476+ };
477+ for g in groups {
478+ let sum = &g["sum"];
479+ let number = |field: &str| sum[field].as_f64().unwrap_or(0.0);
480+ let dims = &g["dimensions"];
481+ let name_of = |field: &str| dims[field].as_str().filter(|s| !s.is_empty()).unwrap_or("(unnamed)").to_owned();
482+ match query.key {
483+ "workers" => add("workers_requests", &name_of("scriptName"), number("requests")),
484+ "workers_cpu" => add("workers_cpu_ms", &name_of("scriptName"), number("cpuTimeUs") / 1000.0),
485+ "d1" => {
486+ let name = name_of("databaseId");
487+ add("d1_rows_read", &name, number("rowsRead"));
488+ add("d1_rows_written", &name, number("rowsWritten"));
489+ }
490+ "queues" => add("queue_operations", &name_of("queueId"), number("billableOperations")),
491+ "do_invocations" => add("do_requests", &name_of("scriptName"), number("requests")),
492+ "do_periodic" => {
493+ let name = name_of("namespaceId");
494+ // activeTime is in microseconds.
495+ add("do_active_seconds", &name, number("activeTime") / 1_000_000.0);
496+ add("do_storage_write_units", &name, number("storageWriteUnits"));
497+ }
498+ "do_sql" => add("do_rows_written", &name_of("namespaceId"), number("rowsWritten")),
499+ "kv" => {
500+ if let Some(metric) = dims["actionType"].as_str().and_then(kv_metric) {
501+ add(metric, &name_of("namespaceId"), number("requests"));
502+ }
503+ }
504+ "artifacts" => add("artifacts_events", &name_of("repositoryName"), g["count"].as_f64().unwrap_or(0.0)),
505+ _ => {}
506+ }
507+ }
508+ Ok(out)
509+}
510+
511+/// A metric's total, and the name that counted most.
512+#[derive(Clone, Debug, PartialEq)]
513+pub(crate) struct Total {
514+ pub value: f64,
515+ pub top_name: Option<String>,
516+ pub top_value: Option<f64>,
517+}
518+
519+/// Every query's answers merged into one total per metric.
520+pub(crate) fn totals(answers: &[Counted]) -> BTreeMap<&'static str, Total> {
521+ let mut merged: Counted = BTreeMap::new();
522+ for answer in answers {
523+ for (metric, names) in answer {
524+ let entry = merged.entry(metric).or_default();
525+ for (name, value) in names {
526+ *entry.entry(name.clone()).or_default() += value;
527+ }
528+ }
529+ }
530+ merged
531+ .into_iter()
532+ .map(|(metric, names)| {
533+ let value = names.values().sum();
534+ let top = names.into_iter().max_by(|a, b| a.1.total_cmp(&b.1));
535+ (metric, Total { value, top_name: top.as_ref().map(|t| t.0.clone()), top_value: top.map(|t| t.1) })
536+ })
537+ .collect()
538+}
539+
540+// ---- Deciding what is a breach --------------------------------------------
541+
542+/// The middle of the week's hours (the mean of the middle two when even).
543+pub(crate) fn median(values: &mut [f64]) -> Option<f64> {
544+ if values.is_empty() {
545+ return None;
546+ }
547+ values.sort_by(f64::total_cmp);
548+ let mid = values.len() / 2;
549+ Some(if values.len() % 2 == 0 { (values[mid - 1] + values[mid]) / 2.0 } else { values[mid] })
550+}
551+
552+/// A breach of one metric, before it is recorded.
553+#[derive(Clone, Debug, PartialEq)]
554+pub(crate) struct Breach {
555+ pub metric: &'static str,
556+ /// `threshold` or `spike`.
557+ pub rule: &'static str,
558+ pub value: f64,
559+ /// What it was held to: the threshold, or the spike's line.
560+ pub limit: f64,
561+ pub severe: bool,
562+}
563+
564+/// Whether an hour's `value` of `metric` is a breach, given the week's
565+/// earlier hours (`history`). The threshold rule wins over the spike rule.
566+pub(crate) fn judge(watch: &Watch, metric: &'static str, value: f64, history: &[f64]) -> Option<Breach> {
567+ let threshold = watch.threshold(metric);
568+ if threshold > 0.0 && value > threshold {
569+ return Some(Breach { metric, rule: "threshold", value, limit: threshold, severe: value >= threshold * watch.severe_factor });
570+ }
571+ if watch.spike_factor <= 0.0 || history.len() < SPIKE_MIN_HOURS {
572+ return None;
573+ }
574+ let mut week = history.to_vec();
575+ let usual = median(&mut week)?;
576+ let line = usual * watch.spike_factor;
577+ let floor = threshold * watch.spike_floor_percent / 100.0;
578+ // A metric with no threshold never spikes: there is no floor to hold it to.
579+ (threshold > 0.0 && value > line && value >= floor).then_some(Breach { metric, rule: "spike", value, limit: line, severe: false })
580+}
581+
582+/// The levels a severe breach of `metric` pauses: those it feeds that
583+/// `AUTO_PAUSE` allows.
584+pub(crate) fn levels_to_pause(watch: &Watch, breach: &Breach) -> Vec<PauseLevel> {
585+ if !breach.severe {
586+ return vec![];
587+ }
588+ metric(breach.metric).map_or(vec![], |m| m.levels.iter().copied().filter(|l| watch.auto_pause.contains(l)).collect())
589+}
590+
591+/// `20,000,000`, or `1.5` for small fractions.
592+pub(crate) fn amount(n: f64) -> String {
593+ if n.fract().abs() > 0.0 && n.abs() < 100.0 {
594+ return format!("{n:.1}");
595+ }
596+ let digits = format!("{:.0}", n.abs());
597+ let mut out = String::new();
598+ for (i, c) in digits.chars().enumerate() {
599+ if i > 0 && (digits.len() - i).is_multiple_of(3) {
600+ out.push(',');
601+ }
602+ out.push(c);
603+ }
604+ if n < 0.0 { format!("-{out}") } else { out }
605+}
606+
607+/// The sentence an alert and sudo carry for a breach.
608+pub(crate) fn describe(breach: &Breach, hour: &str, top: Option<(&str, f64)>) -> String {
609+ let m = metric(breach.metric);
610+ let title = m.map_or(breach.metric, |m| m.title);
611+ let of = m.map_or("source", |m| m.of);
612+ let what = match breach.rule {
613+ "spike" => format!("more than its spike line of {} (the last week's usual hour times PLATFORM_SPIKE_FACTOR)", amount(breach.limit)),
614+ _ => format!(
615+ "over its hourly threshold of {} (PLATFORM_HOURLY_{})",
616+ amount(breach.limit),
617+ m.map_or("?", |m| m.var)
618+ ),
619+ };
620+ let mut text = format!("{title}: {} in the hour from {hour}, {what}.", amount(breach.value));
621+ if let Some((name, value)) = top {
622+ text.push_str(&format!(" Most of it from the {of} {name} ({}).", amount(value)));
623+ }
624+ text
625+}
626+
627+/// Whether the quarter-hour tick at `now` is the hour's watch: the one at
628+/// a quarter past, when the hour before is in Cloudflare's analytics.
629+pub(crate) fn hourly_due(now: u64) -> bool {
630+ let minute = now / 60_000 % 60;
631+ (15..30).contains(&minute)
632+}
633+
634+impl Billing {
635+ /// Each hour: the hour before and the month so far, from Cloudflare,
636+ /// kept and judged; staff emailed and levels paused on a breach.
637+ /// Returns how many metrics were read and how many breached.
638+ pub(crate) async fn watch_platform(&self, keeper: &Keeper) -> Result<(usize, usize)> {
639+ if !keeper.can_read_bill() {
640+ worker::console_log!("platform watch: no Cloudflare token with Account Analytics Read, so nothing is read");
641+ return Ok((0, 0));
642+ }
643+ let now = now_ms();
644+ let watch = Watch::from_env(&self.env);
645+ let hour_window = last_hour(now);
646+ let month_window = month_so_far(now);
647+ let Window::Hour { since: hour, .. } = &hour_window else { unreachable!() };
648+ let hour = hour.clone();
649+ let month = rfc3339(now)[..7].to_owned();
650+ let read = |window: &Window| {
651+ join_all(QUERIES.iter().map(|query| {
652+ let body = query_body(keeper.account(), query, window);
653+ async move {
654+ match keeper.graphql(body).await {
655+ Ok(answer) => counted(query, &answer),
656+ Err(error) => Err(format!("{} ({}): {error}", query.dataset, query.key)),
657+ }
658+ }
659+ }))
660+ };
661+ let (hour_answers, month_answers) = futures_util::future::join(read(&hour_window), read(&month_window)).await;
662+ let keep = |answers: Vec<std::result::Result<Counted, String>>| {
663+ answers
664+ .into_iter()
665+ .filter_map(|answer| answer.map_err(|error| worker::console_error!("platform watch: skipped {error}")).ok())
666+ .collect::<Vec<_>>()
667+ };
668+ let hourly = totals(&keep(hour_answers));
669+ let monthly = totals(&keep(month_answers));
670+ let read_at = rfc3339(now);
671+ let mut writes = vec![];
672+ for (table, period, values) in [("platform_usage", "hour", &hourly), ("platform_usage_month", "month", &monthly)] {
673+ let key = if period == "hour" { hour.as_str() } else { month.as_str() };
674+ for (metric, total) in values.iter() {
675+ writes.push(
676+ self.db
677+ .prepare(format!(
678+ "INSERT INTO {table} ({period}, metric, value, top_name, top_value, read_at) VALUES (?1, ?2, ?3, ?4, ?5, ?6)
679+ ON CONFLICT ({period}, metric) DO UPDATE SET value = ?3, top_name = ?4, top_value = ?5, read_at = ?6"
680+ ))
681+ .bind(&[
682+ key.into(),
683+ (*metric).into(),
684+ total.value.into(),
685+ crate::optional(total.top_name.as_deref()),
686+ total.top_value.map_or(worker::wasm_bindgen::JsValue::NULL, Into::into),
687+ read_at.as_str().into(),
688+ ])?,
689+ );
690+ }
691+ }
692+ if !writes.is_empty() {
693+ self.db.batch(writes).await?;
694+ }
695+ // The week before this hour, for the spike rule.
696+ #[derive(Deserialize)]
697+ struct Past {
698+ metric: String,
699+ value: f64,
700+ }
701+ let past = self
702+ .db
703+ .prepare("SELECT metric, value FROM platform_usage WHERE hour >= ? AND hour < ?")
704+ .bind(&[hour_label(now / HOUR_MS * HOUR_MS - 8 * 24 * HOUR_MS).into(), hour.as_str().into()])?
705+ .all()
706+ .await?
707+ .results::<Past>()?;
708+ let mut breaches = vec![];
709+ for (metric, total) in &hourly {
710+ let history: Vec<f64> = past.iter().filter(|p| p.metric == *metric).map(|p| p.value).collect();
711+ if let Some(breach) = judge(&watch, metric, total.value, &history) {
712+ breaches.push((breach, total.clone()));
713+ }
714+ }
715+ let found = breaches.len();
716+ if found > 0 {
717+ self.on_breaches(&watch, &hour, breaches).await?;
718+ }
719+ Ok((hourly.len(), found))
720+ }
721+
722+ /// Records each breach, pauses what a severe one should, and emails
723+ /// staff about the metrics not emailed in the last 6 hours.
724+ async fn on_breaches(&self, watch: &Watch, hour: &str, breaches: Vec<(Breach, Total)>) -> Result<()> {
725+ let now = now_ms();
726+ let opened_at = rfc3339(now);
727+ let current = self.pause_now().await;
728+ let mut lines = vec![];
729+ let mut ids = vec![];
730+ let mut paused_now: Vec<PauseLevel> = vec![];
731+ for (breach, total) in breaches {
732+ let top = total.top_name.as_deref().zip(total.top_value);
733+ let mut detail = describe(&breach, hour, top);
734+ let pause: Vec<PauseLevel> =
735+ levels_to_pause(watch, &breach).into_iter().filter(|l| !current.is(*l) && !paused_now.contains(l)).collect();
736+ for level in &pause {
737+ self.set_pause(*level, true, &format!("Automatic: {detail}"), "g1t-billing's usage watcher", true).await?;
738+ paused_now.push(*level);
739+ }
740+ if !pause.is_empty() {
741+ let names: Vec<&str> = pause.iter().map(|l| l.as_str()).collect();
742+ detail.push_str(&format!(" Severe (over {}× the threshold): paused {}.", amount(watch.severe_factor), names.join(" and ")));
743+ }
744+ let id = new_id("pal", now);
745+ let paused_text = pause.iter().map(|l| l.as_str()).collect::<Vec<_>>().join(",");
746+ self.db
747+ .prepare(
748+ "INSERT INTO platform_alerts (id, metric, hour, rule, value, threshold, severe, top_name, top_value, detail, paused, opened_at)
749+ VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
750+ )
751+ .bind(&[
752+ id.as_str().into(),
753+ breach.metric.into(),
754+ hour.into(),
755+ breach.rule.into(),
756+ breach.value.into(),
757+ breach.limit.into(),
758+ i32::from(breach.severe).into(),
759+ crate::optional(total.top_name.as_deref()),
760+ total.top_value.map_or(worker::wasm_bindgen::JsValue::NULL, Into::into),
761+ detail.as_str().into(),
762+ crate::optional(Some(paused_text.as_str()).filter(|t| !t.is_empty())),
763+ opened_at.as_str().into(),
764+ ])?
765+ .run()
766+ .await?;
767+ // Emailed at most once per metric every 6 hours; a new pause is
768+ // always said.
769+ #[derive(Deserialize)]
770+ struct Last {
771+ at: Option<String>,
772+ }
773+ let last = self
774+ .db
775+ .prepare("SELECT MAX(emailed_at) AS at FROM platform_alerts WHERE metric = ? AND emailed_at >= ?")
776+ .bind(&[breach.metric.into(), rfc3339(now.saturating_sub(ALERT_EVERY_MS)).into()])?
777+ .first::<Last>(None)
778+ .await?
779+ .and_then(|l| l.at);
780+ if last.is_none() || !pause.is_empty() {
781+ lines.push(detail);
782+ ids.push(id);
783+ }
784+ }
785+ let alert_to = &self.caps.alert_to;
786+ if lines.is_empty() || alert_to.is_empty() {
787+ return Ok(());
788+ }
789+ let subject = if paused_now.is_empty() {
790+ format!("g1t: platform usage breach in the hour from {hour}")
791+ } else {
792+ let names: Vec<&str> = paused_now.iter().map(|l| l.as_str()).collect();
793+ format!("g1t: platform usage breach, {} paused", names.join(" and "))
794+ };
795+ lines.push(
796+ "Ids are Cloudflare's: `node scripts/ops/platform-usage.mjs` names them and shows the last 24 hours. To pause or resume a level: sudo, Costs & margin, Platform pause. Thresholds: PLATFORM_HOURLY_* in services/billing/wrangler.jsonc. See docs/SPEND-GUARDRAILS.md.".to_owned(),
797+ );
798+ match crate::margin::email_staff_page(&self.env, alert_to, &subject, &lines, ("Platform pause", "https://sudo.g1t.sh/costs#platform"), "g1t-billing's usage watcher").await {
799+ Ok(()) => {
800+ let at = rfc3339(now_ms());
801+ for id in ids {
802+ self.db.prepare("UPDATE platform_alerts SET emailed_at = ? WHERE id = ?").bind(&[at.as_str().into(), id.as_str().into()])?.run().await?;
803+ }
804+ }
805+ Err(error) => worker::console_error!("could not email the platform usage breach: {error}"),
806+ }
807+ Ok(())
808+ }
809+}
810+
811+#[cfg(test)]
812+mod tests {
813+ use super::*;
814+
815+ fn watch() -> Watch {
816+ Watch::from_vars(|_| None)
817+ }
818+
819+ #[test]
820+ fn every_metric_has_a_threshold_and_a_query_that_counts_it() {
821+ let w = watch();
822+ for m in METRICS {
823+ assert!(w.threshold(m.key) > 0.0, "{}", m.key);
824+ }
825+ // The thresholds named in docs/SPEND-GUARDRAILS.md.
826+ assert_eq!(w.threshold("kv_lists"), 200_000.0);
827+ assert_eq!(w.threshold("queue_operations"), 5_000_000.0);
828+ assert_eq!(w.threshold("d1_rows_read"), 2_000_000_000.0);
829+ assert_eq!(w.threshold("do_rows_written"), 5_000_000.0);
830+ assert_eq!(w.threshold("workers_requests"), 20_000_000.0);
831+ assert_eq!(w.auto_pause, vec![Schedules, Indexing]);
832+ }
833+
834+ #[test]
835+ fn variables_change_thresholds_and_auto_pause() {
836+ let w = Watch::from_vars(|name| match name {
837+ "PLATFORM_HOURLY_KV_LISTS" => Some("50_000".into()),
838+ "PLATFORM_HOURLY_D1_ROWS_READ" => Some("0".into()),
839+ "AUTO_PAUSE" => Some("compute, renders,nonsense".into()),
840+ "PLATFORM_SEVERE_FACTOR" => Some("0.5".into()),
841+ _ => None,
842+ });
843+ assert_eq!(w.threshold("kv_lists"), 50_000.0);
844+ assert_eq!(w.threshold("d1_rows_read"), 0.0);
845+ assert_eq!(w.auto_pause, vec![Compute, Renders]);
846+ // Never below the threshold itself.
847+ assert_eq!(w.severe_factor, 1.0);
848+ // Empty: nothing pauses by itself.
849+ assert!(Watch::from_vars(|n| (n == "AUTO_PAUSE").then(String::new)).auto_pause.is_empty());
850+ }
851+
852+ #[test]
853+ fn the_hour_read_is_the_one_before_at_a_quarter_past() {
854+ // 2026-10-08 13:17:00 UTC.
855+ let now = 1_791_465_420_000;
856+ assert_eq!(rfc3339(now), "2026-10-08T13:17:00.000Z");
857+ assert!(hourly_due(now));
858+ assert!(!hourly_due(now - 15 * 60_000));
859+ assert!(!hourly_due(now + 15 * 60_000));
860+ assert_eq!(last_hour(now), Window::Hour { since: "2026-10-08T12:00:00Z".into(), until: "2026-10-08T13:00:00Z".into() });
861+ assert_eq!(month_so_far(now), Window::Days { since: "2026-10-01".into(), until: "2026-10-08".into() });
862+ }
863+
864+ #[test]
865+ fn queries_alias_their_dataset_and_filter_by_the_window() {
866+ let d1 = QUERIES.iter().find(|q| q.key == "d1").unwrap();
867+ let hour = query_body("acct", d1, &last_hour(1_791_465_420_000));
868+ let text = hour["query"].as_str().unwrap();
869+ assert!(text.contains("rows: d1AnalyticsAdaptiveGroups(limit: 10000, filter: { datetimeHour_geq: $since, datetimeHour_lt: $until })"), "{text}");
870+ assert!(text.contains("$since: Time!") && text.contains("sum { rowsRead rowsWritten }") && text.contains("dimensions { databaseId }"));
871+ assert_eq!(hour["variables"]["since"], "2026-10-08T12:00:00Z");
872+ let month = query_body("acct", d1, &month_so_far(1_791_465_420_000));
873+ let text = month["query"].as_str().unwrap();
874+ assert!(text.contains("date_geq: $since, date_leq: $until") && text.contains("$since: Date!"), "{text}");
875+ }
876+
877+ fn answer(rows: Value) -> Value {
878+ json!({ "data": { "viewer": { "accounts": [{ "rows": rows }] } }, "errors": null })
879+ }
880+
881+ fn query(key: &str) -> &'static Query {
882+ QUERIES.iter().find(|q| q.key == key).unwrap()
883+ }
884+
885+ #[test]
886+ fn answers_become_metrics_by_name_and_a_refused_dataset_is_skipped() {
887+ let kv = counted(
888+ query("kv"),
889+ &answer(json!([
890+ { "sum": { "requests": 150000 }, "dimensions": { "namespaceId": "ns_a", "actionType": "list" } },
891+ { "sum": { "requests": 90000 }, "dimensions": { "namespaceId": "ns_b", "actionType": "list" } },
892+ { "sum": { "requests": 4000 }, "dimensions": { "namespaceId": "ns_a", "actionType": "read" } },
893+ { "sum": { "requests": 7 }, "dimensions": { "namespaceId": "ns_a", "actionType": "mystery" } },
894+ ])),
895+ )
896+ .unwrap();
897+ let d1 = counted(query("d1"), &answer(json!([{ "sum": { "rowsRead": 10.0, "rowsWritten": 2.0 }, "dimensions": { "databaseId": "db1" } }]))).unwrap();
898+ let cpu = counted(query("workers_cpu"), &answer(json!([{ "sum": { "cpuTimeUs": 5_000_000 }, "dimensions": { "scriptName": "g1t-web" } }]))).unwrap();
899+ let refused = counted(query("do_sql"), &json!({ "data": null, "errors": [{ "message": "unknown field rowsWritten" }] }));
900+ assert!(refused.unwrap_err().contains("unknown field rowsWritten"));
901+ assert!(counted(query("queues"), &json!({ "data": { "viewer": { "accounts": [] } } })).is_err());
902+ let all = totals(&[kv, d1, cpu]);
903+ assert_eq!(all["kv_lists"], Total { value: 240_000.0, top_name: Some("ns_a".into()), top_value: Some(150_000.0) });
904+ assert_eq!(all["kv_reads"].value, 4_000.0);
905+ assert_eq!(all["d1_rows_written"].value, 2.0);
906+ assert_eq!(all["workers_cpu_ms"].value, 5_000.0);
907+ assert!(!all.contains_key("do_rows_written"));
908+ }
909+
910+ #[test]
911+ fn over_the_threshold_is_a_breach_and_five_times_it_is_severe() {
912+ let w = watch();
913+ assert_eq!(judge(&w, "kv_lists", 200_000.0, &[]), None);
914+ let breach = judge(&w, "kv_lists", 240_000.0, &[]).unwrap();
915+ assert_eq!((breach.rule, breach.severe), ("threshold", false));
916+ assert!(levels_to_pause(&w, &breach).is_empty());
917+ let severe = judge(&w, "kv_lists", 1_000_000.0, &[]).unwrap();
918+ assert!(severe.severe);
919+ // KV feeds indexing and renders; AUTO_PAUSE allows indexing only.
920+ assert_eq!(levels_to_pause(&w, &severe), vec![Indexing]);
921+ // Durable Objects feed compute, which is never paused by itself by default.
922+ let dos = judge(&w, "do_rows_written", 30_000_000.0, &[]).unwrap();
923+ assert_eq!(levels_to_pause(&w, &dos), vec![Schedules]);
924+ }
925+
926+ #[test]
927+ fn a_spike_is_ten_times_the_weeks_usual_hour_above_a_floor() {
928+ let w = watch();
929+ let quiet = vec![1_000.0; 168];
930+ // Ten times the usual hour, but under 10% of the 5M threshold: not one.
931+ assert_eq!(judge(&w, "queue_operations", 20_000.0, &quiet), None);
932+ let usual = vec![60_000.0; 168];
933+ let spike = judge(&w, "queue_operations", 700_000.0, &usual).unwrap();
934+ assert_eq!((spike.rule, spike.limit, spike.severe), ("spike", 600_000.0, false));
935+ assert_eq!(judge(&w, "queue_operations", 590_000.0, &usual), None);
936+ // Under a day of history: no spike rule yet.
937+ assert_eq!(judge(&w, "queue_operations", 700_000.0, &usual[..23]), None);
938+ assert_eq!(median(&mut [3.0, 1.0, 2.0, 10.0]), Some(2.5));
939+ assert_eq!(median(&mut []), None);
940+ }
941+
942+ #[test]
943+ fn an_alert_names_where_to_look() {
944+ let breach = Breach { metric: "kv_lists", rule: "threshold", value: 1_250_000.0, limit: 200_000.0, severe: true };
945+ let text = describe(&breach, "2026-10-08T12:00:00Z", Some(("e627b571", 1_200_000.0)));
946+ assert_eq!(
947+ text,
948+ "KV lists: 1,250,000 in the hour from 2026-10-08T12:00:00Z, over its hourly threshold of 200,000 (PLATFORM_HOURLY_KV_LISTS). Most of it from the KV namespace id e627b571 (1,200,000)."
949+ );
950+ assert_eq!(amount(1.5), "1.5");
951+ assert_eq!(amount(5.0), "5");
952+ }
953+
954+ #[test]
955+ fn a_pause_refuses_the_level_a_reservation_waits_on() {
956+ use g1t_contracts::billing::ComputeKind;
957+ assert_eq!(level_for(ComputeKind::Embedding), Indexing);
958+ for kind in [ComputeKind::Agent, ComputeKind::Check, ComputeKind::Workflow, ComputeKind::Queue, ComputeKind::Deploy] {
959+ assert_eq!(level_for(kind), Compute);
960+ }
961+ assert!(pause_refusal(Compute).contains("Work already running finishes"));
962+ assert_eq!(PauseLevel::parse(" indexing "), Some(Indexing));
963+ assert_eq!(PauseLevel::parse("everything"), None);
964+ }
965+}
+32−1
134134 // 00:00 UTC; staff are emailed and can lift it in sudo. Workspaces
135135 // paying with real money are never paused. 0 turns either cap off.
136136 "PLATFORM_DAILY_SPEND_CAP_MICROS": "75000000",
137+ // The platform watch (src/platform.rs, docs/SPEND-GUARDRAILS.md): at
138+ // a quarter past each hour, what Cloudflare counted in the hour
139+ // before, per metric, against these thresholds. Each is a dollar or
140+ // a few an hour at list prices, far above a small alpha's hour; 0
141+ // turns that metric's threshold off. Over one, staff are emailed
142+ // (COSTS_ALERT_EMAIL, once per metric every 6 hours) and sudo shows
143+ // it.
144+ "PLATFORM_HOURLY_WORKERS_REQUESTS": "20000000",
145+ "PLATFORM_HOURLY_WORKERS_CPU_MS": "100000000",
146+ "PLATFORM_HOURLY_D1_ROWS_READ": "2000000000",
147+ "PLATFORM_HOURLY_D1_ROWS_WRITTEN": "5000000",
148+ "PLATFORM_HOURLY_QUEUE_OPERATIONS": "5000000",
149+ "PLATFORM_HOURLY_DO_REQUESTS": "20000000",
150+ "PLATFORM_HOURLY_DO_ROWS_WRITTEN": "5000000",
151+ "PLATFORM_HOURLY_DO_STORAGE_WRITE_UNITS": "5000000",
152+ "PLATFORM_HOURLY_DO_ACTIVE_SECONDS": "3000000",
153+ "PLATFORM_HOURLY_KV_READS": "10000000",
154+ "PLATFORM_HOURLY_KV_WRITES": "200000",
155+ "PLATFORM_HOURLY_KV_DELETES": "200000",
156+ "PLATFORM_HOURLY_KV_LISTS": "200000",
157+ "PLATFORM_HOURLY_ARTIFACTS_EVENTS": "1000000",
158+ // A spike: an hour over 10 times the last week's median hour, and at
159+ // least 10% of its threshold. Emailed, never paused by itself.
160+ "PLATFORM_SPIKE_FACTOR": "10",
161+ "PLATFORM_SPIKE_FLOOR_PERCENT": "10",
162+ // Severe: 5 times a threshold pauses the levels that metric feeds,
163+ // of those AUTO_PAUSE names (compute, schedules, indexing, renders;
164+ // comma-separated, empty for none). Staff resume them in sudo.
165+ "PLATFORM_SEVERE_FACTOR": "5",
166+ "AUTO_PAUSE": "schedules,indexing",
137167 // Cloudflare's fixed subscriptions a month, for sudo's figures only:
138168 // Workers Paid ($5) and Workers for Platforms ($25).
139169 "CLOUDFLARE_FIXED_MONTHLY_MICROS": "30000000"
140170 },
141− // Settling runs every 15 minutes; checking costs daily (keeper::DAILY).
171+ // Settling runs every 15 minutes, and the tick at a quarter past each
172+ // hour also runs the platform watch; checking costs daily (keeper::DAILY).
142173 "triggers": { "crons": ["*/15 * * * *", "17 4 * * *"] },
143174 // Secrets: STRIPE_SECRET_KEY. Without it nothing is charged and the
144175 // runner decides who may start agents some other way.
+21−0
4848 billingClient,
4949 embeddingEstimateMicros,
5050 localRefusal,
51+ platformPaused,
5152 currentMovedPath,
5253 currentWorkspaceSlug,
5354 repoMove,
125126 const SHARED: Set<EntityKind> = new Set(["owner", "language", "integration"]);
126127
127128 /** One compute gate per isolate, so entitlements are kept between calls. */
129+/** What a backfill says while indexing is paused across g1t (billing's `platform_pause`). */
130+const INDEXING_PAUSED = "g1t has paused indexing across the platform for now. Rebuild the context again later.";
131+
128132 let computeGate: ComputeGate | null = null;
129133 function gateFor(billing: ServiceBinding): ComputeGate {
130134 computeGate ??= new ComputeGate(billing);
943947 async backfill(a: { actor: User; workspace: string }): Promise<Result<Backfill>> {
944948 const workspace = a.workspace.toLowerCase();
945949 if (!isMember(a.actor, workspace)) return fail("forbidden", "Only members can rebuild a workspace's context.");
950+ if (await platformPaused(this.env.BILLING, "indexing")) return fail("paused", INDEXING_PAUSED);
946951 const running = await this.backfillRow(workspace);
947952 if (running?.status === "running" && Date.now() - Date.parse(running.startedAt) < BACKFILL_STALE_MS) return ok(running);
948953 const actor = (await this.workspaceActor(workspace)) ?? a.actor;
968973 async runJob(job: Job): Promise<void> {
969974 const actor = await this.workspaceActor(job.workspace);
970975 if (!actor) return;
976+ // Indexing paused across g1t: the job is counted done with the pause
977+ // as its note, so the backfill finishes and can be run again later.
978+ if (await platformPaused(this.env.BILLING, "indexing")) {
979+ if (job.type === "backfill_project") {
980+ await this.db
981+ .prepare(
982+ `UPDATE backfills SET done = done + 1, error = ?,
983+ status = CASE WHEN done + 1 >= projects THEN 'done' ELSE status END,
984+ finished_at = CASE WHEN done + 1 >= projects THEN ? ELSE finished_at END
985+ WHERE workspace = ?`,
986+ )
987+ .bind(INDEXING_PAUSED, now(), job.workspace)
988+ .run();
989+ }
990+ return;
991+ }
971992 if (job.type === "backfill_memory") {
972993 const memories = await memoryReviewClient(this.env.WORK).searchMemories(job.workspace, null, { limit: 100 });
973994 const indexed = await this.indexMemories(job.workspace, memories);
+6−1
1717 * `capture.ts`).
1818 */
1919 import { WorkerEntrypoint } from "cloudflare:workers";
20−import { type ServiceBinding, identityClient, projectsClient, reposClient, workClient } from "@g1t/contracts";
20+import { type ServiceBinding, identityClient, platformPaused, projectsClient, reposClient, workClient } from "@g1t/contracts";
2121 import type { Font } from "satori/standalone";
2222 import resvgWasm from "@resvg/resvg-wasm/index_bg.wasm";
2323 import yogaWasm from "satori/yoga.wasm";
3333 import { cardPng } from "./render.ts";
3434 import { cacheKey } from "./cache.ts";
3535 import { type Shot, screenshotOf, sweep, take } from "./capture.ts";
36+import { pausedCard } from "./paused.ts";
3637 import { parseShot } from "./screenshot.ts";
3738 import { BRAND, type Card, docsCard, resolve } from "./resolve.ts";
3839
4344 REPOS: ServiceBinding;
4445 WORK: ServiceBinding;
4546 PROJECTS: ServiceBinding;
47+ /** Whether rendering is paused across g1t (billing's `platform_pause`). */
48+ BILLING?: ServiceBinding;
4649 /** Browser Rendering, for production screenshots. */
4750 BROWSER: Fetcher;
4851 /** Production screenshots, by app hostname. */
121124 const cache = (caches as unknown as { default: Cache }).default;
122125 const hit = await cache.match(key);
123126 if (hit) return hit;
127+ // Rendering paused across g1t: no card is drawn (paused.ts).
128+ if (await platformPaused(env.BILLING, "renders")) return pausedCard(await cache.match(cacheKey(new URL("/", url))));
124129
125130 const card = await cardFor(url, env);
126131 const failed = card.kind === "brand" && card.failed === true;
+41−0
1+import assert from "node:assert/strict";
2+import { test } from "node:test";
3+
4+import { PAUSE_KEPT_MS, forgetPlatformPause, platformPause, readPause } from "../../../packages/contracts/src/platform.ts";
5+import { STATIC_CARD, pausedCard } from "./paused.ts";
6+
7+test("while renders are paused, a miss gets the cached brand card, kept a minute", async () => {
8+ const brand = new Response("png", { headers: { "content-type": "image/png", "cache-control": "public, max-age=3600" } });
9+ const answer = pausedCard(brand);
10+ assert.equal(answer.status, 200);
11+ assert.equal(answer.headers.get("content-type"), "image/png");
12+ assert.equal(answer.headers.get("cache-control"), "public, max-age=60");
13+ assert.equal(await answer.text(), "png");
14+});
15+
16+test("with no brand card cached, a miss is sent to the static logo", () => {
17+ const answer = pausedCard(undefined);
18+ assert.equal(answer.status, 302);
19+ assert.equal(answer.headers.get("location"), STATIC_CARD);
20+});
21+
22+test("the platform pause is read once per 30 seconds, and nothing is paused when billing cannot say", async () => {
23+ forgetPlatformPause();
24+ let calls = 0;
25+ const billing = {
26+ async fetch() {
27+ calls += 1;
28+ return new Response(JSON.stringify({ compute: false, schedules: true, indexing: false, renders: true }));
29+ },
30+ };
31+ const now = 1_000_000;
32+ assert.equal((await platformPause(billing, now)).renders, true);
33+ assert.equal((await platformPause(billing, now + PAUSE_KEPT_MS - 1)).schedules, true);
34+ assert.equal(calls, 1);
35+ const down = { fetch: async () => new Response("no", { status: 500 }) };
36+ assert.deepEqual(await platformPause(down, now + PAUSE_KEPT_MS), { compute: false, schedules: false, indexing: false, renders: false });
37+ // No binding: never paused.
38+ assert.equal((await platformPause(undefined)).renders, false);
39+ assert.deepEqual(readPause({ renders: "yes", indexing: true }), { compute: false, schedules: false, indexing: true, renders: false });
40+ forgetPlatformPause();
41+});
+25−0
1+/**
2+ * While rendering is paused across g1t (billing's `platform_pause`,
3+ * `renders`; docs/SPEND-GUARDRAILS.md), no card is drawn: a request that
4+ * misses the edge cache gets the brand card the cache already has, or a
5+ * redirect to g1t's static logo. Either is kept only a minute, so cards
6+ * come back soon after rendering is resumed.
7+ */
8+
9+/** g1t's logo on dark, served by the site as a static file. */
10+export const STATIC_CARD = "https://g1t.sh/brand/g1t-logo-on-dark.png";
11+
12+const PAUSED_CACHE_CONTROL = "public, max-age=60";
13+
14+/** What a cache miss is answered with while renders are paused. */
15+export function pausedCard(brand: Response | undefined): Response {
16+ if (brand) {
17+ const response = new Response(brand.body, brand);
18+ response.headers.set("cache-control", PAUSED_CACHE_CONTROL);
19+ return response;
20+ }
21+ return new Response(null, {
22+ status: 302,
23+ headers: { location: STATIC_CARD, "cache-control": PAUSED_CACHE_CONTROL, "access-control-allow-origin": "*" },
24+ });
25+}
+4−1
3131 { "binding": "IDENTITY", "service": "g1t-identity" },
3232 { "binding": "REPOS", "service": "g1t-repos" },
3333 { "binding": "WORK", "service": "g1t-work" },
34− { "binding": "PROJECTS", "service": "g1t-projects" }
34+ { "binding": "PROJECTS", "service": "g1t-projects" },
35+ // Whether rendering is paused across g1t (billing's platform_pause):
36+ // cards fall back to a static image while it is (src/paused.ts).
37+ { "binding": "BILLING", "service": "g1t-billing" }
3538 ],
3639 "observability": { "enabled": true }
3740 }
+9−1
3737 agentEstimateMicros,
3838 eventsClient,
3939 isWaiting,
40+ platformPaused,
4041 issueCapReached,
4142 refusalMessage,
4243 sandboxEstimateMicros,
19621963 async scheduled(): Promise<void> {
19631964 await this.drainWaits();
19641965 await this.advanceAll();
1965− await this.startReady();
1966+ // Schedules paused across g1t (billing's platform_pause, kept 30
1967+ // seconds): the sweep starts no queued agents. Events still start
1968+ // them, through the compute gate, which holds while compute is paused.
1969+ if (await platformPaused(this.env.BILLING, "schedules")) {
1970+ console.log("sweep: schedules are paused across g1t, so no queued agents start");
1971+ } else {
1972+ await this.startReady();
1973+ }
19661974 await this.startBackups().catch((error: unknown) => console.log("backups not started", String(error)));
19671975 }
19681976
+57−0
5353 /// Accounts or workspaces per directory job.
5454 const DIRECTORY_PAGE: u32 = 200;
5555
56+/// Where `park` keeps a backfill page in `meta`, by its key's prefix.
57+const PARKED: &str = "parked_job:";
58+
59+/// The `meta` key a backfill page is parked under while indexing is
60+/// paused: one per kind of walk. None for jobs that are not a backfill's.
61+pub fn parked_key(job: &Job) -> Option<String> {
62+ match job {
63+ Job::Backfill { .. } => Some(format!("{PARKED}backfill")),
64+ Job::Directory { kind, .. } => Some(format!("{PARKED}directory:{kind}")),
65+ _ => None,
66+ }
67+}
68+
5669 /// Work queued for later, one job per message.
5770 #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
5871 #[serde(tag = "type", rename_all = "snake_case")]
734747 Ok(())
735748 }
736749
750+ /// A backfill's next page, set aside while indexing is paused across
751+ /// g1t (billing's `platform_pause`), to go on from where it was when it
752+ /// is resumed (`resume_parked`). Repositories already queued finish.
753+ pub async fn park(&self, job: &Job) -> Result<()> {
754+ let Some(key) = parked_key(job) else { return Ok(()) };
755+ let value = serde_json::to_string(job)?;
756+ store::run_all(
757+ &self.db,
758+ vec![store::prepare(
759+ &self.db,
760+ "INSERT INTO meta (key, value) VALUES (?1, ?2) ON CONFLICT (key) DO UPDATE SET value = ?2",
761+ vec![p(&key), p(&value)],
762+ )?],
763+ )
764+ .await
765+ }
766+
767+ /// Queues the backfill pages `park` set aside. How many it found.
768+ pub async fn resume_parked(&self) -> Result<usize> {
769+ #[derive(Deserialize)]
770+ struct Parked {
771+ value: String,
772+ }
773+ let parked = store::all::<Parked>(
774+ &self.db,
775+ &Sql { text: format!("DELETE FROM meta WHERE key LIKE '{PARKED}%' RETURNING value"), params: vec![] },
776+ )
777+ .await?;
778+ let jobs: Vec<Job> = parked.iter().filter_map(|row| serde_json::from_str(&row.value).ok()).collect();
779+ let found = jobs.len();
780+ if found > 0 {
781+ self.enqueue(jobs).await?;
782+ }
783+ Ok(found)
784+ }
785+
737786 pub async fn on_event(&self, event: &Event) -> Result<()> {
738787 let data = &event.data;
739788 match event.kind.as_str() {
942991 }
943992
944993 #[test]
994+ fn only_a_backfills_pages_are_parked_one_per_walk() {
995+ assert_eq!(parked_key(&Job::Backfill { after: Some("rep_9".into()) }).as_deref(), Some("parked_job:backfill"));
996+ assert_eq!(parked_key(&Job::Directory { kind: "user".into(), after: None }).as_deref(), Some("parked_job:directory:user"));
997+ assert_eq!(parked_key(&Job::Repo { repo_id: "rep_1".into() }), None);
998+ assert_eq!(parked_key(&Job::Drain { repo_id: "rep_1".into() }), None);
999+ }
1000+
1001+ #[test]
9451002 fn a_push_from_nothing_is_a_whole_comparison() {
9461003 assert!(is_null_commit("0000000000000000000000000000000000000000"));
9471004 assert!(is_null_commit(""));
+40−2
3232 use std::collections::HashMap;
3333
3434 use g1t_contracts::User;
35+use g1t_contracts::billing::PauseLevel;
3536 use g1t_contracts::events::Event;
3637 use g1t_kit::{args, reply, rpc_method};
3738 use worker::{Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event};
4142 /// The queue of this service's own jobs.
4243 const JOBS_QUEUE: &str = "g1t-search-jobs";
4344
45+/// How often an isolate looks for backfill pages parked while indexing
46+/// was paused.
47+const RESUME_EVERY_MS: u64 = 5 * 60 * 1000;
48+
49+thread_local! {
50+ static RESUME_CHECKED: std::cell::Cell<u64> = const { std::cell::Cell::new(0) };
51+}
52+
53+/// Whether to look for parked pages now; marks it looked for.
54+fn resume_due(now: u64) -> bool {
55+ RESUME_CHECKED.with(|last| {
56+ if now.saturating_sub(last.get()) < RESUME_EVERY_MS {
57+ return false;
58+ }
59+ last.set(now);
60+ true
61+ })
62+}
63+
4464 pub struct Search {
4565 db: D1Database,
4666 env: Env,
89109 async fn queue(batch: MessageBatch<serde_json::Value>, env: Env, _ctx: Context) -> Result<()> {
90110 let search = Search::new(env)?;
91111 let jobs = batch.queue() == JOBS_QUEUE;
92− if !jobs && let Err(error) = search.ensure_backfill().await {
93− worker::console_error!("search: could not start the backfill: {error}");
112+ // Indexing paused across g1t (billing's `platform_pause`, kept 30
113+ // seconds in the isolate): backfills wait, and events still keep the
114+ // index current. Without a billing binding nothing is ever paused.
115+ let paused = match search.env.service("BILLING") {
116+ Ok(billing) => g1t_kit::pause::paused(&billing, PauseLevel::Indexing).await,
117+ Err(_) => false,
118+ };
119+ if !jobs && !paused {
120+ if let Err(error) = search.ensure_backfill().await {
121+ worker::console_error!("search: could not start the backfill: {error}");
122+ }
123+ // Pages parked while paused, looked for at most every few minutes.
124+ if resume_due(g1t_kit::now_ms()) {
125+ match search.resume_parked().await {
126+ Ok(0) => {}
127+ Ok(found) => worker::console_log!("search: resumed {found} parked backfill pages"),
128+ Err(error) => worker::console_error!("search: could not resume parked backfill pages: {error}"),
129+ }
130+ }
94131 }
95132 for message in batch.messages()? {
96133 let body = message.body().clone();
97134 let outcome = if jobs {
98135 match serde_json::from_value::<Job>(body) {
136+ Ok(job) if paused && index::parked_key(&job).is_some() => search.park(&job).await,
99137 Ok(job) => search.run_job(job).await,
100138 Err(error) => {
101139 worker::console_error!("search: a job could not be read: {error}");
+4−1
2222 "services": [
2323 { "binding": "IDENTITY", "service": "g1t-identity" },
2424 { "binding": "REPOS", "service": "g1t-repos" },
25− { "binding": "WORK", "service": "g1t-work" }
25+ { "binding": "WORK", "service": "g1t-work" },
26+ // Whether indexing is paused across g1t (billing's platform_pause):
27+ // backfills wait while it is. Read at most every 30 seconds.
28+ { "binding": "BILLING", "service": "g1t-billing" }
2629 ],
2730 "queues": {
2831 // Its own jobs: a push's files, a repository compared whole, a page