igneum/sim/horizon/network/propagation.py
2026-10-06 20:52:33 +00:00

257 lines
15 KiB
Python

#!/usr/bin/env python3
"""Propagation and red-share model for the Horizon network lane (lane 5, 6 October 2026).
What it models:
* Poisson block production at `bps` blocks per second, split evenly over the mining nodes.
* A topology: `star` (every miner peers with one hub only, the Devnet 2 run A shape: tools/fleet/box-dn2.sh passes a
single --addpeer=<seed>, the collector read peers=41 at the seed and peers=1 at every pod) or `mesh` (every node
has `degree` random peers, blocks flood by gossip, each node forwards a block once to the peers that have not
announced it).
* One relay hop = inv, request, block: three link traversals (measured in the M21 run, docs/fud-ledger.md M21:
one hop p50 318 to 329 ms on 100-ms proxied links) plus the receiver's own processing (6 to 16 ms measured there).
The link latency is lognormal with median `--link` ms and sigma 0.29 (the cloud devnet's p50 343 / p90 497 ms
spread, docs/bench-log.md "4 October 2026, cloud devnet"). `--hop` sets the hop median directly instead.
* The hub (star only) is a single server: validating a block costs s0 + s1 x mergeset ms (`--hub-s0`, `--hub-s1`);
blocks queue behind it. Run A's hub went from about 6 ms of CPU per accepted block at mergeset 8 to about 400 ms
at mergeset 150 to 200 (A.jsonl cpu_rss against the seed log's Processed lines), so the defaults are 6 and 2.
`--hub-s1 0 --hub-s0 0` is an ideal hub.
* GHOSTDAG colouring: for every block, selected parent = max blue work (hash tiebreak), mergeset = past minus the
selected parent's past, candidates in topological order; a candidate is blue when its anticone inside the blue
set of the new block's chain is at most k (the first k-cluster condition of
vendor/igneum-node/consensus/src/processes/ghostdag/protocol.rs). The second condition (every blue peer's own
anticone bound) is approximated by the candidate's check, so this undercounts reds slightly. Parents capped at
`--parents` by blue work (Kaspa's max_block_parents, consensus/core/src/config/bps.rs). The mergeset size
limit is applied: candidates past the limit are left for a later block (they stay in the anticone).
* Miners join over `--join-spread` seconds with a view of genesis only and mine at once (the fleet's
--enable-unsynced-mining start; run A's pods joined 19:23 to 19:29Z).
* `--schedule a.csv` replaces the constant rate with a per-minute production schedule (blocks per minute).
What it reports: cumulative red share (reds in the mergesets of the final selected chain, over all merged blocks),
mean mergeset size, mean tips at a miner, the natural selected-chain reorg depth at miner 0 (max and p99), the hub's
mean queueing delay and utilisation, mean end-to-end delay to 90% of nodes.
Run (the lock is the main checkout's):
/Users/joshm/Projects/igneum/tools/lock/with-lock.sh run nice -n 19 python3 sim/horizon/network/propagation.py --grid
... --topology star --bps 10 --k 124 --nodes 41 --link 40 --hub-s0 6 --hub-s1 2 --seconds 600 --join-spread 360
"""
import argparse, heapq, math, random, sys, time
def lognormal(rng, median, sigma):
return median * math.exp(rng.gauss(0.0, sigma))
class Dag:
def __init__(self, k, parents_cap, mergeset_cap, rng):
self.k, self.pcap, self.mcap, self.rng = k, parents_cap, mergeset_cap, rng
self.parents = [()]; self.past = [0]; self.sp = [-1]
self.blueset = [1] # bitset: blues of the block's chain incl. itself
self.nblue = [0]; self.nred = [0]
self.blue_work = [0]; self.hash = [0]; self.time = [0.0]; self.miner = [-1]
self.children = [0]
def key(self, b): return (self.blue_work[b], self.hash[b])
def anc(self, a, b): return (self.past[b] >> a) & 1
def add(self, parents, t, miner):
parents = sorted(set(parents), key=self.key, reverse=True)[: self.pcap]
sp = parents[0]
past = 0
for p in parents: past |= self.past[p] | (1 << p)
x = past & ~self.past[sp] & ~(1 << sp)
ms = []
while x:
lsb = x & -x; ms.append(lsb.bit_length() - 1); x ^= lsb
ms.sort(key=self.key)
blueset = self.blueset[sp]; nb = 1; nr = 0; k = self.k
for c in ms[: self.mcap]:
cand = blueset & ~self.past[c] & ~(1 << c)
n = 0; ok = True
while cand:
lsb = cand & -cand; y = lsb.bit_length() - 1; cand ^= lsb
if not self.anc(c, y):
n += 1
if n > k: ok = False; break
if ok: blueset |= (1 << c); nb += 1
else: nr += 1
idx = len(self.parents)
self.parents.append(tuple(parents)); self.past.append(past); self.sp.append(sp)
self.blueset.append(blueset | (1 << idx)); self.nblue.append(nb); self.nred.append(nr)
self.blue_work.append(self.blue_work[sp] + nb); self.hash.append(self.rng.getrandbits(62))
self.time.append(t); self.miner.append(miner); self.children.append(0)
for p in parents: self.children[p] += 1
return idx
def reorg_depth(self, old_tip, new_tip):
"""chain blocks of old_tip's selected chain that are not ancestors (or self) of new_tip"""
if old_tip == new_tip or self.anc(old_tip, new_tip): return 0
d = 0; cur = old_tip
while cur > 0 and not self.anc(cur, new_tip) and cur != new_tip:
d += 1; cur = self.sp[cur]
return d
class Node:
__slots__ = ("id", "known", "tips", "peers", "tip", "reorgs", "joined", "ntips")
def __init__(self, i):
self.id = i; self.known = 1; self.tips = {0}; self.peers = []; self.tip = 0; self.reorgs = []; self.joined = False; self.ntips = []
def simulate(a, rng):
dag = Dag(a.k, a.parents, a.mergeset, rng)
hub = a.topology == "star"
N = a.nodes + (1 if hub else 0)
nodes = [Node(i) for i in range(N)]
miners = list(range(1, N)) if hub else list(range(N))
if hub:
for m in miners: nodes[m].peers = [0]
nodes[0].peers = miners[:]
else:
for i in range(N):
others = [j for j in range(N) if j != i]
for j in rng.sample(others, min(a.degree, len(others))):
if j not in nodes[i].peers: nodes[i].peers.append(j)
if i not in nodes[j].peers: nodes[j].peers.append(i)
joins = {m: rng.uniform(0, a.join_spread) for m in miners}
sigma = 0.29
def hop_delay():
if a.hop: return lognormal(rng, a.hop, sigma) / 1000.0
return (3 * lognormal(rng, a.link, sigma) + a.proc) / 1000.0
ev = [] # (time, seq, kind, payload)
seq = 0
def push(t, kind, payload):
nonlocal seq; seq += 1; heapq.heappush(ev, (t, seq, kind, payload))
# production schedule
sched = None
if a.schedule:
sched = [float(l.split(",")[-1]) for l in open(a.schedule) if l.strip() and not l.startswith("#")]
def rate_at(t):
if sched is None: return a.bps
i = int(t // 60); return (sched[i] if i < len(sched) else sched[-1]) / 60.0
t = 0.0
r0 = rate_at(0.0)
push(rng.expovariate(r0) if r0 > 0 else 1e9, "mine", None)
hub_busy_until = 0.0; hub_wait = 0.0; hub_served = 0; hub_busy_total = 0.0
node_busy_until = [0.0] * N; node_wait = 0.0; node_served = 0; node_queued = [0] * N
arrivals = {} # block -> list of arrival times
produced = 0
def receive(n, b, now):
node = nodes[n]
if (node.known >> b) & 1: return False
node.known |= (1 << b)
for p in dag.parents[b]: node.tips.discard(p)
# b is a tip at this node unless a known child exists (rare: children arrive after parents in this model)
node.tips.add(b)
old = node.tip
if dag.key(b) > dag.key(old):
d = dag.reorg_depth(old, b)
if d > 0: node.reorgs.append(d)
node.tip = b
arrivals.setdefault(b, []).append(now)
return True
def hub_process(b, now):
nonlocal hub_busy_until, hub_wait, hub_served, hub_busy_total
start = max(now, hub_busy_until)
m = dag.nblue[b] + dag.nred[b]
s = (a.hub_s0 + a.hub_s1 * m) / 1000.0
hub_busy_until = start + s; hub_wait += start - now; hub_served += 1; hub_busy_total += s
push(hub_busy_until, "hubdone", b)
while ev:
t, _, kind, payload = heapq.heappop(ev)
if t > a.seconds: break
if kind == "mine":
r = rate_at(t)
push(t + (rng.expovariate(r) if r > 0 else 1e9), "mine", None)
active = [m for m in miners if joins[m] <= t]
if not active: continue
m = rng.choice(active); node = nodes[m]
if not node.joined: node.joined = True
node.ntips.append(len(node.tips))
b = dag.add(list(node.tips), t, m); produced += 1
receive(m, b, t)
for p in node.peers: push(t + hop_delay(), "arrive", (p, b))
elif kind == "arrive":
n, b = payload
if hub and n == 0:
if (nodes[0].known >> b) & 1: continue
nodes[0].known |= (1 << b) # accepted into the queue once
hub_process(b, t)
elif a.node_s0 or a.node_s1:
if (nodes[n].known >> b) & 1 or (node_queued[n] >> b) & 1: continue
node_queued[n] |= (1 << b)
start = max(t, node_busy_until[n]); m = dag.nblue[b] + dag.nred[b]
node_busy_until[n] = start + (a.node_s0 + a.node_s1 * m) / 1000.0
node_wait += start - t; node_served += 1
push(node_busy_until[n], "nodedone", (n, b))
else:
if receive(n, b, t):
for p in nodes[n].peers: push(t + hop_delay(), "arrive", (p, b))
elif kind == "nodedone":
n, b = payload
if receive(n, b, t):
for p in nodes[n].peers:
if p != dag.miner[b]: push(t + hop_delay(), "arrive", (p, b))
elif kind == "hubdone":
b = payload
nodes[0].known &= ~(1 << b); receive(0, b, t)
for p in nodes[0].peers:
if p != dag.miner[b]: push(t + hop_delay(), "arrive", (p, b))
# metrics from the hub's (or node 0's) final selected chain
tip = max(range(len(dag.parents)), key=dag.key)
blues = reds = chain = 0; ms_sizes = []
cur = tip
while cur > 0:
blues += dag.nblue[cur]; reds += dag.nred[cur]; chain += 1; ms_sizes.append(dag.nblue[cur] + dag.nred[cur])
cur = dag.sp[cur]
merged = blues + reds
red_share = reds / merged if merged else 0.0
m0 = nodes[miners[0]]
reorgs = sorted(m0.reorgs)
p99 = reorgs[int(0.99 * (len(reorgs) - 1))] if reorgs else 0
# delay to 90% of nodes
d90 = []
for b, arr in arrivals.items():
if len(arr) >= 0.9 * len(miners):
arr = sorted(arr); d90.append(arr[int(0.9 * len(miners)) - 1] - dag.time[b])
tips_mean = sum(sum(n.ntips) for n in nodes) / max(1, sum(len(n.ntips) for n in nodes))
return dict(produced=produced, merged=merged, red_share=red_share, chain=chain, ms_mean=sum(ms_sizes) / max(1, len(ms_sizes)),
ms_max=max(ms_sizes) if ms_sizes else 0, tips_mean=tips_mean, reorg_max=max(reorgs) if reorgs else 0, reorg_p99=p99,
hub_wait=hub_wait / max(1, hub_served), hub_util=hub_busy_total / max(a.seconds, 1e-9), d90=sum(d90) / max(1, len(d90)),
node_wait=node_wait / max(1, node_served),
blocks_per_s=produced / a.seconds)
def main():
ap = argparse.ArgumentParser()
ap.add_argument("--topology", default="mesh", choices=["mesh", "star"])
ap.add_argument("--bps", type=float, default=10); ap.add_argument("--k", type=int, default=124)
ap.add_argument("--nodes", type=int, default=41); ap.add_argument("--degree", type=int, default=8)
ap.add_argument("--parents", type=int, default=16); ap.add_argument("--mergeset", type=int, default=248)
ap.add_argument("--link", type=float, default=40.0, help="one-way link latency median, ms")
ap.add_argument("--hop", type=float, default=0.0, help="hop delay median, ms (overrides --link)")
ap.add_argument("--proc", type=float, default=10.0, help="receiver processing per hop, ms")
ap.add_argument("--hub-s0", type=float, default=0.0); ap.add_argument("--hub-s1", type=float, default=0.0)
ap.add_argument("--node-s0", type=float, default=0.0, help="every other node's per-block validation cost, ms")
ap.add_argument("--node-s1", type=float, default=0.0, help="plus this many ms per mergeset block")
ap.add_argument("--seconds", type=float, default=300); ap.add_argument("--join-spread", type=float, default=0.0)
ap.add_argument("--schedule", default=None); ap.add_argument("--seed", type=int, default=7)
ap.add_argument("--grid", action="store_true", help="the lane's grid: bps x hop delay, mesh of degree 8, ideal links")
ap.add_argument("--out", default=None)
a = ap.parse_args()
lines = []
def emit(s): print(s, flush=True); lines.append(s)
if a.grid:
from_k = {1: 18, 10: 124, 32: 362}
emit("| bps | k | hop delay median ms | red share | mergeset mean / max | tips mean | reorg max / p99 (miner 0) | delay to 90% of nodes s | blocks/s |")
emit("|---|---|---|---|---|---|---|---|---|")
for bps in (1, 10, 32, 100):
k = from_k.get(bps) or 1074
for hop in (50, 100, 300, 1000):
rng = random.Random(a.seed)
b = argparse.Namespace(**vars(a)); b.bps = bps; b.k = k; b.hop = hop; b.topology = "mesh"; b.hub_s0 = 0; b.hub_s1 = 0
b.mergeset = 180 if bps == 1 else min(512, 2 * k); b.parents = 10 if bps == 1 else 16
b.seconds = min(a.seconds, max(60.0, 12000.0 / bps))
t0 = time.time(); r = simulate(b, rng)
emit(f"| {bps} | {k} | {hop} | {100*r['red_share']:.1f}% | {r['ms_mean']:.1f} / {r['ms_max']} | {r['tips_mean']:.1f} | {r['reorg_max']} / {r['reorg_p99']} | {r['d90']:.2f} | {r['blocks_per_s']:.1f} |")
print(f"# ({b.seconds:.0f} s simulated in {time.time()-t0:.0f} s wall)", file=sys.stderr)
else:
rng = random.Random(a.seed); r = simulate(a, rng)
emit(f"{a.topology} bps={a.bps} k={a.k} nodes={a.nodes} degree={a.degree} link={a.link} hop={a.hop} hub=({a.hub_s0},{a.hub_s1}) node=({a.node_s0},{a.node_s1}) join={a.join_spread} sched={a.schedule} seconds={a.seconds} seed={a.seed}")
emit(f"produced {r['produced']} ({r['blocks_per_s']:.2f}/s) merged {r['merged']} red share {100*r['red_share']:.1f}% chain {r['chain']} mergeset mean {r['ms_mean']:.1f} max {r['ms_max']} tips {r['tips_mean']:.1f} reorg max {r['reorg_max']} p99 {r['reorg_p99']} hub wait {r['hub_wait']:.2f}s util {r['hub_util']:.2f} node wait {r['node_wait']:.2f}s d90 {r['d90']:.2f}s")
if a.out:
with open(a.out, "a") as f: f.write("\n".join(lines) + "\n")
if __name__ == "__main__":
main()