#!/usr/bin/env python3 """Minimal wRPC JSON client for igneumd, standard library only (any Debian python3, no websocket package). The node's wRPC JSON endpoint (--rpclisten-json=127.0.0.1:28610) is a plain WebSocket. Requests are {"id": n, "method": "getBlockDagInfo", "params": {...}}; the answer carries the same id and either "params" (the result) or "error"; notifications arrive as {"method": "blockAddedNotification", "params": {...}}. Field names are camelCase (serde rename_all on the RPC model). This mirrors the Rpc class in tools/observer/observer.mjs. wrpc.py call [json-params] one call, prints the result as JSON wrpc.py watch-blocks subscribe and append, forever: /blocks.tsv recv_ms hash header_timestamp_ms blue_score daa_score parents /chain.tsv recv_ms added removed first_removed_hashes(up to 3, comma separated) /samples.tsv ts_ms sink blue_score daa_score tips peers difficulty headers blocks (every 5 s) wrpc.py sample one samples.tsv line to stdout Environment: IGNEUM_RPC (default ws://127.0.0.1:28610). """ import base64, json, os, socket, struct, sys, time, threading from urllib.parse import urlparse RPC = os.environ.get("IGNEUM_RPC", "ws://127.0.0.1:28610") class WebSocket: """RFC 6455 client: text frames, masking, ping/pong, continuation frames.""" def __init__(self, url, timeout=10.0): u = urlparse(url) self.sock = socket.create_connection((u.hostname, u.port or 80), timeout=timeout) key = base64.b64encode(os.urandom(16)).decode() path = u.path or "/" req = (f"GET {path} HTTP/1.1\r\nHost: {u.hostname}:{u.port}\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n" f"Sec-WebSocket-Key: {key}\r\nSec-WebSocket-Version: 13\r\n\r\n") self.sock.sendall(req.encode()) head = b"" while b"\r\n\r\n" not in head: chunk = self.sock.recv(4096) if not chunk: raise ConnectionError("handshake: connection closed") head += chunk status, _, rest = head.partition(b"\r\n\r\n") if b" 101 " not in status.split(b"\r\n")[0]: raise ConnectionError("handshake refused: " + status.split(b"\r\n")[0].decode(errors="replace")) self.buf = rest self.lock = threading.Lock() def _read(self, n): while len(self.buf) < n: chunk = self.sock.recv(65536) if not chunk: raise ConnectionError("connection closed") self.buf += chunk out, self.buf = self.buf[:n], self.buf[n:] return out def send_text(self, text): data = text.encode() mask = os.urandom(4) n = len(data) if n < 126: hdr = struct.pack("!BB", 0x81, 0x80 | n) elif n < 65536: hdr = struct.pack("!BBH", 0x81, 0x80 | 126, n) else: hdr = struct.pack("!BBQ", 0x81, 0x80 | 127, n) masked = bytes(b ^ mask[i % 4] for i, b in enumerate(data)) with self.lock: self.sock.sendall(hdr + mask + masked) def _send_control(self, opcode, payload=b""): mask = os.urandom(4) with self.lock: self.sock.sendall(struct.pack("!BB", 0x80 | opcode, 0x80 | len(payload)) + mask + bytes(b ^ mask[i % 4] for i, b in enumerate(payload))) def recv_text(self): """Returns the next complete text message (handles ping and continuation).""" message = b"" while True: b0, b1 = self._read(2) fin, opcode = b0 & 0x80, b0 & 0x0F n = b1 & 0x7F if n == 126: n = struct.unpack("!H", self._read(2))[0] elif n == 127: n = struct.unpack("!Q", self._read(8))[0] if b1 & 0x80: mask = self._read(4) payload = bytes(b ^ mask[i % 4] for i, b in enumerate(self._read(n))) else: payload = self._read(n) if opcode == 0x9: self._send_control(0xA, payload); continue if opcode == 0xA: continue if opcode == 0x8: raise ConnectionError("server closed the websocket") message += payload if fin: return message.decode() def close(self): try: self._send_control(0x8) except Exception: pass self.sock.close() class Rpc: def __init__(self, url=RPC): self.ws = WebSocket(url) self.next_id = 0 self.notifications = [] def call(self, method, params=None, timeout=10.0): self.next_id += 1 rid = self.next_id self.ws.send_text(json.dumps({"id": rid, "method": method, "params": params or {}})) deadline = time.time() + timeout while time.time() < deadline: m = json.loads(self.ws.recv_text()) if m.get("id") == rid: if "error" in m and m["error"]: raise RuntimeError(f"{method}: {m['error']}") return m.get("params") if m.get("method"): self.notifications.append(m) raise TimeoutError(method) def next_notification(self): if self.notifications: return self.notifications.pop(0) while True: m = json.loads(self.ws.recv_text()) if m.get("method"): return m def sample_line(rpc): dag = rpc.call("getBlockDagInfo") peers = rpc.call("getConnectedPeerInfo").get("peerInfo", []) try: blue = rpc.call("getSinkBlueScore").get("blueScore") except Exception: blue = "" return "\t".join(str(x) for x in [int(time.time() * 1000), dag.get("sink", ""), blue, dag.get("virtualDaaScore", ""), len(dag.get("tipHashes", []) or []), len(peers), dag.get("difficulty", ""), dag.get("headerCount", ""), dag.get("blockCount", "")]) def watch_blocks(directory): os.makedirs(directory, exist_ok=True) while True: try: rpc = Rpc() rpc.call("subscribe", {"BlockAdded": {}}) rpc.call("subscribe", {"VirtualChainChanged": {"include_accepted_transaction_ids": False}}) sys.stdout.write(f"{time.strftime('%H:%M:%S')} subscribed on {RPC}\n"); sys.stdout.flush() last_sample = 0.0 with open(os.path.join(directory, "blocks.tsv"), "a") as fb, open(os.path.join(directory, "chain.tsv"), "a") as fc, \ open(os.path.join(directory, "samples.tsv"), "a") as fs: while True: now = time.time() if now - last_sample >= 5.0: try: fs.write(sample_line(rpc) + "\n"); fs.flush() except Exception as e: sys.stdout.write(f"sample failed: {e}\n"); sys.stdout.flush() last_sample = now rpc.ws.sock.settimeout(5.0) try: m = rpc.next_notification() except socket.timeout: continue recv_ms = int(time.time() * 1000) p = m.get("params") or {} inner = p.get("BlockAdded") or p.get("VirtualChainChanged") or p if m["method"] == "blockAddedNotification" and inner.get("block"): h = inner["block"].get("header", {}) levels = h.get("parentsByLevel") or [] # array of arrays of hashes; level 0 = direct parents direct = levels[0] if levels and isinstance(levels[0], list) else [] fb.write("\t".join(str(x) for x in [recv_ms, h.get("hash", ""), h.get("timestamp", ""), h.get("blueScore", ""), h.get("daaScore", ""), len(direct)]) + "\n") fb.flush() elif m["method"] == "virtualChainChangedNotification": added = inner.get("addedChainBlockHashes", []) or [] removed = inner.get("removedChainBlockHashes", []) or [] fc.write("\t".join([str(recv_ms), str(len(added)), str(len(removed)), ",".join(removed[:3])]) + "\n") fc.flush() except Exception as e: sys.stdout.write(f"{time.strftime('%H:%M:%S')} rpc lost ({e}); retrying in 5 s\n"); sys.stdout.flush() time.sleep(5) def main(argv): if len(argv) >= 2 and argv[1] == "call": params = json.loads(argv[3]) if len(argv) > 3 else {} print(json.dumps(Rpc().call(argv[2], params), indent=None)) elif len(argv) >= 3 and argv[1] == "watch-blocks": watch_blocks(argv[2]) elif len(argv) >= 2 and argv[1] == "sample": print(sample_line(Rpc())) else: sys.stderr.write(__doc__); sys.exit(2) if __name__ == "__main__": main(sys.argv)