// Electrum (Fulcrum) JSON-RPC over WebSocket for the wallet. One live // connection at a time, chosen by walking the server list in order; the // caller gets a stable `call()` that reconnects transparently on the next // request after a drop. Notifications (headers / scripthash subscriptions) // fan out to `onNotify`. module.exports = function makeElectrum({ WebSocket, log = () => {} }) { const CALL_TIMEOUT_MS = 20000; const CONNECT_TIMEOUT_MS = 10000; const PING_EVERY_MS = 60000; const RECONNECT_MIN_MS = 2000; const RECONNECT_MAX_MS = 60000; class Connection { constructor(url) { this.url = url; this.id = 0; this.pending = new Map(); this.buf = ""; this.closed = false; this.onNotify = null; this.onClose = null; } connect() { return new Promise((resolve, reject) => { const ws = new WebSocket(this.url); this.ws = ws; const fail = (e) => { if (!this.closed) { this.closed = true; reject(e instanceof Error ? e : new Error("electrum ws error: " + this.url)); } }; // A server that accepts the TCP connection and then says nothing // used to hold the whole server walk hostage; give up and move on. const connectTimer = setTimeout(() => { fail(new Error("electrum connect timeout: " + this.url)); try { ws.close(); } catch {} }, CONNECT_TIMEOUT_MS); ws.on("open", () => clearTimeout(connectTimer)); ws.on("close", () => clearTimeout(connectTimer)); ws.on("open", async () => { // A RANGE, not a flat "1.4". Fulcrum only attaches `token_data` to // listunspent results once protocol >= 1.5 is negotiated, and with // a flat 1.4 it silently omits it — which is why imported BCH // wallets showed no CashTokens at all. A [min, max] pair lets a // modern server pick 1.5.3 while an older one still settles on 1.4, // so nothing that worked before stops working. try { const v = await this.call("server.version", ["theseus-bchwallet", ["1.4", "1.5.3"]]); this.serverVersion = Array.isArray(v) ? v[0] : null; this.protocolVersion = Array.isArray(v) ? v[1] : null; resolve(this); } catch (e) { fail(e); this.close(); } }); ws.on("error", fail); ws.on("message", (d) => this._onData(String(d))); ws.on("close", () => { this.closed = true; for (const p of this.pending.values()) p.reject(new Error("electrum connection closed")); this.pending.clear(); if (this.onClose) this.onClose(); }); }); } _onData(chunk) { this.buf += chunk; let nl; while ((nl = this.buf.indexOf("\n")) >= 0) { const line = this.buf.slice(0, nl).trim(); this.buf = this.buf.slice(nl + 1); if (line) this._handleLine(line); } const rest = this.buf.trim(); if (rest) { try { JSON.parse(rest); this._handleLine(rest); this.buf = ""; } catch {} } } _handleLine(line) { let msg; try { msg = JSON.parse(line); } catch { return; } if (msg.id != null && this.pending.has(msg.id)) { const p = this.pending.get(msg.id); this.pending.delete(msg.id); clearTimeout(p.timer); if (msg.error) p.reject(new Error(typeof msg.error === "object" ? (msg.error.message || JSON.stringify(msg.error)) : String(msg.error))); else p.resolve(msg.result); } else if (msg.method && this.onNotify) { this.onNotify(msg.method, msg.params || []); } } call(method, params = []) { if (this.closed) return Promise.reject(new Error("electrum connection closed")); const id = ++this.id; return new Promise((resolve, reject) => { const timer = setTimeout(() => { if (this.pending.has(id)) { this.pending.delete(id); reject(new Error(`electrum timeout: ${method}`)); } }, CALL_TIMEOUT_MS); this.pending.set(id, { resolve, reject, timer }); try { this.ws.send(JSON.stringify({ id, method, params }) + "\n"); } catch (e) { clearTimeout(timer); this.pending.delete(id); reject(e); } }); } close() { this.closed = true; try { this.ws.close(); } catch {} } } class Client { constructor(servers) { this.servers = servers.slice(); this.conn = null; this.connecting = null; this.subscriptions = new Map(); // method+key -> params (replayed on reconnect) this.onNotify = null; this.onServer = null; // (url|null) connection state for the UI this.disposed = false; this._reconnectTimer = null; this._reconnectDelay = RECONNECT_MIN_MS; this._pingTimer = null; } setServers(servers) { this.servers = servers.slice(); this.disconnect(); this._scheduleReconnect(250); // keep the live feed going on the new list } // A wallet with live subscriptions used to go quiet for good when its // server dropped: the connection was only re-opened by the next call(), // and nothing calls while the user is just looking at a balance. Come // back on our own, with backoff, for as long as anything is subscribed. _scheduleReconnect(delay) { if (this.disposed || !this.subscriptions.size || this._reconnectTimer) return; const wait = delay != null ? delay : this._reconnectDelay; this._reconnectTimer = setTimeout(() => { this._reconnectTimer = null; if (this.disposed || (this.conn && !this.conn.closed)) return; this._ensure().then( () => { this._reconnectDelay = RECONNECT_MIN_MS; }, () => { this._reconnectDelay = Math.min(RECONNECT_MAX_MS, this._reconnectDelay * 2); this._scheduleReconnect(); }, ); }, wait); if (this._reconnectTimer.unref) this._reconnectTimer.unref(); } // A half-dead socket (laptop slept, NAT dropped the flow) never fires // "close". Ping it; a ping that fails or times out closes it, which // lands in the reconnect path above. _armPing(c) { clearInterval(this._pingTimer); this._pingTimer = setInterval(() => { if (this.disposed || this.conn !== c || c.closed) { clearInterval(this._pingTimer); return; } c.call("server.ping", []).catch(() => { try { c.close(); } catch {} }); }, PING_EVERY_MS); if (this._pingTimer.unref) this._pingTimer.unref(); } // Stop for good: no reconnects, no pings, and any later call() fails // instead of quietly resurrecting the connection of a removed wallet. dispose() { this.disposed = true; clearTimeout(this._reconnectTimer); this._reconnectTimer = null; clearInterval(this._pingTimer); this._pingTimer = null; this.subscriptions.clear(); this.disconnect(); } get url() { return this.conn && !this.conn.closed ? this.conn.url : null; } // The protocol the live connection settled on. Callers use it to decide // whether listunspent will carry `token_data` (>= 1.5) or whether they // have to classify tokens the expensive way, by fetching each UTXO's // parent transaction. get protocolVersion() { return this.conn && !this.conn.closed ? (this.conn.protocolVersion || null) : null; } // True when the server will report CashTokens on listunspent itself. get hasTokenData() { const v = String(this.protocolVersion || ""); const m = /^(\d+)\.(\d+)/.exec(v); if (!m) return false; const major = Number(m[1]), minor = Number(m[2]); return major > 1 || (major === 1 && minor >= 5); } async _ensure() { if (this.disposed) throw new Error("electrum client disposed"); if (this.conn && !this.conn.closed) return this.conn; if (this.connecting) return this.connecting; this.connecting = (async () => { let lastErr; for (const url of this.servers) { try { const c = await new Connection(url).connect(); if (this.disposed) { c.close(); throw new Error("electrum client disposed"); } c.onNotify = (m, p) => { if (this.onNotify) this.onNotify(m, p); }; c.onClose = () => { if (this.conn !== c) return; this.conn = null; if (this.onServer) this.onServer(null); this._scheduleReconnect(); }; this.conn = c; log("connected", url); if (this.onServer) this.onServer(url); this._armPing(c); // Re-arm subscriptions so a reconnect keeps the live feed, and // hand each answer on as if it were a notification: whatever // changed while we were away (a new tip, a payment) is in that // answer, and the wallet only refreshes when it is told. for (const [method, params] of this.subscriptions.values()) { c.call(method, params).then((result) => { if (!this.onNotify || this.conn !== c) return; this.onNotify(method, method === "blockchain.headers.subscribe" ? [result] : [...params, result]); }, () => {}); } return c; } catch (e) { lastErr = e; log("failed", url, e?.message); if (this.disposed) break; } } throw lastErr || new Error("no electrum server reachable"); })(); try { return await this.connecting; } finally { this.connecting = null; } } async call(method, params = []) { const c = await this._ensure(); return c.call(method, params); } // Remember a subscription so it survives reconnects. async subscribe(method, params = []) { this.subscriptions.set(method + ":" + JSON.stringify(params), [method, params]); return this.call(method, params); } clearSubscriptions() { this.subscriptions.clear(); } disconnect() { clearInterval(this._pingTimer); this._pingTimer = null; if (this.conn) { const c = this.conn; this.conn = null; c.close(); } if (this.onServer) this.onServer(null); } } return { Client }; };