#!/usr/bin/env python3 """Offline analysis of a results// directory (runs on the Mac, standard library only). analyze.py rtt region-pair RTT table from /rtt.tsv analyze.py propagation per-node and per-region block propagation delay from /nodes//blocks.tsv (window t0..t1) analyze.py locks conflicting locks between the checkpoint dumps at heal analyze.py hop block rate and difficulty per hop.log phase (node 1 samples) analyze.py summary summary.md of everything present under analyze.py bench-entry 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//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 "" 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-, genesis 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): ", "", "Block propagation (arrival minus first arrival anywhere): p50 ms, p90 ms, p99 ms, max ms; per region .", "", "Partition , min: minority reorg depth max (getVirtualChainFromBlock from the minority sink at heal), majority ; converged s after heal; conflicting locks: .", "", "Hash-rate steps (hop.sh): : 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 , p99 , max over reorgs (summary.md).", "", "Reading: . 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)