igneum/infra/cloud-devnet/experiments/analyze.py
igneum-labs 8255839d9e Cloud devnet: two Singapore partitions and the hash-rate steps on the 12-node network, results and bench-log entry
partition.sh adapted to private-network mode: the cut sits on the region's gateway (INPUT, OUTPUT and FORWARD
against the far gateways' public IPs over 26611 and the DNAT ports 27001:27099), since a per-node port-26611 rule
leaves the DNAT links up. It now records the locks per side at cut, during, at heal and after convergence, the
first lock after the heal from the journals, and the heal time as each minority node's first chain removal of 5+
blocks (the sink-count criterion is tip churn on a healthy network). hop.sh and partition.sh hold the Mac awake
with caffeinate; analyze.py gains a 10-s hop series (difficulty, block count, 1- and 2-min rates, threads) and an
overshoot table; collect.sh writes hop-series.tsv and hop.md and gzips the journals.

Results 2026-10-04: partition 1 (window still filling) reorg 431/496 on the minority, 2 on the majority, healed in
10 and 14 s; partition 2 (locks active): minority locked nothing during the cut, majority locked every interval at
66.8% to 84.5% of total, 0 conflicting locks over 107 indices, healed in 11 and 15 s, first lock after heal 13 s.
Hash-rate steps x1.42, x0.70, x0.75, x1.32: difficulty overshoots x1.67, x0.66, x0.53, x1.60, settle 751 s,
never in 900 s, 241 s, 646 s. Failures stated in summary.md: the Mac hibernated during the hop (phase 2 ran 94
min), the first partition could not see locks, the script's heal and first-lock figures were artefacts (fixed),
and another agent's v2 rollout restarted every node during the second heal.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-10-04 14:39:54 +00:00

