#!/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. The closing line enforces the # recursive box->agent->box discipline: post the RESULT back in this thread # (box records it and dispatches the next chained step); never DM the next # agent directly. ~166 chars, within the 600-char digest budget. try: import lookup_engine HAS_LOOKUP_ENGINE = True except ImportError: HAS_LOOKUP_ENGINE = False _DEFAULT_CONTRACT_FOOTER = ("Reply: [ACK id] seen | [CLAIM id] mine | " "[RESULT id] done | [DECLINE id] | [NO-ACTION id]. " "Report back here. Box dispatches the next step; " "do not DM the next agent directly.") CONTRACT_FOOTER = lookup_engine.get_contract_footer() if HAS_LOOKUP_ENGINE else _DEFAULT_CONTRACT_FOOTER 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 # --------------------------------------------------------------------------- # Brain workspace — the main loop operates its thinking in the "main-loop # brain" sidechat (opm's account). See bin/brain.py. # User directive 2026-10-04: "we need main loop to operate its brains in # side chat". Safety: brain posts carry [BRAIN], never [JOB]; never # --expect-reply; own messages skipped on read; !loop only from authorized # senders. All enforced in brain.py. # --------------------------------------------------------------------------- _BRAIN_LAST_IDS = {} # (agent, sidechat_name) -> last message_id seen def read_brain_messages(agent, sidechat_name, since_ts): """read_fn adapter for brain.BrainWorkspace. Resolves sidechat_name -> UUID via dm (never hardcoded; UUIDs rotate), reads via the same muse_hybrid primitive the loop uses, returns [{"sender", "text", "ts"}]. Tracks last message_id per (agent, name) in the module cache, which do_check() persists to the watermark JSON (brain.brain_last_ids) across ticks. First run anchors at newest with no backfill (same policy as new_messages for main chat). """ try: import dm uuid = dm.resolve_sidechat_target(sidechat_name, agent) except Exception: return [] if not uuid: return [] try: import muse_hybrid msgs, err = muse_hybrid.get_history(agent, thread_id=uuid, limit=10) except Exception: return [] if err or not msgs: return [] key = (agent, sidechat_name) last_id = _BRAIN_LAST_IDS.get(key) ids = [m.get("message_id") or "msg-%s" % m.get("seq") for m in msgs] if last_id is None: # First run: anchor at newest, no backfill (old !loop commands # must not fire on deploy). _BRAIN_LAST_IDS[key] = ids[-1] if ids else None return [] if last_id in ids: new_msgs = msgs[ids.index(last_id) + 1:] else: # Watermark fell out of the read window: treat all as new. # Safe: brain intake skips own messages and only authorized # !loop senders can act. new_msgs = msgs _BRAIN_LAST_IDS[key] = ids[-1] if ids else last_id now = time.time() return [{"sender": m.get("role", "unknown"), "text": m.get("text", ""), "ts": now} for m in new_msgs] 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 quiet_held = 0 # Brain workspace: the loop thinks in the sidechat, takes !loop # instructions there. Non-fatal: a brain failure must never break # the tick. brain = None try: from brain import BrainWorkspace brain = BrainWorkspace(STATE_FILE, read_fn=read_brain_messages) # Restore persisted brain message IDs so !loop commands posted # between ticks are not missed (module cache is per-process). for k, v in (brain.state.get("brain_last_ids") or {}).items(): try: ag, nm = k.split("|", 1) _BRAIN_LAST_IDS[(ag, nm)] = v except ValueError: pass intake = brain.intake() if intake.get("commands"): log("brain: %d commands, %d acks, %d ignored-senders" % ( intake["commands"], intake.get("acks", 0), intake.get("ignored_senders", 0))) except Exception as e: log("brain init/intake failed (non-fatal): %r" % e) brain = None 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 if brain is not None and brain.ignored(agent): log("%s: skipped (brain !loop ignore active)" % agent) results[agent] = {"ok": True, "new": 0, "ignored": 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 # Brain quiet mode: hold actionable escalations. Reads continue, # digests are composed, but nothing escalates to the agent — the # thinking note in the brain carries the activity instead. if brain is not None and brain.quiet(): is_act, _urg = classify_digest(digest) if is_act: log("%s: quiet mode - digest held (not escalated)" % agent) results[agent] = {"ok": True, "new": total_new, "quiet_held": True} quiet_held += 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 # Brain: post the tick's thinking to the sidechat workspace, then # persist brain state. Non-fatal on failure. if brain is not None: try: seen = {} escalated = [] for a in agents: r = results.get(a, {}) seen[a] = (r.get("new", 0), 0, 0) if r.get("prompted"): wm_a = agents_state.get(a) or {} did = wm_a.get("last_digest_id") if did and wm_a.get("last_digest_actionable"): escalated.append(did) closure_rate = None health = None try: health = digest_health() closure_rate = (health or {}).get("closure_rate") except Exception: pass wm_epochs = {} for a in agents: ts_s = (agents_state.get(a) or {}).get("last_ts") try: if ts_s: dt = datetime.fromisoformat( ts_s.replace("Z", "+00:00")) wm_epochs[a] = dt.timestamp() except Exception: pass brain.post_thinking( {"seen": seen, "escalated": escalated, "skipped_info": quiet_held, "errors": errors, "closure_rate": closure_rate}, cfg={a: bool(enabled.get(a, True)) for a in agents}, watermarks=wm_epochs, health=health, ) # Persist brain message IDs for the next tick. try: brain.state["brain_last_ids"] = { "%s|%s" % k: v for k, v in _BRAIN_LAST_IDS.items() if v is not None} except Exception: pass brain.save() except Exception as e: log("brain post_thinking failed (non-fatal): %r" % e) # 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))