#!/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 0 (success; prompted count in JSON). - 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 re import subprocess import sys from contextlib import contextmanager from datetime import datetime, timezone BASE = "/home/super/Projects/NetVM" BIN = os.path.join(BASE, "bin") if BIN not in sys.path: sys.path.insert(0, 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") STATE_LOCK_FILE = os.path.join(BASE, "self-main-loop-state.lock") DEFAULT_AGENTS = ["muse", "pip", "646", "opm"] # Prompting sidechat per agent — mirrors box-ctl.py NOTIFY_SIDECHATS. DEFAULT_PROMPT_SIDECHAT = { "646": "646 tasks", "opm": "main-loop brain", "pip": "pip tasks", "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 # Word-boundary "operator" that does NOT match hyphenated identities like # "operator-646" (the "-" counts as a boundary for \b, so we exclude it # explicitly). "operator needed" matches; "operator-646" does not. OPERATOR_RE = re.compile(r"(? 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 without synthetic ceremony.""" lines = ["New messages in %s main chat (%d):" % (agent, len(new))] for m in new[:5]: text = (m.get("text") or "").strip().replace("\n", " ") q = " (?)" if "?" in text else "" hot = " (!)" if (OPERATOR_RE.search(text) 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 available.") digest = "\n".join(lines) return digest[:DIGEST_MAX] def get_monitored_sidechats(agent): """Return dict of {thread_uuid: name} for sidechats registered for this agent.""" sc_file = os.path.join(BASE, "job-sidechats.json") if not os.path.exists(sc_file): return {} try: with open(sc_file, "r", encoding="utf-8") as f: sc_data = json.load(f) except Exception: return {} monitored = {} for name, item in sc_data.items(): if not isinstance(item, dict): continue uuid = item.get("thread_uuid") or item.get("uuid") if not uuid: continue item_agent = item.get("agent") if name.endswith(f"@{agent}") or (item_agent == agent and "@" not in name): if name.startswith("pipe-") or name.startswith("test-"): continue monitored[uuid] = name return monitored def send_prompt(sender, agent, sidechat, digest): """Post the digest to the agent's prompting sidechat. Tries fast direct gateway send via muse_hybrid first, falling back to dm.py. """ target_uuid = sidechat try: import dm target_uuid = dm.resolve_sidechat_target(sidechat, agent) except Exception: pass # Try fast direct gateway send if target_uuid and re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", target_uuid.lower()): try: import muse_hybrid send_node = agent if agent in DEFAULT_AGENTS else sender res, err = muse_hybrid.send_message(send_node, digest, thread_id=target_uuid, wait=0) if res and not err: return True, "sent (gateway)" log("Gateway send fallback for %s/%s due to: %s" % (agent, sidechat, err or res)) except Exception as e: log("Gateway send exception for %s/%s: %s" % (agent, sidechat, e)) # Fallback to dm.py send 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 (dm.py)" def check_sidechats(agent, cfg, wm, threads_meta): """Tier 1 & Tier 2 unread check on agent's registered sidechats.""" import muse_hybrid monitored = get_monitored_sidechats(agent) prompt_sidechat = cfg["prompt_sidechat"].get(agent) prompt_uuid = None try: import dm prompt_uuid = dm.resolve_sidechat_target(prompt_sidechat, agent) except Exception: pass sc_watermarks = wm.setdefault("sidechats", {}) threads_by_id = {t.get("session_id"): t for t in (threads_meta or []) if isinstance(t, dict)} new_activity = [] for uuid, sc_name in monitored.items(): # Do not monitor prompt destination to avoid loopback if prompt_uuid and uuid.lower() == prompt_uuid.lower(): continue th_meta = threads_by_id.get(uuid) if not th_meta: continue last_updated = sc_watermarks.get(uuid, {}).get("updated") curr_updated = th_meta.get("updated") # First run: if no watermark, initialize to current newest message without alerting if uuid not in sc_watermarks: msgs, err = muse_hybrid.get_history(agent, thread_id=uuid, limit=3) last_id = msgs[-1]["message_id"] if (msgs and not err and msgs) else None sc_watermarks[uuid] = {"last_id": last_id, "updated": curr_updated} continue # Tier 1 fast check: if updated timestamp unchanged, skip if last_updated and curr_updated and last_updated == curr_updated: continue # Tier 2 check: fetch recent messages msgs, err = muse_hybrid.get_history(agent, thread_id=uuid, limit=10) if err or not msgs: continue last_id = sc_watermarks.get(uuid, {}).get("last_id") formatted = [{ "id": m.get("message_id") or f"msg-{m.get('seq')}", "from": {"name": m.get("role", "unknown")}, "role": m.get("role", "unknown"), "text": m.get("text", "") } for m in msgs] new_msgs = new_messages(formatted, {"last_id": last_id}) incoming = [m for m in new_msgs if m.get("role") != "assistant"] if incoming: new_activity.append({ "sidechat": sc_name, "uuid": uuid, "messages": incoming }) # Advance watermark sc_watermarks[uuid] = { "last_id": formatted[-1]["id"] if formatted else None, "updated": curr_updated } return new_activity def compose_dual_digest(agent, main_new, sc_activity): """Compose concise digest combining main chat and sidechat activity.""" lines = [] if main_new: lines.append("New messages in %s main chat (%d):" % (agent, len(main_new))) for m in main_new[:3]: text = (m.get("text") or "").strip().replace("\n", " ") q = " (?)" if "?" in text else "" hot = " (!)" if (OPERATOR_RE.search(text) or "urgent" in text.lower()) else "" lines.append("- %s: %s%s%s" % (author_of(m), text[:PREVIEW_MAX], q, hot)) if len(main_new) > 3: lines.append("(+%d more)" % (len(main_new) - 3)) if sc_activity: if lines: lines.append("") lines.append("Incoming sidechat activity for %s:" % agent) for item in sc_activity: sc_name = item["sidechat"] msgs = item["messages"] lines.append("- [%s] (%d new):" % (sc_name, len(msgs))) for m in msgs[:2]: text = (m.get("text") or "").strip().replace("\n", " ") lines.append(" * %s: %s" % (author_of(m), text[:PREVIEW_MAX])) lines.append("Check chat via box when available.") return "\n".join(lines)[:DIGEST_MAX] def check_subagents(agent, cfg, threads_meta): """Inspect active child subagents spawned by agent. Returns list of alert strings to deliver to the parent's prompting sidechat: - New output from subagent: prompts parent with snippet and thread view command. - 15-minute stall: warns parent that subagent has had no activity. """ alerts = [] try: import subagent_tracker import muse_hybrid active = subagent_tracker.get_active_sessions(parent=agent) if not active or not threads_meta: return alerts meta_by_id = {t.get("session_id"): t for t in threads_meta if isinstance(t, dict)} now = datetime.now(timezone.utc) for s in active: sid = s.get("session_id") title = s.get("title") or "subagent" last_wm = s.get("last_watermark") t_meta = meta_by_id.get(sid) if not t_meta: continue current_updated = t_meta.get("updated") # 1. New output from subagent if current_updated and current_updated != last_wm: if last_wm is None: # Seed initial watermark so we only trigger on subsequent turns subagent_tracker.update_session(sid, last_watermark=current_updated, last_activity_at=utcnow()) continue history, _ = muse_hybrid.get_history(agent, thread_id=sid, limit=2) snippet = "" if history: for m in reversed(history): if m.get("role") == "assistant": snippet = m.get("text", "").strip()[:180] break subagent_tracker.update_session( sid, last_watermark=current_updated, last_activity_at=utcnow(), timeout_warned=False ) short_id = sid[:8] alert_text = f"[SUBAGENT-UPDATE] Subagent '{title}' ({short_id}) has posted new output." if snippet: alert_text += f"\nSnippet: {snippet}..." alert_text += f"\nRun `box thread view {sid}` to review and aggregate deliverables." alerts.append(alert_text) # 2. 15-minute stall watchdog elif not s.get("timeout_warned"): last_act = s.get("last_activity_at") or s.get("spawned_at") if last_act: try: dt = datetime.fromisoformat(last_act.replace("Z", "+00:00")) idle_sec = (now - dt).total_seconds() if idle_sec >= 900: # 15 minutes subagent_tracker.update_session(sid, timeout_warned=True) short_id = sid[:8] alerts.append( f"[WARN] Subagent '{title}' ({short_id}) has had no activity for 15+ minutes. " f"Check status with `box thread view {sid}` or consider respawning." ) except Exception: pass except Exception as e: log("%s: subagent check error: %s" % (agent, e)) return alerts 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 {} enabled = cfg.get("enabled") or {} results = {} prompted = 0 errors = 0 import muse_hybrid for agent in agents: if not enabled.get(agent, True): results[agent] = {"ok": True, "new": 0, "disabled": True} continue wm = agents_state.get(agent) or {} # 1. Main chat check 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 main_new = new_messages(messages, wm) newest = messages[-1] if messages else None if newest: wm["last_id"] = newest.get("id") wm["last_ts"] = newest.get("ts") # 2. Sidechats dual-monitoring check threads_meta, _ = muse_hybrid.get_threads(agent) sc_activity = check_sidechats(agent, cfg, wm, threads_meta) # 3. Subagent reactive monitoring check subagent_alerts = check_subagents(agent, cfg, threads_meta) agents_state[agent] = wm total_new = len(main_new) + sum(len(a["messages"]) for a in sc_activity) + len(subagent_alerts) if total_new == 0: results[agent] = {"ok": True, "new": 0} continue parts = [] if len(main_new) > 0 or len(sc_activity) > 0: parts.append(compose_dual_digest(agent, main_new, sc_activity)) if subagent_alerts: parts.extend(subagent_alerts) digest = "\n\n".join(parts)[:DIGEST_MAX * 2] 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: prompted += 1 log("%s: prompted %s with %d new (%s)" % (agent, sidechat, total_new, detail)) results[agent] = {"ok": True, "new": total_new, "prompted": sidechat, "detail": detail} else: log("%s: PROMPT SEND FAILED: %s" % (agent, detail)) results[agent] = {"ok": False, "error": detail, "new": total_new} errors += 1 # Save under the state lock with a fresh reload: an enable/disable may # have landed during the slow chat reads; preserve its config changes # and only update the keys this run owns (watermarks, last_run/result). with state_locked(): fresh = load_state() fresh["agents"] = agents_state fresh["last_run"] = utcnow() fresh["last_result"] = {"prompted": prompted, "errors": errors, "agents": results} # Seed defaults on first run; otherwise preserve operator-edited config. if "config" not in fresh: fresh["config"] = {"agents": cfg["agents"], "prompt_sidechat": cfg["prompt_sidechat"], "sender": cfg["sender"], "read_limit": cfg["read_limit"], "enabled": dict(cfg.get("enabled") or {a: True for a in cfg["agents"]})} save_state(fresh) return {"ok": errors == 0, "prompted": prompted, "errors": errors, "new_total": sum(r.get("new", 0) for r in results.values()), "agents": results} def _set_enabled(agent, value): """Enable/disable the loop for one agent (or all if agent is None).""" with state_locked(): st = load_state() cfg = st.get("config") or {} agents = cfg.get("agents") or list(DEFAULT_AGENTS) agents = [a for a in agents if a in DEFAULT_AGENTS] or list(DEFAULT_AGENTS) if agent and agent not in DEFAULT_AGENTS: return {"ok": False, "error": "unknown agent: %s" % agent} enabled = dict(cfg.get("enabled") or {}) targets = [agent] if agent else agents for a in targets: enabled[a] = value cfg["enabled"] = enabled st["config"] = cfg save_state(st) return {"ok": True, "enabled": {a: enabled.get(a, True) for a in targets}} def do_status(): st = load_state() cfg = get_config(st) agents_state = st.get("agents") or {} return {"ok": True, "config": cfg, "enabled": cfg.get("enabled") or {}, "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 == "enable": print(json.dumps(_set_enabled(only_agent, True))) return 0 if action == "disable": print(json.dumps(_set_enabled(only_agent, False))) return 0 if action != "check": print(json.dumps({"ok": False, "error": "usage: check|status|enable|disable [--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 0 if __name__ == "__main__": sys.exit(main(sys.argv))