Make a starved discovery cycle say so, in the log and on /api/health.
For about a year the production RELAYS list did not include the relay carrying
the kind 38000/38172 archive. Every backfill read about thirty events, wrote them
faithfully, reported ok=true, and the nightly build republished an index of eight
mints. Nothing measured the difference between "the cycle completed" and "the
cycle read anything", so nothing went red.
Three signals now do:
- Per-relay attribution. queryRelays() replaces pool.querySync(), which merges
every relay into one deduplicated array and throws away who sent what. It
keeps one subscription per relay over the pool's existing sockets and shares
a single alreadyHaveEvent across them, so an event five relays carry is still
verified once; receivedEvent fires before that check, which is what makes the
per-relay count mean "what this relay contributed". The deadline moved out of
each Subscription's own EOSE timer so `eose` means a frame arrived rather than
something timed out.
- A WARN naming any relay that will not connect, on every cycle, and any relay
that connected and sent nothing, on backfills only. An incremental cycle is
supposed to come back empty.
- BACKFILL_MIN_EVENTS, default 200. Under it, ERROR discovery starvation
suspected and a flag health reports as discovery_starved, forcing 503. Sticky
across incremental cycles so an hourly cycle finding four events cannot clear
what a backfill diagnosed; stored in the database so a restart cannot either.
A fresh database is starved until its first backfill lands. That is intended: it
holds the build's health gate rather than publishing a site made from nothing.
Verified against the live relay set — 1528 events, five relays connected, EOSE on
all five, health 200 — and against an unreachable list, which produces the two
WARN lines, the ERROR, and 503.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
65307ba278
commit
9ffa53094d
@@ -112,6 +112,20 @@ export const config = {
|
||||
iconDir: process.env['ICON_DIR'] ?? path.join(apiRoot, 'data', 'icons'),
|
||||
relays: (process.env['RELAYS']?.split(',').map((r) => r.trim()).filter(Boolean) ??
|
||||
[...DEFAULT_RELAYS]) as string[],
|
||||
/**
|
||||
* The floor a backfill cycle has to clear before it counts as a real read.
|
||||
*
|
||||
* For about a year this deployment's RELAYS list did not include the relay carrying
|
||||
* the kind 38000/38172 archive. Every backfill returned about thirty events, wrote
|
||||
* them, reported ok=true, and the index sat at eight mints while every health signal
|
||||
* stayed green. A backfill asks five relays for the whole history of four kinds; on a
|
||||
* working relay set it comes back with thousands. Anything under this is not a quiet
|
||||
* network, it is a misconfigured one, and it says so in the log and on /api/health.
|
||||
*
|
||||
* Raise it on a deployment that genuinely has more history, lower it for a local
|
||||
* test relay. It is deliberately not zero-able: set it to 1 if you mean "off".
|
||||
*/
|
||||
backfillMinEvents: int('BACKFILL_MIN_EVENTS', 200),
|
||||
probeIntervalMin: int('PROBE_INTERVAL_MIN', 10),
|
||||
discoveryIntervalMin: int('DISCOVERY_INTERVAL_MIN', 60),
|
||||
probeConcurrency: int('PROBE_CONCURRENCY', 8),
|
||||
|
||||
+276
-14
@@ -19,9 +19,10 @@ import {
|
||||
type FedimintAnnouncement,
|
||||
type LnurlAnnouncement,
|
||||
type LnurlFields,
|
||||
type RelayHealth,
|
||||
} from '@cashumints/shared';
|
||||
import { config } from './config.ts';
|
||||
import { getDb, setState, getStateNumber } from './db.ts';
|
||||
import { getDb, setState, getState, getStateNumber } from './db.ts';
|
||||
import type { Sql } from './db-driver.ts';
|
||||
import { log } from './log.ts';
|
||||
import { insertMintIfNew, upsertFedimint, upsertLnurl } from './mints.ts';
|
||||
@@ -39,8 +40,118 @@ export interface DiscoveryResult {
|
||||
newMints: string[];
|
||||
newReviews: number;
|
||||
ok: boolean;
|
||||
/** What each configured relay actually did, in `config.relays` order. */
|
||||
relays: RelayHealth[];
|
||||
/** A backfill that came in under `config.backfillMinEvents`. */
|
||||
starved: boolean;
|
||||
}
|
||||
|
||||
/**
|
||||
* What the last cycle did, kept so /api/health can answer for it.
|
||||
*
|
||||
* Written to the `state` table rather than held in memory, because the question it
|
||||
* answers — "is discovery actually reading anything?" — has to survive the restart that
|
||||
* would otherwise reset it to "no cycle yet, nothing to report". A process that crash
|
||||
* loops would clear an in-memory flag on every attempt.
|
||||
*/
|
||||
export interface DiscoveryReport {
|
||||
at: number;
|
||||
mode: 'backfill' | 'incremental';
|
||||
events: number;
|
||||
ok: boolean;
|
||||
relays: RelayHealth[];
|
||||
starved: boolean;
|
||||
}
|
||||
|
||||
/** `state` key holding the JSON of the above. */
|
||||
const REPORT_KEY = 'last_discovery_report';
|
||||
|
||||
/**
|
||||
* The last cycle's report, or null before any cycle has run.
|
||||
*
|
||||
* A row that will not parse reads as null — the same as no cycle — because the caller
|
||||
* is a health endpoint and "I cannot tell you" must not be dressed up as "fine".
|
||||
*/
|
||||
export async function lastDiscoveryReport(): Promise<DiscoveryReport | null> {
|
||||
const raw = await getState(REPORT_KEY);
|
||||
if (!raw) return null;
|
||||
try {
|
||||
const parsed = JSON.parse(raw) as DiscoveryReport;
|
||||
return Array.isArray(parsed.relays) ? parsed : null;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Per-relay bookkeeping for one cycle.
|
||||
*
|
||||
* Every relay in `config.relays` gets a row up front, including the ones that are never
|
||||
* reached, because a relay that produced no row at all is exactly the one worth naming:
|
||||
* the year-long starvation was a relay list that connected cleanly and simply did not
|
||||
* hold the archive, and the only field that would have shown it is a zero here.
|
||||
*
|
||||
* `events` counts what a relay sent *before* cross-relay deduplication, so five relays
|
||||
* carrying the same 400 events report 400 each rather than 400 once and 0 four times.
|
||||
* Attribution is the whole point; the deduplicated total is reported separately.
|
||||
*/
|
||||
class RelayTally {
|
||||
private readonly rows = new Map<string, { events: number; subs: number; eoses: number; connected: boolean }>();
|
||||
|
||||
constructor(urls: readonly string[]) {
|
||||
for (const url of urls) {
|
||||
this.rows.set(url, { events: 0, subs: 0, eoses: 0, connected: false });
|
||||
}
|
||||
}
|
||||
|
||||
private row(url: string) {
|
||||
let found = this.rows.get(url);
|
||||
if (!found) {
|
||||
found = { events: 0, subs: 0, eoses: 0, connected: false };
|
||||
this.rows.set(url, found);
|
||||
}
|
||||
return found;
|
||||
}
|
||||
|
||||
connected(url: string): void {
|
||||
this.row(url).connected = true;
|
||||
}
|
||||
|
||||
subscribed(url: string): void {
|
||||
this.row(url).subs++;
|
||||
}
|
||||
|
||||
event(url: string): void {
|
||||
this.row(url).events++;
|
||||
}
|
||||
|
||||
eose(url: string): void {
|
||||
this.row(url).eoses++;
|
||||
}
|
||||
|
||||
/** One row per configured relay, in configuration order. */
|
||||
list(): RelayHealth[] {
|
||||
return [...this.rows.entries()].map(([url, row]) => ({
|
||||
url,
|
||||
connected: row.connected,
|
||||
events: row.events,
|
||||
// A cycle asks a relay many questions. It only counts as having reached the end
|
||||
// of the stream if it reached the end of every one of them.
|
||||
eose: row.subs > 0 && row.eoses === row.subs,
|
||||
}));
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Long enough that the relay's own EOSE timer never wins.
|
||||
*
|
||||
* `Subscription` fires `oneose` both when an EOSE frame arrives and when its internal
|
||||
* timer expires, so the two are indistinguishable from the callback. Pushing that timer
|
||||
* out of reach and running the deadline here instead is what makes `eose` in the report
|
||||
* mean "the relay said it was done" rather than "something gave up".
|
||||
*/
|
||||
const NEVER_EOSE_MS = 24 * 60 * 60 * 1000;
|
||||
|
||||
let pool: SimplePool | null = null;
|
||||
|
||||
/**
|
||||
@@ -66,12 +177,100 @@ export function closePool(): void {
|
||||
pool = null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Ask every configured relay one filter, and record what each of them did.
|
||||
*
|
||||
* This replaces `pool.querySync(config.relays, …)`, which answers the same question and
|
||||
* throws the attribution away: it merges five relays into one deduplicated array, so a
|
||||
* relay list where four relays are empty and one carries everything is indistinguishable
|
||||
* from five healthy ones. That indistinguishability is the bug this whole file is being
|
||||
* changed for — a year of ~31-event backfills, `ok=true` every time.
|
||||
*
|
||||
* What it keeps from `querySync`, deliberately:
|
||||
*
|
||||
* - One subscription per relay over the pool's existing sockets, so this is the same
|
||||
* number of connections as before.
|
||||
* - A single `alreadyHaveEvent` shared across all five. `AbstractRelay._onmessage`
|
||||
* consults it *before* `JSON.parse` and signature verification, so an event five
|
||||
* relays all carry is still verified once. Per-relay `querySync` calls would have
|
||||
* verified it five times, which at 500 events a page is real CPU.
|
||||
* - `receivedEvent`, which fires on the way past that check, so the per-relay count is
|
||||
* what the relay sent rather than what was new because of it.
|
||||
*
|
||||
* What it changes: the deadline is run here rather than by each `Subscription`'s own
|
||||
* EOSE timer, so `oneose` firing means an EOSE frame actually arrived. See NEVER_EOSE_MS.
|
||||
*
|
||||
* Never throws. A relay that will not connect is a fact to record, not a reason to
|
||||
* abandon the four that did.
|
||||
*/
|
||||
async function queryRelays(filter: Filter, tally: RelayTally | null): Promise<NostrEvent[]> {
|
||||
const events: NostrEvent[] = [];
|
||||
const known = new Set<string>();
|
||||
const alreadyHaveEvent = (id: string): boolean => {
|
||||
if (known.has(id)) return true;
|
||||
known.add(id);
|
||||
return false;
|
||||
};
|
||||
|
||||
await Promise.all(
|
||||
config.relays.map(async (url) => {
|
||||
let relay;
|
||||
try {
|
||||
// The same connection budget subscribeMap would have used for this maxWait.
|
||||
relay = await getPool().ensureRelay(url, {
|
||||
connectionTimeout: Math.max(MAX_WAIT_MS * 0.8, MAX_WAIT_MS - 1000),
|
||||
});
|
||||
} catch {
|
||||
// Left as connected=false in the tally, which is the whole report this needs.
|
||||
return;
|
||||
}
|
||||
tally?.connected(url);
|
||||
|
||||
await new Promise<void>((resolve) => {
|
||||
let settled = false;
|
||||
let deadline: ReturnType<typeof setTimeout> | undefined;
|
||||
const finish = (): void => {
|
||||
if (settled) return;
|
||||
settled = true;
|
||||
if (deadline !== undefined) clearTimeout(deadline);
|
||||
resolve();
|
||||
};
|
||||
|
||||
try {
|
||||
const sub = relay.subscribe([filter], {
|
||||
onevent: (event) => events.push(event),
|
||||
alreadyHaveEvent,
|
||||
receivedEvent: () => tally?.event(url),
|
||||
oneose: () => {
|
||||
tally?.eose(url);
|
||||
sub.close('closed automatically on eose');
|
||||
},
|
||||
onclose: finish,
|
||||
eoseTimeout: NEVER_EOSE_MS,
|
||||
});
|
||||
tally?.subscribed(url);
|
||||
deadline = setTimeout(() => sub.close('closed on maxWait'), MAX_WAIT_MS);
|
||||
} catch {
|
||||
// The socket went away between ensureRelay and the REQ.
|
||||
finish();
|
||||
}
|
||||
});
|
||||
}),
|
||||
);
|
||||
|
||||
return events;
|
||||
}
|
||||
|
||||
/**
|
||||
* Query one kind, paging backwards with `until` until a page yields nothing new.
|
||||
* Relays cap `limit` independently, so paging is the only way a fresh database
|
||||
* converges to the complete history.
|
||||
*/
|
||||
async function fetchKind(kind: number, since: number | null): Promise<NostrEvent[]> {
|
||||
async function fetchKind(
|
||||
kind: number,
|
||||
since: number | null,
|
||||
tally: RelayTally | null,
|
||||
): Promise<NostrEvent[]> {
|
||||
const seen = new Map<string, NostrEvent>();
|
||||
let until: number | undefined;
|
||||
|
||||
@@ -82,7 +281,7 @@ async function fetchKind(kind: number, since: number | null): Promise<NostrEvent
|
||||
|
||||
let batch: NostrEvent[];
|
||||
try {
|
||||
batch = await getPool().querySync(config.relays, filter, { maxWait: MAX_WAIT_MS });
|
||||
batch = await queryRelays(filter, tally);
|
||||
} catch (err) {
|
||||
log.warn('relay query failed', {
|
||||
kind,
|
||||
@@ -123,6 +322,7 @@ async function fetchKind(kind: number, since: number | null): Promise<NostrEvent
|
||||
async function fetchReviewsForMint(
|
||||
target: ReviewTarget,
|
||||
since: number | null,
|
||||
tally: RelayTally | null,
|
||||
): Promise<NostrEvent[]> {
|
||||
const filters: Filter[] = [];
|
||||
const base: Filter = { kinds: [KIND_REVIEW], limit: QUERY_LIMIT };
|
||||
@@ -178,11 +378,7 @@ async function fetchReviewsForMint(
|
||||
}
|
||||
|
||||
const batches = await Promise.all(
|
||||
filters.map((filter) =>
|
||||
getPool()
|
||||
.querySync(config.relays, filter, { maxWait: MAX_WAIT_MS })
|
||||
.catch(() => [] as NostrEvent[]),
|
||||
),
|
||||
filters.map((filter) => queryRelays(filter, tally).catch(() => [] as NostrEvent[])),
|
||||
);
|
||||
|
||||
return batches.flat();
|
||||
@@ -596,6 +792,7 @@ export async function runDiscovery(backfill: boolean): Promise<DiscoveryResult>
|
||||
const since = lastRun === null ? null : Math.max(0, lastRun - 3600);
|
||||
|
||||
const newMints = new Set<string>();
|
||||
const tally = new RelayTally(config.relays);
|
||||
let newReviews = 0;
|
||||
let events = 0;
|
||||
let ok = true;
|
||||
@@ -610,7 +807,7 @@ export async function runDiscovery(backfill: boolean): Promise<DiscoveryResult>
|
||||
const announcementsByType = new Map<string, NostrEvent[]>();
|
||||
await Promise.all(
|
||||
Object.entries(ANNOUNCEMENT_KINDS).map(async ([type, kind]) => {
|
||||
announcementsByType.set(type, await fetchKind(kind, since));
|
||||
announcementsByType.set(type, await fetchKind(kind, since, tally));
|
||||
}),
|
||||
);
|
||||
|
||||
@@ -625,10 +822,10 @@ export async function runDiscovery(backfill: boolean): Promise<DiscoveryResult>
|
||||
|
||||
// Announcements alone miss mints that only ever appear in a review's `u` tag,
|
||||
// so reviews feed discovery too.
|
||||
const reviews = await fetchKind(KIND_REVIEW, since);
|
||||
const reviews = await fetchKind(KIND_REVIEW, since, tally);
|
||||
// The recent window catches anything a relay dropped from the unbounded query.
|
||||
const recent =
|
||||
since === null ? await fetchKind(KIND_REVIEW, now - RECENT_WINDOW_S) : [];
|
||||
since === null ? await fetchKind(KIND_REVIEW, now - RECENT_WINDOW_S, tally) : [];
|
||||
|
||||
const byId = new Map<string, NostrEvent>();
|
||||
for (const e of [...reviews, ...recent]) byId.set(e.id, e);
|
||||
@@ -689,7 +886,7 @@ export async function runDiscovery(backfill: boolean): Promise<DiscoveryResult>
|
||||
while (cursor < targets.length) {
|
||||
const target = targets[cursor++];
|
||||
if (!target) continue;
|
||||
const found = await fetchReviewsForMint(target, since);
|
||||
const found = await fetchReviewsForMint(target, since, tally);
|
||||
if (found.length > 0) {
|
||||
events += found.length;
|
||||
newReviews += await ingestReviews(found, index);
|
||||
@@ -707,14 +904,79 @@ export async function runDiscovery(backfill: boolean): Promise<DiscoveryResult>
|
||||
log.error('discovery failed', { reason: err instanceof Error ? err.message : String(err) });
|
||||
}
|
||||
|
||||
const mode = backfill ? 'backfill' : 'incremental';
|
||||
const relays = tally.list();
|
||||
|
||||
/*
|
||||
* Name the relay, every time, one line each.
|
||||
*
|
||||
* A relay that would not connect is worth saying on any cycle: the address is wrong,
|
||||
* or it is down, and neither gets better by itself. A relay that connected and sent
|
||||
* nothing is only news on a backfill — an incremental cycle asking for the last hour
|
||||
* of four kinds legitimately comes back empty, and warning about that hourly would
|
||||
* train everyone to skip the line that eventually matters.
|
||||
*/
|
||||
for (const relay of relays) {
|
||||
if (!relay.connected) {
|
||||
log.warn('discovery relay unreachable', { relay: relay.url, mode });
|
||||
continue;
|
||||
}
|
||||
if (backfill && relay.events === 0) {
|
||||
log.warn('discovery relay returned no events', { relay: relay.url, mode });
|
||||
} else if (!relay.eose) {
|
||||
log.warn('discovery relay never reached EOSE', {
|
||||
relay: relay.url,
|
||||
mode,
|
||||
events: relay.events,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
* The floor, and the flag the health endpoint reads.
|
||||
*
|
||||
* Only a backfill is measured against it. A backfill asks for the entire history of
|
||||
* every announcement kind and every review, so on a working relay set it is thousands
|
||||
* of events; an incremental cycle asks for one interval and is supposed to be small.
|
||||
*
|
||||
* The flag is sticky across incremental cycles: an hourly cycle that finds four
|
||||
* events must not clear a starvation a backfill diagnosed, so a non-backfill carries
|
||||
* forward whatever the last backfill concluded.
|
||||
*/
|
||||
let starved: boolean;
|
||||
if (backfill) {
|
||||
starved = events < config.backfillMinEvents;
|
||||
if (starved) {
|
||||
log.error('discovery starvation suspected', {
|
||||
events,
|
||||
floor: config.backfillMinEvents,
|
||||
relays: relays.length,
|
||||
silent: relays.filter((r) => r.events === 0).length,
|
||||
unreachable: relays.filter((r) => !r.connected).length,
|
||||
hint: 'check RELAYS: a relay list missing the announcement archive looks exactly like this',
|
||||
});
|
||||
}
|
||||
} else {
|
||||
// No backfill has ever run in this deployment: nothing has confirmed the relay set
|
||||
// reads anything, and saying "fine" would be the whole original bug.
|
||||
starved = (await lastDiscoveryReport())?.starved ?? true;
|
||||
}
|
||||
|
||||
const report: DiscoveryReport = { at: now, mode, events, ok, relays, starved };
|
||||
// A report that cannot be written is not worth failing a cycle over; the cycle's own
|
||||
// work is already committed, and health degrades on the stale timestamp instead.
|
||||
await setState(REPORT_KEY, JSON.stringify(report)).catch(() => undefined);
|
||||
|
||||
log.info('discovery cycle', {
|
||||
mode: backfill ? 'backfill' : 'incremental',
|
||||
mode,
|
||||
events,
|
||||
new_mints: newMints.size,
|
||||
new_reviews: newReviews,
|
||||
ok,
|
||||
starved,
|
||||
relays: relays.map((r) => `${r.url}=${r.connected ? r.events : 'down'}`).join(' '),
|
||||
ms: Date.now() - started,
|
||||
});
|
||||
|
||||
return { events, newMints: [...newMints], newReviews, ok };
|
||||
return { events, newMints: [...newMints], newReviews, ok, relays, starved };
|
||||
}
|
||||
|
||||
+31
-3
@@ -16,6 +16,7 @@ import {
|
||||
} from '@cashumints/shared';
|
||||
import { config, startedAt } from './config.ts';
|
||||
import { getDb, getStateNumber, getState } from './db.ts';
|
||||
import { lastDiscoveryReport } from './discovery.ts';
|
||||
import { mintByHost, parseEcosystem, type MintRow } from './mints.ts';
|
||||
|
||||
/**
|
||||
@@ -379,28 +380,55 @@ export function resetStatsCache(): void {
|
||||
statsCache = null;
|
||||
}
|
||||
|
||||
/** Health bypasses the stats cache: it is the endpoint you page on. */
|
||||
/**
|
||||
* Health bypasses the stats cache: it is the endpoint you page on.
|
||||
*
|
||||
* Three things can degrade it, and they are three different failures:
|
||||
*
|
||||
* probeStale nothing has checked a mint in three intervals
|
||||
* !discoveryOk the last discovery cycle threw
|
||||
* report.starved the last backfill read less than BACKFILL_MIN_EVENTS
|
||||
*
|
||||
* The third is the one added after the postmortem, and it is the only one that would
|
||||
* have caught a year of the index sitting at eight mints: the cycles were completing,
|
||||
* `ok` was true, the probes were fresh, and the relay list simply did not contain the
|
||||
* relay holding the archive. "Ran without throwing" is not the same claim as "read
|
||||
* anything", and only the second one is worth a green light.
|
||||
*
|
||||
* A deployment with an empty database reports degraded until its first backfill lands,
|
||||
* because until then nothing has confirmed the relay set reads anything at all. That is
|
||||
* intended: it holds `cashumints-web.service` at its health gate rather than letting it
|
||||
* publish a site built from nothing.
|
||||
*/
|
||||
export async function getHealth(): Promise<Health> {
|
||||
const now = Math.floor(Date.now() / 1000);
|
||||
const db = await getDb();
|
||||
|
||||
const [lastProbe, lastDiscovery, discoveryOkRaw, tracked] = await Promise.all([
|
||||
const [lastProbe, lastDiscovery, discoveryOkRaw, tracked, report] = await Promise.all([
|
||||
getStateNumber('last_probe_at'),
|
||||
getStateNumber('last_discovery_at'),
|
||||
getState('last_discovery_ok'),
|
||||
db.get<{ n: number }>('SELECT COUNT(*) AS n FROM mints'),
|
||||
lastDiscoveryReport(),
|
||||
]);
|
||||
|
||||
const discoveryOk = discoveryOkRaw !== '0';
|
||||
const staleAfter = config.probeIntervalMin * 60 * 3;
|
||||
const probeStale = lastProbe === null || now - lastProbe > staleAfter;
|
||||
// No report at all is starvation by default: see the note above.
|
||||
const starved = report?.starved ?? true;
|
||||
|
||||
return {
|
||||
status: probeStale || !discoveryOk ? 'degraded' : 'ok',
|
||||
status: probeStale || !discoveryOk || starved ? 'degraded' : 'ok',
|
||||
uptime_s: now - startedAt,
|
||||
last_probe_at: lastProbe,
|
||||
last_discovery_at: lastDiscovery,
|
||||
mints_tracked: tracked?.n ?? 0,
|
||||
updated_at: now,
|
||||
discovery_relays: report?.relays ?? [],
|
||||
last_discovery_events: report?.events ?? null,
|
||||
last_discovery_mode: report?.mode ?? null,
|
||||
discovery_starved: starved,
|
||||
backfill_min_events: config.backfillMinEvents,
|
||||
};
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user