@@ -0,0 +1,257 @@
|
||||
import { SimplePool } from 'nostr-tools/pool';
|
||||
import type { Event as NostrEvent, Filter } from 'nostr-tools';
|
||||
import {
|
||||
KIND_MINT_ANNOUNCEMENT,
|
||||
KIND_REVIEW,
|
||||
isCashuMintReview,
|
||||
mintPubkeyRef,
|
||||
mintUrlSpellings,
|
||||
mintUrlsFromEvent,
|
||||
normalizeMintUrl,
|
||||
parseRating,
|
||||
} from '@cashumints/shared';
|
||||
import { config } from './config.ts';
|
||||
import { getDb, getStateNumber, setState } from './db.ts';
|
||||
import { log } from './log.ts';
|
||||
import { insertMintIfNew, mintUrlByHost, mintUrlByPubkey } 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(
|
||||
mintUrl: string,
|
||||
mintPubkey: string | null,
|
||||
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)] });
|
||||
}
|
||||
filters.push({ ...base, '#u': mintUrlSpellings(mintUrl) });
|
||||
|
||||
const batches = await Promise.all(
|
||||
filters.map((filter) =>
|
||||
getPool()
|
||||
.querySync(config.relays, filter, { maxWait: MAX_WAIT_MS })
|
||||
.catch(() => [] as NostrEvent[]),
|
||||
),
|
||||
);
|
||||
|
||||
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);
|
||||
if (created) newMints.add(created);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 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): string | 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);
|
||||
if (canonical) return canonical;
|
||||
}
|
||||
|
||||
const pubkey = mintPubkeyRef(event);
|
||||
return pubkey ? mintUrlByPubkey(pubkey) : null;
|
||||
}
|
||||
|
||||
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');
|
||||
// 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 {
|
||||
const announcements = await fetchKind(KIND_MINT_ANNOUNCEMENT, since);
|
||||
events += announcements.length;
|
||||
for (const e of announcements) ingestMintUrls(e, 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;
|
||||
|
||||
for (const e of byId.values()) ingestMintUrls(e, 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()]);
|
||||
|
||||
// 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 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.url, target.pubkey, since);
|
||||
if (found.length > 0) {
|
||||
events += found.length;
|
||||
ingestReviews(found);
|
||||
}
|
||||
}
|
||||
}),
|
||||
);
|
||||
|
||||
setState('discovery_since', String(now));
|
||||
setState('last_discovery_at', String(now));
|
||||
setState('last_discovery_ok', '1');
|
||||
} catch (err) {
|
||||
ok = false;
|
||||
setState('last_discovery_ok', '0');
|
||||
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 };
|
||||
}
|
||||
Reference in New Issue
Block a user