130 lines
4.5 KiB
Python
130 lines
4.5 KiB
Python
|
|
#!/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
|
||
|
|
<machine>\\n<ts>\\n<facts-json>, 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()
|