diff --git a/.env.example b/.env.example index 1912ab2..1e195d3 100644 --- a/.env.example +++ b/.env.example @@ -18,3 +18,10 @@ PORT=3000 # 10002=relay lists). On by default; set to 0 to run a single hose without # writing collections still owned by the legacy firehose processes. # INDEX_LEGACY_KINDS=1 + +# Relay-health prober (node probe.js / npm run probe). Dials every relay in the +# directory, writing online/uptime/latency + NIP-11 auth/payment flags. +# PROBE_INTERVAL=3600000 # sweep cadence in ms (default 1h) +# PROBE_CONCURRENCY=25 # relays probed in parallel +# PROBE_TIMEOUT=7000 # per-relay connect/HTTP timeout in ms +# RUN_ONCE=1 # single sweep then exit (or pass --once) diff --git a/README.md b/README.md index 0b96a00..5d89b3b 100644 --- a/README.md +++ b/README.md @@ -53,6 +53,8 @@ The indexer is being reworked into composable **hoses** (`src/hoses/`), each own Every kind is now hose-owned; the legacy raw-upsert fallback is empty. +Separately, an **active relay-health prober** (`node probe.js` / `npm run probe`) keeps the `relays` directory fresh: it dials every known relay for reachability + latency, reads the NIP-11 info doc for auth/payment requirements, and rolls up uptime. It self-schedules (`PROBE_INTERVAL`, default hourly) or runs a single sweep with `--once`. It's not a hose — it's active, not a subscription. + ## License MIT diff --git a/package.json b/package.json index 8e8e23e..2cf222e 100644 --- a/package.json +++ b/package.json @@ -10,7 +10,8 @@ "test": "node --test", "start": "node index.js", "index": "node -e \"import('./src/indexer.js').then(m => m.runIndexer())\"", - "serve": "node -e \"import('./src/server.js').then(m => m.startServer())\"" + "serve": "node -e \"import('./src/server.js').then(m => m.startServer())\"", + "probe": "node probe.js" }, "dependencies": { "@noble/curves": "^2.2.0", diff --git a/probe.js b/probe.js new file mode 100644 index 0000000..9898d6b --- /dev/null +++ b/probe.js @@ -0,0 +1,14 @@ +// Production entrypoint: start the relay-health prober. +// +// Like serve.js, this avoids relying on the `import.meta.url === main` guard +// (pm2's launcher runs files through a wrapper, so that guard won't fire). Use: +// +// PROBE_INTERVAL=3600000 MONGO_DB=nostr pm2 start probe.js --name relay-prober +import { runProber, proberConfig, USAGE } from './src/prober.js'; + +if (proberConfig().help) { console.log(USAGE); process.exit(0); } + +runProber().catch((e) => { + console.error('[beacon] prober failed to start:', e); + process.exit(1); +}); diff --git a/src/prober.js b/src/prober.js new file mode 100644 index 0000000..89422f2 --- /dev/null +++ b/src/prober.js @@ -0,0 +1,156 @@ +// Active relay-health prober — refreshes the `relays` directory the /relays +// page renders. Unlike the subscription hoses this is *active*: it dials every +// known relay, measures reachability + latency, reads the NIP-11 info document +// for auth/payment requirements, and rolls up uptime. Kind-less, so it is not a +// hose — it's a self-scheduling module of its own. +// +// The relay universe is the `relays` collection itself, continuously seeded by +// the follows/relay-lists hoses' URL harvesting. +import { connect } from './db.js'; +import { ensureRelayDirectoryIndex } from './hoses/relays.js'; +import 'dotenv/config'; + +const RELAY_DIRECTORY = process.env.MONGO_RELAY_DIRECTORY_COLLECTION || 'relays'; + +export const USAGE = `beacon relay-health prober + +Usage: node probe.js [flags] (or: node src/prober.js [flags]) + +Flags (override the matching env var): + --once run a single sweep and exit (env RUN_ONCE=1) + --concurrency relays probed in parallel (env PROBE_CONCURRENCY, default 25) + --timeout per-relay connect/HTTP timeout (env PROBE_TIMEOUT, default 7000) + --interval sweep interval when scheduling (env PROBE_INTERVAL, default 3600000) + -h, --help show this help`; + +/** Parse prober CLI flags into an overlay (only keys the user passed). */ +export function parseArgs(argv = process.argv.slice(2)) { + const out = {}; + const numNext = (i) => { const v = Number(argv[i + 1]); return Number.isFinite(v) ? v : undefined; }; + for (let i = 0; i < argv.length; i++) { + const a = argv[i]; + if (a === '-h' || a === '--help') out.help = true; + else if (a === '--once') out.once = true; + else if (a === '--concurrency') { const v = numNext(i); if (v !== undefined) { out.concurrency = v; i++; } } + else if (a === '--timeout') { const v = numNext(i); if (v !== undefined) { out.timeout = v; i++; } } + else if (a === '--interval') { const v = numNext(i); if (v !== undefined) { out.interval = v; i++; } } + } + return out; +} + +/** Resolve config from env + CLI flags (flags win). Pure — inputs passed in. */ +export function proberConfig(env = process.env, argv = process.argv.slice(2)) { + const flags = parseArgs(argv); + const pos = (v, d) => { const n = Number(v); return Number.isFinite(n) && n > 0 ? n : d; }; + return { + interval: flags.interval ?? pos(env.PROBE_INTERVAL, 3_600_000), + concurrency: flags.concurrency ?? pos(env.PROBE_CONCURRENCY, 25), + timeout: flags.timeout ?? pos(env.PROBE_TIMEOUT, 7_000), + once: flags.once || env.RUN_ONCE === '1' || env.RUN_ONCE === 'true', + help: !!flags.help, + }; +} + +/** wss://… -> https://…, ws://… -> http://… (for the NIP-11 fetch). null if unparseable. */ +export function relayInfoUrl(wsUrl) { + try { + const u = new URL(wsUrl); + if (u.protocol === 'wss:') u.protocol = 'https:'; + else if (u.protocol === 'ws:') u.protocol = 'http:'; + else return null; + return u.toString(); + } catch { return null; } +} + +/** Extract the health-relevant flags from a NIP-11 relay info document. */ +export function nip11Flags(info) { + const lim = (info && typeof info === 'object' && info.limitation) || {}; + return { requiresAuth: !!lim.auth_required, requiresPayment: !!lim.payment_required }; +} + +/** Open a WebSocket; resolve { online, responseTime, error }. Never rejects. */ +export function checkRelay(url, timeoutMs) { + return new Promise((resolve) => { + const start = Date.now(); + let done = false, ws, timer; + const finish = (r) => { if (done) return; done = true; clearTimeout(timer); try { ws?.close(); } catch {} resolve(r); }; + timer = setTimeout(() => finish({ online: false, responseTime: null, error: 'timeout' }), timeoutMs); + try { ws = new WebSocket(url); } catch (e) { return finish({ online: false, responseTime: null, error: String(e?.message || e) }); } + ws.addEventListener('open', () => finish({ online: true, responseTime: Date.now() - start, error: null })); + ws.addEventListener('error', (e) => finish({ online: false, responseTime: null, error: e?.message || 'error' })); + }); +} + +/** Fetch + parse the NIP-11 info doc. Returns flags, or null on any failure. */ +export async function fetchNip11(url, timeoutMs) { + const httpUrl = relayInfoUrl(url); + if (!httpUrl) return null; + const ctrl = new AbortController(); + const timer = setTimeout(() => ctrl.abort(), timeoutMs); + try { + const res = await fetch(httpUrl, { headers: { Accept: 'application/nostr+json' }, signal: ctrl.signal }); + if (!res.ok) return null; + return nip11Flags(await res.json()); + } catch { return null; } + finally { clearTimeout(timer); } +} + +/** Probe one relay and persist the result (uptime rolled up from running totals). */ +export async function probeRelay(db, relay, cfg) { + const conn = await checkRelay(relay, cfg.timeout); + const nip11 = conn.online ? await fetchNip11(relay, cfg.timeout) : null; + await db.collection(RELAY_DIRECTORY).updateOne({ relay }, [ + { + $set: { + lastChecked: '$$NOW', + online: conn.online, + responseTime: conn.responseTime, + lastError: conn.error, + ...(nip11 ? { requiresAuth: nip11.requiresAuth, requiresPayment: nip11.requiresPayment } : {}), + checksTotal: { $add: [{ $ifNull: ['$checksTotal', 0] }, 1] }, + checksOnline: { $add: [{ $ifNull: ['$checksOnline', 0] }, conn.online ? 1 : 0] }, + }, + }, + { $set: { uptime: { $cond: [{ $gt: ['$checksTotal', 0] }, { $multiply: [{ $divide: ['$checksOnline', '$checksTotal'] }, 100] }, null] } } }, + ], { upsert: true }); + return conn; +} + +// Run `fn` over `items` with at most `concurrency` in flight. +async function mapPool(items, concurrency, fn) { + let i = 0; + const worker = async () => { while (i < items.length) { const idx = i++; await fn(items[idx], idx); } }; + await Promise.all(Array.from({ length: Math.max(1, Math.min(concurrency, items.length)) }, worker)); +} + +/** One full sweep over every relay in the directory. Returns { total, online }. */ +export async function sweepOnce(db, cfg) { + const docs = await db.collection(RELAY_DIRECTORY).find({}, { projection: { relay: 1, _id: 0 } }).toArray(); + const relays = [...new Set(docs.map((d) => d.relay).filter(Boolean))]; + let online = 0; + await mapPool(relays, cfg.concurrency, async (relay) => { + try { if ((await probeRelay(db, relay, cfg)).online) online++; } + catch (e) { console.error(`[beacon] probe error ${relay}: ${e.message}`); } + }); + console.log(`[beacon] relay sweep: ${online}/${relays.length} online`); + return { total: relays.length, online }; +} + +export async function runProber() { + const cfg = proberConfig(); + const db = await connect(); + await ensureRelayDirectoryIndex(db); + if (cfg.once) { await sweepOnce(db, cfg); return { stop() {} }; } + console.log(`[beacon] relay-health prober: every ${Math.round(cfg.interval / 60000)}min, concurrency ${cfg.concurrency}, timeout ${cfg.timeout}ms`); + await sweepOnce(db, cfg); + const timer = setInterval(() => sweepOnce(db, cfg).catch((e) => console.error('[beacon] sweep error:', e.message)), cfg.interval); + const stop = () => clearInterval(timer); + process.on('SIGINT', () => { stop(); process.exit(0); }); + process.on('SIGTERM', () => { stop(); process.exit(0); }); + return { stop }; +} + +if (import.meta.url === `file://${process.argv[1]}`) { + if (proberConfig().help) { console.log(USAGE); process.exit(0); } + runProber(); +} diff --git a/test/prober.test.js b/test/prober.test.js new file mode 100644 index 0000000..c5680c4 --- /dev/null +++ b/test/prober.test.js @@ -0,0 +1,57 @@ +// Relay-health prober: pure helpers (config/flags, NIP-11 flags, URL mapping). +// The network/Mongo paths (checkRelay/fetchNip11/sweepOnce) are exercised by a +// manual smoke run, not here. +import test from 'node:test'; +import assert from 'node:assert/strict'; +import { parseArgs, proberConfig, relayInfoUrl, nip11Flags } from '../src/prober.js'; + +test('parseArgs: flags', () => { + assert.deepEqual(parseArgs(['--once']), { once: true }); + assert.deepEqual(parseArgs(['--concurrency', '50']), { concurrency: 50 }); + assert.deepEqual(parseArgs(['--timeout', '3000', '--interval', '600000']), { timeout: 3000, interval: 600000 }); + assert.equal(parseArgs(['--help']).help, true); + assert.equal(parseArgs(['-h']).help, true); + assert.deepEqual(parseArgs([]), {}); +}); + +test('parseArgs: non-numeric value is ignored (not consumed as flag value)', () => { + assert.deepEqual(parseArgs(['--concurrency', 'abc']), {}); +}); + +test('proberConfig: defaults', () => { + const c = proberConfig({}, []); + assert.equal(c.interval, 3_600_000); + assert.equal(c.concurrency, 25); + assert.equal(c.timeout, 7_000); + assert.equal(c.once, false); +}); + +test('proberConfig: env applies, non-positive falls back to default', () => { + const c = proberConfig({ PROBE_CONCURRENCY: '40', PROBE_TIMEOUT: '0', RUN_ONCE: '1' }, []); + assert.equal(c.concurrency, 40); + assert.equal(c.timeout, 7_000); // 0 is rejected -> default + assert.equal(c.once, true); +}); + +test('proberConfig: flags override env', () => { + const c = proberConfig({ PROBE_CONCURRENCY: '40', RUN_ONCE: '1' }, ['--concurrency', '10']); + assert.equal(c.concurrency, 10); // flag beat env + assert.equal(c.once, true); // env still applies where no flag +}); + +test('relayInfoUrl: ws(s) -> http(s), else null', () => { + assert.equal(relayInfoUrl('wss://relay.example.com'), 'https://relay.example.com/'); + assert.equal(relayInfoUrl('wss://relay.example.com/inbox'), 'https://relay.example.com/inbox'); + assert.equal(relayInfoUrl('ws://relay.example.com'), 'http://relay.example.com/'); + assert.equal(relayInfoUrl('https://not-a-relay.com'), null); + assert.equal(relayInfoUrl('garbage'), null); +}); + +test('nip11Flags: reads limitation flags, defaults false', () => { + assert.deepEqual(nip11Flags({ limitation: { auth_required: true, payment_required: false } }), + { requiresAuth: true, requiresPayment: false }); + assert.deepEqual(nip11Flags({ limitation: { payment_required: true } }), + { requiresAuth: false, requiresPayment: true }); + assert.deepEqual(nip11Flags({}), { requiresAuth: false, requiresPayment: false }); + assert.deepEqual(nip11Flags(null), { requiresAuth: false, requiresPayment: false }); +});