// BNS index host. Runs in its own process (Electron utilityProcess), started // by main.js, so snapshot parsing, index builds, the Ariadne pipe and any // electrum traffic never compete with the browser's main thread. // Spec: DESIGN-bns-indexer-service.md ("Theseus"). // // Ariadne's Thread owns BNS indexing on a machine; Theseus reads from it // through the shared source chain (Argus/src/lib/bns-source-chain.js): // 1. Ariadne's indexer over its pipe, server-checked (ariadne-helper.exe): // a snapshot, then pushes // 2. local copies: Ariadne's two files (All users / Just me), Theseus's own // raw copy, the bundled snapshot (the one with the most evidence wins), // and Theseus's own index copy written from the pipe // 3. (main.js) the quick single-name lookup on the gateway, provisional // Theseus's own electrum indexer (the shared core, bns-index-core.js) is a // standby, the last resort: // - started when Ariadne is unhealthy: no pipe a few seconds after launch, // a pipe that fails the server check, a pipe that went silent, or an // index that is not being confirmed while Ariadne says it is not paused. // Warm-started from the richest local copy (normally Ariadne's), so no // download and no cold sync. // - stopped again once Ariadne has been healthy for 90 s, so a flapping // service cannot start and stop it over and over. // Under Economy Ariadne is paused, not unhealthy: no takeover. // With no Ariadne at all (portable, not installed, a pre-0.2 install) it is // the 0.3.70 behaviour: the own indexer runs throughout. // // The only thing Theseus ever does to Ariadne: when "Launch at start" is off, // it runs the indexer's task once at launch. It never writes Ariadne's files. // // Protocol (process.parentPort messages), unchanged for main.js: // main -> host { type: "init", resolverPath, coreLib, chainLib, pipeLib, // userData, bundledSnapshot, torPort } // { type: "go" } network work may start // { type: "tor", port } electrum via Tor (null = direct) // { type: "stop-polling" } // { id, type: "poll" | "rebuild" | "ready" } replies { type: "reply", id, ok } // host -> main { type: "index", builtAt, names: [[key, entry]...], network, source } // { type: "fresh", builtAt } confirmed, nothing changed // { type: "status", ... } where the names come from (Settings) // { type: "reply", id, ok, error? } "use strict"; const fs = require("fs"); const path = require("path"); const { execFile } = require("child_process"); 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); }; } 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 {} }); const LAUNCH_GRACE_MS = 4_000; // the pipe gets this long before the own indexer starts const STALE_MS = 10 * 60_000; // unpaused Ariadne whose index is not confirmed this long: unhealthy const HANDBACK_MS = 90_000; // healthy this long before the own indexer is stopped const CHECK_EVERY_MS = 10_000; let cfg = null, R = null, core = null, chainLib = null, pipeLib = null; const importEsm = (p) => import(pathToFileURL(p).href); // ---- 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); // ---- Ariadne's Thread on this machine --------------------------------------- // Each scope it may be installed in: where its files are, its pipe, and how // to check the pipe's server (Ariadne's helper; the server's image is // Ariadne's own node.exe). The install folder comes from Ariadne's uninstall // key. A pre-0.2 Ariadne has no helper and no pipe: only files (none either). const ARIADNE_KEY = "Software\\Microsoft\\Windows\\CurrentVersion\\Uninstall\\{7E7A5F1C-3B4E-4C8A-9E1D-ARIADNERSLVR}_is1"; function execOut(file, args) { return new Promise((resolve) => execFile(file, args, { windowsHide: true, timeout: 10_000 }, (err, out) => resolve(err ? "" : String(out)))); } async function installLocation(hive) { const out = await execOut("reg.exe", ["query", `${hive}\\${ARIADNE_KEY}`, "/v", "InstallLocation", "/reg:64"]); const m = /InstallLocation\s+REG_\w+\s+(.+)/.exec(out); return m ? m[1].trim() : null; } async function findAriadne() { // Tests point Theseus at a scratch Ariadne: {scope, pipe, stateDir, helper, imagePrefix}. if (process.env.THESEUS_ARIADNE_TEST) { try { return [JSON.parse(process.env.THESEUS_ARIADNE_TEST)]; } catch { return []; } } if (process.platform !== "win32") return []; const scopes = []; const programData = process.env.ProgramData || process.env.PROGRAMDATA || "C:\\ProgramData"; const localAppData = process.env.LOCALAPPDATA || ""; const machineLoc = await installLocation("HKLM"); scopes.push({ scope: "machine", pipe: "\\\\.\\pipe\\ariadne-bns", task: "Ariadne BNS Indexer", stateDir: path.join(programData, "Ariadne"), loc: machineLoc, }); const userLoc = await installLocation("HKCU"); if (userLoc || (localAppData && fs.existsSync(path.join(localAppData, "Ariadne", "index")))) { const sid = (/"(S-1-[\d-]+)"/.exec(await execOut("whoami", ["/user", "/fo", "csv", "/nh"])) || [])[1]; scopes.push({ scope: "user", pipe: sid ? `\\\\.\\pipe\\ariadne-bns-${sid}` : null, task: `Ariadne BNS Indexer (${process.env.USERNAME || ""})`, stateDir: path.join(localAppData, "Ariadne"), loc: userLoc, }); } for (const s of scopes) { if (!s.loc) continue; const exe = path.join(s.loc, "tools", "ariadne-helper.exe"); if (fs.existsSync(exe)) { s.helper = { exe }; s.imagePrefix = path.join(s.loc, "runtime") + path.sep; } } return scopes; } const ariadneCopies = (scopes) => scopes.flatMap((s) => [ path.join(s.stateDir, "index", "bns-name-snapshot.json"), path.join(s.stateDir, "index", "bns-name-snapshot.prev.json"), ]); // "Launch at start" off: the indexer runs on demand, i.e. now. The one thing // Theseus starts on Ariadne's side; its task lets local users run it. function runIndexerOnDemand(scopes) { for (const s of scopes) { if (!s.task || !s.loc) continue; let pol = null; try { pol = JSON.parse(fs.readFileSync(path.join(s.stateDir, "policy.json"), "utf8").replace(/^\uFEFF/, "")); } catch {} if (pol && pol.startAtBoot === false) { execFile("schtasks.exe", ["/Run", "/TN", s.task], { windowsHide: true }, (err) => { log(err ? `Ariadne's indexer (${s.scope}) could not be started on demand` : `Ariadne's indexer (${s.scope}) started on demand (Launch at start is off)`); }); } } } // ---- state ------------------------------------------------------------------ let networkId = "chipnet"; const netSuffix = () => (networkId === "chipnet" ? "" : `.${networkId}`); let scopes = []; let chain = null; // the source chain (Ariadne + local copies) let own = null; // Theseus's own indexer, while it runs let ownReason = null; let goReceived = false; let launchedAt = Date.now(); let healthySince = 0; let lastPublished = null; // which index main has: "chain" | "own" function active() { return own && own.index ? "own" : "chain"; } function publish(reason) { const which = active(); const idx = which === "own" ? own.index : chain && chain.index; if (!idx || !idx.size) return; lastPublished = which; const builtAt = which === "own" ? own.builtAt : chain.builtAt; port.postMessage({ type: "index", builtAt, names: [...idx], network: networkId, reason, source: which === "own" ? "theseus-indexer" : chain.source }); sendStatus(); } function sendStatus() { const cs = chain ? chain.status() : null; port.postMessage({ type: "status", source: active() === "own" ? "theseus-indexer" : cs && cs.source, ariadne: { installed: scopes.filter((s) => s.loc).map((s) => s.scope), pipe: cs ? cs.pipe : null, paused: cs ? cs.paused : false, builtAt: cs ? cs.builtAt : 0, healthy: ariadneHealthy(), }, ownIndexer: own ? { running: true, reason: ownReason, ...own.status() } : { running: false }, }); } // ---- health ----------------------------------------------------------------- function ariadneHealthy() { if (!chain || !chain.pipeConnected) return false; // Economy: Ariadne pauses its network work and stops confirming. That is a // setting, not a failure; taking over would undo it. if (chain.paused) return true; return Date.now() - chain.builtAt < STALE_MS; } function checkHealth() { const now = Date.now(); if (ariadneHealthy()) { healthySince = healthySince || now; if (own && now - healthySince >= HANDBACK_MS) stopOwn(); } else { healthySince = 0; if (!own && now - launchedAt >= LAUNCH_GRACE_MS) { const cs = chain && chain.status(); const why = !scopes.some((s) => s.helper) ? "no Ariadne indexer installed" : cs && cs.pipe && cs.pipe.error ? `Ariadne pipe: ${cs.pipe.error}` : chain && chain.pipeConnected ? "Ariadne's index is not being confirmed" : "Ariadne's indexer is not answering"; startOwn(why); } } } // ---- Theseus's own indexer (standby) ------------------------------------------ function startOwn(reason) { if (own) return own; ownReason = reason; log(`own indexer starting: ${reason}`); own = core.createIndexer({ resolver: R, WebSocket: currentWS(), snapshotDir: cfg.userData, snapshotBase: `bns-name-snapshot${netSuffix()}`, // Ariadne's last snapshot first in line: no download, no cold sync. warmPaths: [...ariadneCopies(scopes), cfg.bundledSnapshot], electrumCacheFile: path.join(cfg.userData, `electrum-servers${netSuffix()}.json`), log: (...a) => log(...a), }); own.on("change", () => { if (own) publish("own"); }); own.on("fresh", ({ builtAt }) => { if (own && active() === "own") port.postMessage({ type: "fresh", builtAt }); }); own.warm(); if (own.index && lastPublished !== "own") publish("own-warm"); // its "change" normally did if (goReceived) own.start(); sendStatus(); return own; } function stopOwn() { if (!own) return; log("Ariadne healthy again: own indexer stopped"); own.stop(); own.removeAllListeners(); own = null; ownReason = null; publish("handback"); } // ---- boot ------------------------------------------------------------------- async function init(m) { cfg = m; await setTor(m.torPort); [R, core, chainLib, pipeLib] = await Promise.all([importEsm(m.resolverPath), importEsm(m.coreLib), importEsm(m.chainLib), importEsm(m.pipeLib)]); networkId = (R.NETWORK && R.NETWORK.id) || "chipnet"; scopes = await findAriadne(); launchedAt = Date.now(); const checked = scopes.filter((s) => s.helper && s.pipe); chain = chainLib.createSourceChain({ resolver: R, // Machine scope first: if both are installed, readers use the machine one. connectPipe: checked.length ? async () => { let last = null; for (const s of checked) { try { return await pipeLib.connectIndexerPipe({ pipeName: s.pipe, helper: s.helper, expect: { scope: s.scope, imagePrefix: s.imagePrefix } }); } catch (e) { last = e; } } throw last || new Error("no pipe"); } : null, localPaths: [ ...ariadneCopies(scopes), path.join(cfg.userData, `bns-name-snapshot${netSuffix()}.json`), path.join(cfg.userData, `bns-name-snapshot${netSuffix()}.prev.json`), cfg.bundledSnapshot, ], ownCopyPath: path.join(cfg.userData, `bns-name-index${netSuffix()}.json`), // The public sources are not the chain's job here: the own indexer // fetches the published snapshot when it runs, and main.js has the // quick single-name lookup (through Tor or the session proxy). publicSnapshotUrls: [], publicIndexerUrl: null, log: (...a) => log(...a), }); chain.on("change", () => { if (active() === "chain") publish("chain"); }); chain.on("fresh", ({ builtAt }) => { if (active() === "chain") port.postMessage({ type: "fresh", builtAt }); }); chain.on("pipe", ({ connected, error }) => { log(connected ? "Ariadne's indexer connected" : `Ariadne's indexer unavailable (${error})`); checkHealth(); sendStatus(); }); chain.on("paused", (p) => { log(p ? "Ariadne is in Economy: sync paused" : "Ariadne resumed syncing"); sendStatus(); }); const t0 = Date.now(); chain.start(); // local copies now (synchronous), the pipe in the background if (chain.index.size) log(`local copy: ${chain.index.size} names from ${chain.source} in ${Date.now() - t0} ms`); if (!lastPublished) publish("local"); // the chain's "change" has normally sent it already runIndexerOnDemand(scopes); // Nothing to wait for (portable Theseus, no Ariadne, or one older than the // pipe): the own indexer at once, as in 0.3.70. if (!checked.length) startOwn("no Ariadne indexer installed"); setTimeout(checkHealth, LAUNCH_GRACE_MS); setInterval(checkHealth, CHECK_EVERY_MS); } let initP = null; // "ready": resolves once there is any index. The pipe gets its grace period; // then the own indexer (published snapshot, electrum as the last resort). async function whenReady() { await initP; if (active() === "chain" && chain.index.size) return true; if (own && own.index) return true; const until = launchedAt + LAUNCH_GRACE_MS; while (Date.now() < until && !chain.index.size) await new Promise((r) => setTimeout(r, 200)); if (chain.index.size && active() === "chain") return true; return startOwn(own ? ownReason : "no index yet").ready(); } async function poll() { await initP; if (own) return own.poll({ force: true }); if (chain.pipeConnected) return chain.requestPoll(); return false; // no live source yet: the health check starts the own indexer } 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": initP = init(m); await initP; break; case "go": goReceived = true; await initP; if (own) own.start(); break; case "tor": await setTor(m.port); if (own) own.setWebSocket(currentWS()); break; case "stop-polling": if (own) own.stop(); if (chain) chain.stop(); break; case "ready": reply(await whenReady()); break; // On-demand poll (a tab's name is not in the index): Ariadne's indexer // when it is live, else the own one; allowed before "go". case "poll": reply(await poll()); break; case "rebuild": await initP; reply(await startOwn(ownReason || "rebuild requested").rebuild()); break; default: reply(false, "unknown request"); } } catch (err) { reply(false, err?.message || String(err)); } });