1#!/usr/bin/env node
nix run .#events: the Worker's telemetry (worker/src/events-schema.ts),
read back from Workers Analytics Engine's SQL API.
nix run .#events every event and its columns nix run .#events -- jev.relay its last 50 rows, last day nix run .#events -- jev.relay --since 7d --limit 500 nix run .#events -- --sql "SELECT …" any query (add FORMAT JSON for JSON; it is sent as is)
An event's rows print one JSON object a line, newest first, with the
schema's field names, timestamp, and weight (_sample_interval: how
many events the row stands for; count with it, not with rows). A double
its event did not know is null, from the known-mask in double1.
The account comes from CLOUDFLARE_ACCOUNT_ID (the flake sets it from nix-facts) and the token from 1Password, read per run from EVENTS_OP_ITEM_REF + "/Token", as the deploy's wrangler reads it. It is held in memory for the request and never printed. Querying needs the token to carry Account Analytics: Read.
The SQL for an event's recent rows, its columns named as the schema names them (the fields' names are checked against the schema, so only the 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}
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}
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}