1#!/usr/bin/env node 2// `nix run .#events`: the Worker's telemetry (worker/src/events-schema.ts), 3// read back from Workers Analytics Engine's SQL API. 4// 5// nix run .#events every event and its columns 6// nix run .#events -- jev.relay its last 50 rows, last day 7// nix run .#events -- jev.relay --since 7d --limit 500 8// nix run .#events -- --sql "SELECT …" any query (add FORMAT JSON 9// for JSON; it is sent as is) 10// 11// An event's rows print one JSON object a line, newest first, with the 12// schema's field names, `timestamp`, and `weight` (_sample_interval: how 13// many events the row stands for; count with it, not with rows). A double 14// its event did not know is null, from the known-mask in double1. 15// 16// The account comes from CLOUDFLARE_ACCOUNT_ID (the flake sets it from 17// nix-facts) and the token from 1Password, read per run from 18// EVENTS_OP_ITEM_REF + "/Token", as the deploy's wrangler reads it. It is 19// held in memory for the request and never printed. Querying needs the 20// token to carry Account Analytics: Read. 21import { execFileSync } from 'node:child_process'; 22import { readFileSync } from 'node:fs'; 23import { pathToFileURL } from 'node:url'; 24import { EVENTS } from '../../worker/src/events-schema.ts'; 25 26const CONFIG = new URL('../../worker/wrangler.json', import.meta.url); 27export const dataset = () => JSON.parse(readFileSync(CONFIG, 'utf8')).analytics_engine_datasets[0].dataset; 28 29const UNITS = { m: 'MINUTE', h: 'HOUR', d: 'DAY' }; 30const quote = (s) => `'${s.replaceAll("'", "''")}'`; 31 32// The SQL for an event's recent rows, its columns named as the schema names 33// them (the fields' names are checked against the schema, so only the 34// schema's text reaches the query). 35export function eventQuery(name, { since = '1d', limit = 50, table = dataset() } = {}) { 36 const def = EVENTS[name]; 37 if (!def) throw new Error(`no event ${name}; known: ${Object.keys(EVENTS).join(', ')}`); 38 const m = /^(\d{1,4})([mhd])$/.exec(since); 39 if (!m) throw new Error(`--since is a number and m, h or d (e.g. 30m, 12h, 7d), not ${since}`); 40 if (!Number.isInteger(limit) || limit < 1 || limit > 10000) throw new Error('--limit is 1 to 10000'); 41 if (!/^[A-Za-z0-9_]+$/.test(table)) throw new Error(`not a dataset name: ${table}`); 42 const columns = [ 43 'timestamp', 44 '_sample_interval AS weight', 45 'double1 AS _known', 46 ...Object.entries(def.blobs).map(([field, { column }]) => `${column} AS ${field}`), 47 ...Object.entries(def.doubles).map(([field, { column }]) => `${column} AS ${field}`), 48 ]; 49 return [ 50 `SELECT ${columns.join(', ')}`, 51 `FROM ${table}`, 52 `WHERE index1 = ${quote(name)} AND timestamp > NOW() - INTERVAL '${m[1]}' ${UNITS[m[2]]}`, 53 `ORDER BY timestamp DESC`, 54 `LIMIT ${limit}`, 55 `FORMAT JSON`, 56 ].join('\n'); 57} 58 59// Rows as eventQuery's JSON returns them, with unknown doubles as null. 60export function decodeRows(name, rows) { 61 const def = EVENTS[name]; 62 return rows.map(({ _known, ...row }) => { 63 const known = Number(_known); 64 for (const [field, { column }] of Object.entries(def.doubles)) { 65 const bit = 1 << (Number(column.slice('double'.length)) - 2); 66 row[field] = known & bit ? Number(row[field]) : null; 67 } 68 row.weight = Number(row.weight); 69 return row; 70 }); 71} 72 73function token() { 74 const ref = process.env.EVENTS_OP_ITEM_REF; 75 if (!ref) throw new Error('EVENTS_OP_ITEM_REF is not set; run this as `nix run .#events`'); 76 let op = 'op'; 77 try { 78 execFileSync('sh', ['-c', 'command -v op.exe'], { stdio: 'ignore' }); 79 op = 'op.exe'; 80 } catch {} 81 const value = execFileSync(op, ['read', `${ref}/Token`], { encoding: 'utf8', stdio: ['ignore', 'pipe', 'inherit'] }) 82 .replaceAll('\r', '') 83 .trim(); 84 if (!value) throw new Error('1Password returned an empty token'); 85 return value; 86} 87 88async function query(sql) { 89 const account = process.env.CLOUDFLARE_ACCOUNT_ID; 90 if (!account) throw new Error('CLOUDFLARE_ACCOUNT_ID is not set; run this as `nix run .#events`'); 91 const res = await fetch(`https://api.cloudflare.com/client/v4/accounts/${account}/analytics_engine/sql`, { 92 method: 'POST', 93 headers: { Authorization: `Bearer ${token()}` }, 94 body: sql, 95 }); 96 const text = await res.text(); 97 if (!res.ok) { 98 const hint = res.status === 403 ? ' (the token needs Account Analytics: Read)' : ''; 99 throw new Error(`the SQL API answered ${res.status}${hint}: ${text.slice(0, 500)}`); 100 } 101 return text; 102} 103 104function describe() { 105 const table = dataset(); 106 console.log(`dataset ${table}: index1 is the event's name; double1 is the known-mask (bit k-2 for double k)\n`); 107 for (const [name, def] of Object.entries(EVENTS)) { 108 console.log(`${name}: ${def.about}`); 109 for (const [field, { column, about }] of [...Object.entries(def.blobs), ...Object.entries(def.doubles)]) { 110 console.log(` ${column.padEnd(9)} ${field}: ${about}`); 111 } 112 console.log(); 113 } 114} 115 116async function main(argv) { 117 const flag = (name) => { 118 const i = argv.indexOf(name); 119 if (i < 0) return undefined; 120 const value = argv[i + 1]; 121 argv.splice(i, 2); 122 if (value === undefined) throw new Error(`${name} needs a value`); 123 return value; 124 }; 125 const sql = flag('--sql'); 126 if (sql) return void process.stdout.write(await query(sql)); 127 const since = flag('--since'); 128 const limit = flag('--limit'); 129 const [name, ...rest] = argv; 130 if (rest.length) throw new Error(`unexpected: ${rest.join(' ')}`); 131 if (!name) return describe(); 132 const q = eventQuery(name, { since, limit: limit === undefined ? undefined : Number(limit) }); 133 const { data } = JSON.parse(await query(q)); 134 for (const row of decodeRows(name, data)) console.log(JSON.stringify(row)); 135} 136 137if (import.meta.url === pathToFileURL(process.argv[1]).href) { 138 main(process.argv.slice(2)).catch((e) => { 139 console.error(`events: ${e.message}`); 140 process.exit(1); 141 }); 142}