From 01fa7cd2bc4a1b86cb470fe58954508ee1898693 Mon Sep 17 00:00:00 2001 From: igneum-labs <337424239+igneum-labs@users.noreply.github.com> Date: Thu, 8 Oct 2026 19:55:43 +0000 Subject: [PATCH] V6-08 build, the fleet side: box-prover.py posts each claim as a lease (igneum_claimSegment) and skips a refusal, the chain-health backpressure gate (blocked, re-executing, discontinuous, lagging, stalled, finality paused, unreachable), TASK=segment|shard|auto with one-shard tasks under the segment line and MAX_SHARDS, the preflight (host and export sha256 against the shipped manifest, --mode id against the node's pinned pair, ldd, nvidia-smi, exit 4 naming the file), the fourth outcome abandoned (lease, restart) and the flow line; prover-outcomes.py reports abandoned, the kind column and the conservation identity per kind; both self-tests known-failed first, the gate runs box-prover's Co-Authored-By: Claude Fable 5.1 --- docs/design/proving-task-protection.md | 17 +- tools/ci/checks.txt | 1 + tools/ci/pre-push.sh | 1 + tools/fleet/box-prover.py | 332 +++++++++++++++++++++++-- tools/fleet/prover-outcomes.py | 106 ++++++-- 5 files changed, 409 insertions(+), 48 deletions(-) diff --git a/docs/design/proving-task-protection.md b/docs/design/proving-task-protection.md index 76f8efb76..6ac19c455 100644 --- a/docs/design/proving-task-protection.md +++ b/docs/design/proving-task-protection.md @@ -63,8 +63,10 @@ behaviour). With it set: `leaseExpiresDaa` so the prover's candidate grouping (every shard "open, unpaid, not in the pool") also requires "not leased by another key". 5. The prover posts the claim: `box-prover.py` calls `igneum_claimSegment` before the export and treats `accepted: false` as a - skip (a `RESULT lease_refused` line, no ledger row: a refused claim is not an eligible job). The lease is the prover's reservation - for the time it takes to prove; the segment deadline stays the chain's. + skip (a `RESULT lease_refused` line, no ledger row: a refused claim is not an eligible job; the segment is held out of the + candidates until the DAA the refusal named, the holder's expiry, so the prover moves to the next candidate on the same pass + instead of asking again). The lease is the prover's reservation for the time it takes to prove; the segment deadline stays + the chain's. `LEASE=off` keeps the 8 October behaviour for a node that predates the method. 6. Counters for the chain-side view, in `igneum_getProvingStatus` under `leases`: `open`, `granted`, `refusedHeld` (a holder's lease refused a second claimant), `refusedDeprioritised`, `landed` (a lease whose record arrived), `abandoned` (expired with no record). Flow: `granted = landed + abandoned + open` at any read (the test asserts it). @@ -116,7 +118,8 @@ and submitted 1,681 segment records into a chain that never paid them, the day's On a hold the prover does not claim, does not export, and prints `RESULT backpressure hold cause= ` once per cause change and `RESULT backpressure released after s` when the gate clears; every pass under a hold counts in `state["backpressure"] = {holds, held_s, by_cause}`. Submitted records keep being polled (the paid state of an accepted record is still -worth reading) and held fresh records keep being offered, because both are free. The node side needs no change: every reading +worth reading) and held fresh records keep being offered, because both are free (`unreachable` is the one cause that skips them: there +is nothing to read). The node side needs no change: every reading is already served. `pause_reason` is reused by name in the design so the app's pause and the prover's hold agree on what "paused" is; adding it to `igneum_getProvingStatus` as `paused` is a one-line follow-up for the node lane, not needed tonight. @@ -150,8 +153,10 @@ mode a task is one shard. Shard mode's pass: the candidate shards are the workli leased by another key, assigned-first then open, newest first and smallest `pgas` among equals (the app's `choose`); one claim (`igneum_claimSegment` over the one block, the lease of section 1), one export of that block, one cut, one `--mode compressed --shard i --out`, one `sign-record`, one `igneum_submitProofRecord`; the paid state is polled through `igneum_getAssignedShards` (`paid` on the -row). Each shard task is a ledger row of `kind: shard` with the same four outcomes and causes, so `prover-outcomes.py` reports both -kinds in one table with the kind beside each count. The planner is not a knob: `proving_v1_segment_blocks` is a consensus param +row). Each shard task is a ledger row of `kind: shard` with the same four outcomes and causes (one cause added for both kinds: +`expired, paid_other`, a submitted record whose shard or segment another key was paid for, the steal of the pipeline record, so the +conservation line counts it as a departure and the waste report names it), so `prover-outcomes.py` reports both kinds in one table +with the kind beside each count. The planner is not a knob: `proving_v1_segment_blocks` is a consensus param (section 0). **Tests, known-failed first.** `box-prover.py --self-test`: `task_mode(total_mib, env)` picks shard under the line and segment over @@ -265,7 +270,7 @@ counters sum to the claims after a mixed sequence. | # | Branch | Change | Test | Clock (UK) | |---|---|---|---|---| -| 1 | `v608-proving-tasks-node` (node fork, off `successor-2.0.1`) | `LeaseRules`, `ClaimOutcome`, `claim_segment` and `open_claims` with leases, the abandon ring, the counters; `igneum_claimSegment`, `igneum_getAssignedShards`, `igneum_getProvingStatus` fields; `MAX_PROOF_BYTES` in consensus-core, the pool and RPC bounds | the five pool tests and the RPC tests above, `cargo test -p igneum-exec` on build-2 or build-3 through `tools/build-remote.sh` | 00:30 | +| 1 | `v608-proving-tasks-node` (node fork, off `successor-2.0.1`, the tip merged at 2b1a247a, 20:55 UK) | `LeaseRules`, `ClaimOutcome`, `claim_segment` and `open_claims` with leases, the abandon ring, the counters; `igneum_claimSegment`, `igneum_getAssignedShards`, `igneum_getProvingStatus` fields; `MAX_PROOF_BYTES` in consensus-core, the pool and RPC bounds | the five pool tests and the RPC tests above, `cargo test -p igneum-exec` on build-2 or build-3 through `tools/build-remote.sh` | 00:30 | | 2 | `v608-proving-tasks` (igneum) | `box-prover.py`: the lease call, the backpressure gate, `TASK`/`MAX_SHARDS` and shard mode, `preflight()`, the `abandoned` outcome, the conservation line, `--self-test`; `prover-outcomes.py`: `abandoned`, the conservation line, the kind column; `tools/ci/pre-push.sh` runs `box-prover.py --self-test` beside the ledger's | the python self-tests | 01:30 | | 3 | both | the node tip merged once more; shas to the node lane and the coordinator | | 02:00 | diff --git a/tools/ci/checks.txt b/tools/ci/checks.txt index fe3169d84..d759f5b18 100644 --- a/tools/ci/checks.txt +++ b/tools/ci/checks.txt @@ -75,6 +75,7 @@ the test map: every automated case of the registry maps to a cell or carries a N the harness map page is generated from tools/ci/test-map.json and current P01 part A, the million-vector driver: a clean run is PASS, one wrong hash or one unanswered nonce is FAIL naming it (self-test, a fake worker) the proving outcome ledger (review B F08): every claimed job ends in one outcome; the report's self-test reads a log and a state file to known numbers +the fleet prover's task protection (V6-08): the backpressure gate, task sizing, shard ordering, preflight verdicts and the flow identity on known-failed-first cases the registry's evidence rules: a PASS names evidence that exists, a touched evidence file moves with its row, stale evidence never reads PASS, a run_status needs the approval (self-test) the public ledger (docs/ledger-public.md) is what docs/fud-ledger.md generates: one row per item, no commit ids, times or team names (self-test first) the ledger page reads both entry heading forms (M1 and AP-F8-1) so no in-house pass row is dropped from /ledger (known-failed first) diff --git a/tools/ci/pre-push.sh b/tools/ci/pre-push.sh index 47f4ad888..3ce01de41 100755 --- a/tools/ci/pre-push.sh +++ b/tools/ci/pre-push.sh @@ -174,6 +174,7 @@ tree_checks() { run "the harness map page is generated from tools/ci/test-map.json and current" node tools/ci/test-map-doc.mjs --check run "P01 part A, the million-vector driver: a clean run is PASS, one wrong hash or one unanswered nonce is FAIL naming it (self-test, a fake worker)" python3 tools/ci/p01-vectors.py --self-test run "the proving outcome ledger (review B F08): every claimed job ends in one outcome; the report's self-test reads a log and a state file to known numbers" python3 tools/fleet/prover-outcomes.py --self-test + run "the fleet prover's task protection (V6-08): the backpressure gate, task sizing, shard ordering, preflight verdicts and the flow identity on known-failed-first cases" python3 tools/fleet/box-prover.py --self-test run "the registry's evidence rules: a PASS names evidence that exists, a touched evidence file moves with its row, stale evidence never reads PASS, a run_status needs the approval (self-test)" bash tools/ci/registry-evidence-check.sh --self-test run "the public ledger (docs/ledger-public.md) is what docs/fud-ledger.md generates: one row per item, no commit ids, times or team names (self-test first)" bash -c 'node tools/ledger/export-public.mjs --self-test && node tools/ledger/export-public.mjs --check' run "the ledger page reads both entry heading forms (M1 and AP-F8-1) so no in-house pass row is dropped from /ledger (known-failed first)" node tools/ledger-page.mjs --self-test diff --git a/tools/fleet/box-prover.py b/tools/fleet/box-prover.py index 521a1afc9..9603d1db1 100644 --- a/tools/fleet/box-prover.py +++ b/tools/fleet/box-prover.py @@ -24,9 +24,26 @@ proving-v1, 272b025) and tools/proving-v1/pc2-segments.ps1 around the four binar 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). +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/ (FLEET_CARD) @@ -40,6 +57,100 @@ EXPORT_FROM = int(os.environ.get("EXPORT_FROM", "27276")) # 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))) @@ -64,7 +175,6 @@ 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: @@ -100,6 +210,38 @@ def miner_rate(n=6): 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 @@ -111,27 +253,159 @@ for _ in range(120): 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, - "outcomes": {"paid": 0, "active": 0, "expired": 0, "cancelled": 0}, "wasted_s": {}, "deadline_misses": 0} -OUTCOMES = ("paid", "expired", "cancelled") -def seg_row(first): return next((x for x in state["segments"] if x["first"] == first), None) -def close(first, outcome, cause=None, tip=None): - """The outcome ledger's close (docstring point 6): one outcome per claimed segment, never a second one.""" - x = seg_row(first) + "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 == "expired" and tip is not None and x.get("deadline") is not None: x["miss_daa"] = max(0, int(tip) - int(x["deadline"])); state["deadline_misses"] += 1 + 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"]) - say(f"RESULT outcome {stamp()} segment {first}..{x.get('last')} {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 paid={o['paid']} active={o['active']} expired={o['expired']} cancelled={o['cancelled']}") + 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") @@ -156,7 +430,7 @@ def candidates(st): 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 + 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): @@ -164,10 +438,11 @@ def candidates(st): 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 + 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}) @@ -177,8 +452,18 @@ 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)]) @@ -200,6 +485,7 @@ while (time.time() - t_run0) / 3600 < RUN_HOURS: 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 @@ -220,11 +506,16 @@ while (time.time() - t_run0) / 3600 < RUN_HOURS: 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 + 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 = {"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"}; state["segments"].append(seg) - state["outcomes"]["active"] = max(0, state["claimed"] - state["outcomes"]["paid"] - state["outcomes"]["expired"] - state["outcomes"]["cancelled"]) + 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() @@ -317,6 +608,7 @@ while (time.time() - t_run0) / 3600 < RUN_HOURS: 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']} deadline_misses={state['deadline_misses']} wasted_s={json.dumps(state['wasted_s'], sort_keys=True)} active_segments={[x['first'] for x in state['segments'] if x.get('outcome') == 'active']}") +_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()}") diff --git a/tools/fleet/prover-outcomes.py b/tools/fleet/prover-outcomes.py index 3dea8aaf6..582f2a192 100755 --- a/tools/fleet/prover-outcomes.py +++ b/tools/fleet/prover-outcomes.py @@ -9,14 +9,22 @@ three questions: paid completions, missed deadlines and wasted work by cause, pl Outcomes: paid (a carrying block paid the segment record), active (claimed and not yet closed: in work, submitted and waiting, or held for a retry), expired (the deadline passed with no paid record: cause unpaid after a submit, held_expired -for a held record, or never_submitted when the log ends past the deadline with no submit), cancelled (the prover gave the -job up: cause disk, export, cut, chain, timeout, shards, statement, sign, refused). Wasted work is the export, cut and -chain seconds of every job that was not paid, by cause. A job row from a state file wins over the same segment in a log. +for a held record), cancelled (the prover gave the job up: cause disk, export, cut, chain, timeout, shards, statement, sign, +refused), abandoned (V6-08 section 6, docs/design/proving-task-protection.md: claimed and never submitted, the lease or the +deadline passed: cause lease, restart, or never_submitted when a log ends past the deadline with no submit). Wasted work is +the export, cut and chain seconds of every job that was not paid, by cause. A job row from a state file wins over the same +segment in a log. A row carries its kind (segment, or shard for a one-shard task of box-prover.py's TASK=shard). + +Flow conservation (the master review's measurement correction, V6-08): the report prints and checks + claims N = paid + expired + cancelled + abandoned + active +and, over the span, opening backlog + arrivals - paid - expired - cancelled - abandoned = closing backlog, per kind and in all. """ import sys, os, json, re, datetime, statistics CAUSES_CANCELLED = ("disk", "export", "cut", "chain", "timeout", "shards", "statement", "sign", "refused") -CAUSES_EXPIRED = ("unpaid", "held_expired", "never_submitted") +CAUSES_EXPIRED = ("unpaid", "held_expired", "paid_other") +CAUSES_ABANDONED = ("lease", "restart", "never_submitted") +OUTCOMES = ("paid", "active", "expired", "cancelled", "abandoned") def ts(s): try: return datetime.datetime.strptime(s, "%Y-%m-%dT%H:%M:%SZ").replace(tzinfo=datetime.timezone.utc).timestamp() @@ -30,9 +38,9 @@ def jobs_from_state(state, label=None): "margin_daa": x.get("margin_daa"), "shards": x.get("shards"), "fresh": x.get("fresh"), "spent_s": x.get("spent_s", round(float(x.get("export_s", 0)) + float(x.get("cut_s", 0)) + float(x.get("chain_s", 0)), 1)), "chain_s": x.get("chain_s"), "end_to_end_s": x.get("end_to_end_s"), "submitted_at": x.get("submitted_at"), "paid_at": x.get("paid_at"), - "paid_wei": x.get("paid_wei"), "closed_at": x.get("closed_at"), "miss_daa": x.get("miss_daa"), "source": "state"} + "paid_wei": x.get("paid_wei"), "closed_at": x.get("closed_at"), "miss_daa": x.get("miss_daa"), "source": "state", "kind": x.get("kind") or "segment"} o = x.get("outcome") - if o in ("paid", "expired", "cancelled", "active"): j["outcome"] = o; j["cause"] = x.get("cause") + if o in OUTCOMES: j["outcome"] = o; j["cause"] = x.get("cause") elif x.get("paid_wei") is not None: j["outcome"] = "paid"; j["cause"] = None elif x.get("failed"): j["outcome"] = "cancelled"; j["cause"] = x["failed"] elif x.get("refused") and not x.get("submitted_at"): j["outcome"] = "cancelled"; j["cause"] = "refused" @@ -64,12 +72,22 @@ def jobs_from_log(text, label=None): if not mm: continue first = int(mm.group(1)); j = {"label": label, "first": first, "last": int(mm.group(2)), "claimed_at": at, "shards": int(mm.group(3)), "fresh": mm.group(4) == "fresh", "margin_daa": int(mm.group(5)), "deadline": int(mm.group(6)) + 1 + int(mm.group(5)), "spent_s": 0.0, "chain_s": None, "end_to_end_s": None, + "submitted_at": None, "paid_at": None, "paid_wei": None, "closed_at": None, "miss_daa": None, "outcome": "active", "cause": None, "source": "log", "kind": "segment"} + jobs[first] = j; order.append(first); continue + if kind == "claim_shard": + # V6-08 section 3: a one-shard task, `RESULT claim_shard block N shard S pgas P deadline D tip T` + mm = re.match(r"block (\d+) shard (\d+) pgas (\d+) deadline (\d+) tip (\d+)", rest) + if not mm: continue + first = int(mm.group(1)); j = {"label": label, "first": first, "last": first, "shard": int(mm.group(2)), "claimed_at": at, "shards": 1, "fresh": True, "kind": "shard", + "margin_daa": int(mm.group(4)) - int(mm.group(5)), "deadline": int(mm.group(4)), "spent_s": 0.0, "chain_s": None, "end_to_end_s": None, "submitted_at": None, "paid_at": None, "paid_wei": None, "closed_at": None, "miss_daa": None, "outcome": "active", "cause": None, "source": "log"} jobs[first] = j; order.append(first); continue if kind == "outcome": - mm = re.match(r"segment (\d+)\.\.(\d+) (\w+)(?: cause=(\w+))? spent=([\d.]+) s(?: miss=(\d+) DAA)?", rest) - if not mm or int(mm.group(1)) not in jobs: continue - j = jobs[int(mm.group(1))]; j.update(outcome=mm.group(3), cause=mm.group(4), spent_s=float(mm.group(5)), closed_at=at, miss_daa=int(mm.group(6)) if mm.group(6) else None); continue + mm = re.match(r"(?:segment (\d+)\.\.\d+|block (\d+) shard \d+) (\w+)(?: cause=(\w+))? spent=([\d.]+) s(?: miss=(\d+) DAA)?", rest) + if not mm: continue + first = int(mm.group(1) or mm.group(2)) + if first not in jobs or mm.group(3) not in OUTCOMES: continue + j = jobs[first]; j.update(outcome=mm.group(3), cause=mm.group(4), spent_s=float(mm.group(5)), closed_at=at, miss_daa=int(mm.group(6)) if mm.group(6) else None); continue if kind == "seg": first = int(a) if a.isdigit() else None if first not in jobs: continue @@ -109,7 +127,7 @@ def jobs_from_log(text, label=None): if j["outcome"] == "active" and j.pop("_refused_at", None): j.update(outcome="cancelled", cause="refused", closed_at=j.get("closed_at") or last_at) # a job never submitted whose deadline the chain passed while the log went on: expired, never_submitted if j["outcome"] == "active" and j["submitted_at"] is None and last_tip is not None and j.get("deadline") and last_tip > j["deadline"] + 50: - j.update(outcome="expired", cause="never_submitted", closed_at=last_at, miss_daa=last_tip - j["deadline"]) + j.update(outcome="abandoned", cause="never_submitted", closed_at=last_at, miss_daa=last_tip - j["deadline"]) j.pop("_refused_at", None) return [jobs[f] for f in order] @@ -117,19 +135,41 @@ def median(xs): xs = [x for x in xs if x is not None] return round(statistics.median(xs), 1) if xs else None +def conservation(jobs): + """V6-08 section 6: the identity claims = paid + expired + cancelled + abandoned + active, per kind and in all, and over + the span opening backlog + arrivals - departures = closing backlog (the backlog is the active count at each end; every + job here arrived inside the span, so the opening backlog is 0 and the closing backlog is the active count).""" + out = {} + for kind in ["all"] + sorted({j.get("kind") or "segment" for j in jobs}): + js = jobs if kind == "all" else [j for j in jobs if (j.get("kind") or "segment") == kind] + c = {k: sum(1 for j in js if j["outcome"] == k) for k in OUTCOMES} + departures = c["paid"] + c["expired"] + c["cancelled"] + c["abandoned"] + out[kind] = {"claims": len(js), **c, "departures": departures, "opening_backlog": 0, "arrivals": len(js), "closing_backlog": c["active"], + "ok": len(js) == departures + c["active"]} + return out + +def conservation_lines(cons): + L = [] + for kind, c in cons.items(): + L.append(f"Flow conservation ({kind}): claims {c['claims']} = paid {c['paid']} + expired {c['expired']} + cancelled {c['cancelled']} + abandoned {c['abandoned']} + active {c['active']}: " + + ("ok" if c["ok"] else f"BROKEN by {c['claims'] - c['departures'] - c['active']}") + + f"; backlog {c['opening_backlog']} + arrivals {c['arrivals']} - departures {c['departures']} = closing backlog {c['closing_backlog']}") + return L + def report(jobs): - tot = {k: sum(1 for j in jobs if j["outcome"] == k) for k in ("paid", "active", "expired", "cancelled")} - paid = [j for j in jobs if j["outcome"] == "paid"]; exp = [j for j in jobs if j["outcome"] == "expired"]; can = [j for j in jobs if j["outcome"] == "cancelled"] + tot = {k: sum(1 for j in jobs if j["outcome"] == k) for k in OUTCOMES} + paid = [j for j in jobs if j["outcome"] == "paid"]; exp = [j for j in jobs if j["outcome"] == "expired"]; can = [j for j in jobs if j["outcome"] == "cancelled"]; aband = [j for j in jobs if j["outcome"] == "abandoned"] to_pay = [ts(j["paid_at"]) - ts(j["submitted_at"]) for j in paid if j.get("paid_at") and j.get("submitted_at") and ts(j["paid_at"]) and ts(j["submitted_at"])] wasted = {} - for j in exp + can: + for j in exp + can + aband: k = j.get("cause") or j["outcome"]; w = wasted.setdefault(k, {"jobs": 0, "seconds": 0.0}); w["jobs"] += 1; w["seconds"] = round(w["seconds"] + float(j.get("spent_s") or 0), 1) t0 = min((ts(j["claimed_at"]) for j in jobs if j.get("claimed_at") and ts(j["claimed_at"])), default=None) t1 = max((ts(j[k]) for j in jobs for k in ("closed_at", "paid_at", "submitted_at", "claimed_at") if j.get(k) and ts(j[k])), default=None) span_h = round((t1 - t0) / 3600, 2) if t0 is not None and t1 is not None and t1 > t0 else None shards_paid = sum(int(j.get("shards") or 0) for j in paid) return { - "jobs": len(jobs), "outcomes": tot, + "jobs": len(jobs), "outcomes": tot, "conservation": conservation(jobs), + "abandoned": {"jobs": tot["abandoned"], "by_cause": {c: sum(1 for j in aband if j.get("cause") == c) for c in CAUSES_ABANDONED if any(j.get("cause") == c for j in aband)}, "wasted_s": round(sum(float(j.get("spent_s") or 0) for j in aband), 1)}, "paid": {"segments": tot["paid"], "shards": shards_paid, "ign": round(sum(int(j.get("paid_wei") or 0) for j in paid) / 1e18, 4), "median_end_to_end_s": median([j.get("end_to_end_s") for j in paid]), "median_time_to_pay_s": median(to_pay), "median_margin_daa": median([j.get("margin_daa") for j in paid])}, "missed_deadlines": {"jobs": tot["expired"], "by_cause": {c: sum(1 for j in exp if j.get("cause") == c) for c in CAUSES_EXPIRED if any(j.get("cause") == c for j in exp)}, @@ -143,7 +183,9 @@ def report(jobs): def text(r): o = r["outcomes"]; p = r["paid"]; m = r["missed_deadlines"]; t = r["throughput"] - L = [f"Outcome ledger: {r['jobs']} jobs: paid {o['paid']}, active {o['active']}, expired {o['expired']}, cancelled {o['cancelled']}", + L = [f"Outcome ledger: {r['jobs']} jobs: paid {o['paid']}, active {o['active']}, expired {o['expired']}, cancelled {o['cancelled']}, abandoned {o['abandoned']}"] + L += conservation_lines(r["conservation"]) + L += [f"Abandoned: {r['abandoned']['jobs']} (" + ", ".join(f"{k} {v}" for k, v in r["abandoned"]["by_cause"].items()) + f"); {r['abandoned']['wasted_s']} s of work", f"Paid completions: {p['segments']} segments ({p['shards']} shards, {p['ign']} IGN); median end to end {p['median_end_to_end_s']} s, median time to pay {p['median_time_to_pay_s']} s, median margin at claim {p['median_margin_daa']} DAA", f"Missed deadlines: {m['jobs']} (" + ", ".join(f"{k} {v}" for k, v in m["by_cause"].items()) + f"); median miss {m['median_miss_daa']} DAA past the deadline, median margin at claim {m['median_margin_daa']} DAA, {m['wasted_s']} s of work", f"Wasted work by cause ({r['wasted_s']} s of {r['spent_s']} s spent):"] @@ -178,6 +220,10 @@ RESULT held 2026-10-08T11:12:20Z segment 140 held for retry until DAA 3441 RESULT held_expired 2026-10-08T11:30:00Z segment 140 deadline passed RESULT claim 2026-10-08T11:31:00Z segment 150..159 (10 shards, fresh) margin=240 tip=3800 candidates=1 rank_by=fnv RESULT pass 9 2026-10-08T11:40:00Z no whole segment inside the margin (worklist 20 entries, tip 3900, mhs 1); waiting +RESULT claim 2026-10-08T11:41:00Z segment 160..169 (10 shards, fresh) margin=240 tip=4000 candidates=1 rank_by=fnv +RESULT outcome 2026-10-08T11:46:00Z segment 160..169 abandoned cause=lease spent=12.0 s ledger paid=1 active=1 expired=2 cancelled=2 abandoned=1 +RESULT claim_shard 2026-10-08T11:47:00Z block 170 shard 2 pgas 5000 deadline 4800 tip 4050 +RESULT outcome 2026-10-08T11:48:00Z block 170 shard 2 paid spent=6.0 s ledger paid=1 active=1 expired=2 cancelled=2 abandoned=1 """ def self_test(): @@ -190,22 +236,38 @@ def self_test(): assert [(j["outcome"], j["cause"]) for j in old] == [("cancelled", "chain"), ("paid", None), ("cancelled", "refused"), ("active", None)], old jobs = jobs_from_log(SELF_LOG) got = [(j["first"], j["outcome"], j["cause"]) for j in jobs] - assert got == [(100, "paid", None), (110, "cancelled", "timeout"), (120, "expired", "unpaid"), (130, "cancelled", "shards"), (140, "expired", "held_expired"), (150, "active", None)], got + assert got == [(100, "paid", None), (110, "cancelled", "timeout"), (120, "expired", "unpaid"), (130, "cancelled", "shards"), (140, "expired", "held_expired"), (150, "active", None), (160, "abandoned", "lease"), (170, "paid", None)], got + assert jobs[7]["kind"] == "shard" and jobs[7]["shard"] == 2 and jobs[7]["deadline"] == 4800 and jobs[7]["margin_daa"] == 750, jobs[7] + # known-failed first: the identity line is BROKEN on a ledger whose counters were edited, ok on the real one + broken = report([dict(j, outcome="nothing") for j in jobs[:1]] + jobs[1:])["conservation"]["all"] + assert not broken["ok"] and broken["claims"] == 8 and broken["departures"] == 6 and broken["active"] == 1, broken + # a never-submitted job past the deadline is abandoned (never_submitted), never expired + late = jobs_from_log("RESULT claim 2026-10-08T10:00:10Z segment 100..109 (10 shards, fresh) margin=240 tip=1000 candidates=3 rank_by=fnv\nRESULT pass 9 2026-10-08T11:40:00Z no whole segment inside the margin (worklist 20 entries, tip 3900, mhs 1); waiting\n") + assert [(j["outcome"], j["cause"], j["miss_daa"]) for j in late] == [("abandoned", "never_submitted", 3900 - 1241)], late assert jobs[0]["deadline"] == 1401 and jobs[0]["spent_s"] == 280.0 and jobs[0]["paid_wei"] == 2 * 10**18 and jobs[0]["end_to_end_s"] == 310.0 assert jobs[2]["miss_daa"] == 2700 - (2401 + 250) and jobs[4]["miss_daa"] is None, (jobs[2], jobs[4]) # the held_expired line's last tip is stale: no miss figure, never a negative one r = report(jobs) - assert r["outcomes"] == {"paid": 1, "active": 1, "expired": 2, "cancelled": 2}, r["outcomes"] - assert r["paid"]["segments"] == 1 and r["paid"]["shards"] == 10 and r["paid"]["ign"] == 2.0 and r["paid"]["median_time_to_pay_s"] == 180.0 + assert r["outcomes"] == {"paid": 2, "active": 1, "expired": 2, "cancelled": 2, "abandoned": 1}, r["outcomes"] + cons = r["conservation"] + assert cons["all"]["ok"] and cons["all"]["claims"] == 8 and cons["all"]["departures"] == 7 and cons["all"]["closing_backlog"] == 1, cons["all"] + assert cons["segment"]["claims"] == 7 and cons["shard"] == {"claims": 1, "paid": 1, "active": 0, "expired": 0, "cancelled": 0, "abandoned": 0, "departures": 1, "opening_backlog": 0, "arrivals": 1, "closing_backlog": 0, "ok": True}, cons["shard"] + assert r["abandoned"] == {"jobs": 1, "by_cause": {"lease": 1}, "wasted_s": 12.0}, r["abandoned"] + assert r["paid"]["segments"] == 2 and r["paid"]["shards"] == 11 and r["paid"]["ign"] == 2.0 and r["paid"]["median_time_to_pay_s"] == 180.0 assert r["missed_deadlines"]["jobs"] == 2 and r["missed_deadlines"]["by_cause"] == {"unpaid": 1, "held_expired": 1} assert r["wasted_by_cause"]["timeout"] == {"jobs": 1, "seconds": 1800.0} and r["wasted_by_cause"]["unpaid"]["seconds"] == 280.0 and r["wasted_by_cause"]["shards"]["seconds"] == 250.0 - assert r["wasted_s"] == 1800.0 + 280.0 + 250.0 + 250.0 and r["spent_s"] == r["wasted_s"] + 280.0 - assert r["throughput"]["span_h"] == 1.51 and r["throughput"]["segments_paid_per_h"] == 0.66 and r["throughput"]["paid_share_of_spent"] == 0.098, r["throughput"] + assert r["wasted_s"] == 1800.0 + 280.0 + 250.0 + 250.0 + 12.0 and r["spent_s"] == r["wasted_s"] + 280.0 + 6.0 + assert r["throughput"]["span_h"] == 1.8 and r["throughput"]["segments_paid_per_h"] == 1.11 and r["throughput"]["paid_share_of_spent"] == 0.099, r["throughput"] assert len(r["active"]) == 1 and r["active"][0]["first"] == 150 # a state row wins over the log's row for the same segment merged = merge(jobs_from_state({"label": "t1", "segments": [{"first": 150, "last": 159, "outcome": "paid", "paid_wei": 10**18, "claimed_at": "2026-10-08T11:31:00Z"}]}), jobs) assert sum(1 for j in merged if j["first"] == 150) == 1 and next(j for j in merged if j["first"] == 150)["outcome"] == "paid" - out = text(r); assert "paid 1, active 1, expired 2, cancelled 2" in out and "timeout: 1 jobs, 1800.0 s" in out - print("RESULT prover-outcomes self-test PASS: 6 log jobs, 4 state rows, merge, report and text as expected") + out = text(r); assert "paid 2, active 1, expired 2, cancelled 2, abandoned 1" in out and "timeout: 1 jobs, 1800.0 s" in out, out + assert "Flow conservation (all): claims 8 = paid 2 + expired 2 + cancelled 2 + abandoned 1 + active 1: ok" in out, out + assert "BROKEN by" in "\n".join(conservation_lines(report([dict(j, outcome="nothing") for j in jobs[:1]] + jobs[1:])["conservation"])) + # a state row of the abandoned kind and a shard kind classify as given + st = jobs_from_state({"label": "s", "segments": [{"first": 9, "last": 9, "kind": "shard", "outcome": "abandoned", "cause": "restart"}]}) + assert [(j["kind"], j["outcome"], j["cause"]) for j in st] == [("shard", "abandoned", "restart")], st + print("RESULT prover-outcomes self-test PASS: 8 log jobs (6 segments, 1 abandoned, 1 shard), 4 state rows, conservation ok and BROKEN, merge, report and text as expected") def merge(state_jobs, log_jobs): seen = {(j.get("label"), j["first"]) for j in state_jobs}