Expand ecash explorer capabilities

Add Fedimint discovery, dual SQLite/Postgres storage, richer review handling, and generated social imagery.
This commit is contained in:
michilis
2026-08-21 02:10:48 +02:00
parent aa1771ea20
commit 6f17b572b1
80 changed files with 7580 additions and 704 deletions
+46 -26
View File
@@ -8,44 +8,62 @@
* a test that actually kills a mint rather than a unit test of a helper.
*
* Runs against a throwaway copy of the seeded database, so the real one is untouched.
* The copy is made with the migrator, which means the check works whichever backend
* holds the real data — and, as a side effect, exercises the migrator on every run.
*
* Set CHECK_DB_URL to run the copy on Postgres instead of a temporary SQLite file:
*
* CHECK_DB_URL=postgres://localhost/cashumints_test pnpm --filter ./api test:offline
*/
import assert from 'node:assert/strict';
import fs from 'node:fs';
import os from 'node:os';
import path from 'node:path';
import { parseDbTarget, resolveDbConfig } from './config.ts';
import { TABLES } from './db-schema.ts';
import { openDb } from './db.ts';
import { migrate } from './migrate.ts';
const source = process.env['DB_PATH'] ?? path.resolve(import.meta.dirname, '..', 'data', 'cashumints.db');
if (!fs.existsSync(source)) {
console.error(`No seeded database at ${source}. Run "pnpm seed" first.`);
const source = resolveDbConfig();
if (source.dialect === 'sqlite' && !fs.existsSync(source.file)) {
console.error(`No seeded database at ${source.file}. Run "pnpm seed" first.`);
process.exit(1);
}
// Point the whole process at a copy before anything opens the real database.
const scratch = fs.mkdtempSync(path.join(os.tmpdir(), 'cashumints-offline-'));
const copy = path.join(scratch, 'test.db');
fs.copyFileSync(source, copy);
process.env['DB_PATH'] = copy;
const target = parseDbTarget(process.env['CHECK_DB_URL'] ?? path.join(scratch, 'test.db'));
// A reused Postgres test database still holds the last run's rows, and a stale mint row
// would be picked as the victim. Start empty either way.
const wipe = await openDb(target);
for (const table of TABLES) await wipe.run(`DELETE FROM ${table}`);
await wipe.close();
await migrate({ from: source, to: target, force: true, dryRun: false });
// Point the whole process at the copy before anything else opens a database.
process.env['DATABASE_URL'] = target.dialect === 'postgres' ? target.url : `sqlite:${target.file}`;
delete process.env['DB_PATH'];
const { getDb, closeDb } = await import('./db.ts');
const { probeMint } = await import('./probe.ts');
const { getMintDetail } = await import('./queries.ts');
const { getMintDetail, listMints, resetStatsCache } = await import('./queries.ts');
const { mintByUrl } = await import('./mints.ts');
const db = getDb();
const db = await getDb();
assert.equal(db.label, target.label, 'the check must not be talking to the real database');
// Pick a mint that is online and has real cached metadata to lose.
const victim = db
.prepare(
`SELECT * FROM mints
WHERE status = 'online' AND info_json IS NOT NULL AND name IS NOT NULL
ORDER BY LENGTH(info_json) DESC LIMIT 1`,
)
.get() as { url: string; host: string; name: string } | undefined;
const victim = await db.get<{ url: string; host: string; name: string }>(
`SELECT url, host, name FROM mints
WHERE status = 'online' AND info_json IS NOT NULL AND name IS NOT NULL
ORDER BY LENGTH(info_json) DESC LIMIT 1`,
);
assert.ok(victim, 'seeded database has no online mint with cached metadata');
console.log(`victim: ${victim.name} (${victim.host})`);
const before = getMintDetail(victim.host);
const before = await getMintDetail(victim.host);
assert.ok(before, 'detail must exist before the mint dies');
assert.equal(before.status, 'online');
assert.ok(before.info, 'must have cached /v1/info before');
@@ -53,18 +71,20 @@ assert.ok(before.nuts.length > 0, 'must have parsed NUTs before');
// Rug it: same row, dead URL. This is the "point a mint row at a dead URL" test.
const DEAD = 'https://this-mint-is-gone.invalid';
db.prepare('UPDATE mints SET url = ? WHERE host = ?').run(DEAD, victim.host);
db.prepare('UPDATE reviews SET mint_url = ? WHERE mint_url = ?').run(DEAD, victim.url);
await db.run('UPDATE mints SET url = ? WHERE host = ?', DEAD, victim.host);
await db.run('UPDATE reviews SET mint_url = ? WHERE mint_url = ?', DEAD, victim.url);
// Three failures is the offline threshold.
for (let i = 0; i < 3; i++) {
const row = mintByUrl(DEAD);
const row = await mintByUrl(DEAD);
assert.ok(row, 'row must survive the URL change');
const result = await probeMint(row);
assert.equal(result.ok, false, 'a dead URL must not probe ok');
}
const after = getMintDetail(victim.host);
resetStatsCache();
const after = await getMintDetail(victim.host);
assert.ok(after, 'the mint page must still resolve when the mint is dead');
assert.equal(after.status, 'offline', 'three failures means offline');
@@ -86,17 +106,17 @@ assert.ok(after.last_online !== null, 'last_online must be set so the banner has
assert.ok(after.last_online <= Math.floor(Date.now() / 1000));
// And it must sink below every online mint without disappearing.
const { listMints } = await import('./queries.ts');
const list = listMints();
const list = await listMints();
const index = list.findIndex((m) => m.host === victim.host);
assert.ok(index >= 0, 'an offline mint must stay listed so people can find and review it');
const lastOnline = list.map((m) => m.status).lastIndexOf('online');
assert.ok(index > lastOnline, 'offline mints rank below every online mint');
closeDb();
await closeDb();
fs.rmSync(scratch, { recursive: true, force: true });
console.log(
`ok, offline mint still serves ${after.nuts.length} NUTs, ${after.review_count} reviews, ` +
`and full cached metadata (offline since ${new Date(after.last_online * 1000).toISOString()})`,
`ok on ${db.dialect}, offline mint still serves ${after.nuts.length} NUTs, ` +
`${after.review_count} reviews, and full cached metadata ` +
`(offline since ${new Date(after.last_online * 1000).toISOString()})`,
);
+248 -25
View File
@@ -6,8 +6,18 @@
* it throws on the first failure and prints a count on success.
*/
import assert from 'node:assert/strict';
import { bayesianScore, compareMints, NEUTRAL_PRIOR_MEAN, normalizeMintUrl, parseNuts, parseLimits, parseRating } from '@cashumints/shared';
import path from 'node:path';
import {
bayesianScore, cleanInviteCodes, compareMints, ecosystemForKind, fedimintKey, fedimintSlug,
federationIdFromKey, getMintWarnings, hasModule, NEUTRAL_PRIOR_MEAN, normalizeMintUrl,
normalizeNetwork, parseFedimintAnnouncement, parseLimits, parseModules, parseNuts, parseRating,
reviewEcosystem,
} from '@cashumints/shared';
import { parseDbTarget } from './config.ts';
import { openDb } from './db.ts';
import { toPgPlaceholders } from './db-postgres.ts';
import { statusForFails } from './probe.ts';
import { LATEST_REVIEWS } from './queries.ts';
let checks = 0;
function check(name: string, fn: () => void): void {
@@ -185,31 +195,244 @@ check('limits come from NUT-04 methods', () => {
assert.equal(parseLimits(undefined), null);
});
// The dedupe rule lives in SQL, so exercise it against a real in-memory database
// using the same statement the API uses.
check('one review per author per mint, newest wins', async () => {
const { default: Database } = await import('better-sqlite3');
const db = new Database(':memory:');
db.exec(`CREATE TABLE reviews (event_id TEXT PRIMARY KEY, mint_url TEXT, pubkey TEXT,
rating INTEGER, created_at INTEGER)`);
const insert = db.prepare('INSERT INTO reviews VALUES (?,?,?,?,?)');
insert.run('e1', 'https://m', 'alice', 1, 100);
insert.run('e2', 'https://m', 'alice', 5, 200); // alice changed her mind
insert.run('e3', 'https://m', 'bob', 3, 150);
/* ---------- fedimint ---------- */
const rows = db
.prepare(
`SELECT pubkey, rating FROM (
SELECT pubkey, rating, ROW_NUMBER() OVER (
PARTITION BY mint_url, pubkey ORDER BY created_at DESC, event_id) rn
FROM reviews) WHERE rn = 1 ORDER BY pubkey`,
)
.all() as { pubkey: string; rating: number }[];
/**
* A real kind 38173 from the relay pool, kept verbatim.
*
* Three things in it disagree with a plain reading of NIP-87 and are exactly why this
* fixture is a copy rather than something written to match the spec: `n` is `bitcoin`
* and not `mainnet`, `modules` are short names and not prose, and `content` carries
* `federation_name` rather than a kind-0 `name`.
*/
const REAL_ANNOUNCEMENT = {
id: 'f1',
pubkey: '141d2053cb29535ad45aa9e865cdec492524f0ec0066496b98b7099daab5d658',
kind: 38173,
created_at: 1_756_224_071,
content: '{"federation_name":"E-Cash Club","meta_external_url":"https://fm.ctrb.io/meta.json"}',
tags: [
['d', 'aeca6cc80ffc530bd2d54b09681f6edb9a415c362e4af2fe3d5e04137006fa21'],
[
'u',
'fed11qvqzggnhwden5te0v9cxjtn9vd3jue3wvfkxjmnyva6kzunyd9skutnwv46z7qqqzc28wumn8ghj7' +
'end9e3hgunz9e5k7tmhwvhszqfq4m9xejq0l3fsh5k4fvyks8mwmwdyzhpk9e909l3atczpxuqxlgss2f35eg',
],
['n', 'bitcoin'],
['modules', 'ln,mint,wallet,meta'],
],
};
assert.equal(rows.length, 2, 'alice must count once');
assert.equal(rows.find((r) => r.pubkey === 'alice')?.rating, 5, 'newest review wins');
db.close();
check('a real fedimint announcement parses into the fields the pages render', () => {
const a = parseFedimintAnnouncement(REAL_ANNOUNCEMENT);
assert.ok(a, 'the announcement every other check builds on must parse');
assert.equal(a.federationId, 'aeca6cc80ffc530bd2d54b09681f6edb9a415c362e4af2fe3d5e04137006fa21');
assert.equal(a.inviteCodes.length, 1);
assert.deepEqual(a.modules, ['ln', 'mint', 'wallet', 'meta']);
// `bitcoin` on the wire is `mainnet` here, so one federation cannot appear on two
// networks depending on which word its operator used.
assert.equal(a.network, 'mainnet');
assert.equal(a.name, 'E-Cash Club');
assert.equal(a.announcerPubkey, REAL_ANNOUNCEMENT.pubkey);
});
// The last check is async, so report after the microtask queue drains.
queueMicrotask(() => console.log(`ok, ${checks} checks passed`));
check('an invite code survives the validator that first rejected every real one', () => {
// `fed1` is the human-readable part and the `1` after it is bech32's separator, so
// every real code starts `fed11`. Anchoring on a bech32 data character after `fed1`
// rejected all fifteen announcements on the network.
const code = REAL_ANNOUNCEMENT.tags[1]?.[1] ?? '';
assert.deepEqual(cleanInviteCodes([code]), [code]);
assert.deepEqual(cleanInviteCodes([code, code.toUpperCase()]), [code], 'deduped case-insensitively');
assert.deepEqual(cleanInviteCodes(['https://mint.example.com', 'fed1', '']), []);
});
check('an announcement with no invite code is not a federation anyone can join', () => {
// This is a real event too: someone published a review as a 38173. It has a valid
// `d` and nothing to join, so there is no row to make from it.
assert.equal(
parseFedimintAnnouncement({
id: 'x', pubkey: 'p', kind: 38173, created_at: 1, content: '[5/5] test',
tags: [['d', '412d2a9338ebeee5957382eb06eac07fa5235087b5a7d5d0a6e18c635394e9ed'], ['n', 'mainnet'], ['k', '38173']],
}),
null,
);
});
check('the fedimint slug and its key round trip', () => {
const id = 'aeca6cc80ffc530bd2d54b09681f6edb9a415c362e4af2fe3d5e04137006fa21';
assert.equal(fedimintSlug(id), 'fed-aeca6cc80ffc530b');
assert.equal(fedimintSlug(id).length, 'fed-'.length + 16);
assert.equal(federationIdFromKey(fedimintKey(id)), id);
// A mint URL is not a federation key, and must never be read as one.
assert.equal(federationIdFromKey('https://mint.example.com'), null);
});
check('versioned modules still satisfy the named rows', () => {
// A federation running lnv2 and no ln can do Lightning. A row reading "not
// supported" beside an lnv2 chip would be false.
const modules = parseModules('lnv2,mintv2,walletv2,meta');
assert.ok(hasModule(modules, 'lightning'));
assert.ok(hasModule(modules, 'mint'));
assert.ok(hasModule(modules, 'wallet'));
assert.ok(!hasModule(parseModules('meta'), 'lightning'));
});
check('both spellings of mainnet arrive as one network', () => {
assert.equal(normalizeNetwork('bitcoin'), 'mainnet');
assert.equal(normalizeNetwork('mainnet'), 'mainnet');
assert.equal(normalizeNetwork('signet'), 'signet');
assert.equal(normalizeNetwork(''), null);
assert.equal(normalizeNetwork('<script>'), null);
});
check('the k tag decides which ecosystem a review belongs to', () => {
const review = (tags: string[][]) => ({ id: 'a', pubkey: 'b', kind: 38000, created_at: 1, content: '', tags });
assert.equal(reviewEcosystem(review([['k', '38172']])), 'cashu');
assert.equal(reviewEcosystem(review([['k', '38173']])), 'fedimint');
// No `k` at all is Cashu, and that is a fact about the past rather than a guess:
// every such event predates this site listing anything else. Most reviews on the
// network are this shape.
assert.equal(reviewEcosystem(review([['u', 'https://mint.example.com']])), 'cashu');
// A kind this build has no ecosystem for is nobody's: it is dropped, not filed.
assert.equal(reviewEcosystem(review([['k', '39999']])), null);
assert.equal(ecosystemForKind(38173), 'fedimint');
});
check('a federation never carries a capability warning it cannot know', () => {
const now = 1_800_000_000;
const base = { type: 'fedimint', status: 'announced', last_online: null, last_review_at: null };
// Announced this week: nothing has had the chance to check it, so nothing is said.
assert.deepEqual(getMintWarnings({ ...base, announced_at: now - 86400 }, { now }), []);
// Announced two months ago and still unconfirmed: the one soft banner.
const stale = getMintWarnings({ ...base, announced_at: now - 60 * 86400 }, { now });
assert.equal(stale[0]?.kind, 'never-confirmed');
assert.equal(stale[0]?.severity, 'warning');
assert.equal(stale.length, 1, 'a federation gets one banner at most');
// Reported down by a real check: critical, and dated from the last time it answered.
const down = getMintWarnings(
{ ...base, status: 'offline', last_online: now - 10 * 86400 },
{ now },
);
assert.equal(down[0]?.kind, 'fedimint-offline');
assert.equal(down[0]?.severity, 'critical');
// And none of the NUT-derived banners can ever appear on one, whatever is passed.
const kinds = new Set([...stale, ...down].map((w) => w.kind));
for (const forbidden of ['melt-only', 'melt-disabled', 'frozen', 'gone', 'offline-long']) {
assert.ok(!kinds.has(forbidden as never), `${forbidden} has no fedimint meaning`);
}
});
check('a cashu mint is unaffected by any of the above', () => {
const now = 1_800_000_000;
// The same input with no `type` and with `type: 'cashu'` must agree, because every
// existing caller passes neither.
const mint = {
status: 'offline',
last_online: now - 40 * 86400,
info: { nuts: { '4': { disabled: true }, '5': { methods: [{ method: 'bolt11' }] } } },
};
const untyped = getMintWarnings(mint, { now });
const typed = getMintWarnings({ ...mint, type: 'cashu' }, { now });
assert.deepEqual(untyped, typed);
assert.equal(untyped[0]?.kind, 'gone');
});
check('announced federations sort between confirmed-up and confirmed-down', () => {
// Not knowing is not the same as knowing otherwise: an announced federation must not
// sit with the rows a check confirmed, nor sink below the ones a check failed on.
const sorted = [
{ status: 'offline', score: 4.9, review_count: 99 },
{ status: 'announced', score: 2.7, review_count: 0 },
{ status: 'online', score: 1.2, review_count: 1 },
].sort(compareMints);
assert.deepEqual(sorted.map((m) => m.status), ['online', 'announced', 'offline']);
});
// Dialect portability. Every statement in the API is written once and run against both
// backends, so the translation and the two spellings that differ get their own checks.
check('placeholders translate to $n without touching quoted text', () => {
assert.equal(
toPgPlaceholders('SELECT ? FROM t WHERE a = ? AND b = ?'),
'SELECT $1 FROM t WHERE a = $2 AND b = $3',
);
// A literal question mark is data, not a placeholder.
assert.equal(toPgPlaceholders("SELECT '?' , ?"), "SELECT '?' , $1");
// Postgres doubles a quote to escape it; the closing quote just opens the next literal.
assert.equal(toPgPlaceholders("SELECT 'it''s ?' , ?"), "SELECT 'it''s ?' , $1");
assert.equal(toPgPlaceholders('SELECT "od?d", ?'), 'SELECT "od?d", $1');
});
check('connection targets are read the way people write them', () => {
assert.equal(parseDbTarget('postgres://u:p@h:5432/db').dialect, 'postgres');
assert.equal(parseDbTarget('postgresql://h/db').dialect, 'postgres');
assert.equal(parseDbTarget('/var/lib/cashumints/x.db').dialect, 'sqlite');
assert.equal(parseDbTarget('sqlite:/var/lib/x.db').file, '/var/lib/x.db');
assert.equal(parseDbTarget('sqlite:///var/lib/x.db').file, '/var/lib/x.db');
assert.equal(parseDbTarget('file:./data/x.db').file, path.resolve('./data/x.db'));
assert.equal(parseDbTarget(':memory:').file, ':memory:');
// A password must never reach a log line.
assert.ok(!parseDbTarget('postgres://u:hunter2@h/db').label.includes('hunter2'));
// An unsupported scheme must be an error, never a relative filename: a silent new
// empty SQLite file looks exactly like a successful boot with all the data gone.
assert.throws(() => parseDbTarget('mysql://h/db'), /Cannot tell what database/);
assert.throws(() => parseDbTarget('redis://h'), /Cannot tell what database/);
assert.throws(() => parseDbTarget('sqlite:'), /no path after the scheme/);
assert.throws(() => parseDbTarget(' '), /Empty database target/);
// A schemeless value is a path, including a bare relative filename.
assert.equal(parseDbTarget('cashumints.db').file, path.resolve('cashumints.db'));
});
/**
* The dedupe rule lives in SQL, so exercise the real statement — the one queries.ts
* ships — against a real database. Set CHECK_DB_URL to run it on Postgres too; that is
* what catches a statement that quietly went SQLite-only.
*/
async function checkDedupe(target: string): Promise<void> {
const db = await openDb(parseDbTarget(target));
try {
await db.run('DELETE FROM reviews');
await db.run(
`INSERT INTO reviews (event_id, mint_url, pubkey, rating, k, created_at)
VALUES (?,?,?,?,?,?), (?,?,?,?,?,?), (?,?,?,?,?,?), (?,?,?,?,?,?)`,
'e1', 'https://m', 'alice', 1, '38172', 100,
'e2', 'https://m', 'alice', 5, '38172', 200, // alice changed her mind
'e3', 'https://m', 'bob', 3, '38172', 150,
// Same author, a different subject: one npub reviewing a mint and a federation
// is two reviews, and the dedupe partitions on the subject as well as the author.
'e4', 'fedimint:aeca', 'alice', 4, '38173', 210,
);
const rows = await db.all<{ pubkey: string; rating: number }>(
`SELECT pubkey, rating FROM (${LATEST_REVIEWS}) AS latest WHERE mint_url = ? ORDER BY pubkey`,
'https://m',
);
assert.equal(rows.length, 2, `${db.dialect}: alice must count once`);
assert.equal(rows.find((r) => r.pubkey === 'alice')?.rating, 5, `${db.dialect}: newest wins`);
const federation = await db.all<{ pubkey: string; rating: number }>(
`SELECT pubkey, rating FROM (${LATEST_REVIEWS}) AS latest WHERE mint_url = ?`,
'fedimint:aeca',
);
assert.equal(federation.length, 1, `${db.dialect}: alice's federation review stands alone`);
assert.equal(federation[0]?.rating, 4, `${db.dialect}: and is not collapsed into her mint one`);
checks++;
} catch (err) {
console.error(`FAIL: one review per author per mint, newest wins (${db.dialect})`);
throw err;
} finally {
await db.close();
}
}
await checkDedupe(':memory:');
const pgUrl = process.env['CHECK_DB_URL'];
if (pgUrl) await checkDedupe(pgUrl);
else console.log('note: set CHECK_DB_URL to a postgres:// database to check it there too');
console.log(`ok, ${checks} checks passed`);
+95 -1
View File
@@ -1,6 +1,8 @@
import { fileURLToPath } from 'node:url';
import path from 'node:path';
import { DEFAULT_RELAYS } from '@cashumints/shared';
import type { Dialect } from './db-driver.ts';
import { redact } from './db-postgres.ts';
const here = path.dirname(fileURLToPath(import.meta.url));
const apiRoot = path.resolve(here, '..');
@@ -12,9 +14,101 @@ function int(name: string, fallback: number): number {
return Number.isFinite(n) && n > 0 ? n : fallback;
}
export interface DbConfig {
dialect: Dialect;
/** SQLite file. Empty when the dialect is postgres. */
file: string;
/** libpq connection string. Empty when the dialect is sqlite. */
url: string;
/** Postgres connections held open. Ignored by SQLite, which has one. */
poolMax: number;
/** Safe to log: a Postgres password is replaced with `***`. */
label: string;
}
export const DEFAULT_DB_FILE = path.join(apiRoot, 'data', 'cashumints.db');
/**
* Read one connection target.
*
* Accepts what a person is likely to type or paste:
*
* postgres://user:pw@host:5432/cashumints postgres
* postgresql://… postgres
* sqlite:/var/lib/cashumints/cashumints.db sqlite, absolute
* sqlite://./data/cashumints.db sqlite, relative to the working directory
* file:./data/cashumints.db sqlite
* /var/lib/cashumints/cashumints.db sqlite, a bare path
* :memory: sqlite, throwaway
*
* Throws on anything else rather than guessing, because the wrong guess here is a
* second empty database that looks like data loss.
*/
export function parseDbTarget(raw: string, poolMax = 10): DbConfig {
const value = raw.trim();
if (!value) throw new Error('Empty database target.');
const sqlite = (file: string): DbConfig => {
if (!file) throw new Error(`Database target has no path after the scheme: ${value}`);
const resolved = file === ':memory:' ? file : path.resolve(file);
return { dialect: 'sqlite', file: resolved, url: '', poolMax, label: `sqlite:${resolved}` };
};
// A scheme is checked before anything else, so an unsupported one is an error rather
// than a filename. Left to a "does it look like a path?" heuristic, `mysql://db/x`
// reads as a relative path and the API starts on a brand new empty SQLite file — data
// loss that announces itself as a successful boot.
const scheme = /^([a-z][a-z0-9+.-]*):/i.exec(value)?.[1]?.toLowerCase();
switch (scheme) {
case undefined:
return sqlite(value); // A bare path, absolute or relative.
case 'postgres':
case 'postgresql':
return { dialect: 'postgres', file: '', url: value, poolMax, label: redact(value) };
case 'sqlite':
case 'sqlite3':
case 'file':
return sqlite(value.replace(/^[a-z0-9+.-]+:(?:\/\/)?/i, ''));
default:
// `:memory:` has no scheme by this reading — the regex needs a letter first.
if (value === ':memory:') return sqlite(value);
throw new Error(
`Cannot tell what database "${value}" means. Use a postgres:// URL, ` +
'a sqlite: path, or a filesystem path.',
);
}
}
/**
* Where this process keeps its data.
*
* `DATABASE_URL` decides the backend. Without it the API stays on SQLite at `DB_PATH`,
* which is what every existing deployment already has, so adding Postgres support
* changed nothing for anyone who does not ask for it.
*/
export function resolveDbConfig(env: NodeJS.ProcessEnv = process.env): DbConfig {
const poolMax = int('DB_POOL_MAX', 10);
const url = env['DATABASE_URL']?.trim();
if (url) return parseDbTarget(url, poolMax);
const file = path.resolve(env['DB_PATH'] ?? DEFAULT_DB_FILE);
return { dialect: 'sqlite', file, url: '', poolMax, label: `sqlite:${file}` };
}
let dbConfig: DbConfig | null = null;
export const config = {
port: int('PORT', 8787),
dbPath: process.env['DB_PATH'] ?? path.join(apiRoot, 'data', 'cashumints.db'),
/**
* Resolved on first use rather than at import, so a process can still redirect itself
* — `test:offline` points at a throwaway copy by setting DATABASE_URL before anything
* opens a connection, and would otherwise rug the live database instead.
*/
get db(): DbConfig {
dbConfig ??= resolveDbConfig();
return dbConfig;
},
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[],
+44
View File
@@ -0,0 +1,44 @@
/**
* The one interface every backend implements.
*
* SQL is written once, in the portable subset both SQLite and Postgres speak, and
* bound with `?` placeholders. The Postgres driver rewrites those to `$1..$n`; the
* SQLite driver passes them straight through. See db-schema.ts for the rules that
* keep a statement portable.
*
* Everything is async because `pg` is async. better-sqlite3 is not, so its driver
* resolves already-settled promises — the cost is a microtask per query, which is
* nothing next to the HTTP and relay work around it.
*/
export type Dialect = 'sqlite' | 'postgres';
/** What a statement can be bound to. `undefined` is rejected by both drivers. */
export type Param = string | number | null;
export interface RunResult {
/** Rows inserted, updated or deleted. `ON CONFLICT DO NOTHING` that did nothing is 0. */
changes: number;
}
/** The read/write surface. A transaction hands back one of these bound to its connection. */
export interface Sql {
all<T>(sql: string, ...params: Param[]): Promise<T[]>;
get<T>(sql: string, ...params: Param[]): Promise<T | undefined>;
run(sql: string, ...params: Param[]): Promise<RunResult>;
}
export interface Db extends Sql {
readonly dialect: Dialect;
/** Connection description safe to log: a Postgres password is never in it. */
readonly label: string;
/** Run DDL. Each element is one statement, no placeholders. */
exec(statements: readonly string[]): Promise<void>;
/**
* Run `fn` inside a transaction. Statements issued through the handle `fn` receives
* are part of it and see its uncommitted writes; anything issued through the Db
* itself waits until the transaction settles. A throw rolls back and re-throws.
*/
transaction<T>(fn: (tx: Sql) => Promise<T>): Promise<T>;
close(): Promise<void>;
}
+147
View File
@@ -0,0 +1,147 @@
import pg from 'pg';
import type { Db, Param, RunResult, Sql } from './db-driver.ts';
const { Pool, types } = pg;
/**
* node-postgres hands back BIGINT and NUMERIC as strings, because either can hold a
* value JavaScript's number cannot. Nothing here can: the counters are row counts, the
* timestamps are unix seconds, and the averages are ratings between 1 and 5. Left
* as-is the strings would flow straight into the JSON the API serves, turning
* `review_count: 12` into `"12"` and quietly breaking every comparison on the way.
*/
types.setTypeParser(20, (v) => Number.parseInt(v, 10)); // int8
types.setTypeParser(1700, (v) => Number.parseFloat(v)); // numeric
/**
* Rewrite `?` placeholders to `$1..$n`.
*
* Quoted text is copied verbatim so a `?` inside a string literal is not renumbered.
* Postgres escapes a quote by doubling it, which needs no special case: the closing
* quote ends one literal and the next character opens another.
*/
export function toPgPlaceholders(sql: string): string {
let out = '';
let quote: string | null = null;
let n = 0;
for (const ch of sql) {
if (quote !== null) {
if (ch === quote) quote = null;
out += ch;
} else if (ch === "'" || ch === '"') {
quote = ch;
out += ch;
} else if (ch === '?') {
out += `$${++n}`;
} else {
out += ch;
}
}
return out;
}
const translated = new Map<string, string>();
function translate(sql: string): string {
let pgSql = translated.get(sql);
if (pgSql === undefined) {
pgSql = toPgPlaceholders(sql);
translated.set(sql, pgSql);
}
return pgSql;
}
/** Bind a Sql surface to anything that can run a query: the pool, or one client. */
function surface(run: (text: string, values: Param[]) => Promise<pg.QueryResult>): Sql {
return {
async all<T>(sql: string, ...params: Param[]): Promise<T[]> {
return (await run(translate(sql), params)).rows as T[];
},
async get<T>(sql: string, ...params: Param[]): Promise<T | undefined> {
return (await run(translate(sql), params)).rows[0] as T | undefined;
},
async run(sql: string, ...params: Param[]): Promise<RunResult> {
return { changes: (await run(translate(sql), params)).rowCount ?? 0 };
},
};
}
/** Strip the password so a connection string is safe to put in a log line. */
export function redact(url: string): string {
try {
const parsed = new URL(url);
if (parsed.password) parsed.password = '***';
return parsed.toString();
} catch {
return 'postgres';
}
}
export class PostgresDb implements Db {
readonly dialect = 'postgres' as const;
readonly label: string;
#pool: pg.Pool;
#sql: Sql;
constructor(connectionString: string, poolMax: number) {
this.#pool = new Pool({ connectionString, max: poolMax });
// A pool client can die between checkouts — a Postgres restart, a killed session, an
// idle timeout on a proxy. Without a listener node-postgres emits that on the pool as
// an unhandled 'error' and takes the process down with it.
this.#pool.on('error', () => undefined);
this.#sql = surface((text, values) => this.#pool.query(text, values));
this.label = redact(connectionString);
}
all<T>(sql: string, ...params: Param[]): Promise<T[]> {
return this.#sql.all<T>(sql, ...params);
}
get<T>(sql: string, ...params: Param[]): Promise<T | undefined> {
return this.#sql.get<T>(sql, ...params);
}
run(sql: string, ...params: Param[]): Promise<RunResult> {
return this.#sql.run(sql, ...params);
}
async exec(statements: readonly string[]): Promise<void> {
// One connection for the lot: concurrent CREATE ... IF NOT EXISTS on the same
// catalog rows deadlocks in Postgres, and this runs on every boot.
const client = await this.#pool.connect();
try {
for (const sql of statements) await client.query(sql);
} finally {
client.release();
}
}
async transaction<T>(fn: (tx: Sql) => Promise<T>): Promise<T> {
const client = await this.#pool.connect();
const tx = surface((text, values) => client.query(text, values));
try {
await client.query('BEGIN');
const result = await fn(tx);
await client.query('COMMIT');
return result;
} catch (err) {
try {
await client.query('ROLLBACK');
} catch {
// The connection is already unusable; release() below discards it.
}
throw err;
} finally {
client.release();
}
}
close(): Promise<void> {
return this.#pool.end();
}
}
+120
View File
@@ -0,0 +1,120 @@
/**
* One schema for both backends.
*
* The types below were picked because SQLite and Postgres agree on all of them:
* TEXT, INTEGER and SMALLINT are literal Postgres types and map to SQLite's INTEGER
* and TEXT affinities, and BIGINT is what unix seconds need in Postgres, where plain
* INTEGER runs out in 2038. `CREATE TABLE/INDEX IF NOT EXISTS` is spelled the same in
* both, so first boot on an empty Postgres database needs no separate setup step.
*
* Rules for keeping a statement portable, learned from the ones that were not:
* - Bind with `?`. The Postgres driver rewrites to `$1..$n`.
* - A subquery in FROM needs an alias. Postgres rejects it without one.
* - Count a condition with `COUNT(*) FILTER (WHERE ...)`, never `SUM(cond)`:
* Postgres has no implicit boolean-to-integer cast.
* - `INSERT OR IGNORE` is SQLite-only. `ON CONFLICT DO NOTHING` works in both and,
* with no conflict target, covers every unique constraint on the table.
*/
export const TABLES = ['mints', 'reviews', 'probes', 'state'] as const;
export type TableName = (typeof TABLES)[number];
export const SCHEMA: readonly string[] = [
`CREATE TABLE IF NOT EXISTS mints (
url TEXT PRIMARY KEY,
host TEXT NOT NULL,
type TEXT NOT NULL DEFAULT 'cashu',
name TEXT,
description TEXT,
icon_url TEXT,
icon_file TEXT,
pubkey TEXT,
info_json TEXT,
ecosystem_json TEXT,
nuts_json TEXT,
version TEXT,
status TEXT NOT NULL DEFAULT 'unknown',
consecutive_fails INTEGER NOT NULL DEFAULT 0,
last_online BIGINT,
last_probe BIGINT,
first_seen BIGINT NOT NULL,
updated_at BIGINT NOT NULL
)`,
`CREATE UNIQUE INDEX IF NOT EXISTS idx_mints_host ON mints(host)`,
`CREATE INDEX IF NOT EXISTS idx_mints_pubkey ON mints(pubkey)`,
`CREATE TABLE IF NOT EXISTS reviews (
event_id TEXT PRIMARY KEY,
mint_url TEXT NOT NULL,
pubkey TEXT NOT NULL,
rating INTEGER,
k TEXT,
created_at BIGINT NOT NULL
)`,
`CREATE INDEX IF NOT EXISTS idx_reviews_mint ON reviews(mint_url, created_at)`,
`CREATE INDEX IF NOT EXISTS idx_reviews_author ON reviews(mint_url, pubkey, created_at)`,
`CREATE TABLE IF NOT EXISTS probes (
mint_url TEXT NOT NULL,
ts BIGINT NOT NULL,
ok SMALLINT NOT NULL,
latency_ms INTEGER
)`,
`CREATE INDEX IF NOT EXISTS idx_probes_mint ON probes(mint_url, ts)`,
`CREATE TABLE IF NOT EXISTS state (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
)`,
];
/** Column order used by the migrator, so source and target line up without guessing. */
export const COLUMNS: Record<TableName, readonly string[]> = {
mints: [
'url', 'host', 'type', 'name', 'description', 'icon_url', 'icon_file', 'pubkey',
'info_json', 'ecosystem_json', 'nuts_json', 'version', 'status', 'consecutive_fails',
'last_online', 'last_probe', 'first_seen', 'updated_at',
],
reviews: ['event_id', 'mint_url', 'pubkey', 'rating', 'k', 'created_at'],
probes: ['mint_url', 'ts', 'ok', 'latency_ms'],
state: ['key', 'value'],
};
/**
* Primary key per table, for the migrator's upsert. `probes` has none — it is an
* append-only log with no natural key — so the migrator replaces it wholesale.
*/
export const PRIMARY_KEY: Record<TableName, string | null> = {
mints: 'url',
reviews: 'event_id',
probes: null,
state: 'key',
};
/**
* Statements that bring an existing database up to the schema above.
*
* `CREATE TABLE IF NOT EXISTS` does nothing to a table that is already there, so a
* database seeded before federations were indexed never grows `mints.type` on its own.
* These run once per boot, after `SCHEMA`, and are the whole of this project's
* migration story.
*
* The `ADD COLUMN`s are spelled without `IF NOT EXISTS` on purpose: Postgres accepts it
* there and SQLite does not, so the portable form is to run the statement and swallow
* the "duplicate column" both backends raise. `applySchemaUpgrades` in db.ts does that.
*
* Order matters, and the index at the end is why this list exists rather than the index
* sitting with its table in `SCHEMA`: on an upgraded database the column is created
* here, so an index over it declared any earlier fails on a column that is not there
* yet — which is exactly how the first version of this failed to boot.
*
* Every statement must be re-runnable and must not need a value: a NOT NULL column
* needs a constant DEFAULT (which is what makes existing rows read as Cashu), and
* anything else is nullable.
*/
export const SCHEMA_UPGRADES: readonly string[] = [
`ALTER TABLE mints ADD COLUMN type TEXT NOT NULL DEFAULT 'cashu'`,
`ALTER TABLE mints ADD COLUMN ecosystem_json TEXT`,
`ALTER TABLE reviews ADD COLUMN k TEXT`,
`CREATE INDEX IF NOT EXISTS idx_mints_type ON mints(type, status)`,
];
+124
View File
@@ -0,0 +1,124 @@
import Database from 'better-sqlite3';
import fs from 'node:fs';
import path from 'node:path';
import type { Db, Param, RunResult, Sql } from './db-driver.ts';
/**
* better-sqlite3 behind the async Db interface.
*
* Every statement runs synchronously; the promises are already settled by the time
* they are returned. Prepared statements are cached by SQL text, which the callers
* rely on: they pass constant strings and would otherwise re-prepare on every request.
*
* The one subtlety is transactions. better-sqlite3's own `db.transaction()` cannot
* wrap an async function, so this driver issues BEGIN/COMMIT itself — and that opens a
* window the synchronous version did not have, where an unrelated `await`ing caller
* could slip a statement between them and have it committed, or rolled back, with the
* transaction. `#openTx` closes the window: statements issued through the Db while a
* transaction is open queue behind it, and only the handle passed to the transaction
* body bypasses the queue. Transactions themselves are serialized for the same reason.
*/
export class SqliteDb implements Db {
readonly dialect = 'sqlite' as const;
readonly label: string;
#db: Database.Database;
#statements = new Map<string, Database.Statement>();
/** Resolves when the transaction in flight settles. Null when none is open. */
#openTx: Promise<unknown> | null = null;
constructor(file: string) {
if (file !== ':memory:') fs.mkdirSync(path.dirname(path.resolve(file)), { recursive: true });
this.#db = new Database(file);
this.#db.pragma('journal_mode = WAL');
this.#db.pragma('synchronous = NORMAL');
this.#db.pragma('busy_timeout = 5000');
this.label = `sqlite:${file}`;
}
#prepare(sql: string): Database.Statement {
let stmt = this.#statements.get(sql);
if (!stmt) {
stmt = this.#db.prepare(sql);
this.#statements.set(sql, stmt);
}
return stmt;
}
/** The synchronous core, shared by the Db surface and the transaction handle. */
#direct: Sql = {
all: <T>(sql: string, ...params: Param[]): Promise<T[]> =>
Promise.resolve(this.#prepare(sql).all(...params) as T[]),
get: <T>(sql: string, ...params: Param[]): Promise<T | undefined> =>
Promise.resolve(this.#prepare(sql).get(...params) as T | undefined),
run: (sql: string, ...params: Param[]): Promise<RunResult> =>
Promise.resolve({ changes: this.#prepare(sql).run(...params).changes }),
};
/** Wait out any open transaction, so a statement never lands inside someone else's. */
async #settled(): Promise<void> {
while (this.#openTx) await this.#openTx.catch(() => undefined);
}
async all<T>(sql: string, ...params: Param[]): Promise<T[]> {
await this.#settled();
return this.#direct.all<T>(sql, ...params);
}
async get<T>(sql: string, ...params: Param[]): Promise<T | undefined> {
await this.#settled();
return this.#direct.get<T>(sql, ...params);
}
async run(sql: string, ...params: Param[]): Promise<RunResult> {
await this.#settled();
return this.#direct.run(sql, ...params);
}
exec(statements: readonly string[]): Promise<void> {
for (const sql of statements) this.#db.exec(sql);
return Promise.resolve();
}
transaction<T>(fn: (tx: Sql) => Promise<T>): Promise<T> {
// Claim the slot synchronously, before the first await: two transactions started in
// the same tick would otherwise both find no transaction open and issue a nested
// BEGIN, which SQLite rejects outright.
const previous = this.#openTx;
const work = (async (): Promise<T> => {
if (previous) await previous.catch(() => undefined);
// IMMEDIATE takes the write lock up front rather than on the first write, so two
// writers fail fast against busy_timeout instead of deadlocking mid-transaction.
this.#db.exec('BEGIN IMMEDIATE');
try {
const result = await fn(this.#direct);
this.#db.exec('COMMIT');
return result;
} catch (err) {
try {
this.#db.exec('ROLLBACK');
} catch {
// Already rolled back by SQLite. The original error is the one that matters.
}
throw err;
}
})();
this.#openTx = work;
return work.finally(() => {
// Only the last transaction in the chain clears the slot; anything queued behind
// this one has already replaced it.
if (this.#openTx === work) this.#openTx = null;
});
}
close(): Promise<void> {
this.#statements.clear();
this.#db.close();
return Promise.resolve();
}
}
+103 -79
View File
@@ -1,101 +1,125 @@
import Database from 'better-sqlite3';
import fs from 'node:fs';
import path from 'node:path';
import { config } from './config.ts';
import { config, type DbConfig } from './config.ts';
import type { Db } from './db-driver.ts';
import { PostgresDb } from './db-postgres.ts';
import { SCHEMA, SCHEMA_UPGRADES } from './db-schema.ts';
import { SqliteDb } from './db-sqlite.ts';
const SCHEMA = `
CREATE TABLE IF NOT EXISTS mints (
url TEXT PRIMARY KEY,
host TEXT NOT NULL,
name TEXT,
description TEXT,
icon_url TEXT,
icon_file TEXT,
pubkey TEXT,
info_json TEXT,
nuts_json TEXT,
version TEXT,
status TEXT NOT NULL DEFAULT 'unknown',
consecutive_fails INTEGER NOT NULL DEFAULT 0,
last_online INTEGER,
last_probe INTEGER,
first_seen INTEGER NOT NULL,
updated_at INTEGER NOT NULL
);
CREATE UNIQUE INDEX IF NOT EXISTS idx_mints_host ON mints(host);
CREATE INDEX IF NOT EXISTS idx_mints_pubkey ON mints(pubkey);
export type { Db, Dialect, Param, RunResult, Sql } from './db-driver.ts';
CREATE TABLE IF NOT EXISTS reviews (
event_id TEXT PRIMARY KEY,
mint_url TEXT NOT NULL,
pubkey TEXT NOT NULL,
rating INTEGER,
created_at INTEGER NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_reviews_mint ON reviews(mint_url, created_at);
CREATE INDEX IF NOT EXISTS idx_reviews_author ON reviews(mint_url, pubkey, created_at);
/**
* Open a connection and make sure the schema is there.
*
* Both backends create their tables on first use, so pointing the API at an empty
* Postgres database is the whole setup: no separate migration step to forget.
*/
export async function openDb(target: DbConfig = config.db): Promise<Db> {
const db: Db =
target.dialect === 'postgres'
? new PostgresDb(target.url, target.poolMax)
: new SqliteDb(target.file);
CREATE TABLE IF NOT EXISTS probes (
mint_url TEXT NOT NULL,
ts INTEGER NOT NULL,
ok INTEGER NOT NULL,
latency_ms INTEGER
);
CREATE INDEX IF NOT EXISTS idx_probes_mint ON probes(mint_url, ts);
try {
await db.exec(SCHEMA);
await applySchemaUpgrades(db);
} catch (err) {
await db.close().catch(() => undefined);
throw err;
}
CREATE TABLE IF NOT EXISTS state (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
);
`;
export type DB = Database.Database;
let db: DB | null = null;
export function getDb(): DB {
if (db) return db;
fs.mkdirSync(path.dirname(config.dbPath), { recursive: true });
fs.mkdirSync(config.iconDir, { recursive: true });
const handle = new Database(config.dbPath);
handle.pragma('journal_mode = WAL');
handle.pragma('synchronous = NORMAL');
handle.pragma('busy_timeout = 5000');
handle.exec(SCHEMA);
db = handle;
return handle;
return db;
}
export function closeDb(): void {
db?.close();
db = null;
/**
* Bring an existing database up to the current schema.
*
* `CREATE TABLE IF NOT EXISTS` is a no-op on a table that already exists, so a
* deployment seeded before federations were indexed would otherwise boot against a
* `mints` table with no `type` column and fail on the first query. Running these here
* means an upgrade is a restart, with no separate step to forget.
*
* An already-present column is the expected outcome on every boot after the first, so
* that error is swallowed and anything else is re-thrown. Both backends are matched on
* the message rather than on a code: SQLite raises `duplicate column name: type` and
* Postgres `column "type" of relation "mints" already exists`, and neither exposes
* anything more structured through the drivers here.
*/
async function applySchemaUpgrades(db: Db): Promise<void> {
for (const statement of SCHEMA_UPGRADES) {
try {
await db.exec([statement]);
} catch (err) {
/*
* Postgres localizes its message text (lc_messages), so a German server says
* neither phrase below and every boot after the first would re-throw. Its error
* codes are locale-proof: 42701 duplicate_column, 42P07 duplicate_table. SQLite
* has no code for this, so its English message is still matched.
*/
const code = (err as { code?: unknown }).code;
const message = (err instanceof Error ? err.message : String(err)).toLowerCase();
const alreadyThere =
code === '42701' ||
code === '42P07' ||
message.includes('duplicate column') ||
message.includes('already exists');
if (!alreadyThere) throw err;
}
}
}
export function getState(key: string): string | null {
const row = getDb().prepare('SELECT value FROM state WHERE key = ?').get(key) as
| { value: string }
| undefined;
let opening: Promise<Db> | null = null;
/**
* The process-wide connection, opened on first use.
*
* Async because `pg` is: there is no synchronous way to reach a Postgres server. The
* promise is memoized, so concurrent callers during startup share one connection
* rather than racing to open several.
*/
export function getDb(): Promise<Db> {
if (!opening) {
fs.mkdirSync(config.iconDir, { recursive: true });
opening = openDb().catch((err: unknown) => {
// A failed open must not be cached, or every later call replays the same error
// against a connection that was never established.
opening = null;
throw err;
});
}
return opening;
}
export async function closeDb(): Promise<void> {
const pending = opening;
opening = null;
if (pending) await pending.then((db) => db.close()).catch(() => undefined);
}
export async function getState(key: string): Promise<string | null> {
const db = await getDb();
const row = await db.get<{ value: string }>('SELECT value FROM state WHERE key = ?', key);
return row?.value ?? null;
}
export function setState(key: string, value: string): void {
getDb()
.prepare('INSERT INTO state (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value')
.run(key, value);
export async function setState(key: string, value: string): Promise<void> {
const db = await getDb();
await db.run(
'INSERT INTO state (key, value) VALUES (?, ?) ON CONFLICT (key) DO UPDATE SET value = excluded.value',
key,
value,
);
}
export function getStateNumber(key: string): number | null {
const raw = getState(key);
export async function getStateNumber(key: string): Promise<number | null> {
const raw = await getState(key);
if (raw === null) return null;
const n = Number.parseInt(raw, 10);
return Number.isFinite(n) ? n : null;
}
/** Delete probe rows older than 90 days. Called once a day. */
export function pruneProbes(): number {
export async function pruneProbes(): Promise<number> {
const cutoff = Math.floor(Date.now() / 1000) - 90 * 24 * 60 * 60;
return getDb().prepare('DELETE FROM probes WHERE ts < ?').run(cutoff).changes;
const db = await getDb();
return (await db.run('DELETE FROM probes WHERE ts < ?', cutoff)).changes;
}
+352 -55
View File
@@ -1,6 +1,8 @@
import { SimplePool } from 'nostr-tools/pool';
import type { Event as NostrEvent, Filter } from 'nostr-tools';
import {
ANNOUNCEMENT_KINDS,
KIND_FEDIMINT_ANNOUNCEMENT,
KIND_MINT_ANNOUNCEMENT,
KIND_REVIEW,
isCashuMintReview,
@@ -8,12 +10,18 @@ import {
mintUrlSpellings,
mintUrlsFromEvent,
normalizeMintUrl,
parseFedimintAnnouncement,
parseRating,
reviewEcosystem,
reviewTargetId,
reviewTargetKind,
type FedimintAnnouncement,
} from '@cashumints/shared';
import { config } from './config.ts';
import { getDb, getStateNumber, setState } from './db.ts';
import { getDb, setState, getStateNumber } from './db.ts';
import type { Sql } from './db-driver.ts';
import { log } from './log.ts';
import { insertMintIfNew, mintUrlByHost, mintUrlByPubkey } from './mints.ts';
import { insertMintIfNew, upsertFedimint } from './mints.ts';
const QUERY_LIMIT = 500;
const MAX_WAIT_MS = 12_000;
@@ -101,18 +109,40 @@ async function fetchKind(kind: number, since: number | null): Promise<NostrEvent
* count in the API disagrees with the list the browser renders on the same page.
*/
async function fetchReviewsForMint(
mintUrl: string,
mintPubkey: string | null,
target: ReviewTarget,
since: number | null,
): Promise<NostrEvent[]> {
const filters: Filter[] = [];
const base: Filter = { kinds: [KIND_REVIEW], limit: QUERY_LIMIT };
if (since !== null) base.since = since;
if (mintPubkey) {
filters.push({ ...base, '#d': [mintPubkey], '#k': [String(KIND_MINT_ANNOUNCEMENT)] });
if (target.type === 'fedimint') {
/*
* A federation is asked for by id, and only by id.
*
* The mint side asks by `#u` as well, and the obvious translation of that is `#u`
* with the invite codes. Two of the five relays in the pool answer such a filter
* with `ERROR: bad req: filter item too large` and nothing else — an invite code is
* 150+ characters — so the query costs two round trips and returns less than
* sending nothing would.
*
* It also buys nothing today: every kind 38000 with `k` = 38173 on the network
* carries the federation id in `d` and no `u` tag at all. A review that arrives
* with only an invite code is still matched, by the global sweep, through
* `byInvite` in the index above; it just is not asked for by name.
*/
if (!target.federationId) return [];
filters.push({
...base,
'#d': [target.federationId],
'#k': [String(KIND_FEDIMINT_ANNOUNCEMENT)],
});
} else {
if (target.pubkey) {
filters.push({ ...base, '#d': [target.pubkey], '#k': [String(KIND_MINT_ANNOUNCEMENT)] });
}
filters.push({ ...base, '#u': mintUrlSpellings(target.url) });
}
filters.push({ ...base, '#u': mintUrlSpellings(mintUrl) });
const batches = await Promise.all(
filters.map((filter) =>
@@ -125,14 +155,137 @@ async function fetchReviewsForMint(
return batches.flat();
}
/** Every mint URL an event points at, normalized and inserted if new. */
function ingestMintUrls(event: NostrEvent, now: number, newMints: Set<string>): void {
for (const raw of mintUrlsFromEvent(event)) {
const created = insertMintIfNew(raw, now);
/**
* One thing the targeted second pass asks the relays about.
*
* A mint is asked for by URL and by the pubkey it publishes; a federation by its id and
* its invite codes. Nothing else about the row is needed, so the pass reads five
* columns rather than the whole table.
*/
interface ReviewTarget {
type: string;
/** Row key: the mint URL, or `fedimint:<id>`. */
url: string;
pubkey: string | null;
federationId: string | null;
}
/**
* Every federation the batch announces, inserted or refreshed.
*
* Unlike mints, a federation's row content comes entirely from these events — there is
* nothing to fetch afterwards — so this both creates rows and keeps existing ones
* current. Deduped by federation id inside the batch first, newest announcement per id
* winning, because the same federation is routinely announced by several npubs and by
* the same npub more than once, and writing each of those would be one UPDATE per event
* to land on the same final state.
*/
async function ingestFedimints(
events: Iterable<NostrEvent>,
now: number,
newRows: Set<string>,
): Promise<void> {
const newest = new Map<string, FedimintAnnouncement>();
for (const event of events) {
const announcement = parseFedimintAnnouncement(event);
if (!announcement) continue;
const existing = newest.get(announcement.federationId);
if (existing && existing.announcedAt >= announcement.announcedAt) continue;
newest.set(announcement.federationId, announcement);
}
for (const announcement of newest.values()) {
const created = await upsertFedimint(announcement, now);
if (created) newRows.add(created);
}
}
/**
* Every mint URL the batch points at, normalized and inserted if new.
*
* Deduped across the whole batch before anything is written. A popular mint's URL
* appears in hundreds of events per cycle, and the old row-at-a-time version paid one
* INSERT round trip for each of them — invisible against a local SQLite file, minutes
* of latency against a Postgres server.
*/
async function ingestMintUrls(
events: Iterable<NostrEvent>,
now: number,
newMints: Set<string>,
): Promise<void> {
const seen = new Set<string>();
for (const event of events) {
// A review that says it is about a federation carries an invite code in `u`, not a
// mint URL, and feeding those to the mint normalizer is how `fed11…` filled the
// skipped-URL log. (A handful of kind 38172 announcements carry one too, which is
// their publisher's mistake and still logs once.)
if (reviewEcosystem(event) !== 'cashu') continue;
for (const raw of mintUrlsFromEvent(event)) seen.add(raw);
}
for (const raw of seen) {
const created = await insertMintIfNew(raw, now);
if (created) newMints.add(created);
}
}
/**
* The mints table as two lookup maps, read once per cycle.
*
* Review resolution asks "which mint is this?" once per event, and the answer only
* changes when a mint is inserted — which happens in ingestMintUrls, before any of
* this runs. So it is a snapshot, not a query per event.
*/
interface MintIndex {
byHost: Map<string, string>;
byPubkey: Map<string, string>;
/**
* Federation id to row key. A federation has no URL, so its `d` tag is the only thing
* a review can be matched on, and it is matched exactly.
*/
byFederation: Map<string, string>;
/** Invite code to row key, for a review that carries `u` but no usable `d`. */
byInvite: Map<string, string>;
}
async function loadMintIndex(): Promise<MintIndex> {
const db = await getDb();
const rows = await db.all<{
url: string;
host: string;
type: string;
pubkey: string | null;
ecosystem_json: string | null;
}>('SELECT url, host, type, pubkey, ecosystem_json FROM mints');
const byHost = new Map<string, string>();
const byPubkey = new Map<string, string>();
const byFederation = new Map<string, string>();
const byInvite = new Map<string, string>();
for (const row of rows) {
byHost.set(row.host, row.url);
// First writer wins, matching the LIMIT 1 the per-event query used.
if (row.pubkey && !byPubkey.has(row.pubkey)) byPubkey.set(row.pubkey, row.url);
if (row.type !== 'fedimint' || !row.ecosystem_json) continue;
try {
const fields = JSON.parse(row.ecosystem_json) as {
federation_id?: string;
invite_codes?: string[];
};
if (fields.federation_id) byFederation.set(fields.federation_id.toLowerCase(), row.url);
for (const code of fields.invite_codes ?? []) byInvite.set(code.toLowerCase(), row.url);
} catch {
// A row whose ecosystem column will not parse simply cannot be matched by id.
// It still has a page and still resolves by nothing else, which is honest.
}
}
return { byHost, byPubkey, byFederation, byInvite };
}
/**
* Resolve which mint a review is about.
*
@@ -141,23 +294,157 @@ function ingestMintUrls(event: NostrEvent, now: number, newMints: Set<string>):
* the `d` tag pubkey against the pubkey a mint publishes in its own /v1/info, which
* catches reviews whose `u` tag is missing or wrong.
*/
function resolveReviewTarget(event: NostrEvent): string | null {
function resolveReviewTarget(event: NostrEvent, index: MintIndex): string | null {
/*
* The `k` tag decides which ecosystem's resolver runs, and it is not merely a hint:
* a federation id and a Cashu mint pubkey are both 64 hex characters in a `d` tag, so
* without this a Fedimint review could in principle be filed against a mint. A review
* with no `k` at all is Cashu, because every one of those predates anything else.
*/
const ecosystem = reviewEcosystem(event);
if (ecosystem === 'fedimint') return resolveFedimintReview(event, index);
// A `k` naming a kind this build has no ecosystem for: not ours to file.
if (ecosystem === null) return null;
for (const raw of mintUrlsFromEvent(event)) {
const normalized = normalizeMintUrl(raw);
if (!normalized) continue;
// The row was inserted moments ago by ingestMintUrls, so this almost always hits.
const canonical = mintUrlByHost(normalized.host);
const canonical = index.byHost.get(normalized.host);
if (canonical) return canonical;
}
const pubkey = mintPubkeyRef(event);
return pubkey ? mintUrlByPubkey(pubkey) : null;
return pubkey ? index.byPubkey.get(pubkey) ?? null : null;
}
/**
* Which federation a review is about.
*
* `d` first and almost always: every kind 38000 with `k` = 38173 seen on the network
* carries the federation id there and nothing else — no `u`, no `a`. The invite code
* fallback is for publishers that write `u` instead, and matches the code exactly
* rather than decoding it, because an invite code is opaque to this codebase.
*
* A review of a federation this site has never seen announced resolves to nothing and
* is dropped, exactly as a review of an unknown mint is. There is no page to put it on.
*/
function resolveFedimintReview(event: NostrEvent, index: MintIndex): string | null {
const d = reviewTargetId(event);
if (d && /^[0-9a-f]{64}$/i.test(d)) {
const found = index.byFederation.get(d.toLowerCase());
if (found) return found;
}
for (const raw of mintUrlsFromEvent(event)) {
const found = index.byInvite.get(raw.trim().toLowerCase());
if (found) return found;
}
return null;
}
/** Postgres caps a statement at 65535 bounds parameters; this keeps every batch clear of it. */
const INSERT_BATCH = 500;
interface ReviewRow {
eventId: string;
mintUrl: string;
pubkey: string;
rating: number | null;
/** The `k` tag verbatim, or null when the event carried none. */
k: string | null;
createdAt: number;
}
/** Which of these event ids the reviews table already holds. */
async function knownEventIds(ids: string[]): Promise<Set<string>> {
const db = await getDb();
const known = new Set<string>();
for (let i = 0; i < ids.length; i += INSERT_BATCH) {
const chunk = ids.slice(i, i + INSERT_BATCH);
const rows = await db.all<{ event_id: string }>(
`SELECT event_id FROM reviews WHERE event_id IN (${chunk.map(() => '?').join(',')})`,
...chunk,
);
for (const row of rows) known.add(row.event_id);
}
return known;
}
/**
* Write a batch of reviews and report how many were genuinely new.
*
* Deduped by event id first, keeping the last spelling seen. That is not just a saving:
* Postgres refuses an ON CONFLICT DO UPDATE that would touch the same row twice in one
* statement, so a duplicate inside a single VALUES list is an error, not a no-op. The
* targeted second pass below queries by `#d` and by `#u` and routinely returns both.
*
* `newReviews` is a health signal — is discovery finding anything? — so it counts rows
* that did not exist, not rows written. `changes` would count every re-seen event and
* hide a stalled crawl behind a busy-looking number.
*/
async function ingestReviews(events: Iterable<NostrEvent>, index: MintIndex): Promise<number> {
const rows = new Map<string, ReviewRow>();
for (const e of events) {
/*
* The `k` tag is what marks a 38000 as being about a Cashu mint rather than a
* federation or something this site does not list. Events without it but with a
* resolvable mint URL are still accepted: the old publisher's own events are the
* reason, and every one of them is about a mint.
*
* A federation review has already proved itself by resolving: `resolveReviewTarget`
* only returns a federation row for an event whose `k` says 38173, so nothing here
* has to re-test that.
*/
const target = resolveReviewTarget(e, index);
if (!target) continue;
const fedimint = reviewTargetKind(e) === KIND_FEDIMINT_ANNOUNCEMENT;
if (!fedimint && !isCashuMintReview(e) && mintUrlsFromEvent(e).length === 0) continue;
const k = reviewTargetKind(e);
rows.set(e.id, {
eventId: e.id,
mintUrl: target,
pubkey: e.pubkey,
rating: parseRating(e),
k: k === null ? null : String(k),
createdAt: e.created_at,
});
}
if (rows.size === 0) return 0;
const batch = [...rows.values()];
const known = await knownEventIds(batch.map((r) => r.eventId));
const db = await getDb();
await db.transaction(async (tx: Sql) => {
for (let i = 0; i < batch.length; i += INSERT_BATCH) {
const chunk = batch.slice(i, i + INSERT_BATCH);
await tx.run(
`INSERT INTO reviews (event_id, mint_url, pubkey, rating, k, created_at)
VALUES ${chunk.map(() => '(?, ?, ?, ?, ?, ?)').join(', ')}
ON CONFLICT (event_id) DO UPDATE SET
mint_url = excluded.mint_url,
rating = excluded.rating,
k = excluded.k`,
...chunk.flatMap((r) => [r.eventId, r.mintUrl, r.pubkey, r.rating, r.k, r.createdAt]),
);
}
});
return batch.filter((r) => !known.has(r.eventId)).length;
}
export async function runDiscovery(backfill: boolean): Promise<DiscoveryResult> {
const started = Date.now();
const now = Math.floor(Date.now() / 1000);
const lastRun = backfill ? null : getStateNumber('discovery_since');
const lastRun = backfill ? null : await getStateNumber('discovery_since');
// Overlap the window by an hour so an event that arrived late is not missed.
const since = lastRun === null ? null : Math.max(0, lastRun - 3600);
@@ -167,9 +454,25 @@ export async function runDiscovery(backfill: boolean): Promise<DiscoveryResult>
let ok = true;
try {
const announcements = await fetchKind(KIND_MINT_ANNOUNCEMENT, since);
events += announcements.length;
for (const e of announcements) ingestMintUrls(e, now, newMints);
/*
* One query per announcement kind, driven by ANNOUNCEMENT_KINDS rather than by a
* list written out here: adding an ecosystem is a line in shared/, not an edit in
* this loop. They run together because they are independent reads against the same
* relay pool, and running them in series doubled the cycle's wall clock for nothing.
*/
const announcementsByType = new Map<string, NostrEvent[]>();
await Promise.all(
Object.entries(ANNOUNCEMENT_KINDS).map(async ([type, kind]) => {
announcementsByType.set(type, await fetchKind(kind, since));
}),
);
const announcements = announcementsByType.get('cashu') ?? [];
const federations = announcementsByType.get('fedimint') ?? [];
events += announcements.length + federations.length;
await ingestMintUrls(announcements, now, newMints);
await ingestFedimints(federations, now, newMints);
// Announcements alone miss mints that only ever appear in a review's `u` tag,
// so reviews feed discovery too.
@@ -182,42 +485,36 @@ export async function runDiscovery(backfill: boolean): Promise<DiscoveryResult>
for (const e of [...reviews, ...recent]) byId.set(e.id, e);
events += byId.size;
for (const e of byId.values()) ingestMintUrls(e, now, newMints);
await ingestMintUrls(byId.values(), now, newMints);
const insert = getDb().prepare(
`INSERT INTO reviews (event_id, mint_url, pubkey, rating, created_at)
VALUES (?, ?, ?, ?, ?)
ON CONFLICT(event_id) DO UPDATE SET
mint_url = excluded.mint_url,
rating = excluded.rating`,
);
// `changes` counts updates as well as inserts, so ask first. The log line is a
// health signal (is discovery finding anything?) and an inflated number hides a
// stalled crawl behind a busy-looking figure.
const known = getDb().prepare('SELECT 1 FROM reviews WHERE event_id = ?');
const ingestReviews = getDb().transaction((list: NostrEvent[]) => {
for (const e of list) {
// The `k` tag is what marks a 38000 as being about a Cashu mint rather than
// some other recommendable kind. Events without it but with a resolvable mint
// URL are still accepted: the old publisher's own events are the reason.
const target = resolveReviewTarget(e);
if (!target) continue;
if (!isCashuMintReview(e) && mintUrlsFromEvent(e).length === 0) continue;
if (!known.get(e.id)) newReviews++;
insert.run(e.id, target, e.pubkey, parseRating(e), e.created_at);
}
});
ingestReviews([...byId.values()]);
// Built after every mint this cycle found has been inserted, so review resolution
// below sees them. Nothing after this point inserts a mint.
const index = await loadMintIndex();
newReviews += await ingestReviews(byId.values(), index);
// Second pass: ask each known mint's reviews by name, which is the only way the
// counts converge (see fetchReviewsForMint). Runs after the global sweep so mints
// discovered in this cycle are included.
const targets = getDb()
.prepare('SELECT url, pubkey FROM mints')
.all() as { url: string; pubkey: string | null }[];
const db = await getDb();
const rows = await db.all<{
url: string;
type: string;
pubkey: string | null;
ecosystem_json: string | null;
}>('SELECT url, type, pubkey, ecosystem_json FROM mints');
const targets: ReviewTarget[] = rows.map((row) => {
let federationId: string | null = null;
if (row.type === 'fedimint' && row.ecosystem_json) {
try {
federationId =
(JSON.parse(row.ecosystem_json) as { federation_id?: string }).federation_id ?? null;
} catch {
// Nothing to ask by. The global sweep above already had its chance.
}
}
return { type: row.type, url: row.url, pubkey: row.pubkey, federationId };
});
const CONCURRENCY = 6;
let cursor = 0;
@@ -226,21 +523,21 @@ export async function runDiscovery(backfill: boolean): Promise<DiscoveryResult>
while (cursor < targets.length) {
const target = targets[cursor++];
if (!target) continue;
const found = await fetchReviewsForMint(target.url, target.pubkey, since);
const found = await fetchReviewsForMint(target, since);
if (found.length > 0) {
events += found.length;
ingestReviews(found);
newReviews += await ingestReviews(found, index);
}
}
}),
);
setState('discovery_since', String(now));
setState('last_discovery_at', String(now));
setState('last_discovery_ok', '1');
await setState('discovery_since', String(now));
await setState('last_discovery_at', String(now));
await setState('last_discovery_ok', '1');
} catch (err) {
ok = false;
setState('last_discovery_ok', '0');
await setState('last_discovery_ok', '0').catch(() => undefined);
log.error('discovery failed', { reason: err instanceof Error ? err.message : String(err) });
}
+108
View File
@@ -0,0 +1,108 @@
/**
* The only check this site can actually run against a Fedimint federation.
*
* A Cashu mint answers `GET /v1/info` over plain HTTPS, which is why probing one is
* fifteen lines. A federation has no such endpoint. Its guardians speak a JSON-RPC
* dialect over websockets, their addresses are bech32m-encoded inside the invite code,
* and confirming one is up means decoding that code, opening a socket to a quorum of
* guardians and agreeing a consensus session with them. That is a Fedimint client, not
* a probe, and writing half of one here would produce a status less trustworthy than
* saying nothing.
*
* What exists instead is fedimint.observer, which already keeps those client
* connections open and publishes the result:
*
* GET https://observer.fedimint.org/api/federations
* [{ "id": "<64 hex>", "name": "...", "invite": "fed11...", "health": "online" }, ...]
*
* So this is a real check, and it is somebody else's. Both halves of that matter and
* both are recorded: a federation whose status came from here carries
* `status_source: "fedimint.observer"` and its page prints that beside the status,
* rather than implying this site opened a socket. A federation the observer does not
* track gets `announced` — Nostr says it exists, nothing says it runs — and never
* `online`, never `offline`, and never a made-up uptime figure.
*
* When the observer itself is unreachable, nothing is written. A federation confirmed
* up an hour ago is not demoted because a third party had a bad minute; the row keeps
* its last real answer and `last_probe` is left alone so `/api/health` can still see
* that the cycle happened.
*/
import { log } from './log.ts';
/** Where the health data comes from, recorded on every row it decides. */
export const OBSERVER_SOURCE = 'fedimint.observer';
/** Empty disables the lookup entirely: every federation then stays `announced`. */
export const OBSERVER_URL =
process.env['FEDIMINT_OBSERVER_URL'] ?? 'https://observer.fedimint.org/api/federations';
/** One request for every federation, so it gets longer than a single mint probe does. */
const OBSERVER_TIMEOUT_MS = 12_000;
export interface ObserverFederation {
id: string;
name: string | null;
/** `online` or `offline` as published. Anything else is treated as unknown. */
health: 'online' | 'offline' | null;
}
/** Federation id (lowercase hex) to what the observer says about it. */
export type ObserverIndex = Map<string, ObserverFederation>;
function readHealth(value: unknown): 'online' | 'offline' | null {
if (value === 'online') return 'online';
if (value === 'offline') return 'offline';
// Anything else is a word this build has no meaning for, and guessing at it is
// exactly what "no fake statuses" rules out.
return null;
}
/**
* Ask the observer about every federation it tracks.
*
* Returns null when the lookup failed, which is deliberately distinguishable from an
* empty map: an empty map means the observer answered and tracks nothing, and the
* caller writes `announced` everywhere; null means it did not answer, and the caller
* writes nothing at all.
*/
export async function fetchObserverIndex(userAgent: string): Promise<ObserverIndex | null> {
if (!OBSERVER_URL.trim()) return new Map();
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), OBSERVER_TIMEOUT_MS);
try {
const res = await fetch(OBSERVER_URL, {
signal: controller.signal,
headers: { Accept: 'application/json', 'User-Agent': userAgent },
redirect: 'follow',
});
if (!res.ok) throw new Error(`HTTP ${res.status}`);
const body: unknown = await res.json();
if (!Array.isArray(body)) throw new Error('not a JSON array');
const index: ObserverIndex = new Map();
for (const entry of body) {
if (!entry || typeof entry !== 'object') continue;
const row = entry as Record<string, unknown>;
const id = typeof row['id'] === 'string' ? row['id'].toLowerCase() : '';
if (!/^[0-9a-f]{64}$/.test(id)) continue;
index.set(id, {
id,
name: typeof row['name'] === 'string' ? row['name'] : null,
health: readHealth(row['health']),
});
}
return index;
} catch (err) {
log.warn('fedimint observer unreachable', {
url: OBSERVER_URL,
reason: err instanceof Error ? err.message : String(err),
});
return null;
} finally {
clearTimeout(timer);
}
}
+39
View File
@@ -0,0 +1,39 @@
/**
* Read a response body up to `maxBytes`, or refuse it.
*
* `res.arrayBuffer()` and `res.json()` buffer however much the server sends before any
* size check can run, and a mint is an untrusted server: the abort timer bounds how
* *long* a read may take, this bounds how *large* it may get, and both are needed. The
* Content-Length header is checked first as a courtesy — a chunked or lying response
* still hits the streaming cap.
*
* Returns null when the body is over the limit or unreadable. The caller treats that
* exactly like a failed fetch.
*/
export async function readBodyBounded(res: Response, maxBytes: number): Promise<Buffer | null> {
const declared = Number(res.headers.get('content-length') ?? '');
if (Number.isFinite(declared) && declared > maxBytes) return null;
const reader = res.body?.getReader();
if (!reader) return null;
const chunks: Uint8Array[] = [];
let size = 0;
try {
for (;;) {
const { done, value } = await reader.read();
if (done) break;
size += value.byteLength;
if (size > maxBytes) {
await reader.cancel().catch(() => undefined);
return null;
}
chunks.push(value);
}
} catch {
// Aborted by the caller's timer, or the connection died mid-body.
return null;
}
return Buffer.concat(chunks);
}
+57 -15
View File
@@ -1,6 +1,8 @@
import fs from 'node:fs/promises';
import path from 'node:path';
import { isFetchableUrl } from '@cashumints/shared';
import { config } from './config.ts';
import { readBodyBounded } from './http.ts';
import { log } from './log.ts';
import type { MintRow } from './mints.ts';
@@ -8,13 +10,54 @@ const EXT_BY_TYPE: Record<string, string> = {
'image/png': '.png',
'image/jpeg': '.jpg',
'image/webp': '.webp',
'image/svg+xml': '.svg',
'image/gif': '.gif',
'image/x-icon': '.ico',
'image/vnd.microsoft.icon': '.ico',
// No SVG on purpose: SVG is markup that can carry script, and these files are
// re-served from the site's own origin, so a cached one opened directly would run a
// mint operator's script there. The sandbox header in server.ts covers files cached
// before this rule existed.
};
const MAX_ICON_BYTES = 512 * 1024;
/** Redirect hops followed, each hop re-checked before it is fetched. */
const MAX_REDIRECTS = 3;
/**
* Fetch an icon with redirects validated hop by hop.
*
* `icon_url` is whatever the mint's /v1/info says it is, so every address on the way —
* the first one and each Location after it — has to pass `isFetchableUrl`, or a public
* URL that 302s to a metadata endpoint walks straight around a check done only once.
*/
async function fetchIconResponse(startUrl: string, signal: AbortSignal): Promise<Response | null> {
let target = startUrl;
for (let hop = 0; hop <= MAX_REDIRECTS; hop++) {
if (!isFetchableUrl(target)) return null;
const res = await fetch(target, {
signal,
redirect: 'manual',
headers: { 'User-Agent': config.userAgent },
});
if (res.status >= 300 && res.status < 400) {
const location = res.headers.get('location');
if (!location) return null;
try {
target = new URL(location, target).toString();
} catch {
return null;
}
continue;
}
return res.ok ? res : null;
}
return null;
}
/**
* Cache a mint's icon to disk so offline mints keep theirs.
@@ -30,27 +73,26 @@ export async function cacheIcon(row: MintRow, iconUrl: string | null): Promise<s
try {
const absolute = new URL(iconUrl, `${row.url}/`).toString();
if (!absolute.startsWith('https://') && !absolute.startsWith('http://')) return row.icon_file;
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), config.probeTimeoutMs);
let res: Response;
let buf: Buffer | null = null;
let ext: string | undefined;
try {
res = await fetch(absolute, {
signal: controller.signal,
headers: { 'User-Agent': config.userAgent },
});
const res = await fetchIconResponse(absolute, controller.signal);
if (!res) return row.icon_file;
const type = (res.headers.get('content-type') ?? '').split(';')[0]?.trim() ?? '';
ext = EXT_BY_TYPE[type];
if (!ext) return row.icon_file;
// The timer stays armed through the body read: the abort is what stops a server
// that sends its headers quickly and then never finishes the body.
buf = await readBodyBounded(res, MAX_ICON_BYTES);
} finally {
clearTimeout(timer);
}
if (!res.ok) return row.icon_file;
const type = (res.headers.get('content-type') ?? '').split(';')[0]?.trim() ?? '';
const ext = EXT_BY_TYPE[type];
if (!ext) return row.icon_file;
const buf = Buffer.from(await res.arrayBuffer());
if (buf.byteLength === 0 || buf.byteLength > MAX_ICON_BYTES) return row.icon_file;
if (!buf || buf.byteLength === 0) return row.icon_file;
const filename = `${row.host}${ext}`;
await fs.writeFile(path.join(config.iconDir, filename), buf);
+40 -16
View File
@@ -17,17 +17,30 @@ function track<T>(work: Promise<T>): Promise<T> {
return work;
}
/*
* One cycle of each kind at a time. A cycle that outlives its interval (a relay socket
* that never closes, a wedged fetch) must not have the next timer tick start a second
* copy beside it — that doubles the relay load and interleaves state writes for no
* new information. The late cycle finishes; the ticks it swallowed are simply skipped.
*/
let probeRunning = false;
let discoveryRunning = false;
async function probeCycle(): Promise<void> {
if (shuttingDown) return;
if (shuttingDown || probeRunning) return;
probeRunning = true;
try {
await track(probeAll());
} catch (err) {
log.error('probe cycle failed', { reason: err instanceof Error ? err.message : String(err) });
} finally {
probeRunning = false;
}
}
async function discoveryCycle(backfill: boolean): Promise<void> {
if (shuttingDown) return;
if (shuttingDown || discoveryRunning) return;
discoveryRunning = true;
try {
const result = await track(runDiscovery(backfill));
// New mints are probed straight away, not on the next cycle, so their page has
@@ -37,16 +50,20 @@ async function discoveryCycle(backfill: boolean): Promise<void> {
log.error('discovery cycle failed', {
reason: err instanceof Error ? err.message : String(err),
});
} finally {
discoveryRunning = false;
}
}
function main(): void {
getDb();
async function main(): Promise<void> {
// Fail loudly here rather than on the first request: a bad DATABASE_URL, an
// unreachable Postgres or an unwritable SQLite directory should stop the boot.
await getDb();
const server = serve({ fetch: createApp().fetch, port: config.port }, (info) => {
log.info('api listening', {
port: info.port,
db: config.dbPath,
db: config.db.label,
relays: config.relays.length,
});
});
@@ -57,8 +74,13 @@ function main(): void {
config.discoveryIntervalMin * 60 * 1000,
);
const pruneTimer = setInterval(() => {
const removed = pruneProbes();
if (removed > 0) log.info('pruned probes', { rows: removed });
void pruneProbes()
.then((removed) => {
if (removed > 0) log.info('pruned probes', { rows: removed });
})
.catch((err: unknown) => {
log.error('prune failed', { reason: err instanceof Error ? err.message : String(err) });
});
}, DAY_MS);
// Startup: probe what we already know, then a full discovery backfill so a fresh
@@ -74,17 +96,16 @@ function main(): void {
clearInterval(discoveryTimer);
clearInterval(pruneTimer);
const stop = (): void => {
// closeDb drains the Postgres pool, so it is awaited before the process leaves.
void closeDb().finally(() => process.exit(0));
};
void inFlight.finally(() => {
closePool();
server.close(() => {
closeDb();
process.exit(0);
});
server.close(stop);
// Do not hang forever on a relay socket that refuses to close.
setTimeout(() => {
closeDb();
process.exit(0);
}, 8000).unref();
setTimeout(stop, 8000).unref();
});
};
@@ -92,4 +113,7 @@ function main(): void {
process.on('SIGTERM', () => shutdown('SIGTERM'));
}
main();
void main().catch((err: unknown) => {
log.error('startup failed', { reason: err instanceof Error ? err.message : String(err) });
process.exit(1);
});
+201
View File
@@ -0,0 +1,201 @@
/**
* Copy the whole database from one backend to the other.
*
* pnpm --filter ./api migrate --to postgres://user:pw@localhost/cashumints
* pnpm --filter ./api migrate --from postgres://… --to ./data/cashumints.db
*
* `--from` defaults to whatever the running configuration points at (DATABASE_URL, or
* DB_PATH's SQLite file), so the common direction — the SQLite file you already have,
* into a fresh Postgres — needs only `--to`.
*
* Both sides accept either spelling, so this is also how you move one SQLite file to
* another, or one Postgres database to another.
*
* The target's schema is created if it is missing. Rows are upserted on the primary
* key, so an interrupted run can simply be repeated. `probes` is the exception: it is
* an append-only log with no unique key, so a target that already holds probe rows is
* left alone unless `--force`, which replaces them wholesale.
*
* Icons are files, not rows. `ICON_DIR` is shared by both backends and nothing here
* touches it.
*/
import { parseDbTarget, resolveDbConfig, type DbConfig } from './config.ts';
import type { Db, Param } from './db-driver.ts';
import { COLUMNS, PRIMARY_KEY, TABLES, type TableName } from './db-schema.ts';
import { openDb } from './db.ts';
import { log } from './log.ts';
/** Rows per statement. 16 columns x 500 stays far under Postgres' 65535 parameter cap. */
const BATCH = 500;
interface Options {
from: DbConfig;
to: DbConfig;
force: boolean;
dryRun: boolean;
}
function usage(): string {
return [
'Usage: migrate --to <target> [--from <source>] [--force] [--dry-run]',
'',
' <target>, <source> postgres://user:pw@host:5432/db',
' /var/lib/cashumints/cashumints.db',
' sqlite:./data/cashumints.db',
'',
' --from defaults to the configured database (DATABASE_URL, else DB_PATH)',
' --force replace the target probes table when it already has rows',
' --dry-run report what would be copied, write nothing',
].join('\n');
}
export function parseArgs(argv: string[]): Options {
let from: string | undefined;
let to: string | undefined;
let force = false;
let dryRun = false;
for (let i = 0; i < argv.length; i++) {
const arg = argv[i];
if (arg === '--from') from = argv[++i];
else if (arg === '--to') to = argv[++i];
else if (arg === '--force') force = true;
else if (arg === '--dry-run') dryRun = true;
else if (arg === '--help' || arg === '-h') throw new Error(usage());
else throw new Error(`Unknown argument "${arg}"\n\n${usage()}`);
}
if (!to) throw new Error(`--to is required.\n\n${usage()}`);
const source = from ? parseDbTarget(from) : resolveDbConfig();
const target = parseDbTarget(to);
if (source.label === target.label) {
throw new Error(`Source and target are the same database (${source.label}).`);
}
return { from: source, to: target, force, dryRun };
}
async function countRows(db: Db, table: TableName): Promise<number> {
const row = await db.get<{ n: number }>(`SELECT COUNT(*) AS n FROM ${table}`);
return row?.n ?? 0;
}
/** `INSERT … ON CONFLICT (pk) DO UPDATE`, or a plain insert for the keyless probes table. */
function insertStatement(table: TableName, rows: number): string {
const cols = COLUMNS[table];
const placeholders = Array.from({ length: rows }, () => `(${cols.map(() => '?').join(', ')})`);
const pk = PRIMARY_KEY[table];
const conflict =
pk === null
? ''
: ` ON CONFLICT (${pk}) DO UPDATE SET ${cols
.filter((c) => c !== pk)
.map((c) => `${c} = excluded.${c}`)
.join(', ')}`;
return `INSERT INTO ${table} (${cols.join(', ')}) VALUES ${placeholders.join(', ')}${conflict}`;
}
async function copyTable(source: Db, target: Db, table: TableName): Promise<number> {
const cols = COLUMNS[table];
const pk = PRIMARY_KEY[table];
// Read the whole table. The largest of them is probes, at roughly 13k rows per mint
// per 90-day retention window — tens of megabytes at the very top end, and paging it
// would need a stable sort key that probes does not have.
const rows = await source.all<Record<string, Param>>(`SELECT ${cols.join(', ')} FROM ${table}`);
if (rows.length === 0) return 0;
// Same primary key twice in one statement is an error in Postgres, not a silent
// overwrite. A source that is itself consistent cannot produce one, but a hand-edited
// file can, and the failure is confusing enough to be worth heading off.
const deduped =
pk === null ? rows : [...new Map(rows.map((r) => [String(r[pk]), r])).values()];
await target.transaction(async (tx) => {
for (let i = 0; i < deduped.length; i += BATCH) {
const chunk = deduped.slice(i, i + BATCH);
await tx.run(
insertStatement(table, chunk.length),
...chunk.flatMap((row) => cols.map((c) => row[c] ?? null)),
);
}
});
return deduped.length;
}
export async function migrate(options: Options): Promise<void> {
log.info('migrating', { from: options.from.label, to: options.to.label });
const source = await openDb(options.from);
let target: Db | null = null;
try {
// openDb creates the schema, so the target may be a database with nothing in it.
target = await openDb(options.to);
const existingProbes = await countRows(target, 'probes');
if (existingProbes > 0 && !options.force && !options.dryRun) {
throw new Error(
`Target already holds ${existingProbes} probe rows. probes has no primary key, so ` +
'they cannot be upserted. Re-run with --force to replace them, or drop the table first.',
);
}
if (options.dryRun) {
for (const table of TABLES) {
log.info('would copy', {
table,
rows: await countRows(source, table),
target_rows: await countRows(target, table),
});
}
log.info('dry run, nothing written');
return;
}
if (existingProbes > 0) {
await target.run('DELETE FROM probes');
log.info('cleared target probes', { rows: existingProbes });
}
for (const table of TABLES) {
const started = Date.now();
const copied = await copyTable(source, target, table);
log.info('copied', { table, rows: copied, ms: Date.now() - started });
}
// Read the counts back from the target rather than trusting the copy loop.
for (const table of TABLES) {
const [before, after] = await Promise.all([
countRows(source, table),
countRows(target, table),
]);
const verdict = after >= before ? 'ok' : 'SHORT';
log.info('verify', { table, source: before, target: after, verdict });
if (after < before) {
throw new Error(`${table}: target has ${after} rows, source has ${before}`);
}
}
log.info('migration complete', { to: options.to.label });
} finally {
await source.close().catch(() => undefined);
if (target) await target.close().catch(() => undefined);
}
}
// Only when run directly, so the functions above stay importable from a test.
if (process.argv[1] && import.meta.filename === process.argv[1]) {
try {
await migrate(parseArgs(process.argv.slice(2)));
process.exit(0);
} catch (err) {
console.error(err instanceof Error ? err.message : String(err));
process.exit(1);
}
}
+165 -27
View File
@@ -1,23 +1,33 @@
import { normalizeMintUrl } from '@cashumints/shared';
import {
fedimintKey, fedimintSlug, normalizeMintUrl,
type FedimintAnnouncement, type FedimintFields,
} from '@cashumints/shared';
import { getDb } from './db.ts';
import { log } from './log.ts';
/**
* URLs already reported as unusable. A rejected value (a Fedimint `fed11…` invite, an
* onion address) recurs in dozens of events per cycle, and logging each occurrence
* buries the lines that matter.
* URLs already reported as unusable. A rejected value (an onion address, a bare label)
* recurs in dozens of events per cycle, and logging each occurrence buries the lines
* that matter.
*
* Fedimint invite codes used to be the loudest entry in here, arriving as `u` tags this
* function could make no sense of. They have their own path now and never reach it.
*/
const reportedSkips = new Set<string>();
export interface MintRow {
url: string;
host: string;
/** 'cashu' or 'fedimint'. Rows written before the column existed default to 'cashu'. */
type: string;
name: string | null;
description: string | null;
icon_url: string | null;
icon_file: string | null;
pubkey: string | null;
info_json: string | null;
/** Type-specific data. `FedimintFields` for a federation, null for a Cashu mint. */
ecosystem_json: string | null;
nuts_json: string | null;
version: string | null;
status: string;
@@ -44,7 +54,10 @@ export interface MintRow {
* If that ever stops holding, the upgrade is to probe both casings once and keep the
* one that answers.
*/
export function insertMintIfNew(rawUrl: string, now = Math.floor(Date.now() / 1000)): string | null {
export async function insertMintIfNew(
rawUrl: string,
now = Math.floor(Date.now() / 1000),
): Promise<string | null> {
const normalized = normalizeMintUrl(rawUrl);
if (!normalized) {
if (!reportedSkips.has(rawUrl)) {
@@ -54,41 +67,166 @@ export function insertMintIfNew(rawUrl: string, now = Math.floor(Date.now() / 10
return null;
}
const result = getDb()
.prepare(
// OR IGNORE covers both unique keys, url and host, in one clause.
`INSERT OR IGNORE INTO mints (url, host, status, first_seen, updated_at)
VALUES (?, ?, 'unknown', ?, ?)`,
)
.run(normalized.url, normalized.host, now, now);
const db = await getDb();
// No conflict target, so this covers both unique keys — url and host — in one clause.
// (`INSERT OR IGNORE` would too, but only SQLite knows that spelling.)
const result = await db.run(
`INSERT INTO mints (url, host, type, status, first_seen, updated_at)
VALUES (?, ?, 'cashu', 'unknown', ?, ?)
ON CONFLICT DO NOTHING`,
normalized.url,
normalized.host,
now,
now,
);
return result.changes > 0 ? normalized.url : null;
if (result.changes > 0) return normalized.url;
/*
* The no-op insert is almost always this exact URL already being tracked. The other
* possibility is a *different* URL whose slug collides (`host/a-b` vs `host/a/b`
* both slug to `host-a-b`): that mint is silently not tracked and any review of it
* files under the other one, which is worth a log line the first time it happens.
*/
const holder = await db.get<{ url: string }>(
'SELECT url FROM mints WHERE host = ?',
normalized.host,
);
if (holder && holder.url !== normalized.url && !reportedSkips.has(normalized.url)) {
reportedSkips.add(normalized.url);
log.warn('slug collision, mint not tracked', {
url: normalized.url,
host: normalized.host,
existing: holder.url,
});
}
return null;
}
/* ---------- fedimint ---------- */
/** Read a federation row's type-specific columns back out. */
export function parseEcosystem(row: Pick<MintRow, 'ecosystem_json'>): FedimintFields | null {
if (!row.ecosystem_json) return null;
try {
return JSON.parse(row.ecosystem_json) as FedimintFields;
} catch {
return null;
}
}
/**
* Insert or refresh a federation from a kind 38173 announcement.
*
* Deduped on the federation id, which is the only identity a federation has: the same
* federation is announced by several people (two different npubs currently announce
* "Bitcoin Principles" with the same `d`), and every one of those is the same thing to
* join. One row, keyed on the id, whoever published it.
*
* Unlike `insertMintIfNew` this does update an existing row, and it has to: an
* announcement is the *only* source of a federation's invite codes, modules and name,
* where a Cashu mint's row is refreshed by probing the mint itself. Older announcements
* are ignored (`announced_at` goes forwards only) so a replayed event from last year
* cannot overwrite this week's invite code.
*
* The status columns are never touched here. Whether a federation is up is a probe's
* answer, and an announcement is not evidence of anything being up.
*
* Returns the row's key when a row was created, null when one was merely updated.
*/
export async function upsertFedimint(
announcement: FedimintAnnouncement,
now = Math.floor(Date.now() / 1000),
): Promise<string | null> {
const db = await getDb();
const url = fedimintKey(announcement.federationId);
const host = fedimintSlug(announcement.federationId);
const fields: FedimintFields = {
federation_id: announcement.federationId,
invite_codes: announcement.inviteCodes,
modules: announcement.modules,
network: announcement.network,
announced_at: announcement.announcedAt,
announcer_pubkey: announcement.announcerPubkey,
// Set by the probe, not by an announcement. Carried over below when a row exists.
status_source: null,
};
const existing = await db.get<MintRow>('SELECT * FROM mints WHERE url = ?', url);
if (!existing) {
await db.run(
`INSERT INTO mints (url, host, type, name, description, icon_url, ecosystem_json,
status, first_seen, updated_at)
VALUES (?, ?, 'fedimint', ?, ?, ?, ?, 'announced', ?, ?)
ON CONFLICT DO NOTHING`,
url,
host,
announcement.name,
announcement.about,
announcement.picture,
JSON.stringify(fields),
now,
now,
);
return url;
}
const previous = parseEcosystem(existing);
if (previous && (previous.announced_at ?? 0) > announcement.announcedAt) return null;
await db.run(
`UPDATE mints SET
name = COALESCE(?, name),
description = COALESCE(?, description),
icon_url = COALESCE(?, icon_url),
ecosystem_json = ?,
updated_at = ?
WHERE url = ?`,
announcement.name,
announcement.about,
announcement.picture,
JSON.stringify({ ...fields, status_source: previous?.status_source ?? null }),
now,
url,
);
return null;
}
/** Every federation row, for the probe cycle and for review resolution. */
export async function fedimintRows(): Promise<MintRow[]> {
const db = await getDb();
return db.all<MintRow>(`SELECT * FROM mints WHERE type = 'fedimint'`);
}
/** Canonical mint URL for a normalized slug, or null if no such mint is tracked. */
export function mintUrlByHost(host: string): string | null {
const row = getDb().prepare('SELECT url FROM mints WHERE host = ?').get(host) as
| { url: string }
| undefined;
export async function mintUrlByHost(host: string): Promise<string | null> {
const db = await getDb();
const row = await db.get<{ url: string }>('SELECT url FROM mints WHERE host = ?', host);
return row?.url ?? null;
}
export function allMintRows(): MintRow[] {
return getDb().prepare('SELECT * FROM mints').all() as MintRow[];
export async function allMintRows(): Promise<MintRow[]> {
const db = await getDb();
return db.all<MintRow>('SELECT * FROM mints');
}
export function mintByHost(host: string): MintRow | undefined {
return getDb().prepare('SELECT * FROM mints WHERE host = ?').get(host) as MintRow | undefined;
export async function mintByHost(host: string): Promise<MintRow | undefined> {
const db = await getDb();
return db.get<MintRow>('SELECT * FROM mints WHERE host = ?', host);
}
export function mintByUrl(url: string): MintRow | undefined {
return getDb().prepare('SELECT * FROM mints WHERE url = ?').get(url) as MintRow | undefined;
export async function mintByUrl(url: string): Promise<MintRow | undefined> {
const db = await getDb();
return db.get<MintRow>('SELECT * FROM mints WHERE url = ?', url);
}
/** Resolve a mint by the pubkey it publishes in /v1/info, for reviews with only a `d` tag. */
export function mintUrlByPubkey(pubkey: string): string | null {
const row = getDb().prepare('SELECT url FROM mints WHERE pubkey = ? LIMIT 1').get(pubkey) as
| { url: string }
| undefined;
export async function mintUrlByPubkey(pubkey: string): Promise<string | null> {
const db = await getDb();
const row = await db.get<{ url: string }>('SELECT url FROM mints WHERE pubkey = ? LIMIT 1', pubkey);
return row?.url ?? null;
}
+163 -19
View File
@@ -1,9 +1,11 @@
import { parseNuts, type MintInfo } from '@cashumints/shared';
import { parseNuts, type FedimintFields, type MintInfo } from '@cashumints/shared';
import { config } from './config.ts';
import { getDb, setState } from './db.ts';
import { fetchObserverIndex, OBSERVER_SOURCE, type ObserverIndex } from './fedimint-observer.ts';
import { readBodyBounded } from './http.ts';
import { cacheIcon } from './icons.ts';
import { log } from './log.ts';
import { allMintRows, mintByUrl, type MintRow } from './mints.ts';
import { allMintRows, mintByUrl, parseEcosystem, type MintRow } from './mints.ts';
export interface ProbeResult {
url: string;
@@ -19,6 +21,13 @@ export function statusForFails(fails: number): 'online' | 'degraded' | 'offline'
return fails < 3 ? 'degraded' : 'offline';
}
/**
* The largest /v1/info accepted. Real payloads are a few KB; the whole document is
* stored in `info_json` and served back on every detail request, so an unbounded one
* is both a memory spike here and a bloated row forever after.
*/
const MAX_INFO_BYTES = 256 * 1024;
async function fetchInfo(url: string): Promise<{ info: MintInfo; latencyMs: number }> {
const started = Date.now();
const controller = new AbortController();
@@ -32,7 +41,10 @@ async function fetchInfo(url: string): Promise<{ info: MintInfo; latencyMs: numb
});
if (!res.ok) throw new Error(`HTTP ${res.status}`);
const info = (await res.json()) as MintInfo;
const body = await readBodyBounded(res, MAX_INFO_BYTES);
if (!body) throw new Error(`info larger than ${MAX_INFO_BYTES} bytes or unreadable`);
const info = JSON.parse(body.toString('utf8')) as MintInfo;
if (!info || typeof info !== 'object') throw new Error('not a JSON object');
return { info, latencyMs: Date.now() - started };
@@ -48,7 +60,7 @@ async function fetchInfo(url: string): Promise<{ info: MintInfo; latencyMs: numb
* what lets an offline mint still render its full page, which is acceptance criterion #1.
*/
export async function probeMint(row: MintRow): Promise<ProbeResult> {
const db = getDb();
const db = await getDb();
const now = Math.floor(Date.now() / 1000);
try {
@@ -57,7 +69,7 @@ export async function probeMint(row: MintRow): Promise<ProbeResult> {
const iconFile = await cacheIcon(row, info.icon_url ?? null);
db.prepare(
await db.run(
`UPDATE mints SET
name = COALESCE(?, name),
description = COALESCE(?, description),
@@ -73,7 +85,6 @@ export async function probeMint(row: MintRow): Promise<ProbeResult> {
last_probe = ?,
updated_at = ?
WHERE url = ?`,
).run(
info.name ?? null,
info.description ?? null,
info.icon_url ?? null,
@@ -88,7 +99,8 @@ export async function probeMint(row: MintRow): Promise<ProbeResult> {
row.url,
);
db.prepare('INSERT INTO probes (mint_url, ts, ok, latency_ms) VALUES (?, ?, 1, ?)').run(
await db.run(
'INSERT INTO probes (mint_url, ts, ok, latency_ms) VALUES (?, ?, 1, ?)',
row.url,
now,
latencyMs,
@@ -102,12 +114,18 @@ export async function probeMint(row: MintRow): Promise<ProbeResult> {
const fails = row.consecutive_fails + 1;
const status = statusForFails(fails);
db.prepare(
await db.run(
`UPDATE mints SET consecutive_fails = ?, status = ?, last_probe = ?, updated_at = ?
WHERE url = ?`,
).run(fails, status, now, now, row.url);
fails,
status,
now,
now,
row.url,
);
db.prepare('INSERT INTO probes (mint_url, ts, ok, latency_ms) VALUES (?, ?, 0, NULL)').run(
await db.run(
'INSERT INTO probes (mint_url, ts, ok, latency_ms) VALUES (?, ?, 0, NULL)',
row.url,
now,
);
@@ -127,6 +145,112 @@ export async function probeMint(row: MintRow): Promise<ProbeResult> {
}
}
/* ---------- fedimint ---------- */
/**
* Write what a check said about one federation.
*
* Three outcomes, and the third is the one that matters most. `online` and `offline`
* are recorded exactly as a mint probe records them, probe row and all, so uptime and
* the sparkline work the same way. A federation the checker does not cover gets
* `announced`, which is a statement about this site's reach and not about the
* federation: no probe row is written for it, because nothing was probed, and
* `last_online` stays null so nothing downstream can render a "last seen" it invented.
*/
async function applyFederationHealth(
row: MintRow,
health: 'online' | 'offline' | null,
now: number,
): Promise<ProbeResult> {
const db = await getDb();
const previous = parseEcosystem(row);
const ecosystem = (source: string | null): string =>
JSON.stringify({ ...(previous ?? {}), status_source: source } as FedimintFields);
if (health === null) {
// Nothing checks this federation. Say so, and leave every column a real check
// would have written alone.
const status = 'announced';
const changed = row.status !== status;
await db.run(
`UPDATE mints SET status = ?, ecosystem_json = ?, last_probe = ?, updated_at = ?
WHERE url = ?`,
status,
ecosystem(null),
now,
now,
row.url,
);
if (changed) log.info('mint state change', { url: row.url, from: row.status, to: status });
return { url: row.url, ok: false, latencyMs: null, status, changed };
}
const ok = health === 'online';
const status = ok ? 'online' : 'offline';
await db.run(
`UPDATE mints SET
status = ?,
consecutive_fails = ?,
ecosystem_json = ?,
last_online = COALESCE(?, last_online),
last_probe = ?,
updated_at = ?
WHERE url = ?`,
status,
ok ? 0 : row.consecutive_fails + 1,
ecosystem(OBSERVER_SOURCE),
ok ? now : null,
now,
now,
row.url,
);
await db.run(
'INSERT INTO probes (mint_url, ts, ok, latency_ms) VALUES (?, ?, ?, NULL)',
row.url,
now,
ok ? 1 : 0,
);
const changed = row.status !== status;
if (changed) {
log.info('mint state change', {
url: row.url,
from: row.status,
to: status,
source: OBSERVER_SOURCE,
});
}
return { url: row.url, ok, latencyMs: null, status, changed };
}
/**
* One lookup for every federation, rather than one request each.
*
* The checker publishes its whole table in a single response, so asking per row would
* be the same answer fetched fifty times. Returns an empty list when the lookup itself
* failed: no answer is not a result, and writing one would be inventing it.
*/
async function checkFederations(rows: MintRow[], now: number): Promise<ProbeResult[]> {
if (rows.length === 0) return [];
const index: ObserverIndex | null = await fetchObserverIndex(config.userAgent);
if (!index) {
log.warn('fedimint check skipped', { federations: rows.length, reason: 'no observer data' });
return [];
}
const results: ProbeResult[] = [];
for (const row of rows) {
const federationId = parseEcosystem(row)?.federation_id ?? '';
results.push(await applyFederationHealth(row, index.get(federationId)?.health ?? null, now));
}
return results;
}
/** Run `tasks` with at most `limit` in flight. */
async function pooled<T>(items: T[], limit: number, fn: (item: T) => Promise<unknown>): Promise<void> {
let cursor = 0;
@@ -139,31 +263,51 @@ async function pooled<T>(items: T[], limit: number, fn: (item: T) => Promise<unk
await Promise.all(workers);
}
export async function probeAll(rows: MintRow[] = allMintRows()): Promise<ProbeResult[]> {
if (rows.length === 0) return [];
export async function probeAll(rows?: MintRow[]): Promise<ProbeResult[]> {
const targets = rows ?? (await allMintRows());
if (targets.length === 0) return [];
const started = Date.now();
const now = Math.floor(Date.now() / 1000);
// Two ecosystems, two entirely different checks: an HTTP fetch per mint, and one
// lookup covering every federation. Split here rather than inside the worker so the
// federations are not each waiting behind a slot in the mint pool.
const mints = targets.filter((row) => row.type !== 'fedimint');
const federations = targets.filter((row) => row.type === 'fedimint');
const results: ProbeResult[] = [];
await pooled(rows, config.probeConcurrency, async (row) => {
results.push(await probeMint(row));
});
const [, federationResults] = await Promise.all([
pooled(mints, config.probeConcurrency, async (row) => {
results.push(await probeMint(row));
}),
checkFederations(federations, now),
]);
results.push(...federationResults);
const online = results.filter((r) => r.ok).length;
const changed = results.filter((r) => r.changed).length;
log.info('probe cycle', {
mints: results.length,
mints: mints.length,
federations: federations.length,
online,
down: results.length - online,
changed,
ms: Date.now() - started,
});
setState('last_probe_at', String(Math.floor(Date.now() / 1000)));
// Only a full sweep counts for /api/health. probeUrls passes a handful of newly
// discovered rows through here, and stamping the health marker for those would let
// a fresh discovery hide hours of failed probe cycles.
if (rows === undefined) {
await setState('last_probe_at', String(Math.floor(Date.now() / 1000)));
}
return results;
}
/** Probe a set of URLs immediately, used when discovery finds new mints. */
/** Probe a set of keys immediately, used when discovery finds new mints or federations. */
export async function probeUrls(urls: string[]): Promise<void> {
const rows = urls.map((u) => mintByUrl(u)).filter((r): r is MintRow => r !== undefined);
const found = await Promise.all(urls.map((u) => mintByUrl(u)));
const rows = found.filter((r): r is MintRow => r !== undefined);
if (rows.length > 0) await probeAll(rows);
}
+183 -88
View File
@@ -8,13 +8,14 @@ import {
type MintInfo,
type MintListItem,
type MintStatus,
type MintType,
type ProbeSample,
type RatingDistribution,
type Stats,
} from '@cashumints/shared';
import { config, startedAt } from './config.ts';
import { getDb, getStateNumber, getState } from './db.ts';
import { mintByHost, type MintRow } from './mints.ts';
import { mintByHost, parseEcosystem, type MintRow } from './mints.ts';
/**
* One review per author per mint, newest wins.
@@ -23,15 +24,18 @@ import { mintByHost, type MintRow } from './mints.ts';
* replaces the earlier one. The old site applied the same rule in memory
* (aggregateReviews, keyed by pubkey); doing it in SQL keeps counts honest and stops
* one npub from moving a mint's average by re-posting.
*
* The `AS ranked` alias is not decoration: Postgres rejects an unaliased subquery in
* FROM, and SQLite does not care either way.
*/
const LATEST_REVIEWS = `
export const LATEST_REVIEWS = `
SELECT mint_url, pubkey, rating, created_at FROM (
SELECT mint_url, pubkey, rating, created_at,
ROW_NUMBER() OVER (
PARTITION BY mint_url, pubkey ORDER BY created_at DESC, event_id
) AS rn
FROM reviews
) WHERE rn = 1
) AS ranked WHERE rn = 1
`;
interface AggRow {
@@ -41,17 +45,26 @@ interface AggRow {
last_review_at: number | null;
}
function aggregates(): Map<string, AggRow> {
const rows = getDb()
.prepare(
`WITH latest AS (${LATEST_REVIEWS})
SELECT mint_url,
COUNT(*) AS review_count,
AVG(rating) AS rating_avg,
MAX(created_at) AS last_review_at
FROM latest GROUP BY mint_url`,
)
.all() as AggRow[];
/**
* Review counts and averages, for one mint or for all of them.
*
* A detail page passes its own URL. Computing the whole table to read one row off it
* is free on a local SQLite file and is not on a Postgres server across a socket.
*/
async function aggregates(mintUrl?: string): Promise<Map<string, AggRow>> {
const db = await getDb();
const where = mintUrl === undefined ? '' : 'WHERE mint_url = ?';
const params = mintUrl === undefined ? [] : [mintUrl];
const rows = await db.all<AggRow>(
`WITH latest AS (${LATEST_REVIEWS})
SELECT mint_url,
COUNT(*) AS review_count,
AVG(rating) AS rating_avg,
MAX(created_at) AS last_review_at
FROM latest ${where} GROUP BY mint_url`,
...params,
);
return new Map(rows.map((r) => [r.mint_url, r]));
}
@@ -64,7 +77,7 @@ function aggregates(): Map<string, AggRow> {
* fail its own stated purpose. `SCORE_PRIOR_MEAN=global` restores literal BACKEND.md
* behaviour, or set any number to pin it.
*/
function priorMean(): number {
async function priorMean(): Promise<number> {
const override = process.env['SCORE_PRIOR_MEAN'];
if (override && override !== 'global') {
@@ -73,10 +86,11 @@ function priorMean(): number {
}
if (override === 'global') {
const row = getDb()
.prepare(`WITH latest AS (${LATEST_REVIEWS}) SELECT AVG(rating) AS avg FROM latest`)
.get() as { avg: number | null };
return row.avg ?? NEUTRAL_PRIOR_MEAN;
const db = await getDb();
const row = await db.get<{ avg: number | null }>(
`WITH latest AS (${LATEST_REVIEWS}) SELECT AVG(rating) AS avg FROM latest`,
);
return row?.avg ?? NEUTRAL_PRIOR_MEAN;
}
return NEUTRAL_PRIOR_MEAN;
@@ -103,6 +117,7 @@ function toListItem(row: MintRow, agg: AggRow | undefined, mean: number, now: nu
host: row.host,
name: row.name,
icon: iconPath(row),
type: row.type as MintType,
status: row.status as MintStatus,
last_online: row.last_online,
review_count: base.review_count,
@@ -113,26 +128,38 @@ function toListItem(row: MintRow, agg: AggRow | undefined, mean: number, now: nu
};
}
export function listMints(limit?: number): MintListItem[] {
/**
* The list, optionally narrowed to one ecosystem.
*
* `type` is filtered in SQL rather than after scoring, because the two list pages ask
* for one ecosystem each and there is no reason to score fifty federations to render
* /mints. Unfiltered still returns everything, so `/api/mints` on its own is the whole
* index with a `type` on every item.
*/
export async function listMints(limit?: number, type?: string): Promise<MintListItem[]> {
const now = Math.floor(Date.now() / 1000);
const agg = aggregates();
const mean = priorMean();
const db = await getDb();
const [agg, mean, rows] = await Promise.all([
aggregates(),
priorMean(),
type
? db.all<MintRow>('SELECT * FROM mints WHERE type = ?', type)
: db.all<MintRow>('SELECT * FROM mints'),
]);
const items = (getDb().prepare('SELECT * FROM mints').all() as MintRow[])
.map((row) => toListItem(row, agg.get(row.url), mean, now))
.sort(compareMints);
const items = rows.map((row) => toListItem(row, agg.get(row.url), mean, now)).sort(compareMints);
return limit && limit > 0 ? items.slice(0, limit) : items;
}
function distribution(mintUrl: string): RatingDistribution {
const rows = getDb()
.prepare(
`WITH latest AS (${LATEST_REVIEWS})
SELECT rating, COUNT(*) AS n FROM latest
WHERE mint_url = ? AND rating IS NOT NULL GROUP BY rating`,
)
.all(mintUrl) as { rating: number; n: number }[];
async function distribution(mintUrl: string): Promise<RatingDistribution> {
const db = await getDb();
const rows = await db.all<{ rating: number; n: number }>(
`WITH latest AS (${LATEST_REVIEWS})
SELECT rating, COUNT(*) AS n FROM latest
WHERE mint_url = ? AND rating IS NOT NULL GROUP BY rating`,
mintUrl,
);
const dist: RatingDistribution = { '1': 0, '2': 0, '3': 0, '4': 0, '5': 0 };
for (const r of rows) {
@@ -142,34 +169,39 @@ function distribution(mintUrl: string): RatingDistribution {
return dist;
}
function reviews90d(mintUrl: string): number {
async function reviews90d(mintUrl: string): Promise<number> {
const cutoff = Math.floor(Date.now() / 1000) - 90 * 24 * 60 * 60;
const row = getDb()
.prepare(
`WITH latest AS (${LATEST_REVIEWS})
SELECT COUNT(*) AS n FROM latest WHERE mint_url = ? AND created_at >= ?`,
)
.get(mintUrl, cutoff) as { n: number };
return row.n;
const db = await getDb();
const row = await db.get<{ n: number }>(
`WITH latest AS (${LATEST_REVIEWS})
SELECT COUNT(*) AS n FROM latest WHERE mint_url = ? AND created_at >= ?`,
mintUrl,
cutoff,
);
return row?.n ?? 0;
}
function uptime30d(mintUrl: string): number | null {
async function uptime30d(mintUrl: string): Promise<number | null> {
const cutoff = Math.floor(Date.now() / 1000) - 30 * 24 * 60 * 60;
const row = getDb()
.prepare('SELECT AVG(ok) AS up, COUNT(*) AS n FROM probes WHERE mint_url = ? AND ts >= ?')
.get(mintUrl, cutoff) as { up: number | null; n: number };
const db = await getDb();
const row = await db.get<{ up: number | null; n: number }>(
'SELECT AVG(ok) AS up, COUNT(*) AS n FROM probes WHERE mint_url = ? AND ts >= ?',
mintUrl,
cutoff,
);
if (!row.n || row.up === null) return null;
if (!row?.n || row.up === null) return null;
return Math.round(row.up * 1000) / 1000;
}
function recentProbes(mintUrl: string): ProbeSample[] {
async function recentProbes(mintUrl: string): Promise<ProbeSample[]> {
const cutoff = Math.floor(Date.now() / 1000) - 48 * 60 * 60;
return getDb()
.prepare(
'SELECT ts, ok, latency_ms FROM probes WHERE mint_url = ? AND ts >= ? ORDER BY ts ASC',
)
.all(mintUrl, cutoff) as ProbeSample[];
const db = await getDb();
return db.all<ProbeSample>(
'SELECT ts, ok, latency_ms FROM probes WHERE mint_url = ? AND ts >= ? ORDER BY ts ASC',
mintUrl,
cutoff,
);
}
function parseInfo(json: string | null): MintInfo | null {
@@ -181,13 +213,23 @@ function parseInfo(json: string | null): MintInfo | null {
}
}
export function getMintDetail(host: string): MintDetail | null {
const row = mintByHost(host);
export async function getMintDetail(host: string): Promise<MintDetail | null> {
const row = await mintByHost(host);
if (!row) return null;
const now = Math.floor(Date.now() / 1000);
const agg = aggregates().get(row.url);
const item = toListItem(row, agg, priorMean(), now);
// One round trip's worth of latency instead of six, which is the difference between
// a local file and a Postgres server on another host.
const [agg, mean, ratingDistribution, reviews, uptime, probes] = await Promise.all([
aggregates(row.url),
priorMean(),
distribution(row.url),
reviews90d(row.url),
uptime30d(row.url),
recentProbes(row.url),
]);
const item = toListItem(row, agg.get(row.url), mean, now);
const info = parseInfo(row.info_json);
let nuts: string[] = [];
@@ -200,8 +242,20 @@ export function getMintDetail(host: string): MintDetail | null {
}
if (nuts.length === 0 && info) nuts = parseNuts(info.nuts);
/*
* Type-specific columns are spread across the payload rather than nested under a key.
*
* That is what keeps a Cashu detail byte for byte what it was: `ecosystem_json` is
* null for a mint, so nothing is added and no consumer sees a new empty object to
* handle. A federation gets `federation_id`, `invite_codes`, `modules`, `network` and
* the rest at the top level, where the page reads them beside `status` and `name`
* without unwrapping anything.
*/
const ecosystem = parseEcosystem(row);
return {
...item,
...(ecosystem ?? {}),
description: row.description,
pubkey: row.pubkey,
info,
@@ -209,58 +263,99 @@ export function getMintDetail(host: string): MintDetail | null {
first_seen: row.first_seen,
last_probe: row.last_probe,
updated_at: row.updated_at,
rating_distribution: distribution(row.url),
reviews_90d: reviews90d(row.url),
uptime_30d: uptime30d(row.url),
probes_recent: recentProbes(row.url),
rating_distribution: ratingDistribution,
reviews_90d: reviews,
uptime_30d: uptime,
probes_recent: probes,
};
}
let statsCache: { at: number; value: Stats } | null = null;
const STATS_TTL_S = 60;
export function getStats(): Stats {
export async function getStats(): Promise<Stats> {
const now = Math.floor(Date.now() / 1000);
if (statsCache && now - statsCache.at < STATS_TTL_S) return statsCache.value;
const db = getDb();
const counts = db
.prepare(
const db = await getDb();
/*
* The four `mints_*` fields are scoped to `type = 'cashu'`, which is what they always
* counted and what every consumer of them still means: the pulse ticker's "mints
* indexed", the /mints page description, the home page. Letting federations quietly
* inflate a number three pages print in a sentence would be a worse kind of breakage
* than a missing field, because nothing would fail — the sentences would just stop
* being true.
*
* COUNT(*) FILTER, not SUM(status = 'online'): Postgres has no implicit cast from
* boolean to integer, so the SQLite spelling is a type error there.
*/
const [counts, federations, reviews, fedimintReviews] = await Promise.all([
db.get<{ total: number; online: number; offline: number; degraded: number }>(
`SELECT
COUNT(*) AS total,
SUM(status = 'online') AS online,
SUM(status = 'offline') AS offline,
SUM(status = 'degraded') AS degraded
FROM mints`,
)
.get() as { total: number; online: number | null; offline: number | null; degraded: number | null };
COUNT(*) AS total,
COUNT(*) FILTER (WHERE status = 'online') AS online,
COUNT(*) FILTER (WHERE status = 'offline') AS offline,
COUNT(*) FILTER (WHERE status = 'degraded') AS degraded
FROM mints WHERE type = 'cashu'`,
),
db.get<{ total: number; online: number; offline: number; announced: number }>(
`SELECT
COUNT(*) AS total,
COUNT(*) FILTER (WHERE status = 'online') AS online,
COUNT(*) FILTER (WHERE status = 'offline') AS offline,
COUNT(*) FILTER (WHERE status = 'announced') AS announced
FROM mints WHERE type = 'fedimint'`,
),
db.get<{ n: number; last: number | null }>(
`WITH latest AS (${LATEST_REVIEWS}) SELECT COUNT(*) AS n, MAX(created_at) AS last FROM latest`,
),
db.get<{ n: number }>(
`WITH latest AS (${LATEST_REVIEWS})
SELECT COUNT(*) AS n FROM latest
WHERE mint_url IN (SELECT url FROM mints WHERE type = 'fedimint')`,
),
]);
const reviews = db
.prepare(`WITH latest AS (${LATEST_REVIEWS}) SELECT COUNT(*) AS n, MAX(created_at) AS last FROM latest`)
.get() as { n: number; last: number | null };
const cashuTotal = counts?.total ?? 0;
const value: Stats = {
mints_total: counts.total,
mints_online: counts.online ?? 0,
mints_offline: counts.offline ?? 0,
mints_degraded: counts.degraded ?? 0,
reviews_total: reviews.n,
last_review_at: reviews.last,
mints_total: cashuTotal,
mints_online: counts?.online ?? 0,
mints_offline: counts?.offline ?? 0,
mints_degraded: counts?.degraded ?? 0,
reviews_total: reviews?.n ?? 0,
last_review_at: reviews?.last ?? null,
updated_at: now,
cashu_total: cashuTotal,
fedimint_total: federations?.total ?? 0,
fedimint_online: federations?.online ?? 0,
fedimint_offline: federations?.offline ?? 0,
fedimint_announced: federations?.announced ?? 0,
fedimint_reviews: fedimintReviews?.n ?? 0,
};
statsCache = { at: now, value };
return value;
}
/** Health bypasses the stats cache: it is the endpoint you page on. */
export function getHealth(): Health {
const now = Math.floor(Date.now() / 1000);
const lastProbe = getStateNumber('last_probe_at');
const lastDiscovery = getStateNumber('last_discovery_at');
const discoveryOk = getState('last_discovery_ok') !== '0';
/** Drop the in-process stats cache. For tests that mutate the database underneath it. */
export function resetStatsCache(): void {
statsCache = null;
}
const tracked = (getDb().prepare('SELECT COUNT(*) AS n FROM mints').get() as { n: number }).n;
/** Health bypasses the stats cache: it is the endpoint you page on. */
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([
getStateNumber('last_probe_at'),
getStateNumber('last_discovery_at'),
getState('last_discovery_ok'),
db.get<{ n: number }>('SELECT COUNT(*) AS n FROM mints'),
]);
const discoveryOk = discoveryOkRaw !== '0';
const staleAfter = config.probeIntervalMin * 60 * 3;
const probeStale = lastProbe === null || now - lastProbe > staleAfter;
@@ -269,7 +364,7 @@ export function getHealth(): Health {
uptime_s: now - startedAt,
last_probe_at: lastProbe,
last_discovery_at: lastDiscovery,
mints_tracked: tracked,
mints_tracked: tracked?.n ?? 0,
updated_at: now,
};
}
+30 -12
View File
@@ -59,10 +59,11 @@ function pad(s: string, width: number): string {
}
async function main(): Promise<void> {
getDb();
const db = await getDb();
log.info('seeding', { db: db.label });
let added = 0;
for (const url of SEED_MINTS) if (insertMintIfNew(url)) added++;
for (const url of SEED_MINTS) if (await insertMintIfNew(url)) added++;
log.info('seed list ingested', { listed: SEED_MINTS.length, added });
await probeAll();
@@ -75,13 +76,14 @@ async function main(): Promise<void> {
await probeAll();
}
const mints = listMints();
const width = { host: 44, status: 9, score: 7, reviews: 8, rating: 7 };
const listed = await listMints();
const width = { host: 44, type: 9, status: 10, score: 7, reviews: 8, rating: 7 };
console.log('');
console.log(
[
pad('MINT', width.host),
pad('TYPE', width.type),
pad('STATUS', width.status),
pad('SCORE', width.score),
pad('REVIEWS', width.reviews),
@@ -89,12 +91,13 @@ async function main(): Promise<void> {
'VERSION',
].join(' '),
);
console.log('-'.repeat(110));
console.log('-'.repeat(120));
for (const m of mints) {
for (const m of listed) {
console.log(
[
pad(m.host, width.host),
pad(m.type, width.type),
pad(m.status, width.status),
pad(m.score.toFixed(3), width.score),
pad(String(m.review_count), width.reviews),
@@ -104,16 +107,31 @@ async function main(): Promise<void> {
);
}
const online = mints.filter((m) => m.status === 'online').length;
const offline = mints.filter((m) => m.status === 'offline').length;
const reviews = mints.reduce((sum, m) => sum + m.review_count, 0);
console.log('-'.repeat(110));
const mints = listed.filter((m) => m.type === 'cashu');
const federations = listed.filter((m) => m.type === 'fedimint');
const count = (rows: typeof listed, status: string) =>
rows.filter((m) => m.status === status).length;
const reviews = listed.reduce((sum, m) => sum + m.review_count, 0);
console.log('-'.repeat(120));
console.log(
`${mints.length} mints (${online} online, ${offline} offline), ${reviews} reviews indexed`,
`${mints.length} mints (${count(mints, 'online')} online, ${count(mints, 'offline')} offline)`,
);
/*
* Federations are counted with `announced` called out rather than folded into the
* offline total. It is the honest word for "nothing checks this one", and rolling it
* into either of the other two would be the first place the site started claiming to
* know something it does not.
*/
console.log(
`${federations.length} federations (${count(federations, 'online')} confirmed up, ` +
`${count(federations, 'offline')} reported down, ` +
`${count(federations, 'announced')} announced only)`,
);
console.log(`${reviews} reviews indexed`);
closePool();
closeDb();
await closeDb();
process.exit(0);
}
+24 -8
View File
@@ -12,22 +12,31 @@ export function createApp(): Hono {
app.use('/api/*', cors());
app.use('/icons/*', cors());
app.get('/api/health', (c) => {
const health = getHealth();
app.get('/api/health', async (c) => {
const health = await getHealth();
c.header('Cache-Control', 'no-store');
return c.json(health, health.status === 'ok' ? 200 : 503);
});
app.get('/api/stats', (c) => c.json(getStats()));
app.get('/api/stats', async (c) => c.json(await getStats()));
app.get('/api/mints', (c) => {
/*
* The whole index, every ecosystem, each item carrying its `type`.
*
* `?type=cashu` / `?type=fedimint` narrows it, which is what the two list pages ask
* for at build time. Unknown values are passed through rather than rejected: they
* return an empty list, which is the honest answer to "show me the ecosystem this
* build does not have".
*/
app.get('/api/mints', async (c) => {
const raw = c.req.query('limit');
const limit = raw ? Number.parseInt(raw, 10) : undefined;
return c.json(listMints(Number.isFinite(limit) ? limit : undefined));
const type = c.req.query('type')?.trim() || undefined;
return c.json(await listMints(Number.isFinite(limit) ? limit : undefined, type));
});
app.get('/api/mints/:host', (c) => {
const detail = getMintDetail(c.req.param('host'));
app.get('/api/mints/:host', async (c) => {
const detail = await getMintDetail(c.req.param('host'));
if (!detail) return c.json({ error: 'not_found', message: 'No mint with that host' }, 404);
return c.json(detail);
});
@@ -38,7 +47,14 @@ export function createApp(): Hono {
serveStatic({
root: path.relative(process.cwd(), config.iconDir) || '.',
rewriteRequestPath: (p) => p.replace(/^\/icons/, ''),
onFound: (_p, c) => c.header('Cache-Control', 'public, max-age=86400'),
onFound: (_p, c) => {
c.header('Cache-Control', 'public, max-age=86400');
// These files came from mint operators. nosniff pins the served type, and the
// sandbox neutralizes anything script-capable (an SVG cached before icons.ts
// stopped accepting them) when the file is opened directly on this origin.
c.header('X-Content-Type-Options', 'nosniff');
c.header('Content-Security-Policy', 'sandbox');
},
}),
);