igneum/infra/cloud-devnet/experiments/analyze.py
igneum-labs 74e7c2bdf5 seed-nodes: igneum-seed-1 live and synced; relay path from the Mac; prices in the account's currency (USD)
igneum-seed-1 (Hetzner cx23, fsn1, 188.245.5.161:26611): built on the VM in 1,530 s, synced to the live devnet
(12,204 blocks, same sink as the live node) through a non-mining relay igneumd on the Mac (the live node's addPeer
RPC is refused in safe mode); the live node and the Windows PC learned the seed's address by peer exchange and dialled
it. seeds.txt written. Hetzner prices corrected to USD (pricing API currency) in the plans, READMEs and scripts;
current-generation types per location (cx23 EU, cpx22 sin, cpx21 US) in the cloud-devnet config.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-10-03 22:27:52 +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}; cx23 in the EU, cpx22 in Singapore, cpx21 in the US 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)