Add Fedimint discovery, dual SQLite/Postgres storage, richer review handling, and generated social imagery.
555 lines
20 KiB
TypeScript
555 lines
20 KiB
TypeScript
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,
|
|
mintPubkeyRef,
|
|
mintUrlSpellings,
|
|
mintUrlsFromEvent,
|
|
normalizeMintUrl,
|
|
parseFedimintAnnouncement,
|
|
parseRating,
|
|
reviewEcosystem,
|
|
reviewTargetId,
|
|
reviewTargetKind,
|
|
type FedimintAnnouncement,
|
|
} from '@cashumints/shared';
|
|
import { config } from './config.ts';
|
|
import { getDb, setState, getStateNumber } from './db.ts';
|
|
import type { Sql } from './db-driver.ts';
|
|
import { log } from './log.ts';
|
|
import { insertMintIfNew, upsertFedimint } from './mints.ts';
|
|
|
|
const QUERY_LIMIT = 500;
|
|
const MAX_WAIT_MS = 12_000;
|
|
/** Backfill page cap. 20 pages x 500 is far more than the network holds today. */
|
|
const MAX_PAGES = 20;
|
|
/** The old site's second filter: a recent window alongside the unbounded one, because
|
|
* some relays answer an unbounded query with an arbitrary slice. See NOTES.md. */
|
|
const RECENT_WINDOW_S = 30 * 24 * 60 * 60;
|
|
|
|
export interface DiscoveryResult {
|
|
events: number;
|
|
newMints: string[];
|
|
newReviews: number;
|
|
ok: boolean;
|
|
}
|
|
|
|
let pool: SimplePool | null = null;
|
|
|
|
function getPool(): SimplePool {
|
|
pool ??= new SimplePool();
|
|
return pool;
|
|
}
|
|
|
|
export function closePool(): void {
|
|
try {
|
|
pool?.close(config.relays);
|
|
} catch {
|
|
// A relay that is already gone throws on close. Nothing to do about it.
|
|
}
|
|
pool = null;
|
|
}
|
|
|
|
/**
|
|
* Query one kind, paging backwards with `until` until a page yields nothing new.
|
|
* Relays cap `limit` independently, so paging is the only way a fresh database
|
|
* converges to the complete history.
|
|
*/
|
|
async function fetchKind(kind: number, since: number | null): Promise<NostrEvent[]> {
|
|
const seen = new Map<string, NostrEvent>();
|
|
let until: number | undefined;
|
|
|
|
for (let page = 0; page < MAX_PAGES; page++) {
|
|
const filter: Filter = { kinds: [kind], limit: QUERY_LIMIT };
|
|
if (since !== null) filter.since = since;
|
|
if (until !== undefined) filter.until = until;
|
|
|
|
let batch: NostrEvent[];
|
|
try {
|
|
batch = await getPool().querySync(config.relays, filter, { maxWait: MAX_WAIT_MS });
|
|
} catch (err) {
|
|
log.warn('relay query failed', {
|
|
kind,
|
|
page,
|
|
reason: err instanceof Error ? err.message : String(err),
|
|
});
|
|
break;
|
|
}
|
|
|
|
let added = 0;
|
|
let oldest = Number.POSITIVE_INFINITY;
|
|
for (const e of batch) {
|
|
oldest = Math.min(oldest, e.created_at);
|
|
if (!seen.has(e.id)) {
|
|
seen.set(e.id, e);
|
|
added++;
|
|
}
|
|
}
|
|
|
|
// Nothing new, or the relays returned less than a full page: history is exhausted.
|
|
if (added === 0 || batch.length < QUERY_LIMIT || !Number.isFinite(oldest)) break;
|
|
until = oldest - 1;
|
|
if (since !== null && until <= since) break;
|
|
}
|
|
|
|
return [...seen.values()];
|
|
}
|
|
|
|
/**
|
|
* Reviews for one specific mint, asked for by name.
|
|
*
|
|
* The global sweep above does NOT converge on its own: relays answer an unbounded
|
|
* `{kinds:[38000]}` with an arbitrary capped slice, but answer a filtered query for one
|
|
* mint in full. Measured against the live network, the global sweep found 43 reviews for
|
|
* mint.minibits.cash/Bitcoin while this targeted query finds 80. Without this pass the
|
|
* count in the API disagrees with the list the browser renders on the same page.
|
|
*/
|
|
async function fetchReviewsForMint(
|
|
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 (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) });
|
|
}
|
|
|
|
const batches = await Promise.all(
|
|
filters.map((filter) =>
|
|
getPool()
|
|
.querySync(config.relays, filter, { maxWait: MAX_WAIT_MS })
|
|
.catch(() => [] as NostrEvent[]),
|
|
),
|
|
);
|
|
|
|
return batches.flat();
|
|
}
|
|
|
|
/**
|
|
* 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.
|
|
*
|
|
* `u` tag first, resolved through the slug so that every spelling of a mint URL (with
|
|
* or without a trailing slash, either path case) lands on the one canonical row. Then
|
|
* 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, 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 = index.byHost.get(normalized.host);
|
|
if (canonical) return canonical;
|
|
}
|
|
|
|
const pubkey = mintPubkeyRef(event);
|
|
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 : 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);
|
|
|
|
const newMints = new Set<string>();
|
|
let newReviews = 0;
|
|
let events = 0;
|
|
let ok = true;
|
|
|
|
try {
|
|
/*
|
|
* 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.
|
|
const reviews = await fetchKind(KIND_REVIEW, since);
|
|
// The recent window catches anything a relay dropped from the unbounded query.
|
|
const recent =
|
|
since === null ? await fetchKind(KIND_REVIEW, now - RECENT_WINDOW_S) : [];
|
|
|
|
const byId = new Map<string, NostrEvent>();
|
|
for (const e of [...reviews, ...recent]) byId.set(e.id, e);
|
|
events += byId.size;
|
|
|
|
await ingestMintUrls(byId.values(), now, newMints);
|
|
|
|
// 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 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;
|
|
await Promise.all(
|
|
Array.from({ length: Math.min(CONCURRENCY, targets.length) }, async () => {
|
|
while (cursor < targets.length) {
|
|
const target = targets[cursor++];
|
|
if (!target) continue;
|
|
const found = await fetchReviewsForMint(target, since);
|
|
if (found.length > 0) {
|
|
events += found.length;
|
|
newReviews += await ingestReviews(found, index);
|
|
}
|
|
}
|
|
}),
|
|
);
|
|
|
|
await setState('discovery_since', String(now));
|
|
await setState('last_discovery_at', String(now));
|
|
await setState('last_discovery_ok', '1');
|
|
} catch (err) {
|
|
ok = false;
|
|
await setState('last_discovery_ok', '0').catch(() => undefined);
|
|
log.error('discovery failed', { reason: err instanceof Error ? err.message : String(err) });
|
|
}
|
|
|
|
log.info('discovery cycle', {
|
|
mode: backfill ? 'backfill' : 'incremental',
|
|
events,
|
|
new_mints: newMints.size,
|
|
new_reviews: newReviews,
|
|
ok,
|
|
ms: Date.now() - started,
|
|
});
|
|
|
|
return { events, newMints: [...newMints], newReviews, ok };
|
|
}
|