180 lines
16 KiB
Python
180 lines
16 KiB
Python
"""The standing fleet (the project lead's ruling, 6 October 2026, 19:50 UK): rented cards that stay up and are never destroyed on a
|
|
job's end. A box is standing when its registry row carries standing=true and a role: live (the live devnet: node, miner,
|
|
prover, voter) or dn2 (Devnet 2: seed, miner or prover). This module is the Mac side of the rule; box-standing.sh is the
|
|
box side (the supervisor loop). Every call goes through lib.box.
|
|
|
|
roster() the standing rows with role, card, price and uptime
|
|
install(label) puts box-standing.sh (and the recovery recipe) on the box and starts the supervisor once
|
|
check() one pass: ssh alive? supervisor up? node synced? exec moving? miner and prover present?
|
|
version against the live manifest (dl.igneum.network/dl/public/igneum-downloads.json,
|
|
files.miner-hive.version); returns one dict per box and writes ~/Desktop/fleet/standing.jsonl
|
|
update(label, pkg) the version follow: the box downloads the manifest's hive package (sha256 checked), unpacks it
|
|
to /opt/igneum/pkg.new, swaps it in and writes /root/fleet/standing.node; the supervisor
|
|
restarts the node on it (data dir kept, seconds)
|
|
rerent(label) the host died (ssh dead for two checks, or the provider says the instance is gone): rent the
|
|
same shape (card, VRAM, provider) with the label suffixed "-r<n>", set it up, install the
|
|
supervisor, mark the old row destroyed; the Devnet 2 roles take their seed from the registry
|
|
loop(every=600) check, then rerent what died and update what is behind, every ten minutes
|
|
weight_check(labels) the 10 percent rule (20:00Z): the labels' share of the live devnet's weight (blue blocks per
|
|
key over the window, from the hub) plus what was removed in the last hour must stay under
|
|
10 percent, or the job is refused; every job that stops or shares a standing miner calls it
|
|
Nothing here destroys a standing box; destroy() in lib.box stays for the one-shot boxes.
|
|
"""
|
|
import os, sys, json, time, datetime, subprocess
|
|
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
|
|
from lib.box import Box, Registry, F, SshError
|
|
HERE = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
|
|
ROOT = os.path.expanduser("~/Desktop/fleet"); MANIFEST = "https://dl.igneum.network/dl/public/igneum-downloads.json"
|
|
def now(): return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
|
def roster():
|
|
reg = Registry.load(); rows = []
|
|
for iid, b in reg.items():
|
|
if not b.get("standing") or b.get("state") == "destroyed": continue
|
|
since = b.get("standing_since") or b.get("rented_at") or now()
|
|
up_h = (datetime.datetime.now(datetime.timezone.utc) - datetime.datetime.fromisoformat(since.replace("Z", "+00:00"))).total_seconds() / 3600
|
|
rows.append({"iid": iid, "label": b["label"], "role": b.get("role", "live"), "card": b.get("card"), "provider": b.get("provider"), "dph": round(float(b.get("dph") or 0), 3), "uptime_h": round(up_h, 1), "since": since})
|
|
return sorted(rows, key=lambda r: (r["role"], r["label"]))
|
|
def manifest():
|
|
try: return json.loads(subprocess.run(["curl", "-fsSL", "-m", "20", MANIFEST], capture_output=True, text=True, timeout=30).stdout)["files"]["miner-hive"]
|
|
except Exception: return {}
|
|
def install(label, role=None, prover=None, mine=None):
|
|
"""Mark the row standing and start the supervisor on the box; idempotent (a running supervisor is left alone)."""
|
|
iid, b = Registry.find(label); role = role or b.get("role", "live"); b_ = Box.from_registry(label)
|
|
prover = "1" if (prover if prover is not None else role == "live" or b.get("dn2_prover")) else "0"
|
|
mine = "1" if (mine if mine is not None else not b.get("dn2_seed")) else "0"
|
|
b_.put([os.path.join(HERE, "box-standing.sh"), os.path.join(HERE, "box-exec-snapshot.sh")], f"{F}/in/")
|
|
hub = next((v for v in Registry.load().values() if v.get("hub") and v.get("state") != "destroyed"), {})
|
|
seed = b.get("dn2_seed_addr", "")
|
|
rc, out, err = b_.run(f"cd {F} && chmod +x in/box-standing.sh in/box-exec-snapshot.sh; pgrep -f '^[b]ash in/box-standing.sh' >/dev/null && echo supervisor-already || {{ ROLE={role} LABEL={label} WALLET={b.get('wallet')} PROVER={prover} MINE={mine} HUB_PEER={hub.get('hub_peer','')} SEED={seed} setsid nohup bash in/box-standing.sh </dev/null >/dev/null 2>&1 & sleep 2; echo supervisor-started; }}", 60)
|
|
Registry.patch(iid, standing=True, role=role, standing_since=b.get("standing_since") or now(), doing=f"standing {role}: node, {'miner, ' if mine == '1' else ''}{'prover' if prover == '1' else 'no prover'} (supervised)")
|
|
return out.strip()
|
|
DIGEST_ALERTS = os.path.join(ROOT, "digest-alerts.json")
|
|
def digest_alert(label, got, expected):
|
|
"""One #incidents line per box per hour when a standing box's handshake digest is not the live one (the seed's 7bd98cc4 at 15:3xZ froze a home miner's tip for 53 minutes)."""
|
|
st = json.load(open(DIGEST_ALERTS)) if os.path.exists(DIGEST_ALERTS) else {}
|
|
if time.time() - st.get(label, 0) < 3600: return
|
|
what = f"Fleet box {label} answers handshakes with consensus digest {got[:8]}, the live object is {expected[:8]}: peers on the live object refuse it and a node whose only peer it is sits on a frozen tip"
|
|
subprocess.run(["node", "/Users/joshm/Projects/igneum/tools/community/discord-hooks.mjs", "incident", "open", "--what", what, "--affected", "nodes that peer only with that box (a home miner reads synced on a frozen tip)", "--doing", "the fleet agent restarts the box's node on the live object; the supervisor's five-minute check found it", "--id", f"digest-{label}-{int(time.time())}", "--live"], capture_output=True, text=True, timeout=60)
|
|
st[label] = time.time(); json.dump(st, open(DIGEST_ALERTS, "w")); print(now(), "digest alert", label, got[:8], "expected", expected[:8], flush=True)
|
|
def check_one(r):
|
|
b = Box.from_registry(r["label"]); row = dict(r); row["t"] = now()
|
|
rc, out, err = b.run("echo alive; pgrep -fc '^bash in/box-standing.sh'; tail -1 /root/fleet/out/standing.log 2>/dev/null | cut -c1-400; echo DIGEST=$(grep -o 'digest: [0-9a-f]*' /root/fleet/node.log | tail -1 | awk '{print substr($2,1,16)}')", 45)
|
|
if rc != 0 or "alive" not in out: row.update(alive=False); return row
|
|
lines = [l for l in out.strip().split("\n") if not l.startswith("DIGEST=")]; row["digest"] = next((l.split("=", 1)[1] for l in out.strip().split("\n") if l.startswith("DIGEST=")), "")
|
|
row.update(alive=True, supervisor=int(lines[1] or 0) if len(lines) > 1 else 0, last=lines[2] if len(lines) > 2 else "")
|
|
kv = dict(x.split("=", 1) for x in row["last"].split() if "=" in x and not x.startswith("miner="))
|
|
row.update(blocks=int(kv.get("blocks", 0) or 0), synced=kv.get("synced") == "true", exec_tip=int(kv.get("exec", 0) or 0), version=kv.get("version", ""), node_pid=kv.get("node", ""), prover=int(kv.get("prover", 0) or 0), bin=kv.get("bin", ""))
|
|
row["bin_sha16"] = os.path.basename(row["bin"]).replace("igneumd-", "") if "igneumd-" in row["bin"] and len(os.path.basename(row["bin"])) == 24 else ""
|
|
return row
|
|
def check(write=True):
|
|
from concurrent.futures import ThreadPoolExecutor
|
|
rows = roster(); m = manifest(); want = m.get("version", "")
|
|
with ThreadPoolExecutor(16) as ex: res = list(ex.map(check_one, rows))
|
|
reg = Registry.load(); expected = open(os.path.join(ROOT, "expected-digest")).read().strip()[:16] if os.path.exists(os.path.join(ROOT, "expected-digest")) else ""
|
|
for r in res:
|
|
r["expected_digest"] = expected; r["digest_ok"] = (not expected) or (not r.get("digest")) or r["digest"][:16] == expected
|
|
if expected and r.get("digest") and not r["digest_ok"]: digest_alert(r["label"], r["digest"], expected)
|
|
r["manifest_version"] = want; wanted = (reg.get(r["iid"]) or {}).get("node_sha16_wanted")
|
|
r["behind"] = bool(wanted and r.get("bin_sha16") and r["bin_sha16"] != wanted) # only an explicit publish sets the want; the package version string is not a node version
|
|
if write:
|
|
with open(os.path.join(ROOT, "standing.jsonl"), "a") as f:
|
|
for r in res: f.write(json.dumps(r) + "\n")
|
|
return res
|
|
def update(label):
|
|
"""The version follow: the manifest's hive package onto the box, the node pointer rewritten, the supervisor restarts it."""
|
|
m = manifest(); b = Box.from_registry(label)
|
|
if not m: raise SshError("no manifest")
|
|
rc, out, err = b.run(f"set -e; cd {F}; curl -fsSL -m 600 -o pkg-{m['version']}.tgz https://dl.igneum.network{m['path']}; echo '{m['sha256']} pkg-{m['version']}.tgz' | sha256sum -c - >/dev/null; rm -rf /opt/igneum/pkg.new; mkdir -p /opt/igneum/pkg.new; tar -C /opt/igneum/pkg.new --strip-components=1 -xzf pkg-{m['version']}.tgz; rm -rf /opt/igneum/pkg.prev; mv /opt/igneum/pkg /opt/igneum/pkg.prev; mv /opt/igneum/pkg.new /opt/igneum/pkg; echo /opt/igneum/pkg/bin/igneumd > {F}/standing.node; pkill -9 -f '^[/]opt/igneum/pkg/bin/igneumd' 2>/dev/null; pkill -9 -f '^[/]opt/igneum/pkg.prev/bin/igneumd' 2>/dev/null; echo updated-to-{m['version']}", 900)
|
|
iid, _ = Registry.find(label); Registry.patch(iid, version_wanted=m["version"], updated_at=now()); return out.strip()
|
|
def rerent(label):
|
|
"""Same shape again on the same provider; the old row is marked destroyed. Returns the new label or None."""
|
|
import fleet
|
|
iid, b = Registry.find(label); n = int(b.get("rerents", 0)) + 1; new = f"{label.split('-r')[0]}-r{n}"
|
|
if b.get("provider") == "runpod":
|
|
import runpod; d = runpod.rent(b.get("gpu_type_id") or f"NVIDIA GeForce {b.get('card')}", new, count=1, disk=40)
|
|
if not d.get("id"): return None
|
|
reg = fleet.load(); reg[str(d["id"])] = {**{k: b[k] for k in ("card", "num_gpus", "provider", "phase", "wallet", "role", "standing") if k in b}, "label": new, "dph": d.get("costPerHr") or b.get("dph"), "state": "renting", "rented_at": now(), "standing_since": now(), "rerents": n, "rerent_of": label}; fleet.save(reg)
|
|
else:
|
|
nid = fleet.rent_card(b.get("card"), new, b.get("archs") or [], min_ram=b.get("vram_mb"), phase=b.get("phase", "2"))
|
|
if not nid: return None
|
|
Registry.patch(str(nid), standing=True, role=b.get("role", "live"), standing_since=now(), rerents=n, rerent_of=label)
|
|
Registry.patch(iid, state="destroyed", destroyed_at=now(), destroyed_reason="host died, re-rented as " + new)
|
|
return new
|
|
# ---- the 10 percent rule (20:00Z): never remove more than 10 percent of the live devnet's 30-day weight in any hour ----
|
|
WEIGHT_LOG = os.path.join(ROOT, "weight-removals.jsonl")
|
|
def weight_shares(blocks=2000):
|
|
"""Vote-key-hash shares over the last `blocks` blocks of the live devnet, from the hub's node (igneum-miner inspect:
|
|
one line per block with voteKeyHash=<64 hex>); the hub is pruned, so the 30-day window is approximated by the largest
|
|
read the RPC serves (2,000 blocks is about 33 minutes at 1 block/s; pass more when the RPC allows). Returns
|
|
({key_hash: share}, blocks_read)."""
|
|
hub = next((v for v in Registry.load().values() if v.get("hub") and v.get("state") != "destroyed"), None)
|
|
if not hub: raise SshError("no hub box")
|
|
b = Box(hub["ssh_host"], hub["ssh_port"], "hub-1", None, None, hub.get("provider"))
|
|
rc, out, err = b.run(f"timeout 170 /opt/igneum/pkg/bin/igneum-miner inspect {blocks} grpc://127.0.0.1:26610 2>&1 | grep -oE 'voteKeyHash=[0-9a-f]{{64}}'", 190)
|
|
keys = [ln.split("=", 1)[1] for ln in out.split() if ln.startswith("voteKeyHash=")]
|
|
n = len(keys) or 1; rows = {}
|
|
for k in keys: rows[k] = rows.get(k, 0) + 1
|
|
return {k: v / n for k, v in rows.items()}, len(keys)
|
|
def vote_key_hash(label):
|
|
"""The box's vote key hash: its miner votes with the key its payout label derives (igneum-miner key-hash <label>)."""
|
|
b = Box.from_registry(label); rc, out, err = b.run(f"/opt/igneum/pkg/bin/igneum-miner key-hash {label} 2>/dev/null | grep -oE '[0-9a-f]{{64}}' | head -1", 30)
|
|
return out.strip()
|
|
def weight_check(labels, limit=0.10, hours=1.0):
|
|
"""Refuses when the labels' weight plus what was removed in the last `hours` exceeds `limit`. Records an allowed removal."""
|
|
shares, nblocks = weight_shares()
|
|
keys = {l: vote_key_hash(l) for l in labels}
|
|
want = sum(shares.get(k, 0) for k in keys.values() if k)
|
|
since = time.time() - hours * 3600; recent = 0.0
|
|
if os.path.exists(WEIGHT_LOG):
|
|
for ln in open(WEIGHT_LOG):
|
|
try: j = json.loads(ln)
|
|
except Exception: continue
|
|
if j.get("ts", 0) >= since: recent += float(j.get("share", 0))
|
|
ok = (want + recent) <= limit
|
|
verdict = {"labels": labels, "keys": {l: k[:16] for l, k in keys.items()}, "share": round(want, 4), "removed_last_hour": round(recent, 4), "limit": limit, "ok": ok, "blocks_read": nblocks, "t": now()}
|
|
if ok:
|
|
with open(WEIGHT_LOG, "a") as f: f.write(json.dumps({"ts": time.time(), "labels": labels, "share": want}) + "\n")
|
|
return verdict
|
|
def table_signed_pct():
|
|
"""The signed share of the frozen voter table at the hub's last lock (the "% of total" of the LOCKED line), or None."""
|
|
hub = next((v for v in Registry.load().values() if v.get("hub") and v.get("state") != "destroyed"), None)
|
|
if not hub: return None
|
|
out = Box(hub["ssh_host"], hub["ssh_port"], "hub-1", None, None, hub.get("provider")).run("grep -E 'Finality: checkpoint [0-9]+ LOCKED:' /root/fleet/node.log | tail -1 | grep -oE '[0-9.]+% of total' | head -1", 30)[1].strip()
|
|
try: return float(out.split("%")[0])
|
|
except Exception: return None
|
|
def table_gate(min_pct=75.0):
|
|
"""The coordinator's rule (6 October 2026, 23:0xZ): nothing of the fleet's moves on a standing box, pool slices included, until the
|
|
table reads over min_pct signed. Returns (ok, pct)."""
|
|
pct = table_signed_pct(); return (pct is not None and pct > min_pct), pct
|
|
def loop(every=600):
|
|
dead = {}
|
|
while True:
|
|
try:
|
|
res = check()
|
|
for r in res:
|
|
if not r["alive"]:
|
|
dead[r["label"]] = dead.get(r["label"], 0) + 1
|
|
if dead[r["label"]] >= 2: print(now(), "re-renting", r["label"], "->", rerent(r["label"]), flush=True); dead.pop(r["label"], None)
|
|
else:
|
|
dead.pop(r["label"], None)
|
|
if r["behind"]: print(now(), "behind:", r["label"], r.get("bin_sha16"), "wanted", (Registry.load().get(r["iid"]) or {}).get("node_sha16_wanted"), "(a publish script moves it; the loop only reports)", flush=True)
|
|
print(now(), "standing check:", len(res), "boxes,", sum(1 for r in res if r["alive"]), "alive,", sum(1 for r in res if r.get("synced")), "synced,", sum(1 for r in res if r["behind"]), "behind,", sum(1 for r in res if not r.get("digest_ok", True)), "off the live digest", flush=True)
|
|
except Exception as e: print(now(), "loop error", str(e)[:200], flush=True)
|
|
time.sleep(every)
|
|
if __name__ == "__main__":
|
|
a = sys.argv[1:]
|
|
if not a or a[0] == "roster":
|
|
for r in roster(): print(f"{r['label']:<12} {r['role']:<5} {str(r['card']):<22} {r['provider']:<7} USD {r['dph']:.3f}/h up {r['uptime_h']:.1f} h")
|
|
rs = roster(); print(f"{len(rs)} standing boxes, USD {sum(r['dph'] for r in rs):.2f}/h, USD {24*sum(r['dph'] for r in rs):.0f}/day")
|
|
elif a[0] == "install": print(install(a[1], *(a[2:3] or [None])))
|
|
elif a[0] == "check":
|
|
for r in check(): print({k: r.get(k) for k in ("label", "role", "alive", "supervisor", "blocks", "synced", "exec_tip", "version", "behind", "prover")})
|
|
elif a[0] == "update": print(update(a[1]))
|
|
elif a[0] == "rerent": print(rerent(a[1]))
|
|
elif a[0] == "loop": loop(int(a[1]) if len(a) > 1 else 600)
|
|
elif a[0] == "table_gate":
|
|
ok, pct = table_gate(float(a[1]) if len(a) > 1 else 75.0); print(json.dumps({"ok": ok, "signed_pct_of_table": pct, "min": float(a[1]) if len(a) > 1 else 75.0})); sys.exit(0 if ok else 1)
|
|
elif a[0] == "weight_check":
|
|
v = weight_check(a[1:]); print(json.dumps(v)); sys.exit(0 if v["ok"] else 1)
|
|
elif a[0] == "weights":
|
|
sh, n = weight_shares(); print("blocks read", n, "voters", len(sh)); print(json.dumps({k[:16]: round(v, 4) for k, v in sorted(sh.items(), key=lambda kv: -kv[1])[:40]}, indent=1))
|