#!/usr/bin/env python3 """chat-state-report.py — per-node chat-state tap for box. Probes each NetVM node's browser via CDP (inside its netns): main chat + side chats, last messages. Signs and POSTs to the board chat-state ingest, mirroring accounts-health.sh: payload is \\n\\n, namespace "health". Usage: chat-state-report.py [--no-post] Cron (on bl, every 5 min): */5 * * * * ~/Projects/NetVM/bin/chat-state-report.py >/dev/null 2>&1 Health key setup: same key as accounts-health (namespace "health", registered in /srv/board/health_signers for machine bl). """ import importlib.util import json import os import subprocess import sys import tempfile import time NETVM_DIR = os.environ.get("NETVM_DIR", os.path.expanduser("~/Projects/NetVM")) CHECK = os.path.join(NETVM_DIR, "bin", "chat-state-check.py") MACHINE = os.environ.get("MUSE_MACHINE", "bl") KEY = os.environ.get("CHATSTATE_KEY", os.path.expanduser("~/.ssh/muse-health")) ENDPOINT = os.environ.get( "CHATSTATE_ENDPOINT", "https://board.muse-dev.online/api/box/chat-state/report") MAX_SIDECHATS = 5 FACTS_CAP = 64 * 1024 POST = "--no-post" not in sys.argv def load_nodes(): spec = importlib.util.spec_from_file_location( "netvm_registry", os.path.join(NETVM_DIR, "bin", "netvm-registry.py")) mod = importlib.util.module_from_spec(spec) spec.loader.exec_module(mod) return mod.load() # node -> {"cdp_port": ...} def run_check(node, port, *args): cmd = ["sudo", "-n", "ip", "netns", "exec", "warp-%s" % node, "python3", CHECK, str(port)] + list(args) try: r = subprocess.run(cmd, capture_output=True, text=True, timeout=150) line = r.stdout.strip().splitlines()[-1] return json.loads(line) except Exception as e: # noqa: BLE001 return {"ok": False, "error": "check_failed", "detail": str(e)[:120]} def main(): nodes = load_nodes() chats = [] now = int(time.time()) for node, rec in sorted(nodes.items()): port = rec.get("cdp_port") if not port: continue base = {"node": node, "checked_at": now} m = run_check(node, port, "main") entry = dict(base) if m.get("ok"): entry.update(m) else: entry.update({"context": "main", "error": m.get("error"), "detail": m.get("detail", "")[:120]}) chats.append(entry) continue chats.append(entry) lst = run_check(node, port, "list") for name in (lst.get("sidechats") or [])[:MAX_SIDECHATS]: s = run_check(node, port, "sidechat", name) sentry = dict(base) if s.get("ok"): sentry.update(s) else: sentry.update({"context": "sidechat", "name": name, "error": s.get("error"), "detail": s.get("detail", "")[:120]}) chats.append(sentry) facts = {"chats": chats, "checked_at": now} facts_json = json.dumps(facts, separators=(",", ":")) if len(facts_json) > FACTS_CAP: # Shed message bodies, keep metadata. for c in chats: c.pop("messages", None) facts["truncated"] = True facts_json = json.dumps(facts, separators=(",", ":")) if not POST: print(json.dumps(json.loads(facts_json), indent=2)) return if not os.path.exists(KEY): print("chat-state-report: %s missing — printing JSON, not posting" % KEY, file=sys.stderr) print(facts_json) return tmp = tempfile.mkdtemp() try: payload = os.path.join(tmp, "payload") with open(payload, "w") as f: f.write("%s\n%d\n%s" % (MACHINE, now, facts_json)) # Fresh temp dir: no stale .sig to worry about (cf. AGENTS.md). subprocess.run(["ssh-keygen", "-Y", "sign", "-f", KEY, "-n", "health", payload], check=True, capture_output=True) with open(payload + ".sig") as f: sig = f.read() body = {"machine": MACHINE, "ts": now, "facts_json": facts_json, "signature": sig} r = subprocess.run( ["curl", "-s", "-X", "POST", ENDPOINT, "-H", "Content-Type: application/json", "--data", json.dumps(body)], capture_output=True, text=True, timeout=90) print(r.stdout[:300]) finally: subprocess.run(["rm", "-rf", tmp]) main()