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>
207 lines
9 KiB
Python
Executable file
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)
|