diff --git a/bin/box-ctl.py b/bin/box-ctl.py index f027092..e18ee63 100755 --- a/bin/box-ctl.py +++ b/bin/box-ctl.py @@ -650,6 +650,30 @@ def act_notify(agent, message, sidechat=None, allow_main_chat=False): out(True, agent=agent, target=target, sent=True) +def act_main_loop(sub): + """Run the main-chat self-monitor loop (self_main_loop.py). + + check: one loop iteration -> JSON {ok, prompted, errors, new_total, agents} + status: watermark / last run / prompt-sidechat config, no chat reads. + """ + if sub not in ("check", "status"): + fail("BAD_NAME", "usage: main-loop check|status") + audit("main-loop-" + sub) + script = BIN / "self_main_loop.py" + try: + r = subprocess.run([sys.executable, str(script), sub], + capture_output=True, text=True, timeout=600) + except subprocess.TimeoutExpired: + fail("LOOP_TIMEOUT", "self_main_loop.py timed out") + try: + data = json.loads(r.stdout.strip()) + except Exception: + fail("LOOP_ERROR", "self_main_loop.py returned non-JSON", + {"stdout": (r.stdout or "")[-300:], "stderr": (r.stderr or "")[-300:]}) + ok = bool(data.pop("ok", False)) + out(ok, **data) + + def act_watchdog_alerts(): """Check for new failed browser relaunches since the watermark. @@ -1147,6 +1171,7 @@ fleet: relay-health cdp-latency identity-audit VM identity audit drift check + main-loop check|status main-chat self-monitor: prompt sidechat on new main msgs timer actions: timer-list @@ -1389,6 +1414,10 @@ def main(argv): if rest: fail("BAD_ARGS", "usage: chrome-errors") act_chrome_errors() + elif action == "main-loop": + if len(rest) != 1 or rest[0] not in ("check", "status"): + fail("BAD_NAME", "usage: main-loop check|status") + act_main_loop(rest[0]) elif action == "dm-log": limit = 50 if rest: diff --git a/bin/self_main_loop.py b/bin/self_main_loop.py new file mode 100755 index 0000000..da1e623 --- /dev/null +++ b/bin/self_main_loop.py @@ -0,0 +1,289 @@ +#!/usr/bin/env python3 +"""self_main_loop.py — main-chat self-monitor loop. + +Failure mode it escapes: main chat is the coordination surface, but if no +agent is actively watching it, questions / job assignments / alerts posted +there go unanswered. + +How it works: + 1. box drives a read of each watched agent's muse.ai Main chat on a timer + (via box-chat.py thread-messages main — existing machinery). + 2. On NEW messages since the per-agent watermark, the loop composes a + concise digest and posts it to that agent's prompting sidechat + (via dm.py send, opm browser as neutral sender — mirrors box notify). + 3. The prompt in the sidechat triggers the operator to DM / use box + against main chat. In such, we escape the failure mode. + +State: /home/super/Projects/NetVM/self-main-loop-watermark.json + { "config": {...}, "agents": {agent: {"last_id","last_ts"}}, + "last_run": iso, "last_result": {...} } +First run per agent starts at the current newest message (no backfill spam). + +Contract: + - No new messages -> silent (no sidechat post), exit 0. + - Digests sent -> exit 1. + - Read/send error -> logged to stderr, exit 2 (timer stays alive). + - Overlap guard: fcntl LOCK_EX|LOCK_NB on a lockfile; a second + concurrent run prints {"ok": false, "skipped": "already running"} + and exits 0. + +stdout is always a single JSON object (box-ctl.py parses it); logs go to +stderr (systemd journal). +""" + +import fcntl +import json +import os +import subprocess +import sys +from datetime import datetime, timezone + +BASE = "/home/super/Projects/NetVM" +BIN = os.path.join(BASE, "bin") +BOX_CHAT = os.path.join(BIN, "box-chat.py") +DM_PY = os.path.join(BIN, "dm.py") +STATE_FILE = os.path.join(BASE, "self-main-loop-watermark.json") +LOCK_FILE = os.path.join(BASE, "self-main-loop.lock") + +DEFAULT_AGENTS = ["muse", "pip", "646", "opm"] +# Prompting sidechat per agent — mirrors box-ctl.py NOTIFY_SIDECHATS. +DEFAULT_PROMPT_SIDECHAT = { + "646": "646 tasks", + "opm": "heartbeat", + "pip": "646-pip-coord", + "muse": "muse tasks", +} +DEFAULT_SENDER = "opm" # neutral sender, mirrors `box notify` +READ_LIMIT = 30 +READ_TIMEOUT = 120 +SEND_TIMEOUT = 180 +DIGEST_MAX = 600 # well under dm.py's 1000-char non-raw truncation +PREVIEW_MAX = 120 + + +def utcnow(): + return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") + + +def log(msg): + print("self-main-loop: %s" % msg, file=sys.stderr) + + +def load_state(): + try: + with open(STATE_FILE) as f: + st = json.load(f) + if isinstance(st, dict): + return st + except (FileNotFoundError, json.JSONDecodeError, ValueError): + pass + return {} + + +def save_state(st): + tmp = STATE_FILE + ".tmp" + with open(tmp, "w") as f: + json.dump(st, f, indent=2) + os.replace(tmp, STATE_FILE) + + +def get_config(st): + cfg = st.get("config") or {} + agents = cfg.get("agents") or list(DEFAULT_AGENTS) + agents = [a for a in agents if a in DEFAULT_AGENTS] + prompt = dict(DEFAULT_PROMPT_SIDECHAT) + prompt.update(cfg.get("prompt_sidechat") or {}) + return { + "agents": agents or list(DEFAULT_AGENTS), + "prompt_sidechat": prompt, + "sender": cfg.get("sender") or DEFAULT_SENDER, + "read_limit": int(cfg.get("read_limit") or READ_LIMIT), + } + + +def read_main_chat(agent, limit): + """Read agent's muse.ai Main chat via box-chat.py. Returns (ok, messages|error).""" + cmd = [sys.executable, BOX_CHAT, "thread-messages", agent, "main", + "--limit", str(limit)] + try: + r = subprocess.run(cmd, capture_output=True, text=True, timeout=READ_TIMEOUT) + except subprocess.TimeoutExpired: + return False, "box-chat.py timed out after %ss" % READ_TIMEOUT + except OSError as e: + return False, "could not exec box-chat.py: %s" % e + try: + data = json.loads(r.stdout.strip()) + except Exception: + return False, "box-chat.py non-JSON output: %s" % (r.stdout or r.stderr)[-200:] + if not data.get("ok"): + return False, data.get("error", "box-chat.py failed") + return True, data.get("messages") or [] + + +def new_messages(messages, wm): + """Messages newer than the watermark (oldest-first list).""" + if not messages: + return [] + last_id = (wm or {}).get("last_id") + if not last_id: + return [] # first run: anchor at newest, no backfill + ids = [m.get("id") for m in messages] + if last_id in ids: + return messages[ids.index(last_id) + 1:] + # Watermark id fell out of the read window: fall back to ts compare. + last_ts = (wm or {}).get("last_ts") or "" + return [m for m in messages if (m.get("ts") or "") > last_ts + and m.get("id") != last_id] + + +def author_of(m): + frm = m.get("from") or {} + return frm.get("name") or "?" + + +def compose_digest(agent, new): + """Concise, actionable digest for the prompting sidechat.""" + lines = ["[main-loop] %d new in %s main chat:" % (len(new), agent)] + for m in new[:5]: + text = (m.get("text") or "").strip().replace("\n", " ") + q = " [?]" if "?" in text else "" + hot = " [!]" if ("operator" in text.lower() or "urgent" in text.lower()) else "" + lines.append("- %s: %s%s%s" % (author_of(m), text[:PREVIEW_MAX], q, hot)) + if len(new) > 5: + lines.append("(+%d more)" % (len(new) - 5)) + lines.append("Check main chat via box when you can.") + digest = "\n".join(lines) + return digest[:DIGEST_MAX] + + +def send_prompt(sender, agent, sidechat, digest): + """Post the digest to the agent's prompting sidechat via dm.py.""" + cmd = [sys.executable, DM_PY, "send", + "--agent", sender, "--to", agent, "--target", sidechat, digest] + try: + r = subprocess.run(cmd, capture_output=True, text=True, timeout=SEND_TIMEOUT) + except subprocess.TimeoutExpired: + return False, "dm.py send timed out after %ss" % SEND_TIMEOUT + except OSError as e: + return False, "could not exec dm.py: %s" % e + if r.returncode != 0: + tail = ((r.stderr or "") + (r.stdout or "")).strip()[-300:] + return False, "dm.py send failed rc=%d: %s" % (r.returncode, tail) + return True, "sent" + + +def do_check(only_agent=None): + st = load_state() + cfg = get_config(st) + agents = [only_agent] if only_agent else cfg["agents"] + if only_agent and only_agent not in DEFAULT_AGENTS: + return {"ok": False, "error": "unknown agent: %s" % only_agent} + + agents_state = st.get("agents") or {} + results = {} + prompted = 0 + errors = 0 + + for agent in agents: + wm = agents_state.get(agent) or {} + ok, payload = read_main_chat(agent, cfg["read_limit"]) + if not ok: + log("%s: READ FAILED: %s" % (agent, payload)) + results[agent] = {"ok": False, "error": payload} + errors += 1 + continue + messages = payload + new = new_messages(messages, wm) + newest = messages[-1] if messages else None + if not new: + # Silent: advance watermark to newest seen, no sidechat post. + if newest: + agents_state[agent] = {"last_id": newest.get("id"), + "last_ts": newest.get("ts")} + results[agent] = {"ok": True, "new": 0} + continue + digest = compose_digest(agent, new) + sidechat = cfg["prompt_sidechat"].get(agent) + if not sidechat: + log("%s: no prompting sidechat configured, skipping" % agent) + results[agent] = {"ok": False, "error": "no prompting sidechat"} + errors += 1 + continue + sent, detail = send_prompt(cfg["sender"], agent, sidechat, digest) + if sent: + agents_state[agent] = {"last_id": new[-1].get("id"), + "last_ts": new[-1].get("ts")} + prompted += 1 + log("%s: prompted %s with %d new" % (agent, sidechat, len(new))) + results[agent] = {"ok": True, "new": len(new), "prompted": sidechat} + else: + # Keep watermark: next iteration retries the same digest. + log("%s: PROMPT SEND FAILED: %s" % (agent, detail)) + results[agent] = {"ok": False, "error": detail, "new": len(new)} + errors += 1 + + st["agents"] = agents_state + st["last_run"] = utcnow() + st["last_result"] = {"prompted": prompted, "errors": errors, + "agents": results} + # Preserve operator-edited config; seed defaults on first run. + if "config" not in st: + st["config"] = {"agents": cfg["agents"], + "prompt_sidechat": cfg["prompt_sidechat"], + "sender": cfg["sender"], + "read_limit": cfg["read_limit"]} + save_state(st) + return {"ok": errors == 0, "prompted": prompted, "errors": errors, + "new_total": sum(r.get("new", 0) for r in results.values()), + "agents": results} + + +def do_status(): + st = load_state() + cfg = get_config(st) + agents_state = st.get("agents") or {} + return {"ok": True, + "config": cfg, + "watermark": agents_state, + "last_run": st.get("last_run"), + "last_result": st.get("last_result")} + + +def main(argv): + action = argv[1] if len(argv) > 1 else "check" + only_agent = None + if "--agent" in argv: + i = argv.index("--agent") + if i + 1 < len(argv): + only_agent = argv[i + 1] + + if action == "status": + print(json.dumps(do_status())) + return 0 + + if action != "check": + print(json.dumps({"ok": False, "error": "usage: check|status [--agent X]"})) + return 2 + + # Overlap guard: timer may fire while a slow run is still going. + try: + lockfh = open(LOCK_FILE, "w") + fcntl.flock(lockfh, fcntl.LOCK_EX | fcntl.LOCK_NB) + except (OSError, IOError): + print(json.dumps({"ok": False, "skipped": "already running"})) + return 0 + + try: + result = do_check(only_agent) + except Exception as e: # never crash the timer + log("UNEXPECTED ERROR: %r" % e) + print(json.dumps({"ok": False, "error": "unexpected: %s" % e})) + return 2 + print(json.dumps(result)) + if not result.get("ok"): + return 2 + return 1 if result.get("prompted") else 0 + + +if __name__ == "__main__": + sys.exit(main(sys.argv))