Skip to content

Commit 64ad298

Browse files
Merge pull request #16 from JavaScriptSolidServer/issue-15-relay-prober
Firehose phase 4: active relay-health prober (unfreezes /relays)
2 parents 887d3ed + 2b96354 commit 64ad298

6 files changed

Lines changed: 238 additions & 1 deletion

File tree

.env.example

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,3 +18,10 @@ PORT=3000
1818
# 10002=relay lists). On by default; set to 0 to run a single hose without
1919
# writing collections still owned by the legacy firehose processes.
2020
# INDEX_LEGACY_KINDS=1
21+
22+
# Relay-health prober (node probe.js / npm run probe). Dials every relay in the
23+
# directory, writing online/uptime/latency + NIP-11 auth/payment flags.
24+
# PROBE_INTERVAL=3600000 # sweep cadence in ms (default 1h)
25+
# PROBE_CONCURRENCY=25 # relays probed in parallel
26+
# PROBE_TIMEOUT=7000 # per-relay connect/HTTP timeout in ms
27+
# RUN_ONCE=1 # single sweep then exit (or pass --once)

README.md

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -53,6 +53,8 @@ The indexer is being reworked into composable **hoses** (`src/hoses/`), each own
5353

5454
Every kind is now hose-owned; the legacy raw-upsert fallback is empty.
5555

56+
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.
57+
5658
## License
5759

5860
MIT

