igneum/tools/fleet/lib/box.py

150 lines
11 KiB
Python

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", 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-<sha16>, 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)} </dev/null >> {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 </dev/null >> {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 >/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)