jevstrudel.git / tools / events / events.mjs
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.

21import { execFileSync } from 'node:child_process';
22import { readFileSync } from 'node:fs';
23import { pathToFileURL } from 'node:url';
24import { EVENTS } from '../../worker/src/events-schema.ts';
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("'", "''")}'`;

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}