package.json

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,8 @@
1010
"test": "node --test",
1111
"start": "node index.js",
1212
"index": "node -e \"import('./src/indexer.js').then(m => m.runIndexer())\"",
13-
"serve": "node -e \"import('./src/server.js').then(m => m.startServer())\""
13+
"serve": "node -e \"import('./src/server.js').then(m => m.startServer())\"",
14+
"probe": "node probe.js"
1415
},
1516
"dependencies": {
1617
"@noble/curves": "^2.2.0",

probe.js

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,14 @@
1+
// Production entrypoint: start the relay-health prober.
2+
//
3+
// Like serve.js, this avoids relying on the `import.meta.url === main` guard
4+
// (pm2's launcher runs files through a wrapper, so that guard won't fire). Use:
5+
//
6+
// PROBE_INTERVAL=3600000 MONGO_DB=nostr pm2 start probe.js --name relay-prober
7+
import { runProber, proberConfig, USAGE } from './src/prober.js';
8+
9+
if (proberConfig().help) { console.log(USAGE); process.exit(0); }
10+
11+
runProber().catch((e) => {
12+
console.error('[beacon] prober failed to start:', e);
13+
process.exit(1);
14+
});

src/prober.js

Lines changed: 156 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,156 @@
1+
// Active relay-health prober — refreshes the `relays` directory the /relays
2+
// page renders. Unlike the subscription hoses this is *active*: it dials every
3+
// known relay, measures reachability + latency, reads the NIP-11 info document
4+
// for auth/payment requirements, and rolls up uptime. Kind-less, so it is not a
5+
// hose — it's a self-scheduling module of its own.
6+
//
7+
// The relay universe is the `relays` collection itself, continuously seeded by
8+
// the follows/relay-lists hoses' URL harvesting.
9+
import { connect } from './db.js';
10+
import { ensureRelayDirectoryIndex } from './hoses/relays.js';
11+
import 'dotenv/config';
12+
13+
const RELAY_DIRECTORY = process.env.MONGO_RELAY_DIRECTORY_COLLECTION || 'relays';
14+
15+
export const USAGE = `beacon relay-health prober
16+
17+
Usage: node probe.js [flags] (or: node src/prober.js [flags])
18+
19+
Flags (override the matching env var):
20+
--once run a single sweep and exit (env RUN_ONCE=1)
21+
--concurrency <n> relays probed in parallel (env PROBE_CONCURRENCY, default 25)
22+
--timeout <ms> per-relay connect/HTTP timeout (env PROBE_TIMEOUT, default 7000)
23+
--interval <ms> sweep interval when scheduling (env PROBE_INTERVAL, default 3600000)
24+
-h, --help show this help`;
25+
26+
/** Parse prober CLI flags into an overlay (only keys the user passed). */
27+
export function parseArgs(argv = process.argv.slice(2)) {
28+
const out = {};
29+
const numNext = (i) => { const v = Number(argv[i + 1]); return Number.isFinite(v) ? v : undefined; };
30+
for (let i = 0; i < argv.length; i++) {
31+
const a = argv[i];
32+
if (a === '-h' || a === '--help') out.help = true;
33+
else if (a === '--once') out.once = true;
34+
else if (a === '--concurrency') { const v = numNext(i); if (v !== undefined) { out.concurrency = v; i++; } }
35+
else if (a === '--timeout') { const v = numNext(i); if (v !== undefined) { out.timeout = v; i++; } }
36+
else if (a === '--interval') { const v = numNext(i); if (v !== undefined) { out.interval = v; i++; } }
37+
}
38+
return out;
39+
}
40+
41+
/** Resolve config from env + CLI flags (flags win). Pure — inputs passed in. */
42+
export function proberConfig(env = process.env, argv = process.argv.slice(2)) {
43+
const flags = parseArgs(argv);
44+
const pos = (v, d) => { const n = Number(v); return Number.isFinite(n) && n > 0 ? n : d; };
45+
return {
46+
interval: flags.interval ?? pos(env.PROBE_INTERVAL, 3_600_000),
47+
concurrency: flags.concurrency ?? pos(env.PROBE_CONCURRENCY, 25),
48+
timeout: flags.timeout ?? pos(env.PROBE_TIMEOUT, 7_000),
49+
once: flags.once || env.RUN_ONCE === '1' || env.RUN_ONCE === 'true',
50+
help: !!flags.help,
51+
};
52+
}
53+
54+
/** wss://… -> https://…, ws://… -> http://… (for the NIP-11 fetch). null if unparseable. */
55+
export function relayInfoUrl(wsUrl) {
56+
try {
57+
const u = new URL(wsUrl);
58+
if (u.protocol === 'wss:') u.protocol = 'https:';
59+
else if (u.protocol === 'ws:') u.protocol = 'http:';
60+
else return null;
61+
return u.toString();
62+
} catch { return null; }
63+
}
64+
65+
/** Extract the health-relevant flags from a NIP-11 relay info document. */
66+
export function nip11Flags(info) {
67+
const lim = (info && typeof info === 'object' && info.limitation) || {};
68+
return { requiresAuth: !!lim.auth_required, requiresPayment: !!lim.payment_required };
69+
}
70+
71+
/** Open a WebSocket; resolve { online, responseTime, error }. Never rejects. */
72+
export function checkRelay(url, timeoutMs) {
73+
return new Promise((resolve) => {
74+
const start = Date.now();
75+
let done = false, ws, timer;
76+
const finish = (r) => { if (done) return; done = true; clearTimeout(timer); try { ws?.close(); } catch {} resolve(r); };
77+
timer = setTimeout(() => finish({ online: false, responseTime: null, error: 'timeout' }), timeoutMs);
78+
try { ws = new WebSocket(url); } catch (e) { return finish({ online: false, responseTime: null, error: String(e?.message || e) }); }
79+
ws.addEventListener('open', () => finish({ online: true, responseTime: Date.now() - start, error: null }));
80+
ws.addEventListener('error', (e) => finish({ online: false, responseTime: null, error: e?.message || 'error' }));
81+
});
82+
}
83+
84+
/** Fetch + parse the NIP-11 info doc. Returns flags, or null on any failure. */
85+
export async function fetchNip11(url, timeoutMs) {
86+
const httpUrl = relayInfoUrl(url);
87+
if (!httpUrl) return null;
88+
const ctrl = new AbortController();
89+
const timer = setTimeout(() => ctrl.abort(), timeoutMs);
90+
try {
91+
const res = await fetch(httpUrl, { headers: { Accept: 'application/nostr+json' }, signal: ctrl.signal });
92+
if (!res.ok) return null;
93+
return nip11Flags(await res.json());
94+
} catch { return null; }
95+
finally { clearTimeout(timer); }
96+
}
97+
98+
/** Probe one relay and persist the result (uptime rolled up from running totals). */
99+
export async function probeRelay(db, relay, cfg) {
100+
const conn = await checkRelay(relay, cfg.timeout);
101+
const nip11 = conn.online ? await fetchNip11(relay, cfg.timeout) : null;
102+
await db.collection(RELAY_DIRECTORY).updateOne({ relay }, [
103+
{
104+
$set: {
105+
lastChecked: '$$NOW',
106+
online: conn.online,
107+
responseTime: conn.responseTime,
108+
lastError: conn.error,
109+
...(nip11 ? { requiresAuth: nip11.requiresAuth, requiresPayment: nip11.requiresPayment } : {}),
110+
checksTotal: { $add: [{ $ifNull: ['$checksTotal', 0] }, 1] },
111+
checksOnline: { $add: [{ $ifNull: ['$checksOnline', 0] }, conn.online ? 1 : 0] },
112+
},
113+
},
114+
{ $set: { uptime: { $cond: [{ $gt: ['$checksTotal', 0] }, { $multiply: [{ $divide: ['$checksOnline', '$checksTotal'] }, 100] }, null] } } },
115+
], { upsert: true });
116+
return conn;
117+
}
118+
119+
// Run `fn` over `items` with at most `concurrency` in flight.
120+
async function mapPool(items, concurrency, fn) {
121+
let i = 0;
122+
const worker = async () => { while (i < items.length) { const idx = i++; await fn(items[idx], idx); } };
123+
await Promise.all(Array.from({ length: Math.max(1, Math.min(concurrency, items.length)) }, worker));
124+
}
125+
126+
/** One full sweep over every relay in the directory. Returns { total, online }. */
127+
export async function sweepOnce(db, cfg) {
128+
const docs = await db.collection(RELAY_DIRECTORY).find({}, { projection: { relay: 1, _id: 0 } }).toArray();
129+
const relays = [...new Set(docs.map((d) => d.relay).filter(Boolean))];
130+
let online = 0;
131+
await mapPool(relays, cfg.concurrency, async (relay) => {
132+
try { if ((await probeRelay(db, relay, cfg)).online) online++; }
133+
catch (e) { console.error(`[beacon] probe error ${relay}: ${e.message}`); }
134+
});
135+
console.log(`[beacon] relay sweep: ${online}/${relays.length} online`);
136+
return { total: relays.length, online };
137+
}
138+
139+
export async function runProber() {
140+
const cfg = proberConfig();
141+
const db = await connect();
142+
await ensureRelayDirectoryIndex(db);
143+
if (cfg.once) { await sweepOnce(db, cfg); return { stop() {} }; }
144+
console.log(`[beacon] relay-health prober: every ${Math.round(cfg.interval / 60000)}min, concurrency ${cfg.concurrency}, timeout ${cfg.timeout}ms`);
145+
await sweepOnce(db, cfg);
146+
const timer = setInterval(() => sweepOnce(db, cfg).catch((e) => console.error('[beacon] sweep error:', e.message)), cfg.interval);
147+
const stop = () => clearInterval(timer);
148+
process.on('SIGINT', () => { stop(); process.exit(0); });
149+
process.on('SIGTERM', () => { stop(); process.exit(0); });
150+
return { stop };
151+
}
152+
153+
if (import.meta.url === `file://${process.argv[1]}`) {
154+
if (proberConfig().help) { console.log(USAGE); process.exit(0); }
155+
runProber();
156+
}

test/prober.test.js

Lines changed: 57 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,57 @@
1+
// Relay-health prober: pure helpers (config/flags, NIP-11 flags, URL mapping).
2+
// The network/Mongo paths (checkRelay/fetchNip11/sweepOnce) are exercised by a
3+
// manual smoke run, not here.
4+
import test from 'node:test';
5+
import assert from 'node:assert/strict';
6+
import { parseArgs, proberConfig, relayInfoUrl, nip11Flags } from '../src/prober.js';
7+
8+
test('parseArgs: flags', () => {
9+
assert.deepEqual(parseArgs(['--once']), { once: true });
10+
assert.deepEqual(parseArgs(['--concurrency', '50']), { concurrency: 50 });
11+
assert.deepEqual(parseArgs(['--timeout', '3000', '--interval', '600000']), { timeout: 3000, interval: 600000 });
12+
assert.equal(parseArgs(['--help']).help, true);
13+
assert.equal(parseArgs(['-h']).help, true);
14+
assert.deepEqual(parseArgs([]), {});
15+
});
16+
17+
test('parseArgs: non-numeric value is ignored (not consumed as flag value)', () => {
18+
assert.deepEqual(parseArgs(['--concurrency', 'abc']), {});
19+
});
20+
21+
test('proberConfig: defaults', () => {
22+
const c = proberConfig({}, []);
23+
assert.equal(c.interval, 3_600_000);
24+
assert.equal(c.concurrency, 25);
25+
assert.equal(c.timeout, 7_000);
26+
assert.equal(c.once, false);
27+
});
28+
29+
test('proberConfig: env applies, non-positive falls back to default', () => {
30+
const c = proberConfig({ PROBE_CONCURRENCY: '40', PROBE_TIMEOUT: '0', RUN_ONCE: '1' }, []);
31+
assert.equal(c.concurrency, 40);
32+
assert.equal(c.timeout, 7_000); // 0 is rejected -> default
33+
assert.equal(c.once, true);
34+
});
35+
36+
test('proberConfig: flags override env', () => {
37+
const c = proberConfig({ PROBE_CONCURRENCY: '40', RUN_ONCE: '1' }, ['--concurrency', '10']);
38+
assert.equal(c.concurrency, 10); // flag beat env
39+
assert.equal(c.once, true); // env still applies where no flag
40+
});
41+
42+
test('relayInfoUrl: ws(s) -> http(s), else null', () => {
43+
assert.equal(relayInfoUrl('wss://relay.example.com'), 'https://relay.example.com/');
44+
assert.equal(relayInfoUrl('wss://relay.example.com/inbox'), 'https://relay.example.com/inbox');
45+
assert.equal(relayInfoUrl('ws://relay.example.com'), 'http://relay.example.com/');
46+
assert.equal(relayInfoUrl('https://not-a-relay.com'), null);
47+
assert.equal(relayInfoUrl('garbage'), null);
48+
});
49+
50+
test('nip11Flags: reads limitation flags, defaults false', () => {
51+
assert.deepEqual(nip11Flags({ limitation: { auth_required: true, payment_required: false } }),
52+
{ requiresAuth: true, requiresPayment: false });
53+
assert.deepEqual(nip11Flags({ limitation: { payment_required: true } }),
54+
{ requiresAuth: false, requiresPayment: true });
55+
assert.deepEqual(nip11Flags({}), { requiresAuth: false, requiresPayment: false });
56+
assert.deepEqual(nip11Flags(null), { requiresAuth: false, requiresPayment: false });
57+
});

0 commit comments

Comments
 (0)