igneum/tools/fleet/box-prover.py
igneum-labs 0c245a6083 fleet: every stop goes through a pid file; pkill and killall are gone from the live tree
The founder's word (20:21 BST 8 Oct 2026): pkill and killall refuse by construction on every box and pod (exit 97). lib/pidkill.sh
(kill_pidfile, kill_children by parent pid, kill_sock_owner through ss -xlp) is the shared form; box-prover.py finds the SP1 server by
the pid owning its socket; lib/box.py, box-kill.sh, kill-node.sh, the rig scripts (box-ember, box-rig, box-matrix, box-floor-v5),
lib/standing.py and publish-2-move.py stop by pid files; box-dn3.sh writes node, miner-loop and miner pids and uses pidkill;
fleet-puller.sh and dn3-kill.sh call fleet-stop.sh. Seventeen devnet-2-era scripts that only ever stopped by name move to
tools/fleet/retired/ with a README.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
2026-10-08 20:32:25 +00:00

625 lines
56 KiB
Python

#!/usr/bin/env python3
"""The fleet's segment-aligned prover for a Linux box (phase 2): a port of app/igneum-app/src/prover.rs (branch
proving-v1, 272b025) and tools/proving-v1/pc2-segments.ps1 around the four binaries. One pass every 15 s:
1. igneum_getProvingStatus (start, n, unproven, tip) and igneum_getAssignedShards [[keyHash], 600], grouped into
whole untouched segments (every block present, every shard open, unpaid, not in this pool) whose deadline
(last DAA + unproven) is at least 240 DAA (or 1.5x the last segment's time) past the tip; ranked by FNV-1a of
(first, key hash) so a fleet of provers spreads over the segments instead of racing for one.
2. igneum_getSegmentStatement [first]: executed and pending; chained with --prev when the previous segment's proof is
in this pool (igneum_getSegmentProofBytes), else fresh. A fresh record refused with "pending until" is held and
offered again every pass until the segment's deadline (the 272b025 behaviour).
3. one export (igneum_exportSegments 0..last), one fixture per block (igneum-prove-export), one host run
(igneum-prove-host --mode chain --chain ... --save-shards [--prev]) on the patched server (HOME=/opt/igneum-floor/home,
SP1_GPU_ELEMENT_THRESHOLD from THRESHOLD), the miner paused for the run when MINER=pause (prove-alone cards).
4. every shard record signed (igneum-miner sign-record) and submitted (igneum_submitProofRecord); the segment record
(sign-segment-record, igneum_submitSegmentRecord) once every shard is accepted and the statement equals the node's.
5. the paid state of every submitted segment polled each pass (igneum_getSegmentRecords); a state file for the
collector: /root/fleet/out/prover-state.json; every event a RESULT line in /root/fleet/out/prover.log.
6. the outcome ledger (review B F08, 8 October 2026): every claimed segment is an eligible job and ends in exactly one
outcome, paid, expired (its deadline passed with no paid record: unpaid after the submit, or held past it) or
cancelled (this prover gave it up for a cause: disk, export, cut, chain, timeout, shards, statement, sign,
refused); until then it is active. Each segment row carries outcome, cause, deadline, the margin at the claim,
the seconds spent (export, cut, chain; wasted unless paid) and the deadline miss in DAA; the state carries the
counters (outcomes, wasted_s by cause, deadline_misses) and every close is a RESULT outcome line.
tools/fleet/prover-outcomes.py reads the state files and the logs into the report (paid completions, missed
deadlines, wasted work by cause, accepted-proof throughput).
7. V6-08 (docs/design/proving-task-protection.md, 8 October 2026): the claim is posted to the node as a lease
(igneum_claimSegment) before any export; a refusal (another key holds the lease, or this key is deprioritised for
abandoned leases) is a RESULT lease_refused line and no ledger row; a worklist row leased by another key is never a
candidate. The chain-health backpressure gate (chain_health): no claim while the executor is blocked, re-executing or
discontinuous, more than LAG_BLOCKS behind the sink, the DAA stalled STALL_S, finality paused FINALITY_PAUSE_S, or the
node unreachable; RESULT backpressure lines on each change. Task sizing: TASK=segment|shard|auto (auto: shard under
SEGMENT_MIN_MIB of card memory), MAX_SHARDS caps a segment; a shard task is one block's one shard proven with
--mode compressed --shard i and submitted as a shard record (claim_shard / outcome lines, kind shard in the ledger).
Preflight (preflight_verdict): the host and export binaries resolve, are executable, ldd clean, their sha256 equal
PREFLIGHT_MANIFEST's when it exists, --mode id prints both program ids and they equal the node's pinned pair; a
mismatch is RESULT preflight_failed naming the file and exit 4. The fourth closing outcome abandoned (a claim whose
lease or deadline passed with no record: cause lease, or restart for active rows found in the state file at start)
and the conservation line claims = paid + expired + cancelled + abandoned + active in every ledger line.
8. `box-prover.py --self-test` runs the pure functions (chain_health, task_mode, shard_candidates, preflight_verdict,
the conservation identity) against known-failed-first cases without a node, a card or the fleet directories.
Env: LABEL (the key label, kept for the box's life), WALLET (payout), THRESHOLD (element threshold or empty),
MINER (keep|pause), RUN_HOURS (default 9), TASK (segment|shard|auto, default auto), SEGMENT_MIN_MIB (16384), MAX_SHARDS (0 = no
cap), LEASE (on|off, default on), LAG_BLOCKS (64), STALL_S (120), FINALITY_PAUSE_S (600), PREFLIGHT_MANIFEST
(/root/fleet/in/prover-manifest.json).
"""
import json, os, sys, time, subprocess, datetime, binascii, urllib.request, signal
# a rig runs one loop per card under /root/fleet/card<n>/ (FLEET_CARD)
F = "/root/fleet"; CARD = os.environ.get("FLEET_CARD"); OUT = f"{F}/card{CARD}/out" if CARD else f"{F}/out"; B = "/opt/igneum/pkg/bin"; FLOOR = "/opt/igneum-floor"
HOST = f"{FLOOR}/bin/igneum-prove-host"; EXPORT = f"{FLOOR}/bin/igneum-prove-export"
EVM = "http://127.0.0.1:26790"; GRPC = "grpc://127.0.0.1:26610"; CHAIN = os.environ.get("CHAIN_NAME", "igneum-devnet") # Devnet 2 signs for igneum-devnet-2
LABEL = os.environ.get("LABEL", "box"); WALLET = os.environ.get("WALLET", "0x" + "19" * 20)
THRESHOLD = os.environ.get("THRESHOLD", ""); MINER = os.environ.get("MINER", "keep"); RUN_HOURS = float(os.environ.get("RUN_HOURS", "9"))
EXPORT_FROM = int(os.environ.get("EXPORT_FROM", "27276"))
# the export directory's own rule (6 October 2026, 20:4xZ, after the hub's node died on a full disk): a size cap and an age cap
# on OUT/segs, each segment's export deleted the moment its record is accepted or paid, and a disk-free check before each
# export that skips with a logged line under 10 percent free. Defaults fit a 100 GB box for a week.
SEGS_CAP_GB = float(os.environ.get("SEGS_CAP_GB", "20")); SEGS_MAX_AGE_H = float(os.environ.get("SEGS_MAX_AGE_H", str(7 * 24))); DISK_MIN_FREE_PCT = float(os.environ.get("DISK_MIN_FREE_PCT", "10"))
# V6-08 knobs (docstring point 7)
TASK = os.environ.get("TASK", "auto"); SEGMENT_MIN_MIB = float(os.environ.get("SEGMENT_MIN_MIB", "16384")); MAX_SHARDS = int(os.environ.get("MAX_SHARDS", "0") or 0)
LEASE = os.environ.get("LEASE", "on") != "off"; LAG_BLOCKS = int(os.environ.get("LAG_BLOCKS", "64")); STALL_S = float(os.environ.get("STALL_S", "120")); FINALITY_PAUSE_S = float(os.environ.get("FINALITY_PAUSE_S", "600"))
PREFLIGHT_MANIFEST = os.environ.get("PREFLIGHT_MANIFEST", "/root/fleet/in/prover-manifest.json")
def hexi(v): return int(v, 16) if isinstance(v, str) and v.startswith("0x") else int(v or 0)
# ---- V6-08 pure functions (no node, no card; exercised by --self-test) ----
def chain_health(ex, st, now_s, last_tip, last_tip_at, lag_blocks=None, stall_s=None, finality_pause_s=None):
"""The backpressure gate (design section 2): the hold cause, or None when the chain is healthy enough to claim into."""
lag_blocks = LAG_BLOCKS if lag_blocks is None else lag_blocks; stall_s = STALL_S if stall_s is None else stall_s; finality_pause_s = FINALITY_PAUSE_S if finality_pause_s is None else finality_pause_s
if not ex or not st: return "unreachable", "no reading from igneum_getExecStatus or igneum_getProvingStatus"
if ex.get("blocked"): return "blocked", str(ex.get("blocked"))[:120]
if ex.get("reexecuting"): return "reexecuting", str(ex.get("reexecuting"))[:120]
if ex.get("recordsContinuous") is False: return "discontinuous", f"continuity break at {ex.get('continuityBreak')}"
sink = ex.get("sinkNumber"); tip_n = hexi(ex.get("executedTip"))
if sink is not None and hexi(sink) - tip_n > lag_blocks: return "lagging", f"executed tip {tip_n} is {hexi(sink) - tip_n} chain blocks behind the sink {hexi(sink)} (over {lag_blocks})"
tip_daa = hexi(st.get("tipDaa"))
if last_tip is not None and tip_daa == last_tip and last_tip_at is not None and now_s - last_tip_at >= stall_s: return "stalled", f"tip DAA {tip_daa} unchanged for {int(now_s - last_tip_at)} s (over {int(stall_s)})"
paused = st.get("pausedSinceMs")
if paused is not None and now_s * 1000 - hexi(paused) > finality_pause_s * 1000: return "finality_paused", f"finality paused since {hexi(paused)} ms, {int(now_s - hexi(paused) / 1000)} s (over {int(finality_pause_s)}): {st.get('finalityReason')}"
return None, ""
def task_mode(task, total_mib, segment_min_mib=None):
"""Design section 3: segment or shard. auto puts a card under the line on shard tasks (the 12 GB tier claimed 313 segments and was paid for none)."""
segment_min_mib = SEGMENT_MIN_MIB if segment_min_mib is None else segment_min_mib
if task in ("segment", "shard"): return task
if total_mib is None: return "segment"
return "shard" if float(total_mib) < float(segment_min_mib) else "segment"
def shard_candidates(work, attempted, key):
"""Design section 3, the app's choose: open, unpaid, not in the pool, not leased by another key, not attempted; assigned
before open, newest first, the smallest pgas among equals."""
rows = []
for w in work or []:
rid = f"{hexi(w.get('number'))}:{int(w.get('shard', 0))}"
if rid in attempted or not w.get("open") or w.get("paid") is not None or (w.get("pool") and len(w["pool"]) > 0) or w.get("leased"): continue
rows.append({"id": rid, "number": hexi(w.get("number")), "hash": w.get("hash"), "shard": int(w.get("shard", 0)), "daa": hexi(w.get("daaScore")), "pgas": hexi(w.get("pgas")), "assigned": bool(w.get("assigned")), "wei": hexi(w.get("shardWei"))})
rows.sort(key=lambda r: (not r["assigned"], -r["number"], r["pgas"], r["shard"]))
return rows
def preflight_verdict(files, manifest, printed_ids, node_ids):
"""Design section 5: (ok, line, detail). files: name -> {path, exists, exec, sha256, ldd_missing}; manifest: the expected set or
None; printed_ids: (shard, aggregator) from --mode id or None; node_ids: (shard, aggregator) from the node, each possibly None."""
for name in ("igneum-prove-host", "igneum-prove-export"):
f = files.get(name) or {}
if not f.get("exists"): return False, f"file={f.get('path')} missing (dangling link or absent)", name
if not f.get("exec"): return False, f"file={f.get('path')} not executable", name
if f.get("ldd_missing"): return False, f"file={f.get('path')} ldd missing {f['ldd_missing']}", name
if manifest and manifest.get(name) and str(manifest[name]).lower().replace("0x", "") != str(f.get("sha256", "")).lower():
return False, f"file={f.get('path')} expected={manifest[name]} got={f.get('sha256')}", name
if not printed_ids or not printed_ids[0] or not printed_ids[1]: return False, f"file={files.get('igneum-prove-host', {}).get('path')} --mode id printed no program ids", "ids"
low = lambda x: (x or "").lower()
if manifest:
for k, got in (("shard_program_id", printed_ids[0]), ("aggregator_id", printed_ids[1])):
if manifest.get(k) and low(manifest[k]) != low(got): return False, f"file={files.get('igneum-prove-host', {}).get('path')} {k} expected={manifest[k]} got={got}", k
for k, got, node in (("shard_program_id", printed_ids[0], node_ids[0]), ("aggregator_id", printed_ids[1], node_ids[1])):
if node and low(node) != low(got): return False, f"file={files.get('igneum-prove-host', {}).get('path')} {k} host={got} node={node}", k
return True, f"manifest={'ok' if manifest else 'absent'} host={files['igneum-prove-host'].get('sha256', '')[:16]} export={files['igneum-prove-export'].get('sha256', '')[:16]} shard={printed_ids[0][:10]} aggregator={printed_ids[1][:10]}", ""
def flow_line(claimed, o):
"""Design section 6: the conservation identity as one line, ok or BROKEN by the difference."""
s = o.get("paid", 0) + o.get("expired", 0) + o.get("cancelled", 0) + o.get("abandoned", 0) + o.get("active", 0)
return f"flow claims {claimed} = paid {o.get('paid', 0)} + expired {o.get('expired', 0)} + cancelled {o.get('cancelled', 0)} + abandoned {o.get('abandoned', 0)} + active {o.get('active', 0)}: " + ("ok" if s == claimed else f"BROKEN by {claimed - s}")
def self_test():
# chain_health, known-failed first: a healthy reading holds nothing (a stub that holds on everything would starve the fleet)
ex = {"blocked": None, "reexecuting": None, "recordsContinuous": True, "sinkNumber": "0x70", "executedTip": "0x6e"}; st = {"tipDaa": "0x1000", "pausedSinceMs": None}
assert chain_health(ex, st, 1000.0, 0xfff, 900.0)[0] is None
assert chain_health(None, st, 1000.0, None, None)[0] == "unreachable" and chain_health(ex, None, 1000.0, None, None)[0] == "unreachable"
assert chain_health(dict(ex, blocked="snapshot refused"), st, 1000.0, None, None)[0] == "blocked"
assert chain_health(dict(ex, reexecuting="from 100"), st, 1000.0, None, None)[0] == "reexecuting"
assert chain_health(dict(ex, recordsContinuous=False, continuityBreak="0x5"), st, 1000.0, None, None)[0] == "discontinuous"
assert chain_health(dict(ex, sinkNumber="0xb0"), st, 1000.0, None, None, lag_blocks=64)[0] == "lagging" and chain_health(dict(ex, sinkNumber="0xae"), st, 1000.0, None, None, lag_blocks=64)[0] is None, "64 behind is not over 64"
assert chain_health(ex, st, 1000.0, 0x1000, 880.0, stall_s=120)[0] == "stalled" and chain_health(ex, st, 1000.0, 0x1000, 881.0, stall_s=120)[0] is None and chain_health(ex, st, 1000.0, 0xfff, 0.0, stall_s=120)[0] is None, "a move clears the stall"
assert chain_health(ex, dict(st, pausedSinceMs=hex(1000 * 1000 - 601_000)), 1000.0, None, None, finality_pause_s=600)[0] == "finality_paused" and chain_health(ex, dict(st, pausedSinceMs=hex(1000 * 1000 - 599_000)), 1000.0, None, None, finality_pause_s=600)[0] is None
# task_mode, known-failed first: the stub that always said segment is what sent the 3060 tier after whole segments
assert task_mode("auto", 12288, 16384) == "shard" and task_mode("auto", 24564, 16384) == "segment" and task_mode("auto", None, 16384) == "segment"
assert task_mode("segment", 12288, 16384) == "segment" and task_mode("shard", 49140, 16384) == "shard"
# shard_candidates: assigned first, newest first, smallest pgas among equals; leased, paid, pooled, closed and attempted never offered
w = lambda n, s, **k: dict({"number": hex(n), "shard": s, "hash": f"0x{n:064x}", "daaScore": hex(10 * n), "pgas": hex(k.pop("pgas", 100)), "open": True, "paid": None, "pool": [], "assigned": False, "shardWei": "0x1"}, **k)
work = [w(10, 0, assigned=True, pgas=300), w(12, 0), w(12, 1, pgas=50), w(11, 0, leased=True), w(9, 0, paid={"wei": "0x1"}), w(8, 0, pool=[{"keyHash": "0x1"}]), w(7, 0, open=False), w(6, 0), w(13, 0, assigned=True, pgas=900)]
got = [c["id"] for c in shard_candidates(work, {"6:0"}, "0xk")]
assert got == ["13:0", "10:0", "12:1", "12:0"], got
# preflight_verdict, known-failed first: a dangling link fails naming the file (the 8 October prover passed on exists alone and failed 629 chains)
files = {"igneum-prove-host": {"path": "/opt/igneum-floor/bin/igneum-prove-host", "exists": False, "exec": False, "sha256": "", "ldd_missing": ""}, "igneum-prove-export": {"path": "/opt/igneum-floor/bin/igneum-prove-export", "exists": True, "exec": True, "sha256": "ee" * 32, "ldd_missing": ""}}
ok, line, which = preflight_verdict(files, None, ("0xaa", "0xbb"), ("0xaa", "0xbb")); assert not ok and "igneum-prove-host missing" in line and which == "igneum-prove-host", line
files["igneum-prove-host"].update(exists=True, exec=True, sha256="ab" * 32)
ok, line, _ = preflight_verdict(files, None, ("0xaa", "0xbb"), ("0xaa", "0xbb")); assert ok and "manifest=absent" in line, line
ok, line, _ = preflight_verdict(files, {"igneum-prove-host": "0x" + "cd" * 32}, ("0xaa", "0xbb"), ("0xaa", "0xbb")); assert not ok and "expected=0x" + "cd" * 32 in line and "got=" + "ab" * 32 in line, line
ok, line, _ = preflight_verdict(files, {"igneum-prove-host": "0x" + "ab" * 32, "shard_program_id": "0xaa", "aggregator_id": "0xbb"}, ("0xaa", "0xbb"), ("0xaa", "0xbb")); assert ok and "manifest=ok" in line, line
ok, line, which = preflight_verdict(files, None, ("0xaa", "0xbb"), ("0xa1", "0xbb")); assert not ok and "host=0xaa node=0xa1" in line and which == "shard_program_id", line
ok, line, _ = preflight_verdict(files, None, None, ("0xaa", "0xbb")); assert not ok and "printed no program ids" in line, line
ok, line, _ = preflight_verdict(files, None, ("0xaa", "0xbb"), (None, None)); assert ok, line # a node without pinned ids: nothing to compare
ok, line, _ = preflight_verdict(dict(files, **{"igneum-prove-export": dict(files["igneum-prove-export"], ldd_missing="libcuda.so.1")}), None, ("0xaa", "0xbb"), ("0xaa", "0xbb")); assert not ok and "ldd missing libcuda.so.1" in line, line
# the conservation identity
assert flow_line(5, {"paid": 1, "expired": 1, "cancelled": 1, "abandoned": 1, "active": 1}).endswith(": ok")
assert flow_line(6, {"paid": 1, "expired": 1, "cancelled": 1, "abandoned": 1, "active": 1}).endswith("BROKEN by 1")
print("RESULT box-prover self-test PASS: chain_health 11 cases, task_mode 5, shard_candidates order and exclusions, preflight_verdict 8 cases, flow line ok and BROKEN")
return 0
if "--self-test" in sys.argv: sys.exit(self_test())
import shutil
def seg_dir_size(d):
return sum(os.path.getsize(os.path.join(r, f)) for r, _, fs in os.walk(d) for f in fs if os.path.exists(os.path.join(r, f)))
def drop_export(first, why):
d = f"{OUT}/segs/seg-{first}"
if os.path.isdir(d): sz = seg_dir_size(d); shutil.rmtree(d, ignore_errors=True); say(f"RESULT export_dropped {stamp()} segment {first} ({sz/1e6:.0f} MB): {why}")
def prune_exports():
root = f"{OUT}/segs"
if not os.path.isdir(root): return
entries = sorted((os.path.join(root, n) for n in os.listdir(root)), key=lambda x: os.path.getmtime(x))
now_ = time.time(); total = sum(seg_dir_size(e) if os.path.isdir(e) else os.path.getsize(e) for e in entries)
for e in entries:
age_h = (now_ - os.path.getmtime(e)) / 3600
if age_h > SEGS_MAX_AGE_H or total > SEGS_CAP_GB * 1e9:
sz = seg_dir_size(e) if os.path.isdir(e) else os.path.getsize(e); shutil.rmtree(e, ignore_errors=True) if os.path.isdir(e) else os.remove(e); total -= sz
say(f"RESULT export_pruned {stamp()} {os.path.basename(e)} ({sz/1e6:.0f} MB, {age_h:.1f} h old): dir {total/1e9:.1f} GB against the {SEGS_CAP_GB:.0f} GB cap, {SEGS_MAX_AGE_H:.0f} h age cap")
def disk_free_pct(path="/"):
st = os.statvfs(path); return 100.0 * st.f_bavail / max(st.f_blocks, 1) # the devnet's exec restart block (ov13.json exec_restart_number) when the node does not report one
MINE = f"{F}/card{CARD}/mine" if CARD else f"{F}/mine"
os.makedirs(f"{OUT}/segs", exist_ok=True); os.makedirs(f"{MINE}/packs", exist_ok=True)
DEV = os.environ.get("IGNEUM_CUDA_DEVICE", "0")
LOG = open(f"{OUT}/prover.log", "a")
def stamp(): return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
def say(s): LOG.write(f"{s}\n"); LOG.flush(); print(s, flush=True)
def rpc(method, params, timeout=60):
body = json.dumps({"jsonrpc": "2.0", "id": 1, "method": method, "params": params}).encode()
try:
with urllib.request.urlopen(urllib.request.Request(EVM, data=body, headers={"Content-Type": "application/json"}), timeout=timeout) as r:
d = json.loads(r.read())
if d.get("error"): say(f"RESULT rpc_error {stamp()} {method}: {str(d['error'])[:200]}"); return None
return d.get("result")
except Exception as e:
say(f"RESULT rpc_fail {stamp()} {method}: {str(e)[:120]}"); return None
def fnv1a(s):
h = 0xcbf29ce484222325
for c in s.encode(): h ^= c; h = (h * 0x100000001b3) & 0xffffffffffffffff
return h
def kill_server():
if CARD: subprocess.run(f"rm -f /tmp/sp1-cuda-{DEV}.sock", shell=True) # a rig: never another card's server
else:
# kills by pid only (founder 8 Oct 2026, by construction): the server is found through the pid that owns its unix socket
# (ss -xlp on /tmp/sp1-cuda-*.sock), written to /root/fleet/pids/sp1-server.pid, then signalled by that pid; never by name
o=subprocess.run("ss -xlp 2>/dev/null | grep -E '/tmp/sp1-cuda-[0-9]*\\.sock' | grep -oE 'pid=[0-9]+' | cut -d= -f2 | sort -u", shell=True, capture_output=True, text=True).stdout.split()
os.makedirs("/root/fleet/pids", exist_ok=True)
for s in o:
try:
pid=int(s); open("/root/fleet/pids/sp1-server.pid","w").write(f"{pid}\n"); os.kill(pid, 15); time.sleep(1)
try: os.kill(pid, 9)
except ProcessLookupError: pass
except Exception: pass
subprocess.run("rm -f /tmp/sp1-cuda-*.sock", shell=True)
MPROC = None
def miner_start():
global MPROC
if MINER == "none": return # the box's own miner loop runs outside (Devnet 2 boxes)
if MPROC and MPROC.poll() is None: return
subprocess.run(f"cd {MINE} && rm -rf packs/devnet && {B}/igneum-miner export-pack {GRPC} packs/devnet > {OUT}/prover-export-pack.log 2>&1", shell=True)
MPROC = subprocess.Popen([f"{B}/igneum-miner", "mine", GRPC, "1", "100000000", LABEL, "--worker", f"{B}/igneum-worker-cuda", "--worker-args", f"--device {DEV} --pack packs/devnet",
"--prepare-packs", "packs/prepare", "--exit-on-seed-change", "--evm-address", WALLET, "--payout-label", LABEL, "--status-secs", "30"],
cwd=MINE, stdout=open(f"{OUT}/prover-miner.log", "a"), stderr=subprocess.STDOUT)
say(f"RESULT miner_start {stamp()} pid={MPROC.pid}")
def miner_stop():
global MPROC
if MPROC: MPROC.terminate(); time.sleep(2); MPROC.kill(); MPROC = None
pass # the worker is the miner's child; MPROC.terminate() above ends it (no kill by name: founder 8 Oct 2026)
def miner_rate(n=6):
try:
vals = [float(l.split(" now=")[1].split()[0]) for l in open(f"{OUT}/prover-miner.log").read().split("\n") if "STATUS" in l and " now=" in l][-n:]
return round(sum(vals) / len(vals), 2) if vals else 0
except Exception: return 0
# V6-08 preflight (docstring point 7), the local half: the binaries, their sha256 against the manifest, ldd, --mode id
import hashlib, re
def _file_facts(path):
real = os.path.realpath(path); f = {"path": path, "real": real, "exists": os.path.isfile(real), "exec": os.access(real, os.X_OK), "sha256": "", "ldd_missing": ""}
if f["exists"]:
h = hashlib.sha256()
with open(real, "rb") as fh:
for chunk in iter(lambda: fh.read(1 << 20), b""): h.update(chunk)
f["sha256"] = h.hexdigest()
try: f["ldd_missing"] = ",".join(sorted({l.split()[0] for l in subprocess.run(["ldd", real], capture_output=True, text=True, timeout=30).stdout.split("\n") if "not found" in l}))
except Exception as e: f["ldd_missing"] = f"ldd failed: {str(e)[:60]}"
return f
def _printed_ids():
try: out = subprocess.run([HOST, "--mode", "id"], capture_output=True, text=True, timeout=120, env=dict(os.environ, HOME=f"{FLOOR}/home")).stdout
except Exception as e: say(f"RESULT preflight {stamp()} --mode id failed: {str(e)[:100]}"); return None
m1 = re.search(r"shard program id (0x[0-9a-fA-F]+)", out); m2 = re.search(r"aggregator id (0x[0-9a-fA-F]+)", out)
return (m1.group(1) if m1 else None, m2.group(1) if m2 else None)
PF_FILES = {"igneum-prove-host": _file_facts(HOST), "igneum-prove-export": _file_facts(EXPORT)}
PF_MANIFEST = None
if os.path.isfile(PREFLIGHT_MANIFEST):
try: PF_MANIFEST = json.load(open(PREFLIGHT_MANIFEST))
except Exception as e: say(f"RESULT preflight_failed {stamp()} file={PREFLIGHT_MANIFEST} does not parse: {str(e)[:80]}"); sys.exit(4)
PF_IDS = _printed_ids() if PF_FILES["igneum-prove-host"]["exists"] else None
_ok, _line, _which = preflight_verdict(PF_FILES, PF_MANIFEST, PF_IDS, (None, None))
if not _ok: say(f"RESULT preflight_failed {stamp()} {_line}"); sys.exit(4)
say(f"RESULT preflight {stamp()} local ok {_line}")
for name, dep in (("sp1-gpu-server", f"{FLOOR}/bin/sp1-gpu-server"), ("igneum-miner", f"{B}/igneum-miner")):
if not os.access(dep, os.X_OK): say(f"RESULT preflight_failed {stamp()} file={dep} missing or not executable ({name})"); sys.exit(4)
try: TOTAL_MIB = float(subprocess.run(["nvidia-smi", "-i", DEV, "--query-gpu=memory.total", "--format=csv,noheader,nounits"], capture_output=True, text=True, timeout=30).stdout.strip().split("\n")[0])
except Exception as e: say(f"RESULT preflight_failed {stamp()} file=nvidia-smi device {DEV} does not answer: {str(e)[:80]}"); sys.exit(4)
MODE = task_mode(TASK, TOTAL_MIB)
say(f"RESULT task_mode {stamp()} {MODE} (TASK={TASK}, card {TOTAL_MIB:.0f} MiB, line {SEGMENT_MIN_MIB:.0f} MiB, MAX_SHARDS={MAX_SHARDS}, lease={'on' if LEASE else 'off'})")
# the key
kh = subprocess.run([f"{B}/igneum-miner", "key-hash", LABEL], capture_output=True, text=True).stdout.strip().split("\n")[-1].strip()
if len(kh) == 64: kh = "0x" + kh
if len(kh) != 66: say(f"RESULT prover_failed key-hash gave '{kh[:40]}'"); sys.exit(2)
say(f"RESULT start {stamp()} label={LABEL} key={kh[:18]} wallet={WALLET[:10]} threshold={THRESHOLD or 'default'} miner={MINER} host={os.path.exists(HOST)}")
# wait for the node to be synced
for _ in range(120):
w = subprocess.run(f"{B}/igneum-miner watch 1 {GRPC} 2>/dev/null | grep -o 'synced=[a-z]*' | tail -1", shell=True, capture_output=True, text=True).stdout.strip()
if w == "synced=true": break
time.sleep(15)
say(f"RESULT node {stamp()} {w}")
# V6-08 preflight, the node half: the host's program ids against the node's pinned pair (what the chain pays for)
_st0 = rpc("igneum_getProvingStatus", []) or {}
_node_ids = ((_st0.get("v1") or {}).get("shardProgramId"), (_st0.get("v1") or {}).get("aggregatorId"))
_ok, _line, _which = preflight_verdict(PF_FILES, PF_MANIFEST, PF_IDS, _node_ids)
if not _ok: say(f"RESULT preflight_failed {stamp()} {_line}"); sys.exit(4)
say(f"RESULT preflight {stamp()} ok {_line} node_shard={str(_node_ids[0])[:10]} node_aggregator={str(_node_ids[1])[:10]} device={DEV}")
# V6-08 section 6: active rows of a previous run in the state file are abandoned by this restart (the log is the ledger across runs)
try:
_prev = json.load(open(f"{OUT}/prover-state.json"))
for x in _prev.get("segments") or []:
if x.get("outcome") == "active":
_sp = round(float(x.get("export_s", 0)) + float(x.get("cut_s", 0)) + float(x.get("chain_s", 0)), 1)
say(f"RESULT outcome {stamp()} " + (f"block {x['first']} shard {x.get('shard')}" if x.get("kind") == "shard" else f"segment {x['first']}..{x.get('last')}") + f" abandoned cause=restart spent={_sp} s ledger restart")
except Exception: pass
miner_start()
state = {"passes": 0, "claimed": 0, "submitted": 0, "paid": 0, "paid_wei": 0, "shards_accepted": 0, "shards_refused": 0, "segment_refused": 0, "held": 0,
"last_segment_s": 0, "segments": [], "started": stamp(), "label": LABEL, "wallet": WALLET, "key": kh, "task_mode": MODE,
"outcomes": {"paid": 0, "active": 0, "expired": 0, "cancelled": 0, "abandoned": 0}, "wasted_s": {}, "deadline_misses": 0,
"backpressure": {"holds": 0, "held_s": 0.0, "by_cause": {}}, "lease_refused": 0,
"preflight": {"host_sha256": PF_FILES["igneum-prove-host"]["sha256"], "export_sha256": PF_FILES["igneum-prove-export"]["sha256"], "shard_program_id": PF_IDS[0], "aggregator_id": PF_IDS[1], "at": stamp()}}
OUTCOMES = ("paid", "expired", "cancelled", "abandoned")
def seg_row(rid): return next((x for x in state["segments"] if x.get("id", x["first"]) == rid), None)
def ledger_tail():
o = state["outcomes"]; return f"ledger paid={o['paid']} active={o['active']} expired={o['expired']} cancelled={o['cancelled']} abandoned={o['abandoned']}"
def close(rid, outcome, cause=None, tip=None):
"""The outcome ledger's close (docstring point 6): one outcome per claimed task, never a second one; abandoned is point 7's fourth close."""
x = seg_row(rid)
if x is None or x.get("outcome") in OUTCOMES: return
spent = round(float(x.get("export_s", 0)) + float(x.get("cut_s", 0)) + float(x.get("chain_s", 0)), 1)
x["outcome"] = outcome; x["closed_at"] = stamp(); x["spent_s"] = spent
if cause: x["cause"] = cause
if outcome in ("expired", "abandoned") and tip is not None and x.get("deadline") is not None: x["miss_daa"] = max(0, int(tip) - int(x["deadline"]))
if outcome == "expired" and "miss_daa" in x: state["deadline_misses"] += 1
if outcome != "paid":
x["wasted_s"] = spent; k = cause or outcome; state["wasted_s"][k] = round(state["wasted_s"].get(k, 0) + spent, 1)
o = state["outcomes"]; o[outcome] = o.get(outcome, 0) + 1; o["active"] = max(0, state["claimed"] - o["paid"] - o["expired"] - o["cancelled"] - o["abandoned"])
what = f"block {x['first']} shard {x.get('shard')}" if x.get("kind") == "shard" else f"segment {x['first']}..{x.get('last')}"
say(f"RESULT outcome {stamp()} {what} {outcome}" + (f" cause={cause}" if cause else "") + f" spent={spent} s"
+ (f" miss={x['miss_daa']} DAA" if "miss_daa" in x else "") + f" {ledger_tail()}")
def sweep_abandoned(tip):
"""V6-08 section 6: an active task never submitted or held whose lease (or deadline) passed is abandoned, cause lease."""
for x in state["segments"]:
rid = x.get("id", x["first"])
if x.get("outcome") != "active" or x.get("submitted_at") or rid in held or rid in submitted: continue
limit = x.get("lease_expires") or x.get("deadline")
if limit is not None and tip > int(limit): close(rid, "abandoned", "lease", tip)
def lease(first, last):
"""V6-08 section 1: the claim posted to the node; None = accepted or no answer (the node may predate the method), else the refusal text."""
if not LEASE: return None, None
r = rpc("igneum_claimSegment", [{"first": first, "last": last, "keyHash": kh}])
if r is None or r.get("accepted") is not False: return None, (r or {}).get("expiresDaa")
return str(r.get("reason") or "refused")[:160], r.get("expiresDaa") # refused: the holder's expiry (or the penalty's end) in expiresDaa
refused_until = {} # V6-08: a lease refusal holds the task out of the candidates until the DAA the refusal named
attempted = set(); submitted = {}; held = {} # held: first -> {body file, deadline, last}
last_seg_secs = 0; t_run0 = time.time()
bp_cause = None; bp_since = 0.0; bp_last_tip = None; bp_last_tip_at = None # the backpressure gate's memory
def export_and_cut(d, start_blk, blocks, row):
"""One export of start_blk..blocks[-1] and one fixture per block in `blocks`; the fixture paths, or None with the row's
export_s / cut_s set (the segment path's steps 3 and 4, shared with the shard task)."""
t = time.time(); body = json.dumps({"jsonrpc": "2.0", "id": 1, "method": "igneum_exportSegments", "params": [hex(start_blk), hex(blocks[-1])]})
subprocess.run(["curl", "-s", "-m", "600", "-X", "POST", EVM, "-H", "Content-Type: application/json", "--data-binary", body, "-o", f"{d}/seq.json"])
try: json.dump(json.load(open(f"{d}/seq.json"))["result"], open(f"{d}/export.json", "w"))
except Exception as e: row["export_s"] = round(time.time() - t, 1); return None, f"export FAILED {str(e)[:100]}"
row["export_s"] = round(time.time() - t, 1); os.remove(f"{d}/seq.json")
t = time.time(); fixtures = []
for b in blocks:
rr = subprocess.run([EXPORT, f"{d}/export.json", str(b), f"{d}/block-{b}.json", "--source", f"fleet {LABEL} live devnet, segment-aligned prover"], capture_output=True, text=True, timeout=600)
if rr.returncode != 0: row["cut_s"] = round(time.time() - t, 1); return None, f"cut {b} FAILED: {(rr.stdout + rr.stderr)[-200:]}"
fixtures.append(f"{d}/block-{b}.json")
row["cut_s"] = round(time.time() - t, 1)
return fixtures, None
def export_start(ex):
start_blk = hexi(ex.get("restartNumber") or ex.get("execRestartNumber") or ex.get("startedAt") or 0)
if not start_blk:
import re as _re; m = _re.search(r"chain block (\d+)", str(ex.get("startedFrom", ""))); start_blk = int(m.group(1)) if m else 0
return start_blk or EXPORT_FROM or 0
attempted_shards = set(); submitted_shards = {}
def hold(cause, reading):
"""V6-08 section 2: one pass under a backpressure hold (a RESULT line on each cause change, the counters, no claim)."""
global bp_cause, bp_since
now_s = time.time()
if cause != bp_cause: say(f"RESULT backpressure {stamp()} hold cause={cause} {reading}"); bp_cause = cause; bp_since = now_s; state["backpressure"]["holds"] += 1; state["backpressure"]["by_cause"][cause] = state["backpressure"]["by_cause"].get(cause, 0) + 1
state["backpressure"]["held_s"] = round(state["backpressure"]["held_s"] + 15, 1); save_state(); time.sleep(15)
def shard_pass(st, tip, cause, reading):
"""V6-08 section 3: one shard of one block per task (the smallest unit the chain pays), for a card under the segment line.
The paid state of submitted records is read under a backpressure hold too (free); nothing is claimed under one."""
global last_seg_secs
work = rpc("igneum_getAssignedShards", [[kh], 600]) or []
by_id = {f"{hexi(w.get('number'))}:{int(w.get('shard', 0))}": w for w in work}
for rid in list(submitted_shards):
w = by_id.get(rid); s = submitted_shards[rid]
if w is not None and w.get("paid"):
if str(w["paid"].get("keyHash", "")).lower() == kh.lower():
wei = hexi(w["paid"].get("wei")); state["paid"] += 1; state["paid_wei"] += wei; x = seg_row(rid); x["paid_wei"] = wei; x["paid_at"] = stamp()
say(f"RESULT paid {stamp()} block {s['number']} shard {s['shard']} wei={wei} ({wei/1e18:.4f} IGN) carrier={hexi(w['paid'].get('carrierNumber'))} after {int(time.time()-s['at'])} s"); close(rid, "paid")
else: say(f"RESULT paid_other {stamp()} block {s['number']} shard {s['shard']} paid to {str(w['paid'].get('keyHash'))[:18]}"); close(rid, "expired", "paid_other", tip)
del submitted_shards[rid]; drop_export(s["number"], "shard record settled")
elif tip > s["deadline"] + 50:
say(f"RESULT unpaid {stamp()} block {s['number']} shard {s['shard']} past its deadline unpaid"); close(rid, "expired", "unpaid", tip); del submitted_shards[rid]; drop_export(s["number"], "shard record expired")
cands = shard_candidates(work, attempted_shards | {k for k, v in refused_until.items() if isinstance(k, str) and v >= tip}, kh); save_state()
if cause:
hold(cause, reading); return
if not cands:
if state["passes"] % 4 == 1: say(f"RESULT pass {state['passes']} {stamp()} no shard to prove (worklist {len(work)} entries, tip {tip}, mhs {miner_rate()}); waiting")
time.sleep(15); return
c = cands[0]; rid = c["id"]
refused, expires = lease(c["number"], c["number"])
if refused: state["lease_refused"] += 1; refused_until[rid] = hexi(expires) if expires else tip + 60; say(f"RESULT lease_refused {stamp()} block {c['number']} shard {c['shard']}: {refused}; back at DAA {refused_until[rid]}"); return
attempted_shards.add(rid)
deadline = c["daa"] + hexi(st.get("recordWindow") or 600)
state["claimed"] += 1
say(f"RESULT claim_shard {stamp()} block {c['number']} shard {c['shard']} pgas {c['pgas']} deadline {deadline} tip {tip} assigned={c['assigned']} candidates={len(cands)}")
row = {"id": rid, "kind": "shard", "first": c["number"], "last": c["number"], "shard": c["shard"], "claimed_at": stamp(), "shards": 1, "fresh": True, "deadline": deadline, "margin_daa": deadline - tip, "outcome": "active", "lease_expires": hexi(expires) if expires else None}
state["segments"].append(row); state["outcomes"]["active"] = max(0, state["claimed"] - sum(state["outcomes"][k] for k in OUTCOMES))
prune_exports(); free = disk_free_pct()
if free < DISK_MIN_FREE_PCT: say(f"RESULT skip {stamp()} block {c['number']}: disk {free:.1f}% free is under the {DISK_MIN_FREE_PCT:.0f}% floor, no export"); close(rid, "cancelled", "disk"); time.sleep(60); return
d = f"{OUT}/segs/seg-{c['number']}"; os.makedirs(d, exist_ok=True); t0 = time.time()
fixtures, err = export_and_cut(d, export_start(rpc("igneum_getExecStatus", []) or {}), [c["number"]], row)
if err: say(f"RESULT seg {c['number']} {err}"); close(rid, "cancelled", "export" if err.startswith("export") else "cut"); return
if MINER == "pause": miner_stop()
kill_server()
env = dict(os.environ, HOME=f"{FLOOR}/home", SP1_PROVER="cuda", RUST_LOG="off")
if THRESHOLD: env["SP1_GPU_ELEMENT_THRESHOLD"] = THRESHOLD
args = [HOST, fixtures[0], "--mode", "compressed", "--shard", str(c["shard"]), "--prover", WALLET, "--out", f"{d}/shard-results.json"]
t = time.time()
try: rr = subprocess.run(args, env=env, capture_output=True, text=True, timeout=1800)
except subprocess.TimeoutExpired: rr = None
row["chain_s"] = round(time.time() - t, 1); kill_server()
if MINER == "pause": miner_start()
open(f"{d}/chain.log", "w").write((rr.stdout if rr else "") + "\n" + (rr.stderr if rr else "TIMEOUT"))
proof_file = f"{d}/block-{c['number']}-shard-{c['shard']}-compressed.bin"
if not rr or rr.returncode != 0 or not os.path.exists(f"{d}/shard-results.json") or not os.path.exists(proof_file):
say(f"RESULT seg {c['number']} chain FAILED {stamp()} rc={rr.returncode if rr else 'timeout'} wall={row['chain_s']} s: {((rr.stderr if rr else '') or '')[-200:].strip()}"); close(rid, "cancelled", "chain" if rr else "timeout"); return
res = json.load(open(f"{d}/shard-results.json")); statement = res.get("statement")
proof_sha = "0x" + hashlib.sha256(open(proof_file, "rb").read()).hexdigest()
say(f"RESULT seg {c['number']} chain {stamp()} 1 shard record, proof {os.path.getsize(proof_file)} bytes, shards {float(res.get('compressed_prove_seconds', 0)):.1f} s, aggregation 0.0 s, wall {row['chain_s']} s")
sg = subprocess.run([f"{B}/igneum-miner", "sign-record", LABEL, CHAIN, c["hash"], str(c["number"]), str(c["shard"]), WALLET, statement, proof_sha], capture_output=True, text=True).stdout.strip().split("\n")[-1]
try: record = json.loads(sg).get("record")
except Exception: record = None
if not record: say(f"RESULT seg {c['number']} shard {c['number']}/{c['shard']} sign FAILED: {sg[:120]}"); state["shards_refused"] += 1; close(rid, "cancelled", "sign"); return
reply = submit("igneum_submitProofRecord", record, proof_file)
e2e = round(time.time() - t0, 1); row["end_to_end_s"] = e2e
if reply and reply.get("accepted"):
state["shards_accepted"] += 1; state["submitted"] += 1; last_seg_secs = e2e; state["last_segment_s"] = e2e; row["submitted_at"] = stamp(); row["shards_accepted"] = 1
submitted_shards[rid] = {"number": c["number"], "shard": c["shard"], "at": time.time(), "deadline": deadline}
say(f"RESULT submitted {stamp()} block {c['number']} shard {c['shard']} record accepted (new={reply.get('new')}), shard part {c['wei']/1e18:.4f} IGN, end to end {e2e} s, mhs {miner_rate()}")
else:
reason = str((reply or {}).get("reason", reply))[:200]; state["shards_refused"] += 1; row["refused"] = reason
say(f"RESULT seg {c['number']} shard {c['number']}/{c['shard']} refused: {reason}; end to end {e2e} s"); close(rid, "cancelled", "refused")
for f in (fixtures[0], f"{d}/export.json"):
try: os.remove(f)
except OSError: pass
save_state()
def save_state():
tmp = f"{OUT}/prover-state.json.{os.getpid()}.tmp" # one tmp per process: two instances racing on one name lost the file (21:4xZ)
json.dump(state, open(tmp, "w"), indent=1); os.replace(tmp, f"{OUT}/prover-state.json")
# one prover per box: a pid file under OUT; a second instance exits at once instead of killing the first's GPU server
_pidf = f"{OUT}/prover.pid"
try:
_old = int(open(_pidf).read().strip()); os.kill(_old, 0); print(f"RESULT refused {stamp()} another box-prover.py runs as pid {_old}; exiting", flush=True); sys.exit(3)
except (FileNotFoundError, ValueError, ProcessLookupError): pass
open(_pidf, "w").write(str(os.getpid()))
def submit(method, record, proof_file):
proof = "0x" + binascii.hexlify(open(proof_file, "rb").read()).decode()
return rpc(method, [{"record": record, "proof": proof}], timeout=180)
def candidates(st):
start = hexi(st["v1"]["start"]); n = max(1, hexi(st["v1"]["segmentBlocks"])); unproven = hexi(st["v1"]["unprovenDaa"]); tip = hexi(st["tipDaa"])
work = rpc("igneum_getAssignedShards", [[kh], 600]) or []
by = {}
for w in work:
num = hexi(w.get("number"))
if num < start: continue
by.setdefault(num, []).append(w)
need = max(240, int(last_seg_secs * 1.5) + 1)
segs = []; seen = set()
for num in sorted(by):
k = (num - start) // n; first = start + k * n; last = first + n - 1
if first in seen or first in attempted or refused_until.get(first, -1) >= tip: continue
seen.add(first)
whole = True; shards = []; last_daa = 0
for b in range(first, last + 1):
if b not in by: whole = False; break
es = {int(e.get("shard", 0)): e for e in by[b]}
for si in sorted(es):
e = es[si]
if not e.get("open") or e.get("paid") is not None or (e.get("pool") and len(e["pool"]) > 0) or e.get("leased"): whole = False; break
shards.append({"number": b, "hash": e.get("hash"), "shard": si})
if b == last: last_daa = hexi(e.get("daaScore"))
if not whole: break
if whole and shards and MAX_SHARDS and len(shards) > MAX_SHARDS: say(f"RESULT skip {stamp()} segment {first}: {len(shards)} shards over MAX_SHARDS {MAX_SHARDS}"); attempted.add(first); continue
if whole and shards:
deadline = last_daa + unproven
if deadline >= tip + 1 + need: segs.append({"first": first, "last": last, "last_daa": last_daa, "deadline": deadline, "shards": shards, "margin": deadline - tip - 1})
segs.sort(key=lambda s: fnv1a(f"{s['first']}:{kh}"))
return segs, start, n, tip, len(work)
while (time.time() - t_run0) / 3600 < RUN_HOURS:
state["passes"] += 1; p = state["passes"]
if MINER != "none" and MPROC and MPROC.poll() is not None and MINER == "keep": say(f"RESULT miner_exit {stamp()} rc={MPROC.returncode}; restarting"); MPROC = None; miner_start()
st = rpc("igneum_getProvingStatus", [])
ex = rpc("igneum_getExecStatus", [])
# V6-08 section 2: the backpressure gate, read before anything is claimed
now_s = time.time()
if st and st.get("tipDaa") is not None and hexi(st["tipDaa"]) != bp_last_tip: bp_last_tip = hexi(st["tipDaa"]); bp_last_tip_at = now_s
cause, reading = chain_health(ex, st, now_s, bp_last_tip, bp_last_tip_at)
if cause == "unreachable": hold(cause, reading); continue
if not cause and bp_cause: say(f"RESULT backpressure {stamp()} released after {int(now_s - bp_since)} s (was {bp_cause})"); bp_cause = None
if not st or not st.get("v1") or not st["v1"].get("active"): say(f"RESULT pass {p} {stamp()} v1 not active or no status"); time.sleep(20); continue
tip = hexi(st["tipDaa"])
sweep_abandoned(tip)
if MODE == "shard":
shard_pass(st, tip, cause, reading); continue
# paid state
for first in list(submitted):
rec = rpc("igneum_getSegmentRecords", [hex(first)])
if rec and rec.get("paid"):
wei = hexi(rec["paid"]["wei"]); state["paid"] += 1; state["paid_wei"] += wei; drop_export(first, "record paid")
say(f"RESULT paid {stamp()} segment {first}..{submitted[first]['last']} wei={wei} ({wei/1e18:.4f} IGN) carrier={hexi(rec['paid'].get('carrierNumber'))} after {int(time.time()-submitted[first]['at'])} s")
for s in state["segments"]:
if s["first"] == first: s["paid_wei"] = wei; s["paid_at"] = stamp()
close(first, "paid"); del submitted[first]
elif rec is not None and tip > submitted[first]["deadline"] + 50:
say(f"RESULT unpaid {stamp()} segment {first} past its deadline unpaid; carried={len(rec.get('carried') or [])} pool={len(rec.get('pool') or [])}"); close(first, "expired", "unpaid", tip); del submitted[first]
# held fresh records offered again
for first in list(held):
h = held[first]
if tip > h["deadline"]: say(f"RESULT held_expired {stamp()} segment {first} deadline passed"); close(first, "expired", "held_expired", tip); del held[first]; continue
rr = submit("igneum_submitSegmentRecord", h["record"], h["proof_file"])
if rr and rr.get("accepted"):
state["submitted"] += 1; submitted[first] = {"last": h["last"], "at": time.time(), "deadline": h["deadline"]}; drop_export(first, "record accepted"); say(f"RESULT submitted {stamp()} segment {first}..{h['last']} record accepted on retry (held {int(time.time()-h['since'])} s)"); del held[first]
state["held"] = len(held)
cands, start, n, tip, entries = candidates(st)
save_state()
if cause: hold(cause, reading); continue # V6-08 section 2: the paid state and the held records were read; nothing is claimed
if not cands:
if p % 4 == 1: say(f"RESULT pass {p} {stamp()} no whole segment inside the margin (worklist {entries} entries, tip {tip}, mhs {miner_rate()}); waiting")
time.sleep(15); continue
picked = None; prev_file = None; expected = ""
for c in cands[:3]:
stmt = rpc("igneum_getSegmentStatement", [hex(c["first"])])
if not stmt or not stmt.get("executed") or (stmt.get("status") or {}).get("status") != "pending":
attempted.add(c["first"]); say(f"RESULT skip {stamp()} segment {c['first']}: executed={stmt and stmt.get('executed')} status={(stmt or {}).get('status')}"); continue
if stmt.get("previous") is None:
if not st["v1"].get("freshRuleActive") and c["first"] >= start + n:
pr = rpc("igneum_getSegmentRecords", [hex(c["first"] - n)])
waiting = any(e.get("verified") and e.get("includedIn") is None for e in (pr or {}).get("pool") or [])
if waiting or (pr and pr.get("paid")): say(f"RESULT skip {stamp()} segment {c['first']}: previous has a record waiting or paid, a fresh chain would be refused (fresh rule off)"); continue
picked = c; expected = stmt.get("publicValuesFresh") or ""; break
if not stmt["previous"].get("proofInPool"): say(f"RESULT skip {stamp()} segment {c['first']}: previous paid, its proof not in this pool"); continue
got = rpc("igneum_getSegmentProofBytes", [stmt["previous"]["first"], stmt["previous"]["keyHash"]])
if not got or not got.get("proof"): continue
prev_file = f"{OUT}/segs/prev-{c['first']}.bin"; h = got["proof"]; open(prev_file, "wb").write(binascii.unhexlify(h[2:] if h.startswith("0x") else h))
picked = c; expected = stmt.get("publicValuesContinuing") or ""; break
if not picked: say(f"RESULT pass {p} {stamp()} {len(cands)} candidates, none usable; waiting"); time.sleep(15); continue
first, last = picked["first"], picked["last"]
# V6-08 section 1: the lease, before the ledger row and before any export; a refusal is not an eligible job and holds the
# segment out of the candidates until the DAA the refusal named (the holder's expiry), one pass lost, nothing exported
refused, expires = lease(first, last)
if refused: state["lease_refused"] += 1; refused_until[first] = hexi(expires) if expires else tip + 60; say(f"RESULT lease_refused {stamp()} segment {first}..{last}: {refused}; back at DAA {refused_until[first]}"); continue
attempted.add(first); state["claimed"] += 1
say(f"RESULT claim {stamp()} segment {first}..{last} ({len(picked['shards'])} shards, {'continuing' if prev_file else 'fresh'}) margin={picked['margin']} tip={tip} candidates={len(cands)} rank_by=fnv")
seg = {"id": first, "kind": "segment", "first": first, "last": last, "claimed_at": stamp(), "shards": len(picked["shards"]), "fresh": prev_file is None,
"deadline": picked["deadline"], "margin_daa": picked["margin"], "outcome": "active", "lease_expires": hexi(expires) if expires else None}; state["segments"].append(seg)
state["outcomes"]["active"] = max(0, state["claimed"] - sum(state["outcomes"][k] for k in OUTCOMES))
prune_exports(); free = disk_free_pct()
if free < DISK_MIN_FREE_PCT: say(f"RESULT skip {stamp()} segment {first}: disk {free:.1f}% free is under the {DISK_MIN_FREE_PCT:.0f}% floor, no export"); close(first, "cancelled", "disk"); time.sleep(60); continue
d = f"{OUT}/segs/seg-{first}"; os.makedirs(d, exist_ok=True); t_seg0 = time.time()
# export
# a node whose EVM restarted at a chain block (0.3.13's exec restart rule) exports from that block, not genesis: the
# blocks below it are not executed and the exporter refuses their zero state roots (6 October 2026, 16:02Z)
ex = rpc("igneum_getExecStatus", []) or {}
start_blk = hexi(ex.get("restartNumber") or ex.get("execRestartNumber") or ex.get("startedAt") or 0)
if not start_blk:
sf = str(ex.get("startedFrom", ""))
import re as _re; m = _re.search(r"chain block (\d+)", sf); start_blk = int(m.group(1)) if m else 0
if not start_blk and EXPORT_FROM: start_blk = EXPORT_FROM
t = time.time(); body = json.dumps({"jsonrpc": "2.0", "id": 1, "method": "igneum_exportSegments", "params": [hex(start_blk), hex(last)]})
seg["export_from"] = start_blk
r = subprocess.run(["curl", "-s", "-m", "600", "-X", "POST", EVM, "-H", "Content-Type: application/json", "--data-binary", body, "-o", f"{d}/seq.json"])
try: json.dump(json.load(open(f"{d}/seq.json"))["result"], open(f"{d}/export.json", "w"))
except Exception as e: say(f"RESULT seg {first} export FAILED {str(e)[:100]}"); seg["export_s"] = round(time.time() - t, 1); close(first, "cancelled", "export"); continue
seg["export_s"] = round(time.time() - t, 1); os.remove(f"{d}/seq.json")
# cut
t = time.time(); fixtures = []
ok = True
for b in range(first, last + 1):
rr = subprocess.run([EXPORT, f"{d}/export.json", str(b), f"{d}/block-{b}.json", "--source", f"fleet {LABEL} live devnet, segment-aligned prover"], capture_output=True, text=True, timeout=600)
if rr.returncode != 0: say(f"RESULT seg {first} cut {b} FAILED: {(rr.stdout + rr.stderr)[-200:]}"); ok = False; break
fixtures.append(f"{d}/block-{b}.json")
if not ok: seg["cut_s"] = round(time.time() - t, 1); close(first, "cancelled", "cut"); continue
seg["cut_s"] = round(time.time() - t, 1)
# chain
if MINER == "pause": miner_stop()
kill_server()
env = dict(os.environ, HOME=f"{FLOOR}/home", SP1_PROVER="cuda", RUST_LOG="off")
if THRESHOLD: env["SP1_GPU_ELEMENT_THRESHOLD"] = THRESHOLD
args = [HOST, "--mode", "chain", "--chain", ",".join(fixtures), "--prover", WALLET, "--save-shards", "--out", f"{d}/chain-results.json"]
if prev_file: args += ["--prev", prev_file]
samp = subprocess.Popen(f"while :; do nvidia-smi -i {DEV} --query-gpu=memory.used,utilization.gpu,power.draw --format=csv,noheader,nounits; sleep 1; done", shell=True, stdout=open(f"{d}/smi.csv", "w"), stderr=subprocess.DEVNULL, start_new_session=True)
t = time.time()
try: rr = subprocess.run(args, env=env, capture_output=True, text=True, timeout=3600)
except subprocess.TimeoutExpired: rr = None
seg["chain_s"] = round(time.time() - t, 1); os.killpg(samp.pid, signal.SIGKILL); kill_server()
if MINER == "pause": miner_start()
try: peak = max(float(l.split(",")[0]) for l in open(f"{d}/smi.csv") if l.strip())
except Exception: peak = 0
seg["peak_mib"] = peak
open(f"{d}/chain.log", "w").write((rr.stdout if rr else "") + "\n" + (rr.stderr if rr else "TIMEOUT"))
if not rr or rr.returncode != 0 or not os.path.exists(f"{d}/chain-results.json"):
say(f"RESULT seg {first} chain FAILED {stamp()} rc={rr.returncode if rr else 'timeout'} wall={seg['chain_s']} s: {((rr.stderr if rr else '') or '')[-200:].strip()}"); seg["failed"] = "chain"; close(first, "cancelled", "chain" if rr else "timeout"); continue
res = json.load(open(f"{d}/chain-results.json"))
recs = [s for blk in res.get("blocks", []) for s in blk.get("shard_records", [])]
say(f"RESULT seg {first} chain {stamp()} {len(recs)} shard records, chain_len {res.get('segment_chain_len')}, proof {res.get('segment_proof_bytes')} bytes, shards {res.get('shard_prove_seconds_total', 0):.1f} s, aggregation {res.get('aggregate_prove_seconds_total', 0):.1f} s, wall {seg['chain_s']} s, peak {peak:.0f} MiB")
# shard records
ok_shards = 0
for rcd in recs:
sg = subprocess.run([f"{B}/igneum-miner", "sign-record", LABEL, CHAIN, rcd["block_hash"], str(rcd["number"]), str(rcd["shard"]), WALLET, rcd["statement"], rcd["proof_sha256"]], capture_output=True, text=True).stdout.strip().split("\n")[-1]
try: record = json.loads(sg).get("record")
except Exception: record = None
if not record: say(f"RESULT seg {first} shard {rcd['number']}/{rcd['shard']} sign FAILED: {sg[:120]}"); state["shards_refused"] += 1; continue
reply = submit("igneum_submitProofRecord", record, rcd["proof_file"])
if reply and reply.get("accepted"): ok_shards += 1; state["shards_accepted"] += 1
else: state["shards_refused"] += 1; say(f"RESULT seg {first} shard {rcd['number']}/{rcd['shard']} refused: {(reply or {}).get('reason', reply)}")
seg["shards_accepted"] = ok_shards
say(f"RESULT seg {first} shards {stamp()} accepted {ok_shards} of {len(recs)}")
if ok_shards != len(recs): seg["failed"] = "shards"; close(first, "cancelled", "shards"); continue
pv = res.get("segment_public_values", "")
strip = lambda h: (h[2:] if h.startswith("0x") else h); strip2 = lambda h: (strip(h)[:472] + strip(h)[536:]) if len(strip(h)) == 680 else strip(h)
if strip2(pv) != strip2(expected):
a, b = strip2(pv), strip2(expected); off = next((i for i in range(min(len(a), len(b))) if a[i] != b[i]), min(len(a), len(b)))
say(f"RESULT seg {first} FAILED: statement differs from the node's at hex offset {off} (lengths {len(a)} vs {len(b)}); ours ...{a[max(0,off-8):off+56]} node ...{b[max(0,off-8):off+56]}"); seg["failed"] = "statement"; close(first, "cancelled", "statement"); continue
last_hash = next(s["hash"] for s in picked["shards"] if s["number"] == last)
sg = subprocess.run([f"{B}/igneum-miner", "sign-segment-record", LABEL, CHAIN, str(first), str(last), last_hash, WALLET, pv, res["segment_proof_sha256"]], capture_output=True, text=True).stdout.strip().split("\n")[-1]
try: record = json.loads(sg).get("record")
except Exception: record = None
if not record: say(f"RESULT seg {first} segment sign FAILED: {sg[:120]}"); seg["failed"] = "sign"; close(first, "cancelled", "sign"); continue
reply = submit("igneum_submitSegmentRecord", record, res["segment_proof_file"])
seg_s = round(time.time() - t_seg0, 1); seg["end_to_end_s"] = seg_s
if reply and reply.get("accepted"):
state["submitted"] += 1; last_seg_secs = seg_s; state["last_segment_s"] = seg_s
stmt2 = rpc("igneum_getSegmentStatement", [hex(first)]) or {}
submitted[first] = {"last": last, "at": time.time(), "deadline": picked["deadline"]}; seg["submitted_at"] = stamp()
say(f"RESULT submitted {stamp()} segment {first}..{last} record accepted (new={reply.get('new')}, chain_len {res.get('segment_chain_len')}), aggregator share {hexi(stmt2.get('aggregatorWei'))/1e18:.4f} IGN, end to end {seg_s} s, mhs {miner_rate()}")
else:
reason = str((reply or {}).get("reason", reply))[:200]; state["segment_refused"] += 1; seg["refused"] = reason
say(f"RESULT segment_refused {stamp()} segment {first}..{last}: {reason}; end to end {seg_s} s")
if "pending until" in reason or "does not chain" in reason:
held[first] = {"record": record, "proof_file": res["segment_proof_file"], "last": last, "deadline": picked["deadline"], "since": time.time()}; say(f"RESULT held {stamp()} segment {first} held for retry until DAA {picked['deadline']}")
else: close(first, "cancelled", "refused")
for b in range(first, last + 1):
try: os.remove(f"{d}/block-{b}.json")
except OSError: pass
try: os.remove(f"{d}/export.json")
except OSError: pass
save_state()
miner_stop(); kill_server(); save_state()
_o = state["outcomes"]; say(f"RESULT ledger {stamp()} paid={_o['paid']} active={_o['active']} expired={_o['expired']} cancelled={_o['cancelled']} abandoned={_o['abandoned']} deadline_misses={state['deadline_misses']} wasted_s={json.dumps(state['wasted_s'], sort_keys=True)} active_segments={[x.get('id', x['first']) for x in state['segments'] if x.get('outcome') == 'active']} lease_refused={state['lease_refused']} backpressure={json.dumps(state['backpressure'], sort_keys=True)}")
say(f"RESULT {flow_line(state['claimed'], _o)}")
say(f"RESULT summary {stamp()} passes={state['passes']} claimed={state['claimed']} submitted={state['submitted']} paid={state['paid']} paid_wei={state['paid_wei']} shards_accepted={state['shards_accepted']} shards_refused={state['shards_refused']} segment_refused={state['segment_refused']}")
say(f"RESULT prover_done {stamp()}")