#!/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 {}) enabled = dict(cfg.get("enabled") or {}) # Default any unlisted agent to enabled. for a in agents: enabled.setdefault(a, True) 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), "enabled": enabled, } 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 {} enabled = cfg.get("enabled") or {} results = {} prompted = 0 errors = 0 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 {} 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"], "enabled": dict(cfg.get("enabled") or {a: True for a in cfg["agents"]})} 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 _set_enabled(agent, value): """Enable/disable the loop for one agent (or all if agent is None).""" 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 1 if result.get("prompted") else 0 if __name__ == "__main__": sys.exit(main(sys.argv))