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