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" # the bracket form: the shell that runs a pgrep/pkill with this pattern never matches itself 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", user="root"): self.host, self.port, self.label, self.iid, self.wallet, self.provider, self.user = host, int(port), label, iid, wallet, provider, user @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"), b.get("ssh_user", "root")) # ---- transport ---- def run(self, cmd, timeout=60): try: r = subprocess.run(["ssh"] + SSH_OPTS + ["-p", str(self.port), f"{self.user}@{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): # the deploy gate: a fleet script (.sh/.py under tools/fleet) goes to a box only when deploy-gate.sh reads clean srcs = [f for f in files if (f.endswith(".sh") or f.endswith(".py")) and "/tools/fleet/" in os.path.abspath(f)] if srcs: g = subprocess.run(["bash", os.path.join(os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "deploy-gate.sh")] + srcs, capture_output=True, text=True, timeout=120) if g.returncode != 0: raise SshError(f"{self.label}: deploy gate RED, push refused: {(g.stdout + g.stderr).strip()[:300]}") for i in range(tries): r = subprocess.run(["scp"] + SSH_OPTS + ["-P", str(self.port)] + list(files) + [f"{self.user}@{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 '[p]ython3 -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)