419 lines
22 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 hop-series <dir> <nodes.tsv> 10-s series over the whole hop run: difficulty (median over
nodes and node 1), block count, trailing 1-min and 2-min
block rate, miner threads; stdout as TSV (hop-series.tsv)
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).")
# the same phases on the 10-s grid with a trailing 2-min rate (close to the spec's 121-block mean at 1 BPS):
# overshoot = the largest excursion of that rate from 60/min after the step, settle = first 10-s point after which
# the trailing 2-min rate stays within 10% of 60 for 3 min
out.append("")
out.append("| phase | selector | threads | rate 2-min min (blocks/min) | rate 2-min max | overshoot vs 60 | difficulty min | difficulty max | difficulty end / start | settled s (2-min rate within 10% for 3 min) |")
out.append("|---|---|---|---|---|---|---|---|---|---|")
series = _series_rows(rows)
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 series if ph["t"] + 120000 <= r["t"] <= t_end] # after the first 2 min the window is inside the phase
full = [r for r in series if ph["t"] <= r["t"] <= t_end]
if not win:
out.append(f"| {i + 1} | {ph['sel']} | {ph['threads']} | n/a | n/a | n/a | n/a | n/a | n/a | n/a |"); continue
rates = [r["rate2"] for r in win]
lo, hi = min(rates), max(rates)
over = max(abs(hi - 60), abs(60 - lo)) / 60 * 100
sign = "+" if abs(hi - 60) >= abs(60 - lo) else "-"
diffs = [r["difficulty"] for r in full if r["difficulty"] > 0]
settled = "never"
for k in range(len(win)):
if win[k]["t"] + 180000 > t_end:
break
if all(abs(r["rate2"] - 60) <= 6 for r in win[k:] if r["t"] <= win[k]["t"] + 180000):
settled = str((win[k]["t"] - ph["t"]) // 1000); break
out.append(f"| {i + 1} | {ph['sel']} | {ph['threads']} | {lo:.1f} | {hi:.1f} | {sign}{over:.0f}% | {min(diffs):.3g} | {max(diffs):.3g} | {diffs[-1] / diffs[0]:.2f} | {settled} |")
out.append("")
out.append("The second table uses the trailing 2-min block rate on a 10-s grid (hop-series.tsv), evaluated from 2 min into the phase so the window lies inside it; overshoot is the largest excursion of that rate from 60/min, signed by the larger side.")
return "\n".join(out)
def _series_rows(rows, step_ms=10000):
"""10-s grid over one node's samples: difficulty, block count, trailing 1-min and 2-min rates in blocks/min."""
if not rows:
return []
t0 = rows[0]["t"] - rows[0]["t"] % step_ms
out = []
j = 0
for t in range(t0, rows[-1]["t"] + step_ms, step_ms):
while j + 1 < len(rows) and rows[j + 1]["t"] <= t:
j += 1
r = rows[j]
def blocks_at(tt):
k = j
while k > 0 and rows[k]["t"] > tt:
k -= 1
return rows[k]["blocks"] if rows[k]["t"] <= tt else None
b60, b120 = blocks_at(t - 60000), blocks_at(t - 120000)
out.append({"t": t, "difficulty": r["difficulty"], "blocks": r["blocks"], "daa": r["daa"], "tips": r["tips"],
"rate1": (r["blocks"] - b60) if b60 is not None else float("nan"),
"rate2": (r["blocks"] - b120) / 2 if b120 is not None else float("nan")})
return out
def hop_series(d, nodes_file):
nodes = read_nodes(nodes_file)
logp = os.path.join(d, "hop.log")
phases = []
if os.path.exists(logp):
for line in open(logp):
f = line.rstrip("\n").split("\t")
if len(f) >= 3:
phases.append((int(f[0]), f[1], f[2]))
per_node = {}
for n in nodes:
p = os.path.join(d, "nodes", n["name"], "samples.tsv")
if os.path.exists(p):
per_node[n["name"]] = read_samples(p)
if not per_node:
return "no samples"
first = nodes[0]["name"] if nodes[0]["name"] in per_node else sorted(per_node)[0]
base = _series_rows(per_node[first])
others = {n: _series_rows(rs) for n, rs in per_node.items()}
# miner threads over time from hop.log: every phase sets the selected nodes; unselected keep theirs
def threads_at(t):
th = {n["name"]: 1 for n in nodes}
for pt, sel, thr in phases:
if pt > t or sel == "end":
continue
for n in nodes:
i, reg = n["index"], n["region"]
hit = (sel == "all" or (sel in ("half", "odd") and i % 2 == 1) or (sel == "even" and i % 2 == 0)
or (sel.startswith("region:") and reg == sel[7:])
or ("-" in sel and sel.replace("-", "").isdigit() and int(sel.split("-")[0]) <= i <= int(sel.split("-")[1])))
if hit:
th[n["name"]] = int(thr)
return sum(th.values())
if phases:
t_lo, t_hi = phases[0][0] - 300000, max(pt for pt, _, _ in phases) + 300000
else:
t_lo, t_hi = base[0]["t"], base[-1]["t"]
lines = ["t_ms\tutc\tphase\tminer_threads_total\tdifficulty_node1\tdifficulty_median\tblocks_node1\tdaa_node1\ttips_node1\trate_1min\trate_2min\tnodes_sampled"]
idx = {n: {r["t"]: r for r in rs} for n, rs in others.items()}
for r in base:
t = r["t"]
if t < t_lo or t > t_hi:
continue
ph = ""
for k, (pt, sel, thr) in enumerate(phases):
if pt <= t and sel != "end":
ph = f"{k + 1}:{sel}:{thr}"
ds = [idx[n][t]["difficulty"] for n in idx if t in idx[n] and idx[n][t]["difficulty"] > 0]
lines.append("\t".join(str(x) for x in [t, time.strftime("%H:%M:%S", time.gmtime(t / 1000)), ph, threads_at(t), f"{r['difficulty']:.1f}",
f"{statistics.median(ds):.1f}" if ds else "", r["blocks"], r["daa"], r["tips"],
fmt(r["rate1"], 0), fmt(r["rate2"], 1), len(ds)]))
return "\n".join(lines)
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] == "hop-series":
print(hop_series(a[2], a[3]))
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)