theseus/bns-indexer.js

332 lines
16 KiB
JavaScript
Raw Normal View History

// BNS name indexer — runs in its own process (Electron utilityProcess),
// started by main.js, so the snapshot parsing, index builds and electrum
// traffic never compete with the browser's main thread.
//
// Boot is two-phase:
// 1. fast lane: load the local snapshot (user cache > bundled), build the
// index from it and publish it. No network. This is what the first tab
// needs, and it is all that runs while the browser is still starting.
// 2. full work, once main says "go" (its first page has loaded): electrum
// server discovery, the Sia snapshot refresh, and the 30 s delta poll.
// A lookup that misses the snapshot before then still gets a live
// delta poll on demand — only that, nothing else.
//
// Protocol (process.parentPort messages):
// main -> indexer { type: "init", resolverPath, userData, bundledSnapshot, torPort }
// { type: "go" } start phase 2
// { type: "tor", port } electrum via Tor (null = direct)
// { type: "stop-polling" }
// { id, type: "poll" | "rebuild" | "ready" } replies { type: "reply", id, ok }
// indexer -> main { type: "index", builtAt, names: [[key, entry]...], network }
// { type: "fresh", builtAt } poll succeeded, nothing changed
// { type: "reply", id, ok, error? }
//
// The index entries are the resolver's own (records, owner, txid, height, …)
// built from the raw beacon transactions — the trust model is unchanged.
"use strict";
const fs = require("fs");
const path = require("path");
const { pathToFileURL } = require("url");
const WebSocket = require("ws");
const port = process.parentPort;
const log = (...a) => console.log("[bns]", ...a);
// Same guard as main.js: ws turns "closed while still connecting" into an
// uncaught error event. Electrum servers that reset mid-handshake are routine.
{
const close = WebSocket.prototype.close;
WebSocket.prototype.close = function (...args) {
if (this.readyState === WebSocket.CONNECTING && this.listenerCount("error") === 0) this.once("error", () => {});
return close.apply(this, args);
};
}
// A flaky network must not take the indexer down; main restarts it if it
// exits, but the index would go cold for a moment.
process.on("uncaughtException", (err) => { try { console.error("[bns] uncaught:", (err && err.stack) || err); } catch {} });
process.on("unhandledRejection", (err) => { try { console.error("[bns] unhandled rejection:", (err && err.message) || err); } catch {} });
let cfg = null;
let R = null;
async function getResolver() {
if (!R) R = await import(pathToFileURL(cfg.resolverPath).href);
return R;
}
// ---- Tor ------------------------------------------------------------------
let torAgent = null;
async function setTor(torPort) {
if (!torPort) { torAgent = null; return; }
const { SocksProxyAgent } = await import("socks-proxy-agent");
torAgent = new SocksProxyAgent(`socks5h://127.0.0.1:${torPort}`);
}
class TorWebSocket extends WebSocket { constructor(url, opts) { super(url, { agent: torAgent, ...opts }); } }
const currentWS = () => (torAgent ? TorWebSocket : WebSocket);
// ---- electrum server pool: hardcoded seed + on-chain discovery, persisted ----
// Bootstrap from the baked-in seed (with pinned IPs), then refresh from the
// on-chain ELECTRUM_LIST_NAME record so the pool can be rotated without a new
// build. The last discovered list is cached to disk and tried first next launch.
let electrumPool = null;
let lastElectrumRefresh = 0;
// Per-network cache files. Chipnet keeps the historical names so existing
// profiles are not invalidated; any other network gets its own files.
let networkId = "chipnet";
const netSuffix = () => (networkId === "chipnet" ? "" : `.${networkId}`);
const electrumFile = () => path.join(cfg.userData, `electrum-servers${netSuffix()}.json`);
const snapshotUserPath = () => path.join(cfg.userData, `bns-name-snapshot${netSuffix()}.json`);
const serverKey = (s) => (typeof s === "string" ? s : s && s.url);
function mergeServers(preferred, rest) {
const seen = new Set(), out = [];
for (const s of [...(preferred || []), ...(rest || [])]) {
const k = serverKey(s);
if (k && !seen.has(k)) { seen.add(k); out.push(s); }
}
return out;
}
async function initElectrumPool() {
const r = await getResolver();
const seed = r.ELECTRUM || r.CHIPNET_ELECTRUM;
let saved = [];
try { if (fs.existsSync(electrumFile())) saved = JSON.parse(fs.readFileSync(electrumFile(), "utf8")); } catch {}
electrumPool = mergeServers(saved, seed); // discovered first, seed always kept
}
async function refreshElectrumPool() {
try {
const { fetchElectrumServers } = await getResolver();
const found = await fetchElectrumServers({ WebSocket: currentWS(), directIP: true, electrum: electrumPool });
if (found && found.length) {
electrumPool = mergeServers(found, electrumPool);
try { fs.writeFileSync(electrumFile(), JSON.stringify(found, null, 2)); } catch {}
}
} catch { /* list unpublished or unreachable — keep the current pool */ }
}
function maybeRefreshElectrum() {
if (!started) return; // phase 2 only
if (Date.now() - lastElectrumRefresh < 30 * 60 * 1000) return;
lastElectrumRefresh = Date.now();
refreshElectrumPool();
}
// ---- index state + publishing ---------------------------------------------
let index = null; // Map name -> entry
let currentSnapshotState = null; // raw { beacon, history, txs, … } behind `index`
let lastSig = "";
// Cheap fingerprint of the snapshot: an index rebuilt from the same history
// at the same heights is the same index, so main isn't sent a copy.
const sigOf = (snap) => `${snap.history.length}:${snap.history.reduce((a, h) => a + (Number(h.height) || 0), 0)}`;
function publish(idx, snap, reason) {
index = idx;
const sig = snap ? sigOf(snap) : `full:${Date.now()}`;
const builtAt = reason === "snapshot" ? 0 : Date.now(); // a snapshot is a floor, not a ceiling
if (sig === lastSig) { port.postMessage({ type: "fresh", builtAt }); return; }
lastSig = sig;
port.postMessage({ type: "index", builtAt, names: [...idx], network: networkId, reason });
}
function readSnapshotFrom(p) {
try {
if (!fs.existsSync(p)) return null;
const parsed = JSON.parse(fs.readFileSync(p, "utf8"));
if (!parsed || !Array.isArray(parsed.history)) return null;
return parsed;
} catch { return null; }
}
// Phase 1: the local snapshot, nothing else.
async function warmFromSnapshot() {
const r = await getResolver();
networkId = r.NETWORK?.id || "chipnet";
if (!r.buildIndexFromSnapshot) return false; // an older resolver-web.js
const snap = readSnapshotFrom(snapshotUserPath()) || readSnapshotFrom(cfg.bundledSnapshot);
if (!snap) return false;
try {
const t0 = Date.now();
const idx = r.buildIndexFromSnapshot({ snapshot: snap });
currentSnapshotState = snap;
publish(idx, snap, "snapshot");
log(`warm-started from snapshot: ${idx.size} names in ${Date.now() - t0} ms @ height=${snap.asOfHeight ?? "?"} root=${snap.root ?? "?"}`);
return true;
} catch (e) {
log("snapshot warm-start failed:", e.message);
return false;
}
}
// ---- continuous background delta refresh ----------------------------------
// One electrum connection per poll: fetch the beacon's history (one call),
// fetch only the tx bodies not already held, rebuild locally, persist.
const POLL_INTERVAL_MS = 30_000;
// Neither the electrum connect nor its WebSocket handshake has a timeout of
// its own; a server that accepts TCP and then stalls would wedge the loop.
const POLL_DEADLINE_MS = 45_000;
let pollInFlight = null, pollTimer = null;
async function pollAndMerge() {
if (pollInFlight) return pollInFlight;
let conn = null, deadline;
const work = (async () => {
const r = await getResolver();
if (!r.connectElectrum || !r.BEACON_SCRIPTHASH || !r.buildIndexFromSnapshot) return rebuild();
if (!electrumPool) await initElectrumPool();
// Base state: memory > user cache > bundled > empty. The "empty" branch
// turns a first start without any snapshot into a full fetch.
const snap = currentSnapshotState
|| readSnapshotFrom(snapshotUserPath())
|| readSnapshotFrom(cfg.bundledSnapshot)
|| { beacon: r.BEACON_SCRIPTHASH, history: [], txs: {} };
const el = conn = await r.connectElectrum({ electrum: electrumPool, WebSocket: currentWS(), directIP: true });
try {
const freshHistory = await el.call("blockchain.scripthash.get_history", [r.BEACON_SCRIPTHASH]);
const known = new Set(snap.history.map((h) => h.tx_hash));
// Merge fresh into snapshot history (dedup by tx_hash, keep fresh height —
// an event that was mempool at snapshot time now has a real height).
const merged = new Map(snap.history.map((h) => [h.tx_hash, h]));
const txs = { ...(snap.txs || {}) };
let added = 0;
for (const h of freshHistory) {
if (!known.has(h.tx_hash)) {
try { txs[h.tx_hash] = await el.call("blockchain.transaction.get", [h.tx_hash, true]); added++; }
catch { /* unreadable — the reduction rules ignore missing txs */ }
}
merged.set(h.tx_hash, { tx_hash: h.tx_hash, height: h.height });
}
const history = [...merged.values()];
currentSnapshotState = { ...snap, beacon: r.BEACON_SCRIPTHASH, history, txs };
const idx = r.buildIndexFromSnapshot({ snapshot: currentSnapshotState });
publish(idx, currentSnapshotState, "poll");
try {
fs.mkdirSync(path.dirname(snapshotUserPath()), { recursive: true });
fs.writeFileSync(snapshotUserPath(), JSON.stringify(currentSnapshotState));
} catch { /* readonly userdata / disk full — skip */ }
if (added > 0) log(`delta-refresh: +${added} new tx${added === 1 ? "" : "s"} (total ${history.length}, index has ${idx.size} names)`);
return true;
} finally { try { el.close(); } catch {} }
})().catch(() => false);
const timeout = new Promise((resolve) => {
deadline = setTimeout(() => { try { conn?.close(); } catch {} resolve(false); }, POLL_DEADLINE_MS);
});
const p = Promise.race([work, timeout]).finally(() => {
clearTimeout(deadline);
if (pollInFlight === p) pollInFlight = null;
maybeRefreshElectrum();
});
pollInFlight = p;
return p;
}
// Full walk over electrum — only for a resolver without the delta primitives,
// or an explicit rebuild request.
let rebuilding = null;
function rebuild() {
if (rebuilding) return rebuilding;
rebuilding = (async () => {
const r = await getResolver();
if (!electrumPool) await initElectrumPool();
const idx = await r.buildIndex({ WebSocket: currentWS(), directIP: true, electrum: electrumPool });
publish(idx, null, "rebuild");
return true;
})().catch(() => false).finally(() => { rebuilding = null; });
return rebuilding;
}
// The published raw snapshot (public, refreshed every 10 min on the VPS):
// one download that can carry a long-closed browser across days of beacon
// events, so the electrum poll after it only fetches what is newer still.
// Merged in now — not saved for the next start — and only ever adds: the
// index is rebuilt here from the transactions, as with any other source.
async function refreshFromPublishedSnapshot() {
try {
const r = await getResolver();
// Two public copies (dl vhost, Sia via the gateway): first one that answers.
const urls = [r.NETWORK?.nameListUrl, r.NETWORK?.nameListMirrorUrl].filter(Boolean);
let pub = null;
for (const url of urls) {
try {
const res = await fetch(url, { redirect: "follow", signal: AbortSignal.timeout(15000) });
if (!res.ok) { log(`published snapshot ${new URL(url).host}: HTTP ${res.status}`); continue; }
pub = await res.json();
break;
} catch (e) { log(`published snapshot ${new URL(url).host}: ${e?.cause?.code || e?.message || e}`); }
}
if (!pub) return;
if (!pub || !Array.isArray(pub.history) || !pub.txs || (r.BEACON_SCRIPTHASH && pub.beacon !== r.BEACON_SCRIPTHASH)) return;
const base = currentSnapshotState || { beacon: r.BEACON_SCRIPTHASH, history: [], txs: {} };
const merged = new Map(base.history.map((h) => [h.tx_hash, h]));
const txs = { ...(base.txs || {}) };
let added = 0;
for (const h of pub.history) {
const mine = merged.get(h.tx_hash);
if (!pub.txs[h.tx_hash] && !txs[h.tx_hash]) continue; // no evidence, skip
if (!txs[h.tx_hash]) { txs[h.tx_hash] = pub.txs[h.tx_hash]; added++; }
// Keep a confirmed height over a mempool (0 / negative) one.
if (!mine || (!(Number(mine.height) > 0) && Number(h.height) > 0)) merged.set(h.tx_hash, { tx_hash: h.tx_hash, height: h.height });
}
const next = { ...base, beacon: base.beacon || pub.beacon, history: [...merged.values()], txs };
if (sigOf(next) === sigOf(base)) return;
currentSnapshotState = next;
publish(r.buildIndexFromSnapshot({ snapshot: next }), next, "published");
try {
fs.mkdirSync(path.dirname(snapshotUserPath()), { recursive: true });
fs.writeFileSync(snapshotUserPath(), JSON.stringify(next));
} catch {}
log(`published snapshot merged: +${added} tx${added === 1 ? "" : "s"} (total ${next.history.length}, ${index.size} names)`);
} catch (e) { log("published snapshot unavailable:", e?.message || e); }
}
// ---- phases -----------------------------------------------------------------
let warmed = null; // promise: phase 1 finished (index may still be null)
let started = false; // phase 2 running
function startFullWork() {
if (started) return;
started = true;
(async () => {
await warmed;
if (!electrumPool) await initElectrumPool();
// Fastest catch-up first: one download of the published snapshot, then
// the electrum delta for whatever is newer than that.
await refreshFromPublishedSnapshot();
lastElectrumRefresh = Date.now();
refreshElectrumPool();
// No snapshot at all (first start of a build without one): the poll's
// "empty" base state does the full fetch.
pollAndMerge();
pollTimer = setInterval(() => pollAndMerge(), POLL_INTERVAL_MS);
})().catch((e) => log("start failed:", e?.message));
}
// "ready": resolves once there is any index — the local snapshot, or (none
// on disk) the published one, one download; the live poll only if that
// failed too, since a first electrum walk fetches every beacon transaction.
let firstFetch = null;
async function whenReady() {
await warmed;
if (index) return true;
firstFetch = firstFetch || refreshFromPublishedSnapshot().finally(() => { firstFetch = null; });
await firstFetch;
if (index) return true;
return pollAndMerge();
}
port.on("message", async (e) => {
const m = e.data || {};
const reply = (ok, error) => { if (m.id != null) port.postMessage({ type: "reply", id: m.id, ok: !!ok, error }); };
try {
switch (m.type) {
case "init":
cfg = m;
await setTor(m.torPort);
warmed = warmFromSnapshot();
break;
case "go": startFullWork(); break;
case "tor": await setTor(m.port); break;
case "stop-polling": if (pollTimer) { clearInterval(pollTimer); pollTimer = null; } break;
case "ready": reply(await whenReady()); break;
// An on-demand poll is allowed before phase 2: it is the "this tab's
// name isn't in the snapshot" case, and costs one electrum round trip.
case "poll": await warmed; reply(await pollAndMerge()); break;
case "rebuild": await warmed; reply(await rebuild()); break;
default: reply(false, "unknown request");
}
} catch (err) { reply(false, err?.message || String(err)); }
});