294 lines
23 KiB
Python
294 lines
23 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.
|
|
|
|
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).
|
|
"""
|
|
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"))
|
|
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 hexi(v): return int(v, 16) if isinstance(v, str) and v.startswith("0x") else int(v or 0)
|
|
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: subprocess.run("pkill -x sp1-gpu-server; sleep 1; 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
|
|
subprocess.run(f"pkill -f '^/opt/igneum/pkg/bin/igneum-worker-cuda --device {DEV} '", shell=True)
|
|
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
|
|
# 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}")
|
|
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}
|
|
attempted = set(); submitted = {}; held = {} # held: first -> {body file, deadline, last}
|
|
last_seg_secs = 0; t_run0 = time.time()
|
|
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: 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): 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:
|
|
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", [])
|
|
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"])
|
|
# 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()
|
|
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 [])}"); 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"); 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 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"]; 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 = {"first": first, "last": last, "claimed_at": stamp(), "shards": len(picked["shards"]), "fresh": prev_file is None}; state["segments"].append(seg)
|
|
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"); 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]}"); 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: 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"; 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"; 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"; 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"; 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']}")
|
|
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()
|
|
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()}")
|