igneum/infra/cloud-devnet/experiments/analyze.py
igneum-labs d8dc8c6ba3 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

308 lines
16 KiB
Python

#!/usr/bin/env python3
"""Offline analysis of a results/<date>/ directory (runs on the Mac, standard library only).
analyze.py rtt <dir> <nodes.tsv> region-pair RTT table from <dir>/rtt.tsv
analyze.py propagation <dir> <nodes.tsv> <t0_ms> <t1_ms> per-node and per-region block propagation delay from
<dir>/nodes/<name>/blocks.tsv (window t0..t1)
analyze.py locks <partition-dir> conflicting locks between the checkpoint dumps at heal
analyze.py hop <dir> block rate and difficulty per hop.log phase (node 1 samples)
analyze.py summary <dir> <nodes.tsv> summary.md of everything present under <dir>
analyze.py bench-entry <dir> <nodes.tsv> a docs/bench-log.md entry (template with the numbers filled)
"""
import glob, json, os, statistics, sys, time
from collections import defaultdict
def read_nodes(path):
nodes = []
for line in open(path):
f = line.rstrip("\n").split("\t")
if len(f) >= 4:
nodes.append({"name": f[0], "index": int(f[1]), "region": f[2], "ip": f[3]})
return nodes
def pct(xs, p):
if not xs:
return float("nan")
xs = sorted(xs)
k = min(len(xs) - 1, max(0, int(round((len(xs) - 1) * p))))
return xs[k]
def fmt(x, nd=1):
return "n/a" if x != x else f"{x:.{nd}f}"
# ---------------------------------------------------------------------------------------------------------------
def rtt(d, nodes_file):
nodes = read_nodes(nodes_file)
by_ip = {n["ip"]: n for n in nodes}
by_name = {n["name"]: n for n in nodes}
pairs = defaultdict(list)
for line in open(os.path.join(d, "rtt.tsv")):
f = line.rstrip("\n").split("\t")
if len(f) < 5 or f[3] == "NA":
continue
src, dst = by_name.get(f[0]), by_ip.get(f[1])
if not src or not dst:
continue
key = tuple(sorted([src["region"], dst["region"]]))
pairs[key].append(float(f[3]))
regions = sorted({n["region"] for n in nodes})
out = ["| region | " + " | ".join(regions) + " |", "|---|" + "---|" * len(regions)]
for a in regions:
row = [a]
for b in regions:
xs = pairs.get(tuple(sorted([a, b])), [])
row.append(fmt(statistics.median(xs)) if xs else "n/a")
out.append("| " + " | ".join(row) + " |")
out.append("")
out.append("Median of the per-pair average RTT in ms (ping, 5 packets). Same-region cells are the intra-location RTT.")
return "\n".join(out)
def propagation(d, nodes_file, t0, t1):
nodes = read_nodes(nodes_file)
t0, t1 = int(t0), int(t1)
seen = defaultdict(dict) # hash -> node -> recv_ms
header_ts = {}
for n in nodes:
p = os.path.join(d, "nodes", n["name"], "blocks.tsv")
if not os.path.exists(p):
continue
for line in open(p):
f = line.rstrip("\n").split("\t")
if len(f) < 3:
continue
try:
recv = int(f[0])
except ValueError:
continue
if recv < t0 or recv > t1:
continue
seen[f[1]][n["name"]] = recv
try:
header_ts[f[1]] = int(f[2])
except ValueError:
pass
total = len(seen)
full = {h: m for h, m in seen.items() if len(m) >= max(2, len(nodes) * 0.8)}
delays = defaultdict(list)
region_delays = defaultdict(list)
ts_vs_first = []
reg = {n["name"]: n["region"] for n in nodes}
for h, m in full.items():
first = min(m.values())
for name, r in m.items():
delays[name].append(r - first)
region_delays[reg[name]].append(r - first)
if h in header_ts:
ts_vs_first.append(first - header_ts[h])
out = [f"Blocks in the window: {total}; blocks seen by at least 80% of nodes: {len(full)} (the join base).", ""]
out.append("| node | region | blocks | p50 ms | p90 ms | max ms |")
out.append("|---|---|---|---|---|---|")
for n in nodes:
xs = delays.get(n["name"], [])
out.append(f"| {n['name']} | {n['region']} | {len(xs)} | {fmt(pct(xs, 0.5), 0)} | {fmt(pct(xs, 0.9), 0)} | {fmt(max(xs) if xs else float('nan'), 0)} |")
out.append("")
out.append("| region | samples | p50 ms | p90 ms | p99 ms |")
out.append("|---|---|---|---|---|")
for r in sorted(region_delays):
xs = region_delays[r]
out.append(f"| {r} | {len(xs)} | {fmt(pct(xs, 0.5), 0)} | {fmt(pct(xs, 0.9), 0)} | {fmt(pct(xs, 0.99), 0)} |")
allx = [x for xs in delays.values() for x in xs]
out.append("")
out.append(f"All nodes: p50 {fmt(pct(allx, 0.5), 0)} ms, p90 {fmt(pct(allx, 0.9), 0)} ms, p99 {fmt(pct(allx, 0.99), 0)} ms, max {fmt(max(allx) if allx else float('nan'), 0)} ms "
f"(arrival at a node minus the first arrival anywhere; 0 for the node that produced or first received the block).")
if ts_vs_first:
out.append(f"First arrival minus the header timestamp: median {fmt(statistics.median(ts_vs_first), 0)} ms (miner clock and template age; negative = miner clock ahead).")
out.append("Clock caveat: chrony on every VM; per-node offsets in nodes/<name>/chrony.txt.")
return "\n".join(out)
def locks(pdir):
dumps = {}
for p in glob.glob(os.path.join(pdir, "checkpoints-*-at-heal.json")):
try:
dumps[os.path.basename(p)[12:-13]] = json.load(open(p))
except Exception:
continue
locked = defaultdict(dict) # index -> hash -> [nodes]
for node, j in dumps.items():
for cp in (j.get("checkpoints") or []):
if str(cp.get("state", "")).lower() == "locked":
locked[int(cp["index"])].setdefault(cp.get("hash"), []).append(node)
conflicts = {i: h for i, h in locked.items() if len(h) > 1}
if not dumps:
return "no checkpoint dumps (node without the finality layer, or the RPC was refused)"
lines = [f"nodes dumped: {', '.join(sorted(dumps))}; locked indices seen: {len(locked)}; conflicting: {len(conflicts)}"]
for i in sorted(conflicts):
lines.append(f" index {i}: " + "; ".join(f"{h[:12]} by {','.join(ns)}" for h, ns in conflicts[i].items()))
if not conflicts:
lines.append(" none: no index locked with two different hashes on the two sides")
return "\n".join(lines)
def read_samples(path):
rows = []
for line in open(path):
f = line.rstrip("\n").split("\t")
if len(f) >= 9:
try:
rows.append({"t": int(f[0]), "sink": f[1], "blue": int(f[2] or 0), "daa": int(f[3] or 0), "tips": int(f[4]),
"peers": int(f[5]), "difficulty": float(f[6] or 0), "headers": int(f[7] or 0), "blocks": int(f[8] or 0)})
except ValueError:
continue
return rows
def hop(d):
logp = os.path.join(d, "hop.log")
if not os.path.exists(logp):
return "no hop.log in this directory"
phases = []
for line in open(logp):
f = line.rstrip("\n").split("\t")
if len(f) >= 3:
phases.append({"t": int(f[0]), "sel": f[1], "threads": f[2]})
cands = sorted(glob.glob(os.path.join(d, "nodes", "*", "samples.tsv")))
if not cands:
return "no samples (run collect.sh first)"
rows = read_samples(cands[0])
out = [f"Samples from {cands[0].split('/')[-2]}; phases from hop.log.", "",
"| phase | selector | threads | length s | blocks/min mean | blocks/min last 3 min | difficulty start | difficulty end | settled s (3 min within 10% of 60/min) |",
"|---|---|---|---|---|---|---|---|---|"]
for i, ph in enumerate(phases):
if ph["sel"] == "end":
break
t_end = phases[i + 1]["t"] if i + 1 < len(phases) else rows[-1]["t"]
win = [r for r in rows if ph["t"] <= r["t"] <= t_end]
if len(win) < 2:
out.append(f"| {i + 1} | {ph['sel']} | {ph['threads']} | {(t_end - ph['t']) // 1000} | n/a | n/a | n/a | n/a | n/a |")
continue
mins = defaultdict(list)
for r in win:
mins[(r["t"] - ph["t"]) // 60000].append(r["blocks"])
per_min = [(m, max(v) - min(v)) for m, v in sorted(mins.items()) if len(v) >= 2]
rates = [x for _, x in per_min]
settled = "never"
for k in range(len(rates) - 2):
if all(abs(rates[k + j] - 60) <= 6 for j in range(3)):
settled = str(per_min[k][0] * 60); break
mean_rate = (win[-1]["blocks"] - win[0]["blocks"]) / max(1, (win[-1]["t"] - win[0]["t"]) / 60000)
last3 = statistics.mean(rates[-3:]) if len(rates) >= 3 else float("nan")
out.append(f"| {i + 1} | {ph['sel']} | {ph['threads']} | {(t_end - ph['t']) // 1000} | {mean_rate:.1f} | {fmt(last3)} | {win[0]['difficulty']:.3g} | {win[-1]['difficulty']:.3g} | {settled} |")
out.append("")
out.append("blocks/min from the node's blockCount (every block, blue or red). Target 60. 'settled' = first minute of three in a row within 10% of target after the step (spec 02 uses a 100-block criterion on the 121-block mean rate; recompute from samples.tsv for the spec's definition).")
return "\n".join(out)
def summary(d, nodes_file):
nodes = read_nodes(nodes_file)
out = [f"# Cloud devnet results, {os.path.basename(d)}", ""]
ver = os.path.join(d, "version.txt")
if os.path.exists(ver):
out.append("Node build: " + " ".join(open(ver).read().split()))
stamp = os.path.join(d, "src.stamp")
if os.path.exists(stamp):
out.append("Source: " + "; ".join(l.strip() for l in open(stamp) if l.strip()))
region_counts = ", ".join("%s x %d" % (r, sum(1 for n in nodes if n["region"] == r)) for r in sorted({n["region"] for n in nodes}))
out.append(f"Nodes: {len(nodes)} ({region_counts})")
out.append("")
out.append("## Final state per node")
out.append("")
out.append("| node | region | blocks | headers | blue | daa | tips | peers | difficulty | synced | miner blocks found |")
out.append("|---|---|---|---|---|---|---|---|---|---|---|")
for n in nodes:
nd = os.path.join(d, "nodes", n["name"])
s = read_samples(os.path.join(nd, "samples.tsv")) if os.path.exists(os.path.join(nd, "samples.tsv")) else []
last = s[-1] if s else {}
synced = "?"
try:
synced = str(json.load(open(os.path.join(nd, "getInfo.json"))).get("isSynced"))
except Exception:
pass
found = "?"
ml = os.path.join(nd, "miner.log")
if os.path.exists(ml):
found = str(sum(1 for l in open(ml, errors="replace") if "accepted" in l.lower() and "block" in l.lower()))
out.append(f"| {n['name']} | {n['region']} | {last.get('blocks', '?')} | {last.get('headers', '?')} | {last.get('blue', '?')} | {last.get('daa', '?')} | {last.get('tips', '?')} | {last.get('peers', '?')} | {last.get('difficulty', 0):.3g} | {synced} | {found} |")
out.append("")
ev = os.path.join(d, "events.log")
if os.path.exists(ev):
out.append("## Events"); out.append(""); out.append("```"); out.append(open(ev).read().rstrip()); out.append("```"); out.append("")
lat = os.path.join(d, "latency")
if os.path.isdir(lat):
for f in ("rtt-by-region.md", "propagation.md"):
p = os.path.join(lat, f)
if os.path.exists(p):
out.append(f"## Latency: {f}"); out.append(""); out.append(open(p).read().rstrip()); out.append("")
for p in sorted(glob.glob(os.path.join(d, "partition-*", "partition.md"))):
out.append("## " + os.path.basename(os.path.dirname(p))); out.append(""); out.append(open(p).read().rstrip()); out.append("")
if os.path.exists(os.path.join(d, "hop.log")):
out.append("## Hash-rate steps (hop.sh)"); out.append(""); out.append(hop(d)); out.append("")
# reorg distribution over the whole run from chain.tsv (every virtualChainChanged removal)
removals = []
for n in nodes:
p = os.path.join(d, "nodes", n["name"], "chain.tsv")
if os.path.exists(p):
for line in open(p):
f = line.rstrip("\n").split("\t")
if len(f) >= 3 and f[2].isdigit() and int(f[2]) > 0:
removals.append(int(f[2]))
out.append("## Reorg depth distribution (all nodes, whole run, virtualChainChanged removals > 0)")
out.append("")
if removals:
hist = defaultdict(int)
for r in removals:
hist[r] += 1
out.append("| depth | count |"); out.append("|---|---|")
for k in sorted(hist):
out.append(f"| {k} | {hist[k]} |")
out.append("")
out.append(f"{len(removals)} reorgs; p50 {pct(removals, 0.5)}, p99 {pct(removals, 0.99)}, max {max(removals)} (spec 03 C1: d is set from this distribution).")
else:
out.append("no reorgs recorded (or no chain.tsv)")
return "\n".join(out)
def bench_entry(d, nodes_file):
nodes = read_nodes(nodes_file)
day = os.path.basename(d)
try:
date_txt = time.strftime("%-d %B %Y", time.strptime(day, "%Y-%m-%d"))
except ValueError:
date_txt = day
ver = " ".join(open(os.path.join(d, "version.txt")).read().split()) if os.path.exists(os.path.join(d, "version.txt")) else "<igneumd version>"
regions = ", ".join("%s x %d" % (r, sum(1 for n in nodes if n["region"] == r)) for r in sorted({n["region"] for n in nodes}))
lines = [f"## {date_txt}, cloud devnet: {len(nodes)} igneumd nodes across regions, CPU trickle miners, latency, partition and hash-rate steps (consensus-engineer)", "",
f"Machines: {len(nodes)} Hetzner Cloud VMs ({regions}; type cpx21 unless noted), Debian 12, chrony. Node: {ver}, built on a builder VM from the source tarball in results/{day}/src.stamp. Network: igneum-devnet-<suffix>, genesis bits <bits>, {len(nodes)} one-thread CPU miners (igneum-miner --engine igneum-pow), one BLS vote key per node, sparse --addpeer mesh (about 4 peers each).", "",
"Inter-region RTT (ms, median of pair averages): <paste results/" + day + "/latency/rtt-by-region.md>", "",
"Block propagation (arrival minus first arrival anywhere): p50 <n> ms, p90 <n> ms, p99 <n> ms, max <n> ms; per region <paste the region table from propagation.md>.", "",
"Partition <region>, <minutes> min: minority reorg depth max <n> (getVirtualChainFromBlock from the minority sink at heal), majority <n>; converged <n> s after heal; conflicting locks: <none | list>.", "",
"Hash-rate steps (hop.sh): <paste the phase table from hop.md>: settle times per step against spec 02 section 2.3 (simulator: x50 settled 62 s, /50 657 s).", "",
"Reorg depth distribution over the run: p50 <n>, p99 <n>, max <n> over <n> reorgs (summary.md).", "",
"Reading: <what this says about gate 3 (checkpoint depth d, the floor rule under a real partition) and the controller under real latency>. Caveats: CPU hash rate only (no GPU), one evening, clocks by chrony.", ""]
return "\n".join(lines)
def main(a):
if len(a) >= 4 and a[1] == "rtt":
print(rtt(a[2], a[3]))
elif len(a) >= 6 and a[1] == "propagation":
print(propagation(a[2], a[3], a[4], a[5]))
elif len(a) >= 3 and a[1] == "locks":
print(locks(a[2]))
elif len(a) >= 3 and a[1] == "hop":
print(hop(a[2]))
elif len(a) >= 4 and a[1] == "summary":
print(summary(a[2], a[3]))
elif len(a) >= 4 and a[1] == "bench-entry":
print(bench_entry(a[2], a[3]))
else:
sys.stderr.write(__doc__); sys.exit(2)
if __name__ == "__main__":
main(sys.argv)