| 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 | |
| 28 | import { ACCOUNT_ID, cloudflareAuth } from "../deploy/cloudflare.mjs"; |
| 29 | |
| 30 | const API = "https://api.cloudflare.com/client/v4"; |
| 31 | |
| 32 | const HELP = `node scripts/ops/platform-usage.mjs [--top N] [--json] [--account <id>] |
| 33 | |
| 34 | Cloudflare usage per Worker, D1 database, queue, Durable Object namespace, |
| 35 | KV 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 | |
| 42 | Needs CLOUDFLARE_API_TOKEN (Account Analytics: Read), or CLOUDFLARE_API_KEY |
| 43 | with CLOUDFLARE_EMAIL. |
| 44 | |
| 45 | Exits 1 when any dataset errored or fell back to fewer fields, when any of |
| 46 | billing's watcher queries errored, or when every dataset answered with no |
| 47 | rows (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 | */ |
| 55 | export 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. */ |
| 68 | export 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). */ |
| 76 | export 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 | */ |
| 92 | export 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 | */ |
| 121 | export 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. */ |
| 196 | export 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. */ |
| 206 | export 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`. */ |
| 212 | export 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. */ |
| 231 | export 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. */ |
| 234 | export 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 | */ |
| 241 | export 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. */ |
| 264 | export 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. */ |
| 269 | export 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 | */ |
| 278 | export 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 | |
| 313 | const 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. */ |
| 316 | export 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. */ |
| 341 | export 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 | |
| 354 | async 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. */ |
| 366 | async 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. */ |
| 377 | async 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. */ |
| 394 | async 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 | |
| 410 | async 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 | |
| 486 | if (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 | } |