bin/box-chat.py + box-chat-cdp.py — read-only bl helper for Box thread oversight
Session: sidechat/box-chat
This commit is contained in:
Executable
+305
@@ -0,0 +1,305 @@
|
||||
#!/usr/bin/env python3
|
||||
"""box-chat-cdp.py — READ-ONLY CDP scraper for box-chat.py.
|
||||
|
||||
Runs INSIDE the agent's netns (via netvm-exec.sh), where the agent's
|
||||
headless Chromium CDP port is reachable on 127.0.0.1. Performs only
|
||||
Runtime.evaluate reads plus benign navigation clicks (panel open, thread
|
||||
switch, restore-to-main). Never sends messages, never touches the
|
||||
composer, never clicks send.
|
||||
|
||||
Usage:
|
||||
box-chat-cdp.py <agent> threads
|
||||
box-chat-cdp.py <agent> messages <thread-id>
|
||||
|
||||
Prints exactly one JSON object to stdout. Exit 0 on success, 1 on
|
||||
failure (stdout still carries {"ok": false, ...}).
|
||||
"""
|
||||
|
||||
import importlib.util
|
||||
import json
|
||||
import sys
|
||||
import time
|
||||
import urllib.request
|
||||
|
||||
|
||||
def load_accounts():
|
||||
path = "/home/super/Projects/NetVM/bin/netvm-registry.py"
|
||||
spec = importlib.util.spec_from_file_location("netvm_registry", path)
|
||||
mod = importlib.util.module_from_spec(spec)
|
||||
spec.loader.exec_module(mod)
|
||||
accounts = {}
|
||||
for node, rec in mod.load().items():
|
||||
accounts[node] = (node, "http://127.0.0.1:%d/json/list" % rec["cdp_port"])
|
||||
return accounts
|
||||
|
||||
|
||||
def connect(agent):
|
||||
accounts = load_accounts()
|
||||
if agent not in accounts:
|
||||
raise RuntimeError("unknown agent in registry: %s" % agent)
|
||||
_node, cdp_url = accounts[agent]
|
||||
with urllib.request.urlopen(cdp_url, timeout=10) as r:
|
||||
targets = json.load(r)
|
||||
pages = [t for t in targets if t.get("type") == "page"]
|
||||
if not pages:
|
||||
raise RuntimeError("no page target on CDP")
|
||||
import websocket
|
||||
ws = websocket.create_connection(
|
||||
pages[0]["webSocketDebuggerUrl"], timeout=60)
|
||||
return ws
|
||||
|
||||
|
||||
def ev(ws, expr, await_p=False):
|
||||
ws.send(json.dumps({
|
||||
"id": 1, "method": "Runtime.evaluate",
|
||||
"params": {"expression": expr, "returnByValue": True,
|
||||
"awaitPromise": await_p},
|
||||
}))
|
||||
resp = json.loads(ws.recv())
|
||||
res = resp.get("result", {})
|
||||
if res.get("subtype") == "error":
|
||||
raise RuntimeError("JS error: %s" % str(res.get("description"))[:200])
|
||||
return res.get("result", {}).get("value")
|
||||
|
||||
|
||||
TITLE_OF = """const titleOf = r => {
|
||||
const s = r.querySelector('span[title]');
|
||||
return s ? s.getAttribute('title').trim()
|
||||
: ((r.innerText||'').split('\\n')[0]||'').trim();
|
||||
};"""
|
||||
|
||||
ENSURE = """(async () => {
|
||||
const sleep = ms => new Promise(r => setTimeout(r, ms));
|
||||
%s
|
||||
if (window.location.pathname === '/thread/new') {
|
||||
window.location.href = '/'; await sleep(4000);
|
||||
}
|
||||
const nav = document.querySelector('[data-testid="hatch-nav-chat"]');
|
||||
if (!nav) return 'NO_CHAT_NAV';
|
||||
if (nav.getAttribute('aria-current') !== 'page') { nav.click(); await sleep(3000); }
|
||||
const panelOpen = () => !!document.querySelector('[data-testid="hatch-chat-compose"]');
|
||||
if (!panelOpen()) {
|
||||
const sw = document.querySelector('[data-testid="hatch-chat-switcher-trigger"]');
|
||||
if (!sw) return 'NO_SWITCHER';
|
||||
sw.click(); await sleep(2500);
|
||||
if (!panelOpen()) return 'PANEL_CLOSED';
|
||||
}
|
||||
return 'OK';
|
||||
})()""" % TITLE_OF
|
||||
|
||||
|
||||
def op_threads(ws):
|
||||
st = ev(ws, ENSURE, await_p=True)
|
||||
if st != "OK":
|
||||
raise RuntimeError("could not reach chat panel: %s" % st)
|
||||
js = """(async () => {
|
||||
const sleep = ms => new Promise(r => setTimeout(r, ms));
|
||||
%s
|
||||
const snap = [...document.querySelectorAll('[data-testid="hatch-thread-row"]')]
|
||||
.map(r => {
|
||||
const spans = [...r.querySelectorAll('span')];
|
||||
return {title: titleOf(r),
|
||||
rel: spans.length ? (spans[spans.length-1].innerText||'').trim() : ''};
|
||||
});
|
||||
const ensurePanel = async () => {
|
||||
if (document.querySelector('[data-testid="hatch-chat-compose"]')) return true;
|
||||
for (let k = 0; k < 3; k++) {
|
||||
const sw = document.querySelector('[data-testid="hatch-chat-switcher-trigger"]');
|
||||
if (!sw) return false;
|
||||
sw.click(); await sleep(3000);
|
||||
if (document.querySelector('[data-testid="hatch-chat-compose"]')) return true;
|
||||
}
|
||||
return false;
|
||||
};
|
||||
const out = [];
|
||||
for (const s of snap) {
|
||||
const isMain = /^main chat$/i.test(s.title);
|
||||
if (isMain) { out.push({id: 'main', kind: 'main', title: s.title, rel: s.rel}); continue; }
|
||||
const panelOk = await ensurePanel();
|
||||
if (!panelOk) { out.push({id: null, kind: 'sidechat', title: s.title, rel: s.rel, error: 'PANEL_CLOSED'}); continue; }
|
||||
const row = [...document.querySelectorAll('[data-testid="hatch-thread-row"]')]
|
||||
.find(r => titleOf(r) === s.title);
|
||||
if (!row) { out.push({id: null, kind: 'sidechat', title: s.title, rel: s.rel, error: 'ROW_GONE'}); continue; }
|
||||
row.click();
|
||||
let id = null;
|
||||
for (let t = 0; t < 20; t++) {
|
||||
await sleep(500);
|
||||
const m = window.location.href.match(/\\/thread\\/([0-9a-fA-F-]{36})/);
|
||||
if (m) { id = m[1]; break; }
|
||||
}
|
||||
if (!id) {
|
||||
const row2 = [...document.querySelectorAll('[data-testid="hatch-thread-row"]')]
|
||||
.find(r => titleOf(r) === s.title);
|
||||
if (row2) {
|
||||
row2.click();
|
||||
for (let t = 0; t < 12; t++) {
|
||||
await sleep(500);
|
||||
const m = window.location.href.match(/\\/thread\\/([0-9a-fA-F-]{36})/);
|
||||
if (m) { id = m[1]; break; }
|
||||
}
|
||||
}
|
||||
}
|
||||
out.push({id, kind: 'sidechat', title: s.title, rel: s.rel});
|
||||
}
|
||||
return {threads: out};
|
||||
})()""" % TITLE_OF
|
||||
return ev(ws, js, await_p=True)
|
||||
|
||||
|
||||
def op_messages(ws, thread_id):
|
||||
st = ev(ws, ENSURE, await_p=True)
|
||||
if st != "OK":
|
||||
raise RuntimeError("could not reach chat panel: %s" % st)
|
||||
js = """(async () => {
|
||||
const sleep = ms => new Promise(r => setTimeout(r, ms));
|
||||
%s
|
||||
const THREAD = %s;
|
||||
const ensurePanel = async () => {
|
||||
if (document.querySelector('[data-testid="hatch-chat-compose"]')) return true;
|
||||
for (let k = 0; k < 3; k++) {
|
||||
const sw = document.querySelector('[data-testid="hatch-chat-switcher-trigger"]');
|
||||
if (!sw) return false;
|
||||
sw.click(); await sleep(3000);
|
||||
if (document.querySelector('[data-testid="hatch-chat-compose"]')) return true;
|
||||
}
|
||||
return false;
|
||||
};
|
||||
const uuidOf = () => {
|
||||
const m = window.location.href.match(/\\/thread\\/([0-9a-fA-F-]{36})/);
|
||||
return m ? m[1].toLowerCase() : null;
|
||||
};
|
||||
let landed = false;
|
||||
if (THREAD === 'main') {
|
||||
await ensurePanel();
|
||||
const row = [...document.querySelectorAll('[data-testid="hatch-thread-row"]')]
|
||||
.find(r => /^main chat$/i.test(titleOf(r)));
|
||||
if (row) { row.click(); await sleep(2500); }
|
||||
landed = /muse\\.ai\\/?$/.test(window.location.href) && !uuidOf();
|
||||
} else {
|
||||
// already there?
|
||||
if (uuidOf() === THREAD.toLowerCase()) landed = true;
|
||||
// click rows until the URL carries our thread uuid (SPA navigation,
|
||||
// keeps the JS context alive unlike location.href assignment)
|
||||
for (let i = 0; i < 12 && !landed; i++) {
|
||||
await ensurePanel();
|
||||
const rows = [...document.querySelectorAll('[data-testid="hatch-thread-row"]')];
|
||||
if (!rows.length) break;
|
||||
const row = rows[i %% rows.length];
|
||||
row.click();
|
||||
for (let t = 0; t < 14; t++) {
|
||||
await sleep(500);
|
||||
if (uuidOf() === THREAD.toLowerCase()) { landed = true; break; }
|
||||
if (uuidOf()) break; // navigated somewhere else; try next row
|
||||
}
|
||||
}
|
||||
if (!landed) {
|
||||
// fallback: SPA history navigation (no full page load)
|
||||
window.history.pushState({}, '', '/thread/' + THREAD);
|
||||
window.dispatchEvent(new PopStateEvent('popstate'));
|
||||
await sleep(5000);
|
||||
landed = uuidOf() === THREAD.toLowerCase();
|
||||
}
|
||||
}
|
||||
if (!landed) return {error: 'THREAD_NOT_FOUND'};
|
||||
await sleep(2500);
|
||||
const sc = document.getElementById('hatch-chat-scroll');
|
||||
if (sc) {
|
||||
for (let i = 0; i < 3; i++) { sc.scrollTop = 0; await sleep(1500); }
|
||||
sc.scrollTop = sc.scrollHeight; await sleep(800);
|
||||
}
|
||||
const els = [...document.querySelectorAll('[data-message-id]')];
|
||||
const messages = els.map(m => {
|
||||
const id = m.getAttribute('data-message-id');
|
||||
const ps = [...m.querySelectorAll('p')].map(p => (p.innerText||'').trim()).filter(Boolean);
|
||||
let text = ps.join('\\n');
|
||||
if (!text) text = (m.innerText||'').replace(/^(Assistant message:|User message:)\\s*/, '').trim();
|
||||
const t = m.querySelector('time');
|
||||
return {id,
|
||||
author: id.indexOf('assistant-msg') === 0 ? 'assistant' : 'user',
|
||||
text,
|
||||
ts: t ? (t.getAttribute('datetime') || t.innerText || null) : null};
|
||||
});
|
||||
return {url: window.location.href, messages};
|
||||
})()""" % (TITLE_OF, json.dumps(thread_id))
|
||||
data = ev(ws, js, await_p=True)
|
||||
if not isinstance(data, dict) or "messages" not in data:
|
||||
if isinstance(data, dict) and data.get("error") == "THREAD_NOT_FOUND":
|
||||
raise RuntimeError("THREAD_NOT_FOUND: no such thread for this agent")
|
||||
raise RuntimeError("unexpected messages payload")
|
||||
return data
|
||||
|
||||
|
||||
def restore_main(ws):
|
||||
try:
|
||||
ev(ws, """(() => {
|
||||
%s
|
||||
const row = [...document.querySelectorAll('[data-testid="hatch-thread-row"]')]
|
||||
.find(r => /^main chat$/i.test(titleOf(r)));
|
||||
if (row) row.click();
|
||||
return 'ok';
|
||||
})()""" % TITLE_OF)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def main(argv):
|
||||
if len(argv) < 3:
|
||||
print(json.dumps({"ok": False, "code": "BAD_ARGS",
|
||||
"error": "usage: box-chat-cdp.py <agent> threads|messages [thread-id]"}))
|
||||
return 1
|
||||
agent, op = argv[1], argv[2]
|
||||
ws = None
|
||||
try:
|
||||
ws = connect(agent)
|
||||
# CDP evaluate can rarely resolve null when the page is mid-navigation;
|
||||
# one retry covers the flake (observed 2026-10-04).
|
||||
data, last_err = None, None
|
||||
for attempt in range(2):
|
||||
try:
|
||||
if op == "threads":
|
||||
data = op_threads(ws)
|
||||
ok = isinstance(data, dict) and isinstance(
|
||||
data.get("threads"), list)
|
||||
elif op == "messages":
|
||||
if len(argv) < 4:
|
||||
raise RuntimeError("messages requires thread-id")
|
||||
data = op_messages(ws, argv[3])
|
||||
ok = isinstance(data, dict) and isinstance(
|
||||
data.get("messages"), list)
|
||||
else:
|
||||
raise RuntimeError("unknown op: %s" % op)
|
||||
if ok:
|
||||
break
|
||||
last_err = "empty CDP result"
|
||||
data = None
|
||||
except RuntimeError as e:
|
||||
last_err = str(e)
|
||||
data = None
|
||||
time.sleep(3)
|
||||
if data is None:
|
||||
raise RuntimeError(last_err or "CDP returned no usable data")
|
||||
if op == "threads":
|
||||
print(json.dumps({"ok": True, "agent": agent,
|
||||
"threads": data["threads"]}))
|
||||
else:
|
||||
print(json.dumps({"ok": True, "agent": agent, "url": data["url"],
|
||||
"messages": data["messages"]}))
|
||||
return 0
|
||||
except Exception as e:
|
||||
print(json.dumps({"ok": False, "code": "CDP_ERROR",
|
||||
"error": str(e)[:300]}))
|
||||
return 1
|
||||
finally:
|
||||
if ws is not None:
|
||||
try:
|
||||
restore_main(ws)
|
||||
except Exception:
|
||||
pass
|
||||
try:
|
||||
ws.close()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main(sys.argv))
|
||||
Executable
+380
@@ -0,0 +1,380 @@
|
||||
#!/usr/bin/env python3
|
||||
"""box-chat.py — allowlisted bl helper for Box thread oversight (READ-ONLY).
|
||||
|
||||
The board server (VM) never scrapes browsers directly. All thread reads go
|
||||
through this helper, invoked as:
|
||||
|
||||
/home/super/Projects/NetVM/bin/box-chat.py thread-list <agent> [--kind K] [--limit N]
|
||||
/home/super/Projects/NetVM/bin/box-chat.py thread-messages <agent> <thread-id> [--limit N] [--before MSGID]
|
||||
|
||||
Security properties (mirrors box-ctl.py):
|
||||
- Fixed verb set; every argument validated before acting.
|
||||
- <agent> must be a known node (muse, pip, 646, opm); anything else exits
|
||||
before any netns/SSH/CDP work.
|
||||
- <thread-id> must match ^[a-zA-Z0-9-]{1,64}$ or be the literal "main".
|
||||
- READ-ONLY by construction: the CDP driver (box-chat-cdp.py) only runs
|
||||
Runtime.evaluate reads plus benign navigation clicks. No sends, no
|
||||
composer interaction, no shell=True anywhere. All subprocess calls use
|
||||
argv lists.
|
||||
- Every invocation audit-logged to box-chat.jsonl with caller identity.
|
||||
|
||||
Output: JSON to stdout ({"ok": true, ...} or {"ok": false, ...}),
|
||||
exit 0 on success, nonzero on failure.
|
||||
"""
|
||||
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import subprocess
|
||||
import sys
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from pathlib import Path
|
||||
|
||||
NETVM_ROOT = Path("/home/super/Projects/NetVM")
|
||||
BIN = NETVM_ROOT / "bin"
|
||||
CDP_DRIVER = BIN / "box-chat-cdp.py"
|
||||
NETVM_EXEC = BIN / "netvm-exec.sh"
|
||||
CHAT_LOG = NETVM_ROOT / "box-chat.jsonl"
|
||||
DM_LOG = NETVM_ROOT / "dm-log.jsonl"
|
||||
|
||||
VALID_AGENTS = {"muse", "pip", "646", "opm"}
|
||||
VALID_KINDS = {"main", "sidechat", "dm", "all"}
|
||||
AGENT_RE = re.compile(r"^[a-z0-9-]{1,64}$")
|
||||
THREAD_RE = re.compile(r"^[a-zA-Z0-9-]{1,64}$")
|
||||
MSGID_RE = re.compile(r"^[a-zA-Z0-9-]{1,128}$")
|
||||
|
||||
REL_MONTHS = {m: i + 1 for i, m in enumerate(
|
||||
["jan", "feb", "mar", "apr", "may", "jun",
|
||||
"jul", "aug", "sep", "oct", "nov", "dec"])}
|
||||
|
||||
|
||||
def utcnow():
|
||||
return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
|
||||
|
||||
|
||||
def out(ok, **kw):
|
||||
payload = {"ok": ok}
|
||||
payload.update(kw)
|
||||
print(json.dumps(payload))
|
||||
|
||||
|
||||
def fail(code, error, detail=None, exit_code=1):
|
||||
payload = {"ok": False, "code": code, "error": error}
|
||||
if detail is not None:
|
||||
payload["detail"] = detail
|
||||
print(json.dumps(payload))
|
||||
sys.exit(exit_code)
|
||||
|
||||
|
||||
def audit(action, agent=None, name=None, kind=None):
|
||||
"""Append invocation record to the bl-side log."""
|
||||
try:
|
||||
entry = {
|
||||
"ts": utcnow(),
|
||||
"action": action,
|
||||
"agent": agent,
|
||||
"name": name,
|
||||
"kind": kind,
|
||||
"caller": os.environ.get("BOX_CALLER", "unknown"),
|
||||
}
|
||||
with open(CHAT_LOG, "a") as f:
|
||||
f.write(json.dumps(entry) + "\n")
|
||||
except Exception:
|
||||
pass # audit failure must not break the action
|
||||
|
||||
|
||||
def check_agent(agent):
|
||||
if not agent or agent not in VALID_AGENTS:
|
||||
fail("BAD_AGENT", "agent must be one of %s" % sorted(VALID_AGENTS),
|
||||
{"field": "agent", "value": agent})
|
||||
return agent
|
||||
|
||||
|
||||
def check_thread_id(tid):
|
||||
if tid != "main" and not (tid and THREAD_RE.match(tid)):
|
||||
fail("BAD_THREAD", "thread-id must be 'main' or match ^[a-zA-Z0-9-]{1,64}$",
|
||||
{"field": "thread-id", "value": tid})
|
||||
return tid
|
||||
|
||||
|
||||
def check_limit(val, default):
|
||||
if val is None:
|
||||
return default
|
||||
try:
|
||||
n = int(val)
|
||||
except (TypeError, ValueError):
|
||||
fail("BAD_LIMIT", "limit must be an integer 1-200", {"value": val})
|
||||
if not 1 <= n <= 200:
|
||||
fail("BAD_LIMIT", "limit must be an integer 1-200", {"value": val})
|
||||
return n
|
||||
|
||||
|
||||
def run_cdp(agent, op, args, timeout):
|
||||
"""Run the CDP driver inside the agent's netns. argv only, no shell."""
|
||||
cmd = [str(NETVM_EXEC), agent, "--", sys.executable,
|
||||
str(CDP_DRIVER), agent, op] + args
|
||||
try:
|
||||
r = subprocess.run(cmd, capture_output=True, text=True,
|
||||
timeout=timeout)
|
||||
except subprocess.TimeoutExpired:
|
||||
fail("CDP_TIMEOUT", "CDP driver timed out in netns",
|
||||
{"agent": agent, "op": op})
|
||||
if r.returncode != 0 and not r.stdout.strip():
|
||||
fail("CDP_ERROR", "CDP driver failed",
|
||||
{"stderr": (r.stderr or "")[-300:]})
|
||||
try:
|
||||
data = json.loads(r.stdout.strip())
|
||||
except Exception:
|
||||
fail("CDP_ERROR", "CDP driver returned non-JSON",
|
||||
{"stdout": r.stdout[-300:], "stderr": (r.stderr or "")[-300:]})
|
||||
if not data.get("ok"):
|
||||
code = data.get("code", "CDP_ERROR")
|
||||
if "THREAD_NOT_FOUND" in str(data.get("error", "")):
|
||||
code = "THREAD_NOT_FOUND"
|
||||
fail(code, data.get("error", "cdp failed"))
|
||||
return data
|
||||
|
||||
|
||||
def parse_rel_time(rel):
|
||||
"""Panel relative times ('just now', '4m', '2h', '3d', 'Oct 3') ->
|
||||
(approx ISO8601, True). Returns (None, False) when unparseable."""
|
||||
now = datetime.now(timezone.utc)
|
||||
rel = (rel or "").strip().lower()
|
||||
if rel in ("just now", "now"):
|
||||
return now.strftime("%Y-%m-%dT%H:%M:%SZ"), True
|
||||
m = re.match(r"^(\d+)\s*m(in(ute)?s?)?$", rel)
|
||||
if m:
|
||||
return (now - timedelta(minutes=int(m.group(1)))).strftime("%Y-%m-%dT%H:%M:%SZ"), True
|
||||
m = re.match(r"^(\d+)\s*h((ou)?rs?)?$", rel)
|
||||
if m:
|
||||
return (now - timedelta(hours=int(m.group(1)))).strftime("%Y-%m-%dT%H:%M:%SZ"), True
|
||||
m = re.match(r"^(\d+)\s*d(ays?)?$", rel)
|
||||
if m:
|
||||
return (now - timedelta(days=int(m.group(1)))).strftime("%Y-%m-%dT%H:%M:%SZ"), True
|
||||
m = re.match(r"^([a-z]{3})\s+(\d{1,2})$", rel)
|
||||
if m and m.group(1) in REL_MONTHS:
|
||||
dt = now.replace(month=REL_MONTHS[m.group(1)], day=int(m.group(2)),
|
||||
hour=12, minute=0, second=0, microsecond=0)
|
||||
if dt > now:
|
||||
dt = dt.replace(year=dt.year - 1)
|
||||
return dt.strftime("%Y-%m-%dT%H:%M:%SZ"), True
|
||||
return None, False
|
||||
|
||||
|
||||
def dm_conversations(agent):
|
||||
"""DM conversations involving <agent>, synthesized from dm-log.jsonl.
|
||||
|
||||
Real timestamps and counts; bodies are not stored in the log (by design)
|
||||
— message text for these threads is the wire tag, and full bodies live
|
||||
in the target chat's messages.
|
||||
"""
|
||||
convos = {}
|
||||
try:
|
||||
with open(DM_LOG) as f:
|
||||
for line in f:
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
try:
|
||||
e = json.loads(line)
|
||||
except Exception:
|
||||
continue
|
||||
if e.get("type") not in ("sent", "send_done"):
|
||||
continue
|
||||
frm, to = e.get("agent"), e.get("to")
|
||||
if not frm or not to:
|
||||
continue
|
||||
if agent not in (frm, to):
|
||||
continue
|
||||
key = tuple(sorted([frm, to]))
|
||||
c = convos.setdefault(key, {"count": 0, "last_ts": "",
|
||||
"entries": []})
|
||||
c["count"] += 1
|
||||
if e.get("ts", "") > c["last_ts"]:
|
||||
c["last_ts"] = e["ts"]
|
||||
c["entries"].append(e)
|
||||
except FileNotFoundError:
|
||||
return []
|
||||
out = []
|
||||
for (a, b), c in sorted(convos.items()):
|
||||
out.append({
|
||||
"id": "dm-%s-%s" % (a, b),
|
||||
"kind": "dm",
|
||||
"title": "dm:%s:%s" % (a, b),
|
||||
"participants": [a, b],
|
||||
"last_message_at": c["last_ts"] or None,
|
||||
"last_message_approx": False,
|
||||
"message_count": c["count"],
|
||||
})
|
||||
return out
|
||||
|
||||
|
||||
def dm_thread_messages(agent, thread_id):
|
||||
"""Messages for a dm-<a>-<b> thread, from dm-log.jsonl (complete log)."""
|
||||
parts = thread_id[3:].split("-")
|
||||
if len(parts) != 2:
|
||||
fail("THREAD_NOT_FOUND", "no such DM thread: %s" % thread_id)
|
||||
a, b = parts
|
||||
if agent not in (a, b):
|
||||
fail("THREAD_NOT_FOUND", "no such DM thread: %s" % thread_id)
|
||||
msgs = []
|
||||
try:
|
||||
with open(DM_LOG) as f:
|
||||
for line in f:
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
try:
|
||||
e = json.loads(line)
|
||||
except Exception:
|
||||
continue
|
||||
if e.get("type") not in ("sent", "send_done"):
|
||||
continue
|
||||
frm, to = e.get("agent"), e.get("to")
|
||||
if not frm or not to:
|
||||
continue
|
||||
if tuple(sorted([frm, to])) != (a, b):
|
||||
continue
|
||||
sender = e.get("agent")
|
||||
msgs.append({
|
||||
"id": str(e.get("id", "")),
|
||||
"from": {"role": "agent", "name": sender},
|
||||
"text": "[from:%s] [id:%s] -> %s" % (
|
||||
sender, e.get("id"), e.get("target", "?")),
|
||||
"ts": e.get("ts"),
|
||||
"target": e.get("target"),
|
||||
"verified": e.get("verified"),
|
||||
})
|
||||
except FileNotFoundError:
|
||||
pass
|
||||
msgs.sort(key=lambda m: m.get("ts") or "")
|
||||
# de-dupe send_done/sent pairs on DM id, keep the richer record
|
||||
seen = {}
|
||||
for m in msgs:
|
||||
prev = seen.get(m["id"])
|
||||
if prev is None or (m.get("verified") and not prev.get("verified")):
|
||||
seen[m["id"]] = m
|
||||
msgs = sorted(seen.values(), key=lambda m: m.get("ts") or "")
|
||||
return msgs
|
||||
|
||||
|
||||
def act_thread_list(agent, kind, limit):
|
||||
threads = []
|
||||
if kind in ("main", "sidechat", "all"):
|
||||
data = run_cdp(agent, "threads", [], timeout=240)
|
||||
for t in data.get("threads", [])[:limit]:
|
||||
if kind != "all" and t.get("kind") != kind:
|
||||
continue
|
||||
ts, approx = parse_rel_time(t.get("rel", ""))
|
||||
threads.append({
|
||||
"id": t.get("id"),
|
||||
"kind": t.get("kind"),
|
||||
"title": t.get("title"),
|
||||
"participants": [agent, "human"],
|
||||
"last_message_at": ts,
|
||||
"last_message_approx": approx,
|
||||
"message_count": None, # list is a panel scrape; counts need a thread open
|
||||
})
|
||||
if kind in ("dm", "all"):
|
||||
threads.extend(dm_conversations(agent))
|
||||
audit("thread-list", agent=agent, kind=kind)
|
||||
out(True, agent=agent, kind=kind, threads=threads,
|
||||
fetched_at=utcnow())
|
||||
|
||||
|
||||
def act_thread_messages(agent, thread_id, limit, before):
|
||||
if thread_id.startswith("dm-"):
|
||||
msgs = dm_thread_messages(agent, thread_id)
|
||||
thread = {"id": thread_id, "kind": "dm",
|
||||
"title": "dm:%s" % thread_id[3:].replace("-", ":"),
|
||||
"participants": thread_id[3:].split("-")}
|
||||
else:
|
||||
data = run_cdp(agent, "messages", [thread_id], timeout=150)
|
||||
raw = data.get("messages", [])
|
||||
thread = {"id": thread_id,
|
||||
"kind": "main" if thread_id == "main" else "sidechat",
|
||||
"participants": [agent, "human"]}
|
||||
msgs = [{
|
||||
"id": m.get("id"),
|
||||
"from": ({"role": "agent", "name": agent}
|
||||
if m.get("author") == "assistant"
|
||||
else {"role": "human", "name": "human"}),
|
||||
"text": m.get("text", ""),
|
||||
"ts": m.get("ts"),
|
||||
} for m in raw]
|
||||
if before:
|
||||
if not MSGID_RE.match(before):
|
||||
fail("BAD_CURSOR", "before must match ^[a-zA-Z0-9-]{1,128}$",
|
||||
{"value": before})
|
||||
idx = next((i for i, m in enumerate(msgs) if m["id"] == before), None)
|
||||
if idx is None:
|
||||
fail("BAD_CURSOR", "no message with that id in loaded window",
|
||||
{"value": before})
|
||||
msgs = msgs[:idx]
|
||||
older = len(msgs)
|
||||
if len(msgs) > limit:
|
||||
msgs = msgs[-limit:]
|
||||
next_before = msgs[0]["id"] if older > len(msgs) and msgs else None
|
||||
audit("thread-messages", agent=agent, name=thread_id)
|
||||
out(True, agent=agent, thread=thread, messages=msgs,
|
||||
next_before=next_before, loaded_count=older, fetched_at=utcnow())
|
||||
|
||||
|
||||
USAGE = """usage: box-chat.py <action> [args]
|
||||
|
||||
thread-list <agent> [--kind main|sidechat|dm|all] [--limit N]
|
||||
thread-messages <agent> <thread-id> [--limit N] [--before MSGID]
|
||||
|
||||
read-only. <agent> is one of: muse, pip, 646, opm."""
|
||||
|
||||
|
||||
def parse_flags(rest, names):
|
||||
"""Parse [--flag value] pairs; returns (positionals, {flag: value})."""
|
||||
pos, flags = [], {}
|
||||
i = 0
|
||||
while i < len(rest):
|
||||
tok = rest[i]
|
||||
if tok.startswith("--") and tok[2:] in names:
|
||||
if i + 1 >= len(rest):
|
||||
fail("BAD_ARGS", "flag %s needs a value" % tok)
|
||||
flags[tok[2:]] = rest[i + 1]
|
||||
i += 2
|
||||
elif tok.startswith("--"):
|
||||
fail("BAD_ARGS", "unknown flag: %s" % tok)
|
||||
else:
|
||||
pos.append(tok)
|
||||
i += 1
|
||||
return pos, flags
|
||||
|
||||
|
||||
def main(argv):
|
||||
if len(argv) < 2:
|
||||
print(USAGE, file=sys.stderr)
|
||||
sys.exit(2)
|
||||
action = argv[1]
|
||||
|
||||
if action == "thread-list":
|
||||
pos, flags = parse_flags(argv[2:], {"kind", "limit"})
|
||||
if len(pos) != 1:
|
||||
fail("BAD_ARGS", "usage: thread-list <agent> [--kind K] [--limit N]")
|
||||
agent = check_agent(pos[0])
|
||||
kind = flags.get("kind", "all")
|
||||
if kind not in VALID_KINDS:
|
||||
fail("BAD_ARGS", "kind must be one of %s" % sorted(VALID_KINDS))
|
||||
act_thread_list(agent, kind, check_limit(flags.get("limit"), 50))
|
||||
elif action == "thread-messages":
|
||||
pos, flags = parse_flags(argv[2:], {"limit", "before"})
|
||||
if len(pos) != 2:
|
||||
fail("BAD_ARGS",
|
||||
"usage: thread-messages <agent> <thread-id> [--limit N] [--before MSGID]")
|
||||
agent = check_agent(pos[0])
|
||||
tid = check_thread_id(pos[1])
|
||||
act_thread_messages(agent, tid, check_limit(flags.get("limit"), 50),
|
||||
flags.get("before"))
|
||||
else:
|
||||
print(USAGE, file=sys.stderr)
|
||||
fail("BAD_ARGS", "unknown action: %s" % action)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main(sys.argv)
|
||||
Reference in New Issue
Block a user