// wRPC JSON client for igneumd (same shape as tools/observer/observer.mjs). Node 22, no dependencies. // Method names are the lowerCamelCase of the node's RpcApiOps (getBlockTemplate, submitBlock, ...). export class Rpc { constructor(url, { timeoutMs = 10_000 } = {}) { this.url = url; this.id = 0; this.pending = new Map(); this.ws = null; this.open = false; this.timeoutMs = timeoutMs; this.onNotification = () => { }; } connect() { return new Promise((resolve) => { const ws = new WebSocket(this.url); this.ws = ws; ws.onopen = () => { this.open = true; resolve(true); }; ws.onmessage = (e) => { let m; try { m = JSON.parse(e.data); } catch { return; } if (m.id !== undefined && m.id !== null && this.pending.has(m.id)) { const p = this.pending.get(m.id); this.pending.delete(m.id); m.error ? p.reject(new Error(typeof m.error === 'string' ? m.error : (m.error.message || JSON.stringify(m.error)))) : p.resolve(m.params); } else if (m.method) this.onNotification(m.method, m.params); }; // A refused connection may fire error without close in Node's WebSocket: resolve false either way, with a timer. ws.onerror = () => { if (!this.open) resolve(false); }; ws.onclose = () => { const wasOpen = this.open; this.open = false; for (const p of this.pending.values()) p.reject(new Error('rpc closed')); this.pending.clear(); if (!wasOpen) resolve(false); }; setTimeout(() => { if (!this.open) { try { ws.close(); } catch { } resolve(false); } }, 3000); }); } close() { try { this.ws && this.ws.close(); } catch { } } call(method, params = {}, timeoutMs = this.timeoutMs) { return new Promise((resolve, reject) => { if (!this.open) return reject(new Error('rpc not connected')); const id = ++this.id; this.pending.set(id, { resolve, reject }); this.ws.send(JSON.stringify({ id, method, params })); setTimeout(() => { if (this.pending.has(id)) { this.pending.delete(id); reject(new Error(`${method} timed out`)); } }, timeoutMs); }); } /// Sends a raw frame (for malformed-input cases) and resolves with the first response or an error. raw(text, timeoutMs = 3000) { return new Promise((resolve) => { if (!this.open) return resolve({ error: 'rpc not connected' }); const id = ++this.id; this.pending.set(id, { resolve: (p) => resolve({ ok: p }), reject: (e) => resolve({ error: String(e.message || e) }) }); try { this.ws.send(text.replace('__ID__', String(id))); } catch (e) { this.pending.delete(id); return resolve({ error: String(e) }); } setTimeout(() => { if (this.pending.has(id)) { this.pending.delete(id); resolve({ timeout: true }); } }, timeoutMs); }); } } /// Connects with retries (the node takes a few seconds to open its listeners). export async function connectRpc(url, { attempts = 60, waitMs = 500 } = {}) { for (let i = 0; i < attempts; i++) { const rpc = new Rpc(url); if (await rpc.connect()) { try { await rpc.call('getInfo'); return rpc; } catch { rpc.close(); } } await new Promise(r => setTimeout(r, waitMs)); } throw new Error(`could not reach ${url}`); } /// Unwraps the node's SubmitBlockResponse into a short string: "accepted" or "rejected:". export function submitReport(res) { const r = res && res.report; if (!r) return `odd:${JSON.stringify(res)}`; if (r === 'success' || r.type === 'success') return 'accepted'; if (r.type === 'reject') return `rejected:${JSON.stringify(r.reason ?? r.reject ?? r)}`; return `odd:${JSON.stringify(r)}`; }