Files
ethereum-rpc-docker/node-facts.sh
T

290 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__))
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:
for line in open(os.path.join(BASE, ".env")):
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
FIXTURES = os.environ.get("NODE_FACTS_FIXTURES") # test mode: answer HTTP from <dir>/<method-or-path>.json, no network
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()