diff --git a/bin/chat-state-check.py b/bin/chat-state-check.py new file mode 100755 index 0000000..2612cd6 --- /dev/null +++ b/bin/chat-state-check.py @@ -0,0 +1,281 @@ +#!/usr/bin/env python3 +"""chat-state-check.py — one-shot CDP chat-state probe for a single node. + +Runs INSIDE the node's netns (CDP listens on 127.0.0.1 there). + +Usage: chat-state-check.py [name] + modes: main | list | sidechat + +Prints one JSON object to stdout. + +Honest limits (see CHATSTATE_SPEC.md): only the DOM-visible message window +is captured (React virtualization); author attribution is best-effort and +may be null; checked_at is the report time, not per-message times. +""" +import base64 +import json +import sys +import time +import urllib.request + +import websocket + +MSG_N = 10 +MSG_WIDTH = 500 + + +def connect(port): + with urllib.request.urlopen( + "http://127.0.0.1:%s/json/list" % port, timeout=5) as r: + targets = json.load(r) + pages = [t for t in targets if t.get("type") == "page"] + if not pages: + return None + return websocket.create_connection( + pages[0]["webSocketDebuggerUrl"], timeout=15) + + +def ev1(ws, expr, await_p=False, reads=30): + """Runtime.evaluate that skips CDP event chatter while awaiting ours.""" + ws.send(json.dumps({ + "id": 1, "method": "Runtime.evaluate", + "params": {"expression": expr, "returnByValue": True, + "awaitPromise": await_p}})) + for _ in range(reads): + resp = json.loads(ws.recv()) + if resp.get("id") != 1: + continue + if "error" in resp: + raise RuntimeError("CDP evaluate failed: %s" % resp["error"]) + res = resp.get("result", {}).get("result", {}) + return res.get("value") + raise RuntimeError("CDP evaluate: no response") + + +def read_current(ws): + title = ev1(ws, "document.title") or "" + url = ev1(ws, "window.location.href") or "" + msgs = ev1(ws, """(() => { + const ps = [...document.querySelectorAll('p')].slice(-%d) + .map(p => (p.innerText||'').slice(0,%d)).filter(t => t.trim()); + return ps; + })()""" % (MSG_N, MSG_WIDTH)) or [] + return {"title": title, "url": url, + "messages": [{"text": t} for t in msgs]} + + +def _wait_for(ws, expr, timeout_s=12): + """Poll a JS truthiness expression until true or timeout.""" + deadline = time.time() + timeout_s + while time.time() < deadline: + try: + if ev1(ws, expr): + return True + except Exception: # noqa: BLE001 + pass + time.sleep(1) + return False + + +def ensure_sidebar_open(ws): + """The toggle closes an open sidebar — only click when closed. + + Detected structurally via the 'New side chat' button (text matching + 'Side chats' is unreliable: main-chat messages can contain the phrase). + Waits for React to render the panel after opening. + """ + is_open = ev1(ws, """(() => { + return !!document.querySelector('button[aria-label="New side chat"]'); + })()""") + if not is_open: + ev1(ws, """(() => { + const btn = [...document.querySelectorAll('button')].find(b => + (b.textContent||'').includes('Open chat and side chats')); + if (btn) btn.click(); + })()""") + _wait_for(ws, """(() => { + return !!document.querySelector( + 'button[aria-label="New side chat"]'); + })()""") + # The row list renders a beat after the panel; wait for it. + _wait_for(ws, """(() => { + return document.querySelectorAll( + 'button[aria-label="More thread actions"]').length > 0; + })()""", timeout_s=10) + + +_SIDEBAR_JS = """(() => { + const anchor = document.querySelector('button[aria-label="New side chat"]'); + if (!anchor) return 'NOANCHOR'; + let sec = anchor.parentElement; + for (let i = 0; i < 8 && sec; i++) { + if (sec.querySelectorAll( + 'button[aria-label="More thread actions"]').length) break; + sec = sec.parentElement; + } + if (!sec) return 'NOSEC'; + return JSON.stringify({ok: true}); +})()""" + + +def _sidebar_section(ws): + """Return True when the side-chat list section is addressable.""" + raw = ev1(ws, _SIDEBAR_JS) + try: + return json.loads(raw or "").get("ok", False) + except (ValueError, TypeError): + return False + + +_LIST_JS = """(() => { + const T = (el) => (el.innerText || '').trim(); + const TS = + /^(\\d+[smhd]|just now|Yesterday|Today|[A-Z][a-z]{2} \\d{1,2}.*)$/; + const Y = (el) => el.getBoundingClientRect().top; + const h3s = [...document.querySelectorAll('h3')]; + const head = h3s.find(h => h.textContent.trim() === 'Side chats'); + if (!head) return '[]'; + // Scroll the list to top so the section's rows sit above the next + // (sticky) header; without this, rows below the fold are missed. + let sc = head.parentElement; + for (let i = 0; i < 8 && sc; i++) { + if (sc.scrollHeight > sc.clientHeight + 10) break; + sc = sc.parentElement; + } + if (sc) sc.scrollTop = 0; + const y0 = Y(head); + const isStopText = (t) => t === 'Unread updates' || t === 'Chats' || + t === 'Unread chats'; + let stopY = null; + for (const h of h3s) { + const y = Y(h); + if (y > y0 + 5 && (stopY === null || y < stopY)) stopY = y; + } + // Non-H3 section headers ("Unread updates", ...) also bound the list. + for (const el of document.querySelectorAll('*')) { + if (el.children.length !== 0) continue; + if (!isStopText((el.textContent || '').trim())) continue; + const y = Y(el); + if (y > y0 + 5 && (stopY === null || y < stopY)) stopY = y; + } + const names = []; + for (const el of document.querySelectorAll('div')) { + const cn = (el.className || '').toString(); + if (cn.indexOf('nav-row') === -1) continue; + const y = Y(el); + if (y <= y0 + 5) continue; + if (stopY !== null && y >= stopY - 5) continue; + const m = T(el).match(/^([^\\n]+)\\n\\n(.+)$/); + if (m && TS.test(m[2].trim()) && names.indexOf(m[1].trim()) === -1) { + names.push(m[1].trim()); + if (names.length >= 8) break; + } + } + return JSON.stringify(names); +})()""" + + +def _list_once(ws): + raw = ev1(ws, _LIST_JS) + try: + return json.loads(raw or "[]") + except (ValueError, TypeError): + return [] + + +def list_sidechats(ws): + ensure_sidebar_open(ws) + if not _sidebar_section(ws): + return [] + # React renders rows progressively after navigation; poll until the + # list stabilizes instead of trusting the first paint. + prev = None + for _ in range(4): + names = _list_once(ws) + if names and names == prev: + return names + prev = names + time.sleep(2) + return prev or [] + + +def open_sidechat(ws, name): + """Open the side chat by clicking its row div inside the sidebar section. + + Clicking the row (not a text search over the whole document) avoids + hitting message text that happens to contain the chat name. + """ + ensure_sidebar_open(ws) + b64 = base64.b64encode(name.encode()).decode() + return ev1(ws, """(async () => { + const nm = atob('%s'); + const T = (el) => (el.innerText || '').trim(); + const target = [...document.querySelectorAll('div')].find(el => { + const cn = (el.className || '').toString(); + if (cn.indexOf('nav-row') === -1) return false; + const m = T(el).match(/^([^\\n]+)\\n\\n(.+)$/); + return m && m[1].trim() === nm; + }); + if (!target) return 'NOTFOUND'; + target.click(); + await new Promise(r => setTimeout(r, 4000)); + const href = window.location.href; + if (href.indexOf('/thread/') === -1) return 'NOTFOUND'; + return href; + })()""" % b64, await_p=True) + + +def navigate_main(ws): + ws.send(json.dumps({"id": 2, "method": "Page.navigate", + "params": {"url": "https://muse.ai/"}})) + for _ in range(20): + resp = json.loads(ws.recv()) + if resp.get("id") == 2: + break + time.sleep(6) + + +def main(): + if len(sys.argv) < 3: + print(json.dumps({"ok": False, "error": "usage"})) + sys.exit(1) + port, mode = sys.argv[1], sys.argv[2] + name = sys.argv[3] if len(sys.argv) > 3 else "" + try: + ws = connect(port) + except Exception as e: # noqa: BLE001 + print(json.dumps({"ok": False, "error": "cdp_connect_failed", + "detail": str(e)[:120]})) + return + if ws is None: + print(json.dumps({"ok": False, "error": "no_page"})) + return + try: + if mode == "list": + print(json.dumps({"ok": True, + "sidechats": list_sidechats(ws)})) + elif mode == "sidechat": + href = open_sidechat(ws, name) + if href == "NOTFOUND": + print(json.dumps({"ok": False, "error": "sidechat_not_found", + "name": name})) + else: + cur = read_current(ws) + cur.update({"ok": True, "context": "sidechat", "name": name}) + print(json.dumps(cur)) + else: + navigate_main(ws) + cur = read_current(ws) + cur.update({"ok": True, "context": "main", "name": "Main chat"}) + print(json.dumps(cur)) + except Exception as e: # noqa: BLE001 + print(json.dumps({"ok": False, "error": "probe_failed", + "detail": str(e)[:200]})) + finally: + try: + ws.close() + except Exception: # noqa: BLE001 + pass + + +main() diff --git a/bin/chat-state-report.py b/bin/chat-state-report.py new file mode 100755 index 0000000..9152b9d --- /dev/null +++ b/bin/chat-state-report.py @@ -0,0 +1,129 @@ +#!/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()