#!/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 import time 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) # Deterministic per-agent jittered sleeps — de-correlates within-run # traffic across the shared egress IP (see rate_limiter.py). try: from rate_limiter import effective_interval as _effective_interval HAS_RATE_LIMITER = True except ImportError: HAS_RATE_LIMITER = False 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") FOLLOWUPS_FILE = os.path.join(BASE, "followups.json") DEFAULT_AGENTS = ["muse", "pip", "646", "opm", "dev"] # 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", "dev": "dev-audit-channel", } 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 INTER_NODE_SLEEP_BASE = 7.0 # s between per-agent passes; jittered per agent 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-") or name.startswith("onboarding-"): continue if "heartbeat" in name.lower(): continue if any(name.startswith(p) for p in ("box-http-health", "box-service-health", "box-deep-health", "canary")): continue monitored[uuid] = name return monitored def classify_digest(digest): """Return (actionable, urgent) for a composed digest. Actionable = digest carries a (?) question or (!) operator-needed marker. Urgent = digest carries the (!) marker. """ urgent = " (!)" in digest actionable = urgent or " (?)" in digest return actionable, urgent def make_digest_id(agent): """Stable digest ID: ml-- (UTC).""" ts = datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S") return "ml-%s-%s" % (agent, ts) # In-band response-contract footer for ACTIONABLE digests. The verbs are matched # by response-harvester.py to resolve followups: ACK/CLAIM acknowledge (nudge # suppression), RESULT/DECLINE/NO-ACTION close. ~94 chars, well under budget. CONTRACT_FOOTER = ("Reply: [ACK id] seen | [CLAIM id] mine | " "[RESULT id] done | [DECLINE id] | [NO-ACTION id]") def send_prompt(sender, agent, sidechat, digest): """Post the digest to the agent's prompting sidechat. Actionable digests (carrying (?) or (!)) go via dm.py with --expect-reply so the followup machinery tracks them; the digest gets a [JOB ] tag which dm.py auto-extracts into the followup record. Informational digests try the fast gateway first, falling back to dm.py without followup flags. Returns (sent, detail, digest_id, actionable). """ actionable, urgent = classify_digest(digest) digest_id = make_digest_id(agent) if actionable else None if actionable: # Embed [JOB id] for dm.py's job_id extraction -> followup record, # and append the response-contract footer. Reserve space for both so # the tagged digest stays within DIGEST_MAX and dm.py never truncates # the footer (or the job id). head = "[JOB %s]\n" % digest_id room = DIGEST_MAX - len(head) - len(CONTRACT_FOOTER) - 1 body = digest if len(digest) <= room else digest[:room].rstrip() tagged = head + body + "\n" + CONTRACT_FOOTER timeout = 1800 if urgent else 3600 cmd = [sys.executable, DM_PY, "send", "--agent", sender, "--to", agent, "--target", sidechat, "--expect-reply", "--reply-timeout", str(timeout), "--reply-nudges", "2", tagged] 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, digest_id, True except OSError as e: return False, "could not exec dm.py: %s" % e, digest_id, True 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), digest_id, True return True, "sent (dm.py, expect-reply)", digest_id, True # Informational: fast gateway first, dm.py fallback without followup flags. target_uuid = sidechat try: import dm target_uuid = dm.resolve_sidechat_target(sidechat, agent) except Exception: pass # Fast direct gateway send for main chat or resolved UUID is_main = (sidechat == "main" or target_uuid == "main") if is_main: 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=None, wait=25) if res and not err: return True, "sent (gateway main)", None, False log("Gateway send main fallback for %s: %s" % (agent, err or res)) except Exception as e: log("Gateway send main exception for %s: %s" % (agent, e)) elif 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)", None, False 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, None, False except OSError as e: return False, "could not exec dm.py: %s" % e, None, False 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), None, False return True, "sent (dm.py)", None, False 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 = [] for m in new_msgs: if m.get("role") == "assistant": continue text = m.get("text", "") if "[JOB " in text or "Heartbeat check" in text or "[from:super]" in text: continue incoming.append(m) 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 first_agent = True for agent in agents: if not first_agent: # Stagger per-node passes: timer staggering doesn't help within # a run. Deterministic per-agent jitter keeps runs reproducible # while drifting each node's phase apart. sleep_s = (INTER_NODE_SLEEP_BASE * _effective_interval(agent, 1.0) if HAS_RATE_LIMITER else INTER_NODE_SLEEP_BASE) log("%s: inter-node sleep %.1fs" % (agent, sleep_s)) time.sleep(sleep_s) first_agent = False 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 sidechat = cfg["prompt_sidechat"].get(agent) is_main = (sidechat == "main") if is_main: # Side chats drive the main chat if not sc_activity and not subagent_alerts: results[agent] = {"ok": True, "new": 0} continue parts = [] if sc_activity: parts.append(compose_dual_digest(agent, [], sc_activity)) if subagent_alerts: parts.extend(subagent_alerts) digest = "\n\n".join(parts)[:DIGEST_MAX * 2] else: 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) if not parts: results[agent] = {"ok": True, "new": 0} continue digest = "\n\n".join(parts)[:DIGEST_MAX * 2] if not sidechat: log("%s: no prompting sidechat configured, skipping" % agent) results[agent] = {"ok": False, "error": "no prompting sidechat"} errors += 1 continue sent, detail, digest_id, actionable = send_prompt(cfg["sender"], agent, sidechat, digest) if sent: prompted += 1 if digest_id: wm["last_digest_id"] = digest_id wm["last_digest_actionable"] = actionable 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}} # Digest protocol job-id prefix: actionable main-loop digests embed # [JOB ml--]. Informational digests create no # followup at all, so every ml- followup is actionable by construction # and informational digests are excluded from all rates. DIGEST_JOB_PREFIX = "ml-" # Reply verbs that acknowledge without closing (nudge-suppressed). ACK_VERBS = {"ACK", "CLAIM"} # Reply verbs that close the digest. CLOSE_VERBS = {"RESULT", "DECLINE", "NO-ACTION"} def _followup_job_id(rec): tags = rec.get("tags") or {} return tags.get("job_id") or rec.get("job_id") or "" def _followup_is_stale(rec, now): if rec.get("status") == "escalated": return True if rec.get("status") == "pending": try: dl = datetime.fromisoformat( (rec.get("deadline") or "").replace("Z", "+00:00")) return dl < now except (ValueError, TypeError): return False return False def digest_health(): """Digest loop-closure metrics from followups.json. Counts only actionable digests (job_id starting with 'ml-'). Reads the existing followups.json - no new state file, keeping self-main-loop-watermark.json as the single status source. """ now = datetime.now(timezone.utc) health = {"delivered": 0, "acked": 0, "closed": 0, "stale": 0, "closure_rate": None, "by_agent": {}} try: with open(FOLLOWUPS_FILE) as f: data = json.load(f) except (FileNotFoundError, json.JSONDecodeError, ValueError): return health records = data.values() if isinstance(data, dict) else data for rec in records: if not isinstance(rec, dict): continue job_id = _followup_job_id(rec) if not job_id.startswith(DIGEST_JOB_PREFIX): continue # not a main-loop digest: excluded from all rates agent = rec.get("recipient") or "?" per = health["by_agent"].setdefault( agent, {"delivered": 0, "acked": 0, "closed": 0, "stale": 0}) health["delivered"] += 1 per["delivered"] += 1 outcome = (rec.get("outcome") or "").upper() status = rec.get("status") or "" is_acked = status == "acknowledged" or outcome in ACK_VERBS # resolved counts as closed unless the outcome was only an ACK/CLAIM # (legacy pre-protocol resolutions have no outcome field). is_closed = status == "resolved" and outcome not in ACK_VERBS if is_acked: health["acked"] += 1 per["acked"] += 1 if is_closed: health["closed"] += 1 per["closed"] += 1 if _followup_is_stale(rec, now): health["stale"] += 1 per["stale"] += 1 if health["delivered"]: health["closure_rate"] = round( health["closed"] / health["delivered"], 4) for per in health["by_agent"].values(): per["closure_rate"] = (round(per["closed"] / per["delivered"], 4) if per["delivered"] else None) return health 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"), "digest_health": digest_health()} 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))