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>
308 lines
16 KiB
Python
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)
|