igneum/infra/cloud-devnet/node/wrpc.py
igneum-labs d6c69bcb0c infra: cloud devnet (20 nodes), rented GPU bench and seed node scripts; plans for both
infra/cloud-devnet: hcloud (doctl variant) create, builder-VM provision from a git-archive source tarball,
systemd units for igneumd --devnet-suffix with a sparse --addpeer mesh and a CPU trickle miner per node,
stdlib wRPC client, experiments (latency, partition, hop, collect, observer hookup), README with the command
sequence and the Hetzner API prices of 3 Oct 2026.
infra/gpu-bench: RunPod image recipes (CUDA 12.8, ROCm), bundle, run.sh (vectors gate, 10-min raw, sweep,
inline shortcut ratio, nvcc/NVRTC/OpenCL recompile timings, results row, intake upload), bench-log template.
infra/seed-nodes: create-seed (persistent IPv4, firewall), provision on the VM, health check, addPeer from the
Mac over grpcurl, seeds.txt; igneum-seed-1 created at 188.245.5.161 (Hetzner cx23, fsn1).
docs/plans/cloud-devnet.md and docs/plans/seed-nodes.md.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-10-03 22:01:50 +00:00

207 lines
9 KiB
Python
Executable file

#!/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 <method> [json-params] one call, prints the result as JSON
wrpc.py watch-blocks <dir> subscribe and append, forever:
<dir>/blocks.tsv recv_ms hash header_timestamp_ms blue_score daa_score parents
<dir>/chain.tsv recv_ms added removed first_removed_hashes(up to 3, comma separated)
<dir>/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)