diff --git a/tools/fleet/lib/__init__.py b/tools/fleet/lib/__init__.py new file mode 100644 index 00000000..619410fe --- /dev/null +++ b/tools/fleet/lib/__init__.py @@ -0,0 +1,20 @@ +"""tools/fleet/lib: the ONE tested library for operations on a rented box (the project lead's standing rule, CLAUDE.md e1d0b23, +6 October 2026). Every fleet playbook (canary, wave, matrix, recovery, devnet2-gate) calls these; a watcher reads the +chain-side fact (height, peers, exec tip, paid records), never a reported MH/s or a process name. + + from lib import Box + b = Box.from_registry("dn2-1") # or Box(host, port, label) + b.run("hostname") # ssh, with ServerAlive and a timeout; (rc, out, err) + b.put([local, ...], "/root/fleet/in/") # scp with three tries + b.install_payload(url_or_path, sha256) # the node binary into /root/fleet/in/igneumd-, checked on the box + b.start_node(binary, override, peers, devnet_suffix=None, verifier=None, extra=[]) # anchored kill of the old one, then start + b.wait_synced(limit_s) # the watch line's synced=true (chain-side) + b.height(), b.peers(), b.exec_tip(), b.exec_status(), b.proving_status(), b.state_root(height), b.paid_segments() + b.start_miner(label, wallet, device=0, pack_dir=...) # one igneum-miner with the CUDA worker; returns the pid + b.start_prover(label, wallet, hub_peer, threshold="", miner="keep") # box-prover.sh + b.stop_all() # every stage process (anchored), never sshd + b.destroy() # the provider call plus the registry row + +Process patterns are anchored on the executable's path or use -x (tools/ci/pgrep-self-match-check.sh). +""" +from .box import Box, Registry, SshError # noqa: F401 diff --git a/tools/fleet/lib/box.py b/tools/fleet/lib/box.py new file mode 100644 index 00000000..cb8f362c --- /dev/null +++ b/tools/fleet/lib/box.py @@ -0,0 +1,145 @@ +import json, os, subprocess, time, datetime, hashlib, shlex +ROOT = os.path.expanduser("~/Desktop/fleet"); REG = os.path.join(ROOT, "boxes.json") +KEY = os.path.expanduser("~/.ssh/igneum-fleet") +SSH_OPTS = ["-i", KEY, "-o", "StrictHostKeyChecking=no", "-o", "UserKnownHostsFile=/dev/null", "-o", "LogLevel=ERROR", + "-o", "ConnectTimeout=20", "-o", "ServerAliveInterval=10", "-o", "ServerAliveCountMax=3", "-o", "BatchMode=yes"] +B = "/opt/igneum/pkg/bin"; F = "/root/fleet" +NODE_PATH_PATTERN = "^/root/fleet/in/igneumd|^/opt/igneum/pkg/bin/igneumd" +def now(): return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") +class SshError(RuntimeError): pass +def hexi(v): + if v is None: return 0 + return int(v, 16) if isinstance(v, str) and v.startswith("0x") else int(v or 0) + +class Registry: + """~/Desktop/fleet/boxes.json: one row per rented instance; writes are atomic and locked.""" + @staticmethod + def load(): + if not os.path.exists(REG): return {} + for i in range(6): + try: return json.load(open(REG)) + except json.JSONDecodeError: time.sleep(0.2 * (i + 1)) + raise + @staticmethod + def patch(iid, **fields): + import fcntl + with open(REG + ".lock", "w") as lk: + fcntl.flock(lk, fcntl.LOCK_EX) + reg = Registry.load(); reg.setdefault(str(iid), {}).update(fields) + tmp = REG + f".tmp.{os.getpid()}"; json.dump(reg, open(tmp, "w"), indent=1); os.replace(tmp, REG) + return reg[str(iid)] + @staticmethod + def find(label): + for iid, b in Registry.load().items(): + if b.get("label") == label and b.get("state") != "destroyed": return iid, b + raise KeyError(label) + +class Box: + def __init__(self, host, port, label="box", iid=None, wallet=None, provider="vast"): + self.host, self.port, self.label, self.iid, self.wallet, self.provider = host, int(port), label, iid, wallet, provider + @classmethod + def from_registry(cls, label): + iid, b = Registry.find(label) + if not b.get("ssh_host") or not b.get("ssh_port"): raise SshError(f"{label}: no ssh endpoint in the registry") + return cls(b["ssh_host"], b["ssh_port"], label, iid, b.get("wallet"), b.get("provider", "vast")) + # ---- transport ---- + def run(self, cmd, timeout=60): + try: + r = subprocess.run(["ssh"] + SSH_OPTS + ["-p", str(self.port), f"root@{self.host}", cmd], capture_output=True, text=True, timeout=timeout, stdin=subprocess.DEVNULL) + except subprocess.TimeoutExpired: return 124, "", "timeout" + return r.returncode, r.stdout or "", r.stderr or "" + def check(self, cmd, timeout=60): + rc, out, err = self.run(cmd, timeout) + if rc != 0: raise SshError(f"{self.label}: rc {rc}: {err.strip()[:200]}") + return out + def put(self, files, dest, tries=3, timeout=900): + for i in range(tries): + r = subprocess.run(["scp"] + SSH_OPTS + ["-P", str(self.port)] + list(files) + [f"root@{self.host}:{dest}"], capture_output=True, text=True, timeout=timeout) + if r.returncode == 0: return True + time.sleep(5) + raise SshError(f"{self.label}: scp failed {tries} times: {r.stderr.strip()[:120]}") + def alive(self): return self.run("true", 30)[0] == 0 + # ---- payload ---- + def install_payload(self, src, sha256): + """The node binary onto the box as /root/fleet/in/igneumd-, verified there; src is a local file or an http(s) URL.""" + local = src + if src.startswith("http"): + local = os.path.join(ROOT, f"payload-{sha256[:16]}"); subprocess.run(["curl", "-fsSL", "-o", local, src], check=True, timeout=600) + got = hashlib.sha256(open(local, "rb").read()).hexdigest() + if got != sha256: raise SshError(f"payload sha256 is {got[:16]}, not {sha256[:16]}") + name = f"igneumd-{sha256[:16]}" + self.check(f"mkdir -p {F}/in {F}/out"); self.put([local], f"{F}/in/{name}.new") + self.check(f"cd {F}/in && [ \"$(sha256sum {name}.new | cut -c1-64)\" = {sha256} ] && mv {name}.new {name} && chmod +x {name}") + return f"{F}/in/{name}" + # ---- node ---- + def stop_node(self): + self.run(f"pkill -f '{NODE_PATH_PATTERN}'; sleep 3; pkill -9 -f '{NODE_PATH_PATTERN}' 2>/dev/null; true", 40) + def start_node(self, binary, override_local, peers, devnet_suffix=None, verifier=None, extra=(), appdir=None, log=None, unsynced_mining=False): + """Writes the override file, kills the old node (anchored on its path), starts the new one; returns the digest it prints.""" + appdir = appdir or (f"{F}/dn2" if devnet_suffix else f"{F}/node"); log = log or (f"{F}/dn2-node.log" if devnet_suffix else f"{F}/node.log") + self.put([override_local], f"{F}/override-live.json") + self.stop_node() + flags = ["--devnet"] + ([f"--devnet-suffix={devnet_suffix}"] if devnet_suffix else []) + [f"--appdir={appdir}", "--rpclisten=0.0.0.0:26610", "--evm-rpclisten=127.0.0.1:26790", "--listen=0.0.0.0:26611"] + flags += [f"--addpeer={p}" for p in peers] + [f"--override-params-file={F}/override-live.json", "--nodnsseed", "--disable-upnp", "--nologfiles", "--yes"] + (["--enable-unsynced-mining"] if unsynced_mining else []) + list(extra) + env = f"IGNEUM_PROOF_VERIFIER={shlex.quote(verifier)} " if verifier else "" + self.check(f"cd {F} && {env}setsid nohup {shlex.quote(binary)} {' '.join(shlex.quote(x) for x in flags)} > {log} 2>&1 & sleep 8; echo started", 60) + out = self.run(f"grep -o 'digest: [0-9a-f]*' {log} | tail -1 | awk '{{print $2}}'", 30)[1].strip() + return out + def watch(self): + out = self.run(f"{B}/igneum-miner watch 1 grpc://127.0.0.1:26610 2>/dev/null | grep -o 'blocks=[0-9]*.*synced=[a-z]*' | tail -1", 40)[1].strip() + return dict(x.split("=", 1) for x in out.split() if "=" in x) if out else {} + def height(self): return int(self.watch().get("blocks", 0)) + def daa(self): return int(self.watch().get("daa", 0)) + def peers(self): return int(self.watch().get("peers", 0)) + def synced(self): return self.watch().get("synced") == "true" + def wait_synced(self, limit_s=1200, poll=15): + t0 = time.time() + while time.time() - t0 < limit_s: + if self.synced(): return True + time.sleep(poll) + return False + # ---- exec and proving (chain-side facts) ---- + def rpc(self, method, params=None, timeout=20): + body = json.dumps({"jsonrpc": "2.0", "id": 1, "method": method, "params": params or []}) + out = self.run(f"curl -s -m {timeout} -X POST -H 'Content-Type: application/json' --data {shlex.quote(body)} http://127.0.0.1:26790/", timeout + 20)[1] + try: return json.loads(out).get("result") + except json.JSONDecodeError: return None + def exec_status(self): return self.rpc("igneum_getExecStatus") or {} + def exec_tip(self): return hexi(self.exec_status().get("executedTip")) + def proving_status(self): return self.rpc("igneum_getProvingStatus") or {} + def paid_segments(self): return int((self.proving_status().get("v1") or {}).get("paidSegments") or 0) + def paid_shards(self): return int(self.proving_status().get("paidShards") or 0) + def state_root(self, height): + b = self.rpc("eth_getBlockByNumber", [hex(height), False]) or {} + return b.get("stateRoot") + def node_version(self): + return self.run(f"v=$(pgrep -fa '{NODE_PATH_PATTERN}' | head -1 | awk '{{print $2}}'); [ -n \"$v\" ] && echo \"$(sha256sum \"$v\" | cut -c1-16) $($v --version 2>&1 | head -1)\"", 30)[1].strip() + def rejects_since(self, since_utc, log=None): + log = log or f"{F}/node.log" + out = self.run(f"awk -v s='{since_utc}' 'substr($0,1,19) >= s' {log} 2>/dev/null | grep -cE 'got reject message|PoW rejected|block rejected|invalid block'", 30)[1].strip() + return int(out or 0) + def max_reorg_since(self, since_utc, log=None): + log = log or f"{F}/node.log" + out = self.run(f"awk -v s='{since_utc}' 'substr($0,1,19) >= s' {log} 2>/dev/null | grep -oE 'selected-chain reorg: [0-9]+ chain blocks' | grep -oE '[0-9]+ chain' | awk '{{if ($1>m) m=$1}} END {{print m+0}}'", 30)[1].strip() + return int(out or 0) + # ---- miner and prover ---- + def start_miner(self, label, wallet, device=0, grpc="grpc://127.0.0.1:26610", pack="devnet", log=None): + log = log or f"{F}/out/miner-{device}.log" + self.check(f"cd {F}/mine 2>/dev/null || mkdir -p {F}/mine/packs && cd {F}/mine; [ -d packs/{pack} ] || {B}/igneum-miner export-pack {grpc} packs/{pack} >/dev/null 2>&1; setsid nohup {B}/igneum-miner mine {grpc} 1 100000000 {shlex.quote(label)} --worker {B}/igneum-worker-cuda --worker-args '--device {device} --pack packs/{pack}' --prepare-packs packs/prepare-{device} --exit-on-seed-change --evm-address {wallet} --payout-label {shlex.quote(label)} --status-secs 30 > {log} 2>&1 & echo $!", 90) + return int(self.run("pgrep -x igneum-miner | tail -1", 20)[1].strip() or 0) + def stop_miners(self): self.run("pkill -x igneum-miner; pkill -f '^/opt/igneum/pkg/bin/igneum-worker-cuda'; true", 30) + def start_prover(self, label, wallet, hub_peer="", threshold="", miner="keep"): + self.check(f"cd {F} && chmod +x in/box-prover.sh && LABEL={shlex.quote(label)} WALLET={wallet} HUB_PEER={hub_peer} THRESHOLD={threshold} MINER={miner} setsid nohup in/box-prover.sh /dev/null 2>&1 & echo started", 60) + def stop_all(self): + self.run(f"bash {F}/in/box-kill.sh >/dev/null 2>&1; pkill -f '^python3 -u /root/fleet/in/box-prover.py'; pkill -x igneum-miner; pkill -f '^/opt/igneum/pkg/bin/igneum-worker-cuda'; pkill -x sp1-gpu-server; true", 60) + # ---- provider ---- + def destroy(self): + import sys; sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) + if self.provider == "runpod": + import runpod; runpod.destroy(self.iid) + else: + import vast; vast.destroy(self.iid) + b = Registry.load().get(str(self.iid), {}); t1 = now() + h = (datetime.datetime.strptime(t1, "%Y-%m-%dT%H:%M:%SZ") - datetime.datetime.strptime(b.get("rented_at", t1), "%Y-%m-%dT%H:%M:%SZ")).total_seconds() / 3600 if b.get("rented_at") else 0 + Registry.patch(self.iid, state="destroyed", destroyed_at=t1, hours=round(h, 2), cost_usd=round(h * (b.get("dph") or 0), 2)) + return round(h * (b.get("dph") or 0), 2) diff --git a/tools/fleet/lib/test_box.py b/tools/fleet/lib/test_box.py new file mode 100644 index 00000000..91a0eb2c --- /dev/null +++ b/tools/fleet/lib/test_box.py @@ -0,0 +1,39 @@ +"""Exercises the library against a live Devnet 2 box (the chain the fleet may touch freely): transport, the chain-side +readers, a miner start and stop, a node restart on the current binary and object with the digest read back. + python3 -m lib.test_box dn2-3 (run from tools/fleet) +The known-finished case is a PASS line; the known-failed case is a label that does not exist (a SshError). +""" +import sys, os, time +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +from lib import Box, Registry, SshError +def main(label): + try: Box.from_registry("no-such-box-" + label) + except (KeyError, SshError) as e: print("ok known-failed: a missing label raises", type(e).__name__) + else: print("FAIL known-failed case did not raise"); return 1 + b = Box.from_registry(label); fails = 0 + def chk(name, cond, detail=""): + nonlocal fails + print(("ok " if cond else "FAIL") + f" {name} {detail}"); fails += 0 if cond else 1 + chk("alive", b.alive()) + w = b.watch(); chk("watch has blocks and daa", "blocks" in w and "daa" in w, str(w)[:100]) + h0 = b.height(); chk("height > 0", h0 > 0, h0) + chk("peers >= 1", b.peers() >= 1, b.peers()) + chk("exec tip within 60 of the height", abs(b.exec_tip() - h0) <= 60, f"exec {b.exec_tip()} height {h0}") + chk("state root at height-20 reads", bool(b.state_root(max(1, h0 - 20)))) + chk("proving status answers", "tipDaa" in b.proving_status()) + v = b.node_version(); chk("node version reads", "igneumd" in v, v) + chk("rejects since 1970 is an int", isinstance(b.rejects_since("1970-01-01 00:00:00", f"/root/fleet/dn2-node.log"), int)) + b.stop_miners(); pid = b.start_miner(label + "-libtest", b.wallet, pack="dn2", log="/root/fleet/out/libtest-miner.log"); time.sleep(20) + chk("miner started (pid)", pid > 0, pid) + rc, out, _ = b.run("pgrep -c -x igneum-miner", 20); chk("one miner process", out.strip() == "1", out.strip()) + b.stop_miners(); rc, out, _ = b.run("pgrep -c -x igneum-miner", 20); chk("miner stopped", out.strip() == "0", out.strip()) + seed_peer = next((x.get("dn2_peer") for x in Registry.load().values() if x.get("dn2_seed") and x.get("state") != "destroyed"), "") + binary = b.run("ls /root/fleet/in/igneumd-0313 /root/fleet/in/igneumd-gate 2>/dev/null | head -1", 20)[1].strip() + override = os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "devnet2-override.json") + d = b.start_node(binary, override, [seed_peer] if seed_peer else [], devnet_suffix=2, verifier="/opt/igneum-floor/bin/igneum-prove-host", unsynced_mining=True) + chk("node restarted, digest read", d.startswith("4a0b8726"), d[:16]) + chk("height resumes within 60 s", (time.sleep(45) or b.height()) >= h0, b.height()) + b.start_miner(label, b.wallet, pack="dn2", log="/root/fleet/out/dn2-miner.log") # the box back as it was + print(("PASS" if fails == 0 else "FAIL") + f" lib.test_box on {label}: {fails} failures") + return 1 if fails else 0 +if __name__ == "__main__": sys.exit(main(sys.argv[1] if len(sys.argv) > 1 else "dn2-3"))