Skip to content

Commit

Spend guardrails documented, with the owner's dashboard checklist, and scripts/ops/platform-usage.mjs prints the month and the last 24 hours per script, queue, database and namespace

syntaqxcommitted Parent78542f6Browse files
4 files+687−00/4 viewed
+1−0
1313 | The runner's images | `services/runner/base/Dockerfile`, `services/runner/Dockerfile`, `services/runner/base.json`, `scripts/build-runner.mjs`, `scripts/deploy/image.mjs` |
1414 | The workflows | `.g1t/workflows/deploy.yml`, `.g1t/workflows/runner-base.yml` |
1515 | The old entry point | `scripts/deploy.sh`, now a wrapper |
16+| Spend guardrails (platform pause, hourly usage watch) | [SPEND-GUARDRAILS.md](SPEND-GUARDRAILS.md) |
1617
1718 ## The manifest
1819
+196−0
1+# Spend guardrails
2+
3+Cloudflare has no hard spending cap. A loop in a Worker, a queue that
4+retries forever or a cron that lists KV every second shows up on the bill
5+weeks later unless something here catches it. This is what catches it,
6+how to stop it, and what the owner has to set up by hand in Cloudflare's
7+dashboard. Internal: the public side is the "When g1t pauses work for
8+everyone" section of `apps/docs/src/content/docs/guides/usage-and-billing.md`
9+and "How they are enforced" in `guides/guardrails.md`.
10+
11+## What catches what
12+
13+| Guardrail | Catches | Where |
14+| --- | --- | --- |
15+| A workspace's limits, caps and spike pause | One workspace's agents, sandboxes and builds running away | `services/billing/src/limits.rs`, `compute.rs`; reserved through `ComputeGate.admit` (`packages/contracts/src/compute.ts`) |
16+| The comped budget and the daily breaker | g1t's own spend on agents: comped accounts, trials, pools; $75 a day pauses hosted-model agent runs g1t pays for | `services/billing/src/budget.rs`; [BILLING_OPERATIONS.md](BILLING_OPERATIONS.md) |
17+| The daily reconciliation | What Cloudflare billed against what g1t counted, a day later | `services/billing/src/costs.rs`, `margin.rs` |
18+| **The hourly platform watch** | Platform cost no workspace's limit covers: Workers requests and CPU, D1 rows, Queue operations, Durable Objects, KV, Artifacts. Within the hour. | `services/billing/src/platform.rs` |
19+| **The platform pause** | A staff (or automatic) brake on whole kinds of work across g1t | `platform.rs`; read by every service through `g1t_kit::pause` (Rust) or `platformPaused` (`packages/contracts/src/platform.ts`) |
20+| **The models proxy's run cap** | An agent spending past its run's cap by calling the model proxy itself (a prompt-injected `curl`), around the sandbox's `--max-budget-usd` | `services/models/src/spend.ts`, `run-spend.ts` |
21+| Cloudflare's own notifications | Anything else, by product, as a last line | Set up by hand: [the checklist](#manual-steps-in-cloudflares-dashboard) |
22+
23+## The hourly platform watch
24+
25+At a quarter past each hour (billing's `*/15` cron, gated to the tick at
26+`:15`), billing reads two windows from Cloudflare's GraphQL Analytics API:
27+the hour before, and the month so far. It uses the costs reconciliation's
28+token, `CLOUDFLARE_BILLING_TOKEN` (else `CLOUDFLARE_USAGE_TOKEN`), which
29+needs **Account Analytics Read**. Without either it reads nothing and sudo
30+says so; the pause still works.
31+
32+| Metric | Dataset and field | Named by |
33+| --- | --- | --- |
34+| `workers_requests` | `workersInvocationsAdaptive` `sum.requests` | script |
35+| `workers_cpu_ms` | `workersInvocationsAdaptive` `sum.cpuTimeUs` / 1000 | script |
36+| `d1_rows_read`, `d1_rows_written` | `d1AnalyticsAdaptiveGroups` `sum.rowsRead`, `sum.rowsWritten` | database id |
37+| `queue_operations` | `queueMessageOperationsAdaptiveGroups` `sum.billableOperations` | queue id |
38+| `do_requests` | `durableObjectsInvocationsAdaptiveGroups` `sum.requests` | script |
39+| `do_active_seconds`, `do_storage_write_units` | `durableObjectsPeriodicGroups` `sum.activeTime` (µs), `sum.storageWriteUnits` | namespace id |
40+| `do_rows_written` | `durableObjectsPeriodicGroups` `sum.rowsWritten` | namespace id |
41+| `kv_reads`, `kv_writes`, `kv_deletes`, `kv_lists` | `kvOperationsAdaptiveGroups` `sum.requests` by `actionType` | namespace id |
42+| `artifacts_events` | `artifactsEventsAdaptiveGroups` `count` | repository |
43+
44+Each metric whose field is less certain has a query of its own, so a
45+dataset or field GraphQL refuses is logged (`platform watch: skipped …`)
46+and the others still count. Workers Logs has no dataset here; its volume
47+follows Workers requests, and log sampling is set per Worker.
48+
49+Each hour is kept in billing's D1 (`platform_usage`, migration 0050) with
50+the script, queue, database or namespace that counted most; the month so
51+far in `platform_usage_month`; breaches in `platform_alerts`.
52+
53+### Thresholds
54+
55+Hourly, in `services/billing/wrangler.jsonc` `vars`. Each is about a dollar
56+to a few dollars an hour at list prices, far above a small alpha's normal
57+hour. `0` turns a metric's threshold (and its spike rule) off.
58+
59+| Variable | Default | About |
60+| --- | --- | --- |
61+| `PLATFORM_HOURLY_WORKERS_REQUESTS` | 20,000,000 | $6/hour |
62+| `PLATFORM_HOURLY_WORKERS_CPU_MS` | 100,000,000 | $2/hour |
63+| `PLATFORM_HOURLY_D1_ROWS_READ` | 2,000,000,000 | $2/hour |
64+| `PLATFORM_HOURLY_D1_ROWS_WRITTEN` | 5,000,000 | $5/hour |
65+| `PLATFORM_HOURLY_QUEUE_OPERATIONS` | 5,000,000 | $2/hour |
66+| `PLATFORM_HOURLY_DO_REQUESTS` | 20,000,000 | $3/hour |
67+| `PLATFORM_HOURLY_DO_ROWS_WRITTEN` | 5,000,000 | $5/hour |
68+| `PLATFORM_HOURLY_DO_STORAGE_WRITE_UNITS` | 5,000,000 | $5/hour |
69+| `PLATFORM_HOURLY_DO_ACTIVE_SECONDS` | 3,000,000 | about $5/hour at 128 MB |
70+| `PLATFORM_HOURLY_KV_READS` | 10,000,000 | $5/hour |
71+| `PLATFORM_HOURLY_KV_WRITES` | 200,000 | $1/hour |
72+| `PLATFORM_HOURLY_KV_DELETES` | 200,000 | $1/hour |
73+| `PLATFORM_HOURLY_KV_LISTS` | 200,000 | $1/hour |
74+| `PLATFORM_HOURLY_ARTIFACTS_EVENTS` | 1,000,000 | |
75+
76+The rules, in order:
77+
78+1. **Threshold**: an hour over its threshold is a breach.
79+2. **Spike**: otherwise, an hour over `PLATFORM_SPIKE_FACTOR` (10) times the
80+ median hour of the week before, and at least `PLATFORM_SPIKE_FLOOR_PERCENT`
81+ (10%) of its threshold, is a breach. It needs a day of history first.
82+3. **Severe**: a threshold breach at `PLATFORM_SEVERE_FACTOR` (5) times the
83+ threshold or more pauses the levels that metric feeds, of those
84+ `AUTO_PAUSE` names (default `schedules,indexing`; `compute` and `renders`
85+ are off unless added). Spikes never pause.
86+
87+| Metric | Feeds |
88+| --- | --- |
89+| Workers | schedules, indexing, renders |
90+| D1, Queues | schedules, indexing |
91+| Durable Objects, Artifacts | compute, schedules |
92+| KV | indexing, renders |
93+
94+On a breach, staff are emailed at `COSTS_ALERT_EMAIL` through billing's
95+`EMAIL` binding, once per metric every 6 hours (a new automatic pause is
96+always emailed), and sudo's Costs & margin shows it under **Platform
97+pause**. The alert names the top script, queue, database or namespace.
98+Ids are Cloudflare's; `node scripts/ops/platform-usage.mjs` names them.
99+
100+To change a threshold: edit the variable, push, and billing redeploys. The
101+next hour reads with it.
102+
103+### Looking by hand
104+
105+```sh
106+node scripts/ops/platform-usage.mjs # month so far and the last 24 hours
107+node scripts/ops/platform-usage.mjs --json # the same, as JSON
108+```
109+
110+It needs `CLOUDFLARE_API_TOKEN` (or `CLOUDFLARE_API_KEY` and `CLOUDFLARE_EMAIL`) with Account
111+Analytics Read, as `scripts/deploy/cloudflare.mjs` `cloudflareAuth` reads them, and names ids
112+with the D1, Queues, KV and Durable Objects listings when the token can
113+read them.
114+
115+## The platform pause
116+
117+Four levels, each independent, kept in billing's `platform_pause` table.
118+
119+| Level | Stops | Where it is checked |
120+| --- | --- | --- |
121+| `compute` | Every reservation through billing's `reserve` except embeddings: agent runs, checks, the merge queue, workflow jobs on g1t's machines, deploy builds. For every workspace and plan, and where payments are not set up. Refused with `paused`. | `services/billing/src/compute.rs` `reserve` |
122+| `schedules` | Actions' cron-triggered runs (skipped, not run late), and the runner's sweep that starts queued agents (they stay queued) | `services/actions/src/plan.rs` `on_minute`; `services/runner/src/index.ts` `scheduled` |
123+| `indexing` | Embeddings (refused in `reserve`), context backfills (**Rebuild** refused, queued jobs closed with a note), search backfills (pages parked in search's `meta` and resumed where they were) | `reserve`; `services/context/src/index.ts`; `services/search/src/lib.rs` |
124+| `renders` | Social cards: a cache miss gets the cached brand card or a redirect to `https://g1t.sh/brand/g1t-logo-on-dark.png`, kept a minute | `services/og/src/index.ts`, `paused.ts` |
125+
126+Work already running finishes at every level.
127+
128+**Reads are cheap.** Every caller keeps the flags 30 seconds in its isolate:
129+billing itself (`pause_now`), Rust services through `g1t_kit::pause`, and
130+TypeScript services through `platformPaused`. A change reaches everything
131+within about 30 seconds, and there is never a D1 read per request.
132+
133+**When the flag cannot be read, nothing is paused** (fail open), and the
134+failure is kept for the same 30 seconds. A pause is a brake someone pulls
135+on purpose; failing closed would turn a billing outage into a platform
136+outage. The other guardrails (workspace limits, the breaker, the model
137+proxy's run cap) do not depend on it.
138+
139+### Pausing and resuming
140+
141+1. Open sudo, **Costs & margin**, **Platform pause** (`https://sudo.g1t.sh/costs#platform`).
142+2. On the level's card, write why, and choose **Pause** or **Resume**.
143+
144+Every change is in the audit log (`platform_paused`, `platform_resumed`),
145+with who and why. While any level is paused, every sudo page shows a red
146+**Platform pause** bar. A level the usage watcher paused says so and stays
147+paused until staff resume it: fix or understand the cause first.
148+
149+Without sudo (billing's RPC, through a service binding):
150+`admin_set_pause` with `{ "level": "schedules", "paused": false, "note": "…", "by": "you@flagon.io" }`.
151+
152+## The models proxy's run cap
153+
154+A run's model cap was only enforced inside the sandbox
155+(`--max-budget-usd`). Now the proxy holds it too. When a run starts, the
156+runner gives its model session token (`g1tm_`) the run's cap
157+(`cap_model_sessions` on integrations). The proxy counts each answer's cost
158+from its token usage in one `RunSpend` Durable Object per session, so
159+requests fanned out across isolates cannot each spend the cap. At the cap it
160+answers `402` with `run_cap_reached`. At most 16 answers are counted in
161+flight at once, which bounds the overshoot. A session without a cap gets a
162+$100 backstop. A run token reaches only the message and model routes, and is
163+closed when the run ends.
164+
165+## Manual steps in Cloudflare's dashboard
166+
167+These are the owner's, once per account. None can be set from code.
168+
169+- [ ] **Billing → Billable Usage notifications** (Manage Account → Billing →
170+ Notifications, or Notifications → Add → "Usage Based Billing"). Add
171+ one per product g1t uses: Workers (requests and CPU), Workers KV, D1,
172+ Queues, Durable Objects, R2, Workers Logs, Containers, Browser
173+ Rendering, Workers AI and Vectorize. Set each threshold near this
174+ doc's hourly threshold times about 24 times 3 (a day at a third of the
175+ watch's line), and send them to hey@flagon.io.
176+- [ ] **Budget alerts** (Manage Account → Billing → Budget alerts, where the
177+ account has them): one for the whole account's monthly usage, at the
178+ month's expected bill and at twice it.
179+- [ ] **Notifications → Destinations**: add hey@flagon.io (and a webhook,
180+ if one is set up for paging) so the alerts above reach someone.
181+- [ ] **Account API token for billing**: check `CLOUDFLARE_BILLING_TOKEN`
182+ (or `CLOUDFLARE_USAGE_TOKEN`) has Account Analytics Read. sudo's
183+ Platform pause says when the watch cannot read.
184+- [ ] Where to see them: Notifications → History for what fired; Billing →
185+ Billable Usage for the month so far by product.
186+
187+## Deploy order for these changes
188+
189+1. Billing (migration 0050 runs first, then the Worker): the pause, the
190+ watch, `platform_pause`, `admin_platform_guard`, `admin_set_pause`.
191+2. Search and og, with their new `BILLING` binding; actions, context and
192+ runner. Before billing has `platform_pause`, they read nothing paused.
193+3. Sudo.
194+4. For the run cap: integrations (migration 0006), then models (its
195+ `RunSpend` Durable Object migration), then runner. Any order works; a
196+ session without a cap gets the $100 backstop.
+377−0
1+#!/usr/bin/env node
2+// What is the platform using on Cloudflare? Asks Cloudflare's GraphQL
3+// Analytics API, for the month so far (UTC, from the 1st) and the last 24
4+// hours, how much each Worker, D1 database, queue, Durable Object namespace,
5+// KV namespace and Artifacts namespace did, and prints the top of each with
6+// totals. Workers Logs are added when the account's schema has a dataset for
7+// them.
8+//
9+// Read-only: one GraphQL query per dataset and window, all at once, and a few
10+// REST listings to put names on database, queue and namespace ids. A dataset
11+// or field Cloudflare refuses is a note in the report, never the end of it.
12+//
13+// CLOUDFLARE_API_TOKEN=<token with Account Analytics: Read> \
14+// node scripts/ops/platform-usage.mjs [--top 10] [--json] [--account <id>]
15+//
16+// --account (or CLOUDFLARE_ACCOUNT_ID) reads another account than g1t's.
17+// --json prints one object with snake_case keys and every row, not just the top.
18+// Names come from the D1, Queues, KV and Durable Objects listings when the
19+// token may read them (D1: Read, Queues: Read, Workers KV Storage: Read,
20+// Workers Scripts: Read); otherwise the ids are printed.
21+
22+import { ACCOUNT_ID, cloudflareAuth } from "../deploy/cloudflare.mjs";
23+
24+const API = "https://api.cloudflare.com/client/v4";
25+
26+const HELP = `node scripts/ops/platform-usage.mjs [--top N] [--json] [--account <id>]
27+
28+Cloudflare usage per Worker, D1 database, queue, Durable Object namespace,
29+KV namespace and Artifacts namespace, month to date (UTC) and the last 24 hours.
30+
31+ --top N rows shown per dataset in the tables (default 10)
32+ --json one JSON object, snake_case keys, every row
33+ --account <id> the Cloudflare account (default CLOUDFLARE_ACCOUNT_ID, else g1t's)
34+ --help this text
35+
36+Needs CLOUDFLARE_API_TOKEN (Account Analytics: Read), or CLOUDFLARE_API_KEY
37+with CLOUDFLARE_EMAIL.`;
38+
39+/**
40+ * The datasets the report asks for. Each has one or more variants, tried in
41+ * order: when Cloudflare refuses a field, the next variant asks for less.
42+ * `names` says which REST listing puts a name on the first dimension's ids.
43+ */
44+export const DATASETS = [
45+ {
46+ key: "workers",
47+ label: "Workers invocations, by script",
48+ dataset: "workersInvocationsAdaptive",
49+ variants: [
50+ { sum: ["requests", "errors", "cpuTimeUs"], dims: ["scriptName"] },
51+ { sum: ["requests", "errors"], dims: ["scriptName"] },
52+ ],
53+ },
54+ {
55+ key: "d1",
56+ label: "D1 rows, by database",
57+ dataset: "d1AnalyticsAdaptiveGroups",
58+ names: "d1",
59+ variants: [
60+ { sum: ["rowsRead", "rowsWritten", "readQueries", "writeQueries"], dims: ["databaseId"] },
61+ { sum: ["rowsRead", "rowsWritten"], dims: ["databaseId"] },
62+ ],
63+ },
64+ {
65+ key: "queues",
66+ label: "Queue operations, by queue",
67+ dataset: "queueMessageOperationsAdaptiveGroups",
68+ names: "queues",
69+ variants: [{ sum: ["billableOperations"], dims: ["queueId"] }],
70+ },
71+ {
72+ key: "durable_objects",
73+ label: "Durable Object requests, by script",
74+ dataset: "durableObjectsInvocationsAdaptiveGroups",
75+ variants: [
76+ { sum: ["requests", "errors"], dims: ["scriptName"] },
77+ { sum: ["requests"], dims: ["scriptName"] },
78+ ],
79+ },
80+ {
81+ key: "durable_objects_periodic",
82+ label: "Durable Object time and storage, by namespace",
83+ dataset: "durableObjectsPeriodicGroups",
84+ names: "durable_objects",
85+ variants: [
86+ { sum: ["activeTime", "cpuTime", "storageReadUnits", "storageWriteUnits"], dims: ["namespaceId"] },
87+ { sum: ["activeTime", "storageWriteUnits"], dims: ["namespaceId"] },
88+ // Periodic groups may filter by the minute rather than by datetime.
89+ { sum: ["activeTime"], dims: ["namespaceId"], time: "datetimeMinute" },
90+ ],
91+ },
92+ {
93+ key: "kv",
94+ label: "KV operations, by namespace and action",
95+ dataset: "kvOperationsAdaptiveGroups",
96+ names: "kv",
97+ variants: [{ sum: ["requests"], dims: ["namespaceId", "actionType"] }],
98+ },
99+ {
100+ key: "artifacts",
101+ label: "Artifacts events, by namespace and type",
102+ dataset: "artifactsEventsAdaptiveGroups",
103+ variants: [{ count: true, sum: ["durationMs"], dims: ["repositoryNamespace", "eventType"] }],
104+ },
105+ {
106+ key: "workers_logs",
107+ label: "Workers Logs events, by script",
108+ // Which dataset holds Workers Logs is read from the schema (logsDataset).
109+ dataset: null,
110+ optional: true,
111+ variants: [
112+ { count: true, dims: ["scriptName"] },
113+ { count: true, dims: [] },
114+ ],
115+ },
116+];
117+
118+/** The two windows: the UTC month so far, and the last 24 hours. */
119+export function windows(now = new Date()) {
120+ const end = new Date(now.getTime());
121+ const monthStart = new Date(Date.UTC(end.getUTCFullYear(), end.getUTCMonth(), 1));
122+ return [
123+ { key: "month_to_date", label: "Month to date (UTC)", start: monthStart.toISOString(), end: end.toISOString() },
124+ { key: "last_24h", label: "Last 24 hours", start: new Date(end.getTime() - 24 * 3600 * 1000).toISOString(), end: end.toISOString() },
125+ ];
126+}
127+
128+/** Which account field holds Workers Logs, from the account type's field names; null when none does. */
129+export function logsDataset(fieldNames) {
130+ const known = ["workersObservabilityEventsAdaptiveGroups", "workersLogsEventsAdaptiveGroups", "workersLogsAdaptiveGroups"];
131+ return known.find((name) => fieldNames.includes(name)) ?? fieldNames.find((name) => /^workers.*(logs|observability).*groups$/i.test(name)) ?? null;
132+}
133+
134+/** One dataset's query, for one variant. The rows come back under `rows`. */
135+export function buildQuery(dataset, variant) {
136+ const time = variant.time ?? "datetime";
137+ const fields = [
138+ variant.count ? "count" : "",
139+ variant.sum?.length ? `sum { ${variant.sum.join(" ")} }` : "",
140+ variant.dims.length ? `dimensions { ${variant.dims.join(" ")} }` : "",
141+ ].filter(Boolean);
142+ return `query PlatformUsage($accountTag: String!, $start: Time!, $end: Time!) {
143+ viewer {
144+ accounts(filter: { accountTag: $accountTag }) {
145+ rows: ${dataset}(limit: 10000, filter: { ${time}_geq: $start, ${time}_leq: $end }) {
146+ ${fields.join("\n ")}
147+ }
148+ }
149+ }
150+}`;
151+}
152+
153+/** camelCase to snake_case, for the JSON report's keys. */
154+export const snake = (name) => name.replace(/[A-Z]/g, (letter) => `_${letter.toLowerCase()}`);
155+
156+/** The metric names a variant reports, snake_case: `count` first, then its sums. */
157+export const metricsOf = (variant) => [...(variant.count ? ["count"] : []), ...(variant.sum ?? []).map(snake)];
158+
159+/**
160+ * A GraphQL answer to one dataset's query as rows, one per name, summed and
161+ * sorted by the first metric (largest first), with totals per metric. Throws
162+ * with Cloudflare's message when the answer has errors or no such dataset.
163+ */
164+export function parseGroups(body, variant, names = {}) {
165+ if (body?.errors?.length) throw new Error(body.errors.map((error) => error.message).join("; ").slice(0, 400));
166+ const groups = body?.data?.viewer?.accounts?.[0]?.rows;
167+ if (!Array.isArray(groups)) throw new Error("no rows in the answer");
168+ const metrics = metricsOf(variant);
169+ const byName = new Map();
170+ for (const group of groups) {
171+ const parts = variant.dims.map((dim, at) => {
172+ const value = group.dimensions?.[dim];
173+ const text = value == null || value === "" ? "(none)" : String(value);
174+ return at === 0 ? (names[text] ?? text) : text;
175+ });
176+ const name = parts.join(" / ") || "(all)";
177+ const row = byName.get(name) ?? { name, ...Object.fromEntries(metrics.map((metric) => [metric, 0])) };
178+ if (variant.count) row.count += Number(group.count ?? 0);
179+ for (const field of variant.sum ?? []) row[snake(field)] += Number(group.sum?.[field] ?? 0);
180+ byName.set(name, row);
181+ }
182+ const rows = sortRows([...byName.values()], metrics[0]);
183+ return { metrics, rows, totals: totalsOf(rows, metrics) };
184+}
185+
186+/** Rows largest first by one metric, ties by name. */
187+export function sortRows(rows, metric) {
188+ return [...rows].sort((a, b) => (b[metric] ?? 0) - (a[metric] ?? 0) || a.name.localeCompare(b.name));
189+}
190+
191+/** Each metric summed over all rows. */
192+export function totalsOf(rows, metrics) {
193+ return Object.fromEntries(metrics.map((metric) => [metric, rows.reduce((total, row) => total + (row[metric] ?? 0), 0)]));
194+}
195+
196+/**
197+ * The report from every dataset's outcome in every window. An outcome is
198+ * `{ body, variant }` (a GraphQL answer) or `{ error }` or `{ skipped }`;
199+ * whatever cannot be read becomes a note and the other datasets still count.
200+ */
201+export function assemble({ account, now, windowList, outcomes, names = {} }) {
202+ return {
203+ account,
204+ generated_at: now.toISOString(),
205+ windows: windowList.map((window) => {
206+ const notes = [];
207+ const datasets = [];
208+ for (const spec of DATASETS) {
209+ const outcome = outcomes[window.key]?.[spec.key];
210+ if (!outcome) continue;
211+ if (outcome.skipped) {
212+ notes.push(`${spec.label}: ${outcome.skipped}`);
213+ continue;
214+ }
215+ try {
216+ if (outcome.error) throw outcome.error;
217+ const parsed = parseGroups(outcome.body, outcome.variant, names[spec.names] ?? {});
218+ datasets.push({ key: spec.key, label: spec.label, dataset: outcome.dataset ?? spec.dataset, ...parsed });
219+ } catch (error) {
220+ notes.push(`${spec.label} (${outcome.dataset ?? spec.dataset ?? "no dataset"}): ${String(error.message ?? error).split("\n")[0]}`);
221+ }
222+ }
223+ return { key: window.key, label: window.label, start: window.start, end: window.end, datasets, notes };
224+ }),
225+ };
226+}
227+
228+const number = (value) => (Number.isInteger(value) ? value.toLocaleString("en-US") : value.toLocaleString("en-US", { maximumFractionDigits: 2 }));
229+
230+/** The report as plain tables: the top rows of each dataset, and a totals line. */
231+export function format(report, top = 10) {
232+ const lines = [`Cloudflare usage for account ${report.account}, ${report.generated_at}`];
233+ for (const window of report.windows) {
234+ lines.push("", `== ${window.label}: ${window.start} to ${window.end}`);
235+ for (const set of window.datasets) {
236+ lines.push("", `${set.label} (${set.dataset})`);
237+ const shown = set.rows.slice(0, top);
238+ const cells = [["name", ...set.metrics], ...shown.map((row) => [row.name, ...set.metrics.map((m) => number(row[m]))]), ["total", ...set.metrics.map((m) => number(set.totals[m]))]];
239+ const widths = cells[0].map((_, at) => Math.max(...cells.map((row) => String(row[at]).length)));
240+ const render = (row) => " " + row.map((cell, at) => (at === 0 ? String(cell).padEnd(widths[at]) : String(cell).padStart(widths[at]))).join(" ");
241+ lines.push(render(cells[0]));
242+ if (!shown.length) lines.push(" (nothing in this window)");
243+ for (const row of cells.slice(1, -1)) lines.push(render(row));
244+ if (set.rows.length > shown.length) lines.push(` ... and ${set.rows.length - shown.length} more`);
245+ lines.push(render(cells.at(-1)));
246+ }
247+ if (window.notes.length) {
248+ lines.push("", "Notes:");
249+ for (const note of window.notes) lines.push(` ${note}`);
250+ }
251+ }
252+ return lines.join("\n");
253+}
254+
255+/** The command line: flags and the account. */
256+export function parseArgs(argv, env = process.env) {
257+ const option = (name) => {
258+ const at = argv.indexOf(name);
259+ return at >= 0 && argv[at + 1] && !argv[at + 1].startsWith("--") ? argv[at + 1] : null;
260+ };
261+ return {
262+ help: argv.includes("--help") || argv.includes("-h"),
263+ json: argv.includes("--json"),
264+ top: Math.max(1, Number(option("--top")) || 10),
265+ account: option("--account") || env.CLOUDFLARE_ACCOUNT_ID || ACCOUNT_ID,
266+ };
267+}
268+
269+async function graphql(auth, account, query, variables = {}) {
270+ const response = await fetch(`${API}/graphql`, {
271+ method: "POST",
272+ headers: { ...auth, "content-type": "application/json", "user-agent": "g1t-ops" },
273+ body: JSON.stringify({ query, variables: { accountTag: account, ...variables } }),
274+ });
275+ const body = await response.json().catch(() => ({ errors: [{ message: `HTTP ${response.status}, not JSON` }] }));
276+ if (!response.ok && !body.errors?.length) body.errors = [{ message: `HTTP ${response.status}` }];
277+ return body;
278+}
279+
280+/** The account type's field names, to skip datasets the schema lacks; null when it cannot be read. */
281+async function schemaFields(auth, account) {
282+ try {
283+ const body = await graphql(auth, account, `{ __type(name: "account") { fields { name } } }`);
284+ const fields = body.data?.__type?.fields;
285+ return Array.isArray(fields) ? fields.map((field) => field.name) : null;
286+ } catch {
287+ return null;
288+ }
289+}
290+
291+/** One dataset in one window: each variant in turn until one is answered. */
292+async function ask(auth, account, dataset, spec, window) {
293+ let last = null;
294+ for (const variant of spec.variants) {
295+ try {
296+ const body = await graphql(auth, account, buildQuery(dataset, variant), { start: window.start, end: window.end });
297+ if (!body.errors?.length) return { body, variant, dataset };
298+ last = { body, variant, dataset };
299+ } catch (error) {
300+ last = { error, dataset };
301+ }
302+ }
303+ return last;
304+}
305+
306+/** A REST listing as id to name; empty when the token may not read it. */
307+async function listing(auth, path, id, name) {
308+ const out = {};
309+ try {
310+ for (let page = 1; page <= 10; page++) {
311+ const response = await fetch(`${API}${path}${path.includes("?") ? "&" : "?"}per_page=100&page=${page}`, { headers: { ...auth, "user-agent": "g1t-ops" } });
312+ const body = await response.json();
313+ if (!response.ok || !Array.isArray(body.result)) break;
314+ for (const item of body.result) if (item[id]) out[item[id]] = item[name] ?? item[id];
315+ if (body.result.length < 100) break;
316+ }
317+ } catch {
318+ // Ids stand in for names.
319+ }
320+ return out;
321+}
322+
323+async function main() {
324+ const options = parseArgs(process.argv.slice(2));
325+ if (options.help) {
326+ console.log(HELP);
327+ return 0;
328+ }
329+ const auth = cloudflareAuth();
330+ if (!auth) {
331+ console.error(`Set CLOUDFLARE_API_TOKEN to a token with Account Analytics: Read on account ${options.account}, or CLOUDFLARE_API_KEY and CLOUDFLARE_EMAIL.`);
332+ return 2;
333+ }
334+ const now = new Date();
335+ const windowList = windows(now);
336+ const account = options.account;
337+ const base = `/accounts/${account}`;
338+ const [fields, d1, queues, kv, durable] = await Promise.all([
339+ schemaFields(auth, account),
340+ listing(auth, `${base}/d1/database`, "uuid", "name"),
341+ listing(auth, `${base}/queues`, "queue_id", "queue_name"),
342+ listing(auth, `${base}/storage/kv/namespaces`, "id", "title"),
343+ listing(auth, `${base}/workers/durable_objects/namespaces`, "id", "name"),
344+ ]);
345+ const names = { d1, queues, kv, durable_objects: durable };
346+ const outcomes = {};
347+ const jobs = [];
348+ for (const window of windowList) {
349+ outcomes[window.key] = {};
350+ for (const spec of DATASETS) {
351+ const dataset = spec.dataset ?? (fields ? logsDataset(fields) : null);
352+ if (!dataset) {
353+ if (!spec.optional) outcomes[window.key][spec.key] = { skipped: "no dataset" };
354+ continue;
355+ }
356+ if (fields && !fields.includes(dataset)) {
357+ if (!spec.optional) outcomes[window.key][spec.key] = { skipped: `${dataset} is not in this account's schema` };
358+ continue;
359+ }
360+ jobs.push(ask(auth, account, dataset, spec, window).then((outcome) => (outcomes[window.key][spec.key] = outcome)));
361+ }
362+ }
363+ await Promise.all(jobs);
364+ const report = assemble({ account, now, windowList, outcomes, names });
365+ console.log(options.json ? JSON.stringify(report, null, 2) : format(report, options.top));
366+ return 0;
367+}
368+
369+if (process.argv[1]?.replaceAll("\\", "/").endsWith("scripts/ops/platform-usage.mjs")) {
370+ main().then(
371+ (code) => process.exit(code),
372+ (error) => {
373+ console.error(`platform-usage: ${error.message}`);
374+ process.exit(1);
375+ },
376+ );
377+}
+113−0
1+// The platform usage report, from GraphQL answers as Cloudflare gives them,
2+// without the network.
3+
4+import assert from "node:assert/strict";
5+import { test } from "node:test";
6+
7+import { DATASETS, assemble, buildQuery, format, logsDataset, parseArgs, parseGroups, snake, windows } from "./platform-usage.mjs";
8+
9+const NOW = new Date("2026-10-08T15:30:00.000Z");
10+const answer = (rows) => ({ data: { viewer: { accounts: [{ rows }] } } });
11+const variantOf = (key, at = 0) => DATASETS.find((spec) => spec.key === key).variants[at];
12+
13+test("the windows are the UTC month so far and the last 24 hours", () => {
14+ const [month, day] = windows(NOW);
15+ assert.deepEqual(month, { key: "month_to_date", label: "Month to date (UTC)", start: "2026-10-01T00:00:00.000Z", end: "2026-10-08T15:30:00.000Z" });
16+ assert.equal(day.key, "last_24h");
17+ assert.equal(day.start, "2026-10-07T15:30:00.000Z");
18+ assert.equal(day.end, NOW.toISOString());
19+ // Just after midnight on the 1st, the month has only begun.
20+ assert.equal(windows(new Date("2026-11-01T00:05:00Z"))[0].start, "2026-11-01T00:00:00.000Z");
21+});
22+
23+test("a query asks for the variant's sums and dimensions under one alias", () => {
24+ const query = buildQuery("kvOperationsAdaptiveGroups", variantOf("kv"));
25+ assert.match(query, /rows: kvOperationsAdaptiveGroups\(limit: 10000, filter: \{ datetime_geq: \$start, datetime_leq: \$end \}\)/);
26+ assert.match(query, /sum \{ requests \}/);
27+ assert.match(query, /dimensions \{ namespaceId actionType \}/);
28+ assert.match(query, /accounts\(filter: \{ accountTag: \$accountTag \}\)/);
29+ const artifacts = buildQuery("artifactsEventsAdaptiveGroups", variantOf("artifacts"));
30+ assert.match(artifacts, /\bcount\b/);
31+ const minute = buildQuery("durableObjectsPeriodicGroups", variantOf("durable_objects_periodic", 2));
32+ assert.match(minute, /datetimeMinute_geq: \$start, datetimeMinute_leq: \$end/);
33+});
34+
35+test("Workers Logs are read from whichever dataset the schema has, or skipped", () => {
36+ assert.equal(logsDataset(["workersInvocationsAdaptive", "workersObservabilityEventsAdaptiveGroups"]), "workersObservabilityEventsAdaptiveGroups");
37+ assert.equal(logsDataset(["workersLogsSomethingGroups"]), "workersLogsSomethingGroups");
38+ assert.equal(logsDataset(["workersInvocationsAdaptive", "kvOperationsAdaptiveGroups"]), null);
39+});
40+
41+test("groups are summed per name, named from the listing, sorted largest first, with totals", () => {
42+ const body = answer([
43+ { sum: { rowsRead: 10, rowsWritten: 1, readQueries: 2, writeQueries: 1 }, dimensions: { databaseId: "aaa" } },
44+ { sum: { rowsRead: 500, rowsWritten: 20, readQueries: 9, writeQueries: 3 }, dimensions: { databaseId: "bbb" } },
45+ { sum: { rowsRead: 5, rowsWritten: 0, readQueries: 1, writeQueries: 0 }, dimensions: { databaseId: "aaa" } },
46+ ]);
47+ const parsed = parseGroups(body, variantOf("d1"), { bbb: "g1t-repos" });
48+ assert.deepEqual(parsed.metrics, ["rows_read", "rows_written", "read_queries", "write_queries"]);
49+ assert.deepEqual(parsed.rows.map((row) => [row.name, row.rows_read]), [["g1t-repos", 500], ["aaa", 15]]);
50+ assert.deepEqual(parsed.totals, { rows_read: 515, rows_written: 21, read_queries: 12, write_queries: 4 });
51+ const kv = parseGroups(
52+ answer([
53+ { sum: { requests: 3 }, dimensions: { namespaceId: "n1", actionType: "write" } },
54+ { sum: { requests: 40 }, dimensions: { namespaceId: "n1", actionType: "read" } },
55+ ]),
56+ variantOf("kv"),
57+ { n1: "SESSIONS" },
58+ );
59+ assert.deepEqual(kv.rows.map((row) => row.name), ["SESSIONS / read", "SESSIONS / write"]);
60+ assert.throws(() => parseGroups({ errors: [{ message: "unknown field cpuTimeUs" }] }, variantOf("workers")), /cpuTimeUs/);
61+ assert.equal(snake("billableOperations"), "billable_operations");
62+});
63+
64+test("a dataset that errors or is missing is a note, and the others still report", () => {
65+ const windowList = windows(NOW);
66+ const outcomes = {
67+ month_to_date: {
68+ workers: { body: answer([{ sum: { requests: 7, errors: 0, cpuTimeUs: 1200 }, dimensions: { scriptName: "web" } }]), variant: variantOf("workers"), dataset: "workersInvocationsAdaptive" },
69+ d1: { body: { errors: [{ message: "unknown field \"readQueries\"" }] }, variant: variantOf("d1"), dataset: "d1AnalyticsAdaptiveGroups" },
70+ queues: { skipped: "queueMessageOperationsAdaptiveGroups is not in this account's schema" },
71+ kv: { error: new Error("fetch failed") },
72+ artifacts: { body: answer([{ count: 4, sum: { durationMs: 80 }, dimensions: { repositoryNamespace: "g1t", eventType: "pull" } }]), variant: variantOf("artifacts"), dataset: "artifactsEventsAdaptiveGroups" },
73+ },
74+ last_24h: {},
75+ };
76+ const report = assemble({ account: "acct", now: NOW, windowList, outcomes });
77+ const month = report.windows[0];
78+ assert.deepEqual(month.datasets.map((set) => set.key), ["workers", "artifacts"]);
79+ assert.equal(month.notes.length, 3);
80+ assert.match(month.notes.join("\n"), /D1 rows.*readQueries/);
81+ assert.match(month.notes.join("\n"), /not in this account's schema/);
82+ assert.match(month.notes.join("\n"), /fetch failed/);
83+ assert.deepEqual(report.windows[1].datasets, []);
84+ const text = format(report, 10);
85+ assert.match(text, /Workers invocations, by script/);
86+ assert.match(text, /web\s+7\s+0\s+1,200/);
87+ assert.match(text, /Notes:/);
88+});
89+
90+test("the JSON report is one snake_case object with every row", () => {
91+ const windowList = windows(NOW);
92+ const rows = Array.from({ length: 15 }, (_, at) => ({ sum: { requests: at + 1 }, dimensions: { scriptName: `s${at}` } }));
93+ const outcomes = { month_to_date: { durable_objects: { body: answer(rows), variant: variantOf("durable_objects", 1), dataset: "durableObjectsInvocationsAdaptiveGroups" } }, last_24h: {} };
94+ const report = JSON.parse(JSON.stringify(assemble({ account: "acct", now: NOW, windowList, outcomes })));
95+ assert.deepEqual(Object.keys(report), ["account", "generated_at", "windows"]);
96+ assert.deepEqual(Object.keys(report.windows[0]), ["key", "label", "start", "end", "datasets", "notes"]);
97+ const set = report.windows[0].datasets[0];
98+ assert.deepEqual(Object.keys(set), ["key", "label", "dataset", "metrics", "rows", "totals"]);
99+ assert.equal(set.rows.length, 15);
100+ assert.equal(set.rows[0].name, "s14");
101+ assert.equal(set.totals.requests, 120);
102+ const keys = JSON.stringify(report).match(/"([^"]+)":/g).map((key) => key.slice(1, -2));
103+ for (const key of keys) assert.match(key, /^[a-z0-9_]+$/, key);
104+ // The tables show only the top.
105+ assert.match(format(assemble({ account: "acct", now: NOW, windowList, outcomes }), 10), /\.\.\. and 5 more/);
106+});
107+
108+test("the command line reads --json, --top and the account", () => {
109+ assert.deepEqual(parseArgs(["--json", "--top", "3", "--account", "abc"], {}), { help: false, json: true, top: 3, account: "abc" });
110+ assert.equal(parseArgs([], { CLOUDFLARE_ACCOUNT_ID: "env" }).account, "env");
111+ assert.equal(parseArgs(["--help"], {}).help, true);
112+ assert.equal(parseArgs([], {}).top, 10);
113+});