291 lines
20 KiB
Bash
Executable File
291 lines
20 KiB
Bash
Executable File
#!/usr/bin/env python3
|
|
"""node-facts.sh [node_path ...] [--pretty] (hand-maintained; ansible_showcase/node-facts-plan.md)
|
|
|
|
One JSON 'node facts' document per node, same schema (v1) for every client family: own version, peer
|
|
version histogram + newer_than_us, protocol upgrade schedule the running config knows, chain signals,
|
|
head. Modelled on avalanchego's Info API; simulated best-effort for the others. Read-only: compose
|
|
files + docker ps/inspect + HTTP through the local traefik (like sync-status.sh). Parallel over nodes
|
|
and over the per-node calls. Every field present; null + reason instead of silence.
|
|
Exit 0 always once .env is readable; exit 2 if .env/COMPOSE_FILE is missing (one error line).
|
|
"""
|
|
import concurrent.futures as cf, datetime, json, os, re, socket, subprocess, sys, time, urllib.request, ssl
|
|
|
|
BASE = os.path.dirname(os.path.abspath(__file__))
|
|
FIXTURES = os.environ.get("NODE_FACTS_FIXTURES") # test mode: answer HTTP from <dir>/<method-or-path>.json, env from <dir>/env, no network
|
|
INFRA = {"base", "rpc", "monitoring", "backup-http", "drpc", "drpc-free", "squidfura", "squidfura-backend", "oracle", "extra"}
|
|
NOW = lambda: datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
|
SEMVER = re.compile(r"(\d+)\.(\d+)\.(\d+)")
|
|
CTX = ssl.create_default_context()
|
|
|
|
|
|
def env():
|
|
e = {}
|
|
try:
|
|
envfile = os.path.join(FIXTURES, "env") if FIXTURES else os.path.join(BASE, ".env") # fixture mode: no dotfile needed
|
|
for line in open(envfile):
|
|
line = line.strip()
|
|
if line and not line.startswith("#") and "=" in line:
|
|
k, v = line.split("=", 1); e[k] = v.strip().strip('"').strip("'")
|
|
except Exception:
|
|
return None
|
|
return e
|
|
|
|
|
|
|
|
|
|
def _fixture(url, body):
|
|
key = (body or {}).get("method") if body else url.split("/eth/", 1)[-1].replace("/", "_") if "/eth/" in url else os.path.basename(url)
|
|
f = os.path.join(FIXTURES, f"{key}.json")
|
|
if os.path.exists(f): return 1, 200, open(f).read(), None
|
|
return 1, None, None, f"fixture missing: {key}"
|
|
|
|
|
|
def http(url, body=None, timeout=5.0):
|
|
"""-> (ms, status, text|None, err|None)"""
|
|
if FIXTURES: return _fixture(url, body)
|
|
t0 = time.time()
|
|
try:
|
|
req = urllib.request.Request(url, data=(json.dumps(body).encode() if body is not None else None),
|
|
headers={"content-type": "application/json", "user-agent": "node-facts/1"})
|
|
with urllib.request.urlopen(req, timeout=timeout, context=CTX) as r:
|
|
return int((time.time() - t0) * 1000), r.status, r.read().decode(errors="replace"), None
|
|
except urllib.error.HTTPError as ex:
|
|
return int((time.time() - t0) * 1000), ex.code, None, f"http {ex.code}"
|
|
except Exception as ex:
|
|
return int((time.time() - t0) * 1000), None, None, f"{type(ex).__name__}: {str(ex)[:80]}"
|
|
|
|
|
|
def rpc(url, method, params=None, timeout=5.0):
|
|
ms, st, txt, err = http(url, {"jsonrpc": "2.0", "id": 1, "method": method, "params": params if params is not None else []}, timeout)
|
|
if err: return ms, None, err
|
|
try:
|
|
d = json.loads(txt)
|
|
except Exception:
|
|
return ms, None, "non-json response"
|
|
if "error" in d and d["error"]: return ms, None, f"rpc error: {str(d['error'].get('message', d['error']))[:80]}"
|
|
return ms, d.get("result"), None
|
|
|
|
|
|
def rest(url, timeout=5.0):
|
|
ms, st, txt, err = http(url, None, timeout)
|
|
if err: return ms, None, err
|
|
try: return ms, json.loads(txt), None
|
|
except Exception: return ms, None, "non-json response"
|
|
|
|
|
|
def semver(s):
|
|
m = SEMVER.search(s or ""); return tuple(int(x) for x in m.groups()) if m else None
|
|
|
|
|
|
def docker_ps():
|
|
out = subprocess.run(["docker", "ps", "--format", '{{.Names}}\t{{.Image}}\t{{.Label "com.docker.compose.service"}}\t{{.Label "com.docker.compose.project.config_files"}}'],
|
|
capture_output=True, text=True, timeout=20).stdout
|
|
rows = []
|
|
for line in out.splitlines():
|
|
p = line.split("\t")
|
|
if len(p) == 4: rows.append({"container": p[0], "image": p[1], "service": p[2], "files": p[3]})
|
|
return rows
|
|
|
|
|
|
FAMILY = [("avalanchego", "avalanchego"), ("kona-node", "kona"), ("kona", "kona"), ("op-node", "op-node"), ("conduit-op-reth", "op-reth"), ("op-reth", "op-reth"),
|
|
("op-geth", "op-geth"), ("op-erigon", "erigon"), ("lighthouse", "beacon"), ("prysm", "beacon"), ("nimbus", "beacon"), ("teku", "beacon"), ("lodestar", "beacon"),
|
|
("grandine", "beacon"), ("client-go", "geth"), ("bnb-chain/bsc", "geth"), ("bsc", "geth"), ("/geth", "geth"), ("geth", "geth"), ("bor", "bor"), ("erigon", "erigon"),
|
|
("nethermind", "nethermind"), ("besu", "besu"), ("reth-bsc", "reth"), ("bera-reth", "reth"), ("reth", "reth"), ("nitro", "nitro"), ("cometbft", "cometbft"),
|
|
("gaiad", "cometbft"), ("external-node", "zksync-en"), ("maru", "maru"), ("arc-consensus", "arc"), ("arc-execution", "arc")]
|
|
|
|
|
|
def family_of(image):
|
|
low = image.lower()
|
|
for sub, fam in FAMILY:
|
|
if sub in low: return fam
|
|
return "unknown"
|
|
|
|
|
|
def labels_of(compose_text):
|
|
"""traefik stripprefix prefixes found in the compose file -> set of paths like 'linea-mainnet', 'linea-mainnet/node'"""
|
|
return set(re.findall(r"stripprefix\.prefixes=/([^\"'\s,]+)", compose_text))
|
|
|
|
|
|
def node_facts(env_, node, containers, pretty=False):
|
|
proto = "http" if env_.get("NO_SSL") else "https"; domain = env_.get("DOMAIN") or "0.0.0.0"
|
|
compose = os.path.join(BASE, node + ".yml")
|
|
try: text = open(compose).read()
|
|
except Exception as ex:
|
|
return {"schema": 1, "at": NOW(), "host": socket.gethostname().split(".")[0], "node_path": node, "network": None, "client": None,
|
|
"version": None, "peers": None, "schedule": None, "signals": None, "head": None, "timing_ms": None, "error": f"compose file unreadable: {ex}"}
|
|
# services defined in THIS compose file (top-level keys under services:) -> containers by compose service label
|
|
svc_names = set(re.findall(r"^ ([A-Za-z0-9_.-]+):\s*$", text, re.M)) - {"services", "volumes", "networks", "x-upstreams"}
|
|
mine = [c for c in containers if c["service"] in svc_names]
|
|
paths = labels_of(text)
|
|
client_paths = sorted(p for p in paths if "/" not in p); node_paths = sorted(p for p in paths if p.endswith("/node"))
|
|
# beacon REST route (base node template): rule PathPrefix(`/<short>/eth`) with stripprefix /<short> -> URL base/<short>/eth/...
|
|
eth_paths = sorted(set(re.findall(r"PathPrefix\(`/([^`]+?)/eth`\)", text)))
|
|
base = f"{proto}://{domain}"
|
|
# execution-side container = the one whose service does not end with -node / known CL images
|
|
exec_c = next((c for c in mine if family_of(c["image"]) not in ("kona", "op-node", "beacon", "maru", "arc") and not c["service"].endswith("-node")), None) or (mine[0] if mine else None)
|
|
cl_c = next((c for c in mine if family_of(c["image"]) in ("kona", "op-node", "beacon", "maru")), None)
|
|
fam = family_of(exec_c["image"]) if exec_c else "unknown"
|
|
net = os.path.basename(node).split("-")[0] if "/" in node else node
|
|
m = re.search(r"\n\s+network:\s*([\w-]+)|chain:\s*\"?([\w-]+)", text)
|
|
network = (m.group(1) or m.group(2)) if m else net
|
|
def imgtag(c): return (c["image"].rsplit(":", 1) + [""])[:2] if c else ("", "")
|
|
client = {"name": fam, "image": imgtag(exec_c)[0], "tag": imgtag(exec_c)[1], "container": exec_c["container"] if exec_c else None,
|
|
"consensus": ({"name": family_of(cl_c["image"]), "image": imgtag(cl_c)[0], "tag": imgtag(cl_c)[1], "container": cl_c["container"]} if cl_c else None)}
|
|
url = f"{base}/{client_paths[0]}" if client_paths else None
|
|
timing = {}
|
|
R = {"version": None, "peers": None, "schedule": None, "signals": None, "head": None}
|
|
def nul(src, err): return {"source": src, "error": err}
|
|
|
|
# ---- per-family collectors (each returns dict for one field) ----
|
|
def v_rpc():
|
|
if not url: return {"reported": None, "db_version": None, **nul("web3_clientVersion", "no traefik route for the client")}
|
|
ms, r, e = rpc(url, "web3_clientVersion"); timing["version"] = ms
|
|
return {"reported": r.strip().splitlines()[0][:120] if isinstance(r, str) and r.strip() else None, "db_version": None, **nul("web3_clientVersion", e)}
|
|
def p_admin():
|
|
if not url: return {"count": None, "versions": None, "newer_than_us": None, **nul("admin_peers", "no traefik route for the client")}
|
|
ms, r, e = rpc(url, "admin_peers", timeout=10); timing["peers"] = ms
|
|
if e or not isinstance(r, list): return {"count": None, "versions": None, "newer_than_us": None, **nul("admin_peers", e or "admin namespace disabled or unexpected result")}
|
|
hist = {}
|
|
for p in r:
|
|
name = (p.get("name") or "?"); ver = name.split("/")[1] if "/" in name else name
|
|
hist[ver] = hist.get(ver, 0) + 1
|
|
return {"count": len(r), "versions": hist, "newer_than_us": None, **nul("admin_peers", None)}
|
|
def s_nodeinfo():
|
|
if not url: return {"upgrades": None, **nul("admin_nodeInfo", "no traefik route for the client")}
|
|
ms, r, e = rpc(url, "admin_nodeInfo"); timing["schedule"] = ms
|
|
cfg = (((r or {}).get("protocols") or {}).get("eth") or {}).get("config") if isinstance(r, dict) else None
|
|
if e or not isinstance(cfg, dict): return {"upgrades": None, **nul("admin_nodeInfo", e or "admin namespace disabled or no chain config in nodeInfo")}
|
|
ups = {}
|
|
for k, v in cfg.items():
|
|
if k.endswith("Time") and isinstance(v, int): ups[k[:-4]] = {"time": datetime.datetime.fromtimestamp(v, datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ"), "block": None}
|
|
elif k.endswith("Block") and isinstance(v, int): ups[k[:-5]] = {"time": None, "block": v}
|
|
return {"upgrades": ups, **nul("admin_nodeInfo", None)}
|
|
def h_eth():
|
|
if not url: return {"height": None, "timestamp": None, "syncing": None, **nul("eth_syncing+eth_getBlockByNumber", "no traefik route for the client")}
|
|
ms1, s, e1 = rpc(url, "eth_syncing"); ms2, b, e2 = rpc(url, "eth_getBlockByNumber", ["latest", False]); timing["head"] = ms1 + ms2
|
|
h = {"height": int(b["number"], 16) if isinstance(b, dict) and b.get("number") else None,
|
|
"timestamp": datetime.datetime.fromtimestamp(int(b["timestamp"], 16), datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") if isinstance(b, dict) and b.get("timestamp") else None,
|
|
"syncing": (s is not False) if (s is not None and not e1) else None}
|
|
return {**h, **nul("eth_syncing+eth_getBlockByNumber", e2 or e1)}
|
|
def avax():
|
|
# Info API is not routed through traefik; try the client route's host with /ext/info replaced, then bridge IP
|
|
cand = []
|
|
if url: cand.append(url.rsplit("/", 1)[0] + "/" + client_paths[0].split("/")[0] + "-info") # optional future route
|
|
try:
|
|
ip = subprocess.run(["docker", "inspect", exec_c["container"], "--format", "{{range .NetworkSettings.Networks}}{{.IPAddress}} {{end}}"], capture_output=True, text=True, timeout=10).stdout.split()
|
|
if ip: cand.append(f"http://{ip[0]}:9650")
|
|
except Exception: pass
|
|
info = None; used = None; err = "avalanche info API not reachable (needs /ext/info route or host->bridge access)"
|
|
for c in cand:
|
|
ms, r, e = rpc(c + "/ext/info", "info.getNodeVersion", {}); timing["version"] = ms
|
|
if not e and isinstance(r, dict): info, used = c + "/ext/info", c; break
|
|
err = e or err
|
|
if not info:
|
|
R["version"] = {"reported": None, "db_version": None, **nul("info.getNodeVersion", err)}
|
|
R["peers"] = {"count": None, "versions": None, "newer_than_us": None, **nul("info.peers", err)}
|
|
R["schedule"] = {"upgrades": None, **nul("info.upgrades", err)}; R["signals"] = None
|
|
R["head"] = h_eth(); return
|
|
def call(m):
|
|
ms, r, e = rpc(info, m, {}, timeout=10); return ms, r, e
|
|
with cf.ThreadPoolExecutor(4) as ex:
|
|
fv, fp, fu, fa = ex.submit(call, "info.getNodeVersion"), ex.submit(call, "info.peers"), ex.submit(call, "info.upgrades"), ex.submit(call, "info.acps")
|
|
ms, r, e = fv.result(); timing["version"] = ms
|
|
R["version"] = {"reported": (r or {}).get("version") if isinstance(r, dict) else None, "db_version": (r or {}).get("databaseVersion") if isinstance(r, dict) else None,
|
|
"vm_versions": (r or {}).get("vmVersions") if isinstance(r, dict) else None, **nul("info.getNodeVersion", e)}
|
|
ms, r, e = fp.result(); timing["peers"] = ms
|
|
if not e and isinstance(r, dict):
|
|
hist = {}
|
|
for p in r.get("peers") or []: hist[p.get("version") or "?"] = hist.get(p.get("version") or "?", 0) + 1
|
|
R["peers"] = {"count": len(r.get("peers") or []), "versions": hist, "newer_than_us": None, **nul("info.peers", None)}
|
|
else: R["peers"] = {"count": None, "versions": None, "newer_than_us": None, **nul("info.peers", e or "unexpected result")}
|
|
ms, r, e = fu.result(); timing["schedule"] = ms
|
|
R["schedule"] = {"upgrades": ({k[:-4]: {"time": v, "block": None} for k, v in r.items() if k.endswith("Time")} if isinstance(r, dict) else None), **nul("info.upgrades", e)}
|
|
ms, r, e = fa.result()
|
|
R["signals"] = ({"acps": {k: {"supporters": len(v.get("supporters") or []), "objectors": len(v.get("objectors") or []), "abstainers": len(v.get("abstainers") or []), "supportWeight": v.get("supportWeight")}
|
|
for k, v in (r.get("acps") or {}).items()}, "source": "info.acps", "error": None} if not e and isinstance(r, dict) else {"acps": None, "source": "info.acps", "error": e})
|
|
R["head"] = h_eth()
|
|
def opnode():
|
|
nurl = f"{base}/{node_paths[0]}" if node_paths else None
|
|
C = {"version": None, "peers": None, "schedule": None, "head": None}
|
|
if not nurl:
|
|
C["peers"] = {"count": None, "versions": None, "newer_than_us": None, **nul("opp2p_peers", "no /<node>/node route")}
|
|
C["schedule"] = {"upgrades": None, **nul("optimism_rollupConfig", "no /<node>/node route")}; R["consensus"] = C; return
|
|
with cf.ThreadPoolExecutor(4) as ex:
|
|
fp, fc, fs, fv = ex.submit(rpc, nurl, "opp2p_peers", [True], 10), ex.submit(rpc, nurl, "optimism_rollupConfig"), ex.submit(rpc, nurl, "optimism_syncStatus"), ex.submit(rpc, nurl, "optimism_version")
|
|
ms, r, e = fv.result(); timing["consensus_version"] = ms
|
|
C["version"] = {"reported": r if isinstance(r, str) else None, **nul("optimism_version", e)}
|
|
ms, r, e = fp.result(); timing["consensus_peers"] = ms
|
|
if not e and isinstance(r, dict):
|
|
hist = {}
|
|
for p in (r.get("peers") or {}).values():
|
|
ua = p.get("userAgent") or "?"; hist[ua] = hist.get(ua, 0) + 1
|
|
C["peers"] = {"count": len(r.get("peers") or {}), "versions": hist, "newer_than_us": None, **nul("opp2p_peers", None if any(semver(k) for k in hist) else "userAgent carries no version")}
|
|
else: C["peers"] = {"count": None, "versions": None, "newer_than_us": None, **nul("opp2p_peers", e or "unexpected result")}
|
|
ms, r, e = fc.result(); timing["consensus_schedule"] = ms
|
|
C["schedule"] = {"upgrades": ({k[:-5]: {"time": datetime.datetime.fromtimestamp(v, datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ"), "block": None}
|
|
for k, v in r.items() if k.endswith("_time") and isinstance(v, int)} if isinstance(r, dict) else None), **nul("optimism_rollupConfig", e)}
|
|
ms, r, e = fs.result()
|
|
if not e and isinstance(r, dict) and isinstance(r.get("unsafe_l2"), dict):
|
|
u = r["unsafe_l2"]; C["head"] = {"height": u.get("number"), "timestamp": datetime.datetime.fromtimestamp(u["timestamp"], datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") if u.get("timestamp") else None,
|
|
"syncing": None, "unsafe_minus_safe": (u.get("number") or 0) - ((r.get("safe_l2") or {}).get("number") or 0), **nul("optimism_syncStatus", None)}
|
|
else: C["head"] = {"height": None, "timestamp": None, "syncing": None, **nul("optimism_syncStatus", e or "unexpected result")}
|
|
R["consensus"] = C
|
|
def beacon():
|
|
eurl = f"{base}/{eth_paths[0]}" if eth_paths else None # eth/v1/... is appended below
|
|
if not eurl:
|
|
R["consensus_version"] = None; return
|
|
with cf.ThreadPoolExecutor(4) as ex:
|
|
fv, fp, fs, fc = ex.submit(rest, eurl + "/eth/v1/node/version"), ex.submit(rest, eurl + "/eth/v1/node/peers", 10), ex.submit(rest, eurl + "/eth/v1/node/syncing"), ex.submit(rest, eurl + "/eth/v1/config/spec")
|
|
ms, r, e = fv.result(); timing["consensus_version"] = ms
|
|
cv = {"reported": ((r or {}).get("data") or {}).get("version") if isinstance(r, dict) else None, **nul("eth/v1/node/version", e)}
|
|
ms, r, e = fp.result()
|
|
cp = ({"count": sum(1 for x in ((r or {}).get("data") or []) if (x.get("state") or "connected") == "connected"), "versions": None, "newer_than_us": None, **nul("eth/v1/node/peers", "beacon peers API carries no agent version")} if not e and isinstance(r, dict)
|
|
else {"count": None, "versions": None, "newer_than_us": None, **nul("eth/v1/node/peers", e or "unexpected result")})
|
|
ms, r, e = fs.result()
|
|
d = (r or {}).get("data") if isinstance(r, dict) else None
|
|
ch = ({"head_slot": int(d.get("head_slot")) if d.get("head_slot") is not None else None, "syncing": d.get("is_syncing"), **nul("eth/v1/node/syncing", None)} if isinstance(d, dict) else {"head_slot": None, "syncing": None, **nul("eth/v1/node/syncing", e)})
|
|
ms, r, e = fc.result()
|
|
spec = (r or {}).get("data") if isinstance(r, dict) else None
|
|
cs = ({"upgrades": {k[:-11]: {"epoch": int(v)} for k, v in spec.items() if k.endswith("_FORK_EPOCH") and str(v).isdigit() and int(v) < 2**62}, **nul("eth/v1/config/spec", None)} if isinstance(spec, dict) else {"upgrades": None, **nul("eth/v1/config/spec", e)})
|
|
R["consensus"] = {"version": cv, "peers": cp, "head": ch, "schedule": cs}
|
|
|
|
# ---- dispatch ----
|
|
t0 = time.time()
|
|
if fam == "avalanchego":
|
|
avax()
|
|
else:
|
|
with cf.ThreadPoolExecutor(4) as ex:
|
|
fv, fp, fs, fh = ex.submit(v_rpc), ex.submit(p_admin), ex.submit(s_nodeinfo), ex.submit(h_eth)
|
|
R["version"], R["peers"], R["schedule"], R["head"] = fv.result(), fp.result(), fs.result(), fh.result()
|
|
if fam in ("reth", "op-reth") and R["schedule"].get("upgrades") is None:
|
|
R["schedule"]["error"] = "reth exposes no chain config over RPC"
|
|
if fam in ("nitro", "zksync-en", "cometbft", "maru", "arc", "unknown") and R["peers"].get("count") is None:
|
|
R["peers"]["error"] = f"no peer API via RPC for {fam}"
|
|
if node_paths and (cl_c is None or family_of(cl_c["image"]) in ("kona", "op-node")): opnode()
|
|
if eth_paths: beacon()
|
|
# newer_than_us
|
|
mine_v = semver((R["version"] or {}).get("reported")); hist = (R["peers"] or {}).get("versions")
|
|
if mine_v and hist:
|
|
R["peers"]["newer_than_us"] = sum(n for v, n in hist.items() if (semver(v) or (0,)) > mine_v)
|
|
doc = {"schema": 1, "at": NOW(), "host": socket.gethostname().split(".")[0], "node_path": node, "network": network, "client": client,
|
|
"version": R["version"], "peers": R["peers"], "schedule": R["schedule"], "signals": R.get("signals"), "head": R["head"],
|
|
"consensus": R.get("consensus"), "timing_ms": {**timing, "total": int((time.time() - t0) * 1000)}}
|
|
return doc
|
|
|
|
|
|
def main():
|
|
args = [a for a in sys.argv[1:] if not a.startswith("--")]; pretty = "--pretty" in sys.argv
|
|
e = env()
|
|
if not e or not e.get("COMPOSE_FILE"):
|
|
print(json.dumps({"schema": 1, "at": NOW(), "host": None, "node_path": None, "error": "cannot read .env or COMPOSE_FILE empty"})); sys.exit(2)
|
|
nodes = args or [f[:-4] for f in e["COMPOSE_FILE"].split(":") if f.endswith(".yml") and os.path.basename(f)[:-4] not in INFRA and not os.path.basename(f).startswith(("drpc", "squidfura", "oracle"))]
|
|
try: containers = docker_ps() if not FIXTURES else json.load(open(os.path.join(FIXTURES, "docker_ps.json")))
|
|
except Exception: containers = []
|
|
with cf.ThreadPoolExecutor(int(os.environ.get("MAX_WORKERS", "8"))) as ex:
|
|
for doc in ex.map(lambda n: node_facts(e, n, containers), nodes):
|
|
print(json.dumps(doc, indent=2 if pretty else None), flush=True)
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|