Skip to content
494 linesCodeBlameRaw

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

Merge platform pause and the hourly usage watcher: staff can pause compute, schedules, indexing or renders for everyone, the watcher emails on a breach and is never blind quietly, and the models proxy holds each run to its cap (billing 0051, integrations 0006)1#!/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// It is also the check of billing's hourly watcher (services/billing/src/
23// platform.rs): it asks Cloudflare billing's own queries, field for field,
24// over the last full hour, and exits 1 naming every dataset that errored,
25// fell back to fewer fields, or (all of them) answered with no rows. Run it
26// once after deploying a change to the watcher's queries.
27
28import { ACCOUNT_ID, cloudflareAuth } from "../deploy/cloudflare.mjs";
29
30const API = "https://api.cloudflare.com/client/v4";
31
32const HELP = `node scripts/ops/platform-usage.mjs [--top N] [--json] [--account <id>]
33
34Cloudflare usage per Worker, D1 database, queue, Durable Object namespace,
35KV namespace and Artifacts namespace, month to date (UTC) and the last 24 hours.
36
37 --top N rows shown per dataset in the tables (default 10)
38 --json one JSON object, snake_case keys, every row
39 --account <id> the Cloudflare account (default CLOUDFLARE_ACCOUNT_ID, else g1t's)
40 --help this text
41
42Needs CLOUDFLARE_API_TOKEN (Account Analytics: Read), or CLOUDFLARE_API_KEY
43with CLOUDFLARE_EMAIL.
44
45Exits 1 when any dataset errored or fell back to fewer fields, when any of
46billing's watcher queries errored, or when every dataset answered with no
47rows (the wrong account, or a token that cannot see it).`;
48
49/**
50 * Billing's hourly watcher queries, as `QUERIES` in
51 * services/billing/src/platform.rs has them (a test keeps the two the same):
52 * the dataset, what it selects, what it groups by, and the filter field for
53 * an hour.
54 */
55export const WATCHER_QUERIES = [
56 { key: "workers", dataset: "workersInvocationsAdaptive", select: "sum { requests }", dimensions: "scriptName", hourFilter: "datetime" },
57 { key: "workers_cpu", dataset: "workersInvocationsAdaptive", select: "sum { cpuTimeUs }", dimensions: "scriptName", hourFilter: "datetime" },
58 { key: "d1", dataset: "d1AnalyticsAdaptiveGroups", select: "sum { rowsRead rowsWritten }", dimensions: "databaseId", hourFilter: "datetimeHour" },
59 { key: "queues", dataset: "queueMessageOperationsAdaptiveGroups", select: "sum { billableOperations }", dimensions: "queueId", hourFilter: "datetime" },
60 { key: "do_invocations", dataset: "durableObjectsInvocationsAdaptiveGroups", select: "sum { requests }", dimensions: "scriptName", hourFilter: "datetime" },
61 { key: "do_periodic", dataset: "durableObjectsPeriodicGroups", select: "sum { activeTime storageWriteUnits }", dimensions: "namespaceId", hourFilter: "datetime" },
62 { key: "do_sql", dataset: "durableObjectsPeriodicGroups", select: "sum { rowsWritten }", dimensions: "namespaceId", hourFilter: "datetime" },
63 { key: "kv", dataset: "kvOperationsAdaptiveGroups", select: "sum { requests }", dimensions: "namespaceId actionType", hourFilter: "datetime" },
64 { key: "artifacts", dataset: "artifactsEventsAdaptiveGroups", select: "count", dimensions: "repositoryName", hourFilter: "datetime" },
65];
66
67/** The last full hour (UTC) before `now`, as the watcher reads it. */
68export function lastFullHour(now) {
69 const until = new Date(Math.floor(now.getTime() / 3_600_000) * 3_600_000);
70 const since = new Date(until.getTime() - 3_600_000);
71 const label = (date) => `${date.toISOString().slice(0, 13)}:00:00Z`;
72 return { since: label(since), until: label(until) };
73}
74
75/** One watcher query over an hour, as billing sends it (its account variable named as this script names it). */
76export function watcherQuery(query) {
77 return `query ($accountTag: String!, $since: Time!, $until: Time!) {
78 viewer { accounts(filter: { accountTag: $accountTag }) {
79 rows: ${query.dataset}(limit: 10000, filter: { ${query.hourFilter}_geq: $since, ${query.hourFilter}_lt: $until }) {
80 ${query.select}
81 dimensions { ${query.dimensions} }
82 }
83 } }
84}`;
85}
86
87/**
88 * Whether the schema checks out: every problem in the report's datasets and
89 * in the watcher's queries (`watcher`: key to `{ error }` or `{ rows }`).
90 * An empty list is a pass.
91 */
92export function verdict(report, watcher = {}) {
93 const problems = [];
94 for (const window of report.windows) {
95 for (const error of window.errors ?? []) problems.push(`${window.label}: ${error}`);
96 }
97 let answered = 0;
98 let rows = 0;
99 for (const query of WATCHER_QUERIES) {
100 const outcome = watcher[query.key];
101 if (!outcome) continue;
102 if (outcome.error) problems.push(`watcher query ${query.key} (${query.dataset}): ${outcome.error}`);
103 else {
104 answered += 1;
105 rows += outcome.rows;
106 }
107 }
108 const reportRows = report.windows.flatMap((w) => w.datasets).reduce((sum, set) => sum + set.rows.length, 0);
109 const reportAnswered = report.windows.flatMap((w) => w.datasets).length;
110 if (answered + reportAnswered > 0 && rows + reportRows === 0) {
111 problems.push("every dataset answered with no rows: likely the wrong account, or a token that cannot see its analytics");
112 }
113 return problems;
114}
115
116/**
117 * The datasets the report asks for. Each has one or more variants, tried in
118 * order: when Cloudflare refuses a field, the next variant asks for less.
119 * `names` says which REST listing puts a name on the first dimension's ids.
120 */
121export const DATASETS = [
122 {
123 key: "workers",
124 label: "Workers invocations, by script",
125 dataset: "workersInvocationsAdaptive",
126 variants: [
127 { sum: ["requests", "errors", "cpuTimeUs"], dims: ["scriptName"] },
128 { sum: ["requests", "errors"], dims: ["scriptName"] },
129 ],
130 },
131 {
132 key: "d1",
133 label: "D1 rows, by database",
134 dataset: "d1AnalyticsAdaptiveGroups",
135 names: "d1",
136 variants: [
137 { sum: ["rowsRead", "rowsWritten", "readQueries", "writeQueries"], dims: ["databaseId"] },
138 { sum: ["rowsRead", "rowsWritten"], dims: ["databaseId"] },
139 ],
140 },
141 {
142 key: "queues",
143 label: "Queue operations, by queue",
144 dataset: "queueMessageOperationsAdaptiveGroups",
145 names: "queues",
146 variants: [{ sum: ["billableOperations"], dims: ["queueId"] }],
147 },
148 {
149 key: "durable_objects",
150 label: "Durable Object requests, by script",
151 dataset: "durableObjectsInvocationsAdaptiveGroups",
152 variants: [
153 { sum: ["requests", "errors"], dims: ["scriptName"] },
154 { sum: ["requests"], dims: ["scriptName"] },
155 ],
156 },
157 {
158 key: "durable_objects_periodic",
159 label: "Durable Object time and storage, by namespace",
160 dataset: "durableObjectsPeriodicGroups",
161 names: "durable_objects",
162 variants: [
163 { sum: ["activeTime", "cpuTime", "storageReadUnits", "storageWriteUnits"], dims: ["namespaceId"] },
164 { sum: ["activeTime", "storageWriteUnits"], dims: ["namespaceId"] },
165 // Periodic groups may filter by the minute rather than by datetime.
166 { sum: ["activeTime"], dims: ["namespaceId"], time: "datetimeMinute" },
167 ],
168 },
169 {
170 key: "kv",
171 label: "KV operations, by namespace and action",
172 dataset: "kvOperationsAdaptiveGroups",
173 names: "kv",
174 variants: [{ sum: ["requests"], dims: ["namespaceId", "actionType"] }],
175 },
176 {
177 key: "artifacts",
178 label: "Artifacts events, by namespace and type",
179 dataset: "artifactsEventsAdaptiveGroups",
180 variants: [{ count: true, sum: ["durationMs"], dims: ["repositoryNamespace", "eventType"] }],
181 },
182 {
183 key: "workers_logs",
184 label: "Workers Logs events, by script",
185 // Which dataset holds Workers Logs is read from the schema (logsDataset).
186 dataset: null,
187 optional: true,
188 variants: [
189 { count: true, dims: ["scriptName"] },
190 { count: true, dims: [] },
191 ],
192 },
193];
194
195/** The two windows: the UTC month so far, and the last 24 hours. */
196export function windows(now = new Date()) {
197 const end = new Date(now.getTime());
198 const monthStart = new Date(Date.UTC(end.getUTCFullYear(), end.getUTCMonth(), 1));
199 return [
200 { key: "month_to_date", label: "Month to date (UTC)", start: monthStart.toISOString(), end: end.toISOString() },
201 { key: "last_24h", label: "Last 24 hours", start: new Date(end.getTime() - 24 * 3600 * 1000).toISOString(), end: end.toISOString() },
202 ];
203}
204
205/** Which account field holds Workers Logs, from the account type's field names; null when none does. */
206export function logsDataset(fieldNames) {
207 const known = ["workersObservabilityEventsAdaptiveGroups", "workersLogsEventsAdaptiveGroups", "workersLogsAdaptiveGroups"];
208 return known.find((name) => fieldNames.includes(name)) ?? fieldNames.find((name) => /^workers.*(logs|observability).*groups$/i.test(name)) ?? null;
209}
210
211/** One dataset's query, for one variant. The rows come back under `rows`. */
212export function buildQuery(dataset, variant) {
213 const time = variant.time ?? "datetime";
214 const fields = [
215 variant.count ? "count" : "",
216 variant.sum?.length ? `sum { ${variant.sum.join(" ")} }` : "",
217 variant.dims.length ? `dimensions { ${variant.dims.join(" ")} }` : "",
218 ].filter(Boolean);
219 return `query PlatformUsage($accountTag: String!, $start: Time!, $end: Time!) {
220 viewer {
221 accounts(filter: { accountTag: $accountTag }) {
222 rows: ${dataset}(limit: 10000, filter: { ${time}_geq: $start, ${time}_leq: $end }) {
223 ${fields.join("\n ")}
224 }
225 }
226 }
227}`;
228}
229
230/** camelCase to snake_case, for the JSON report's keys. */
231export const snake = (name) => name.replace(/[A-Z]/g, (letter) => `_${letter.toLowerCase()}`);
232
233/** The metric names a variant reports, snake_case: `count` first, then its sums. */
234export const metricsOf = (variant) => [...(variant.count ? ["count"] : []), ...(variant.sum ?? []).map(snake)];
235
236/**
237 * A GraphQL answer to one dataset's query as rows, one per name, summed and
238 * sorted by the first metric (largest first), with totals per metric. Throws
239 * with Cloudflare's message when the answer has errors or no such dataset.
240 */
241export function parseGroups(body, variant, names = {}) {
242 if (body?.errors?.length) throw new Error(body.errors.map((error) => error.message).join("; ").slice(0, 400));
243 const groups = body?.data?.viewer?.accounts?.[0]?.rows;
244 if (!Array.isArray(groups)) throw new Error("no rows in the answer");
245 const metrics = metricsOf(variant);
246 const byName = new Map();
247 for (const group of groups) {
248 const parts = variant.dims.map((dim, at) => {
249 const value = group.dimensions?.[dim];
250 const text = value == null || value === "" ? "(none)" : String(value);
251 return at === 0 ? (names[text] ?? text) : text;
252 });
253 const name = parts.join(" / ") || "(all)";
254 const row = byName.get(name) ?? { name, ...Object.fromEntries(metrics.map((metric) => [metric, 0])) };
255 if (variant.count) row.count += Number(group.count ?? 0);
256 for (const field of variant.sum ?? []) row[snake(field)] += Number(group.sum?.[field] ?? 0);
257 byName.set(name, row);
258 }
259 const rows = sortRows([...byName.values()], metrics[0]);
260 return { metrics, rows, totals: totalsOf(rows, metrics) };
261}
262
263/** Rows largest first by one metric, ties by name. */
264export function sortRows(rows, metric) {
265 return [...rows].sort((a, b) => (b[metric] ?? 0) - (a[metric] ?? 0) || a.name.localeCompare(b.name));
266}
267
268/** Each metric summed over all rows. */
269export function totalsOf(rows, metrics) {
270 return Object.fromEntries(metrics.map((metric) => [metric, rows.reduce((total, row) => total + (row[metric] ?? 0), 0)]));
271}
272
273/**
274 * The report from every dataset's outcome in every window. An outcome is
275 * `{ body, variant }` (a GraphQL answer) or `{ error }` or `{ skipped }`;
276 * whatever cannot be read becomes a note and the other datasets still count.
277 */
278export function assemble({ account, now, windowList, outcomes, names = {} }) {
279 return {
280 account,
281 generated_at: now.toISOString(),
282 windows: windowList.map((window) => {
283 const notes = [];
284 const errors = [];
285 const datasets = [];
286 for (const spec of DATASETS) {
287 const outcome = outcomes[window.key]?.[spec.key];
288 if (!outcome) continue;
289 if (outcome.skipped) {
290 notes.push(`${spec.label}: ${outcome.skipped}`);
291 continue;
292 }
293 if (outcome.fellBack) {
294 const note = `${spec.label} (${outcome.dataset ?? spec.dataset}): fell back to fewer fields: ${outcome.fellBack}`;
295 notes.push(note);
296 errors.push(note);
297 }
298 try {
299 if (outcome.error) throw outcome.error;
300 const parsed = parseGroups(outcome.body, outcome.variant, names[spec.names] ?? {});
301 datasets.push({ key: spec.key, label: spec.label, dataset: outcome.dataset ?? spec.dataset, ...parsed });
302 } catch (error) {
303 const note = `${spec.label} (${outcome.dataset ?? spec.dataset ?? "no dataset"}): ${String(error.message ?? error).split("\n")[0]}`;
304 notes.push(note);
305 errors.push(note);
306 }
307 }
308 return { key: window.key, label: window.label, start: window.start, end: window.end, datasets, notes, errors };
309 }),
310 };
311}
312
313const number = (value) => (Number.isInteger(value) ? value.toLocaleString("en-US") : value.toLocaleString("en-US", { maximumFractionDigits: 2 }));
314
315/** The report as plain tables: the top rows of each dataset, and a totals line. */
316export function format(report, top = 10) {
317 const lines = [`Cloudflare usage for account ${report.account}, ${report.generated_at}`];
318 for (const window of report.windows) {
319 lines.push("", `== ${window.label}: ${window.start} to ${window.end}`);
320 for (const set of window.datasets) {
321 lines.push("", `${set.label} (${set.dataset})`);
322 const shown = set.rows.slice(0, top);
323 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]))]];
324 const widths = cells[0].map((_, at) => Math.max(...cells.map((row) => String(row[at]).length)));
325 const render = (row) => " " + row.map((cell, at) => (at === 0 ? String(cell).padEnd(widths[at]) : String(cell).padStart(widths[at]))).join(" ");
326 lines.push(render(cells[0]));
327 if (!shown.length) lines.push(" (nothing in this window)");
328 for (const row of cells.slice(1, -1)) lines.push(render(row));
329 if (set.rows.length > shown.length) lines.push(` ... and ${set.rows.length - shown.length} more`);
330 lines.push(render(cells.at(-1)));
331 }
332 if (window.notes.length) {
333 lines.push("", "Notes:");
334 for (const note of window.notes) lines.push(` ${note}`);
335 }
336 }
337 return lines.join("\n");
338}
339
340/** The command line: flags and the account. */
341export function parseArgs(argv, env = process.env) {
342 const option = (name) => {
343 const at = argv.indexOf(name);
344 return at >= 0 && argv[at + 1] && !argv[at + 1].startsWith("--") ? argv[at + 1] : null;
345 };
346 return {
347 help: argv.includes("--help") || argv.includes("-h"),
348 json: argv.includes("--json"),
349 top: Math.max(1, Number(option("--top")) || 10),
350 account: option("--account") || env.CLOUDFLARE_ACCOUNT_ID || ACCOUNT_ID,
351 };
352}
353
354async function graphql(auth, account, query, variables = {}) {
355 const response = await fetch(`${API}/graphql`, {
356 method: "POST",
357 headers: { ...auth, "content-type": "application/json", "user-agent": "g1t-ops" },
358 body: JSON.stringify({ query, variables: { accountTag: account, ...variables } }),
359 });
360 const body = await response.json().catch(() => ({ errors: [{ message: `HTTP ${response.status}, not JSON` }] }));
361 if (!response.ok && !body.errors?.length) body.errors = [{ message: `HTTP ${response.status}` }];
362 return body;
363}
364
365/** The account type's field names, to skip datasets the schema lacks; null when it cannot be read. */
366async function schemaFields(auth, account) {
367 try {
368 const body = await graphql(auth, account, `{ __type(name: "account") { fields { name } } }`);
369 const fields = body.data?.__type?.fields;
370 return Array.isArray(fields) ? fields.map((field) => field.name) : null;
371 } catch {
372 return null;
373 }
374}
375
376/** One dataset in one window: each variant in turn until one is answered. */
377async function ask(auth, account, dataset, spec, window) {
378 let last = null;
379 let first = null;
380 for (const variant of spec.variants) {
381 try {
382 const body = await graphql(auth, account, buildQuery(dataset, variant), { start: window.start, end: window.end });
383 if (!body.errors?.length) return first ? { body, variant, dataset, fellBack: first } : { body, variant, dataset };
384 first ??= body.errors.map((e) => e.message).join("; ");
385 last = { body, variant, dataset };
386 } catch (error) {
387 last = { error, dataset };
388 }
389 }
390 return last;
391}
392
393/** A REST listing as id to name; empty when the token may not read it. */
394async function listing(auth, path, id, name) {
395 const out = {};
396 try {
397 for (let page = 1; page <= 10; page++) {
398 const response = await fetch(`${API}${path}${path.includes("?") ? "&" : "?"}per_page=100&page=${page}`, { headers: { ...auth, "user-agent": "g1t-ops" } });
399 const body = await response.json();
400 if (!response.ok || !Array.isArray(body.result)) break;
401 for (const item of body.result) if (item[id]) out[item[id]] = item[name] ?? item[id];
402 if (body.result.length < 100) break;
403 }
404 } catch {
405 // Ids stand in for names.
406 }
407 return out;
408}
409
410async function main() {
411 const options = parseArgs(process.argv.slice(2));
412 if (options.help) {
413 console.log(HELP);
414 return 0;
415 }
416 const auth = cloudflareAuth();
417 if (!auth) {
418 console.error(`Set CLOUDFLARE_API_TOKEN to a token with Account Analytics: Read on account ${options.account}, or CLOUDFLARE_API_KEY and CLOUDFLARE_EMAIL.`);
419 return 2;
420 }
421 const now = new Date();
422 const windowList = windows(now);
423 const account = options.account;
424 const base = `/accounts/${account}`;
425 const [fields, d1, queues, kv, durable] = await Promise.all([
426 schemaFields(auth, account),
427 listing(auth, `${base}/d1/database`, "uuid", "name"),
428 listing(auth, `${base}/queues`, "queue_id", "queue_name"),
429 listing(auth, `${base}/storage/kv/namespaces`, "id", "title"),
430 listing(auth, `${base}/workers/durable_objects/namespaces`, "id", "name"),
431 ]);
432 const names = { d1, queues, kv, durable_objects: durable };
433 const outcomes = {};
434 const jobs = [];
435 for (const window of windowList) {
436 outcomes[window.key] = {};
437 for (const spec of DATASETS) {
438 const dataset = spec.dataset ?? (fields ? logsDataset(fields) : null);
439 if (!dataset) {
440 if (!spec.optional) outcomes[window.key][spec.key] = { skipped: "no dataset" };
441 continue;
442 }
443 if (fields && !fields.includes(dataset)) {
444 if (!spec.optional) outcomes[window.key][spec.key] = { skipped: `${dataset} is not in this account's schema` };
445 continue;
446 }
447 jobs.push(ask(auth, account, dataset, spec, window).then((outcome) => (outcomes[window.key][spec.key] = outcome)));
448 }
449 }
450 // Billing's own watcher queries, over the last full hour.
451 const hour = lastFullHour(now);
452 const watcher = {};
453 jobs.push(
454 ...WATCHER_QUERIES.map(async (query) => {
455 try {
456 const body = await graphql(auth, account, watcherQuery(query), { since: hour.since, until: hour.until });
457 watcher[query.key] = body.errors?.length
458 ? { error: body.errors.map((e) => e.message).join("; ") }
459 : { rows: body.data?.viewer?.accounts?.[0]?.rows?.length ?? 0 };
460 } catch (error) {
461 watcher[query.key] = { error: String(error.message ?? error) };
462 }
463 }),
464 );
465 await Promise.all(jobs);
466 const report = assemble({ account, now, windowList, outcomes, names });
467 const problems = verdict(report, watcher);
468 if (options.json) {
469 console.log(JSON.stringify({ ...report, watcher_hour: hour, watcher, problems }, null, 2));
470 } else {
471 console.log(format(report, options.top));
472 console.log("", `== Billing's watcher queries, ${hour.since} to ${hour.until}`);
473 for (const query of WATCHER_QUERIES) {
474 const outcome = watcher[query.key];
475 console.log(` ${query.key.padEnd(15)} ${outcome?.error ? `ERROR ${outcome.error}` : `${outcome?.rows ?? 0} rows`}`);
476 }
477 }
478 if (problems.length) {
479 console.error("", `platform-usage: ${problems.length} problem(s); the watcher may be blind on these:`);
480 for (const problem of problems) console.error(` ${problem}`);
481 return 1;
482 }
483 return 0;
484}
485
486if (process.argv[1]?.replaceAll("\\", "/").endsWith("scripts/ops/platform-usage.mjs")) {
487 main().then(
488 (code) => process.exit(code),
489 (error) => {
490 console.error(`platform-usage: ${error.message}`);
491 process.exit(1);
492 },
493 );
494}

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