chat-state tap: per-node main+sidechat reporter for box
chat-state-check.py: polls each node Chromium via CDP, captures main chat + up to 5 side chats (last 10 DOM-visible messages each, 500 chars). Structural React selectors; picks the most common non-empty list across repeated renders. chat-state-report.py: signs and POSTs the report to the board /api/box/chat-state/report ingest (namespace health, 5-min timer). Session: sidechat/box-console
This commit is contained in:
Executable
+281
@@ -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 <cdp_port> <mode> [name]
|
||||
modes: main | list | sidechat <name>
|
||||
|
||||
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()
|
||||
Executable
+129
@@ -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
|
||||
<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()
|
||||
Reference in New Issue
Block a user