826 lines
31 KiB
Python
Executable File
826 lines
31 KiB
Python
Executable File
#!/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 <agent> 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": "main",
|
|
"opm": "main",
|
|
"pip": "main",
|
|
"muse": "main",
|
|
"dev": "main",
|
|
}
|
|
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"(?<![\w-])operator(?![\w-])", re.IGNORECASE)
|
|
|
|
|
|
@contextmanager
|
|
def state_locked():
|
|
"""Exclusive lock for load-modify-save on the state file.
|
|
|
|
Prevents the timer's check (mid-run save) from clobbering an
|
|
enable/disable change that landed during the slow chat reads.
|
|
Held only for the brief critical section, never across network I/O.
|
|
"""
|
|
fh = open(STATE_LOCK_FILE, "w")
|
|
try:
|
|
fcntl.flock(fh, fcntl.LOCK_EX)
|
|
yield
|
|
finally:
|
|
fcntl.flock(fh, fcntl.LOCK_UN)
|
|
fh.close()
|
|
|
|
|
|
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)
|
|
user_prompt = cfg.get("prompt_sidechat") or {}
|
|
for a, sc in user_prompt.items():
|
|
# Migrate obsolete defaults
|
|
if a == "pip" and sc in ("646-pip-coord", "heartbeat"):
|
|
continue
|
|
if a == "opm" and sc in ("heartbeat",):
|
|
continue
|
|
prompt[a] = sc
|
|
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. Tries fast headless gateway (muse_hybrid) first, falling back to box-chat.py."""
|
|
# Fast path: muse_hybrid over WebSocket Noise frame inside isolated netns
|
|
try:
|
|
import muse_hybrid
|
|
msgs, err = muse_hybrid.get_history(agent, thread_id=None, limit=limit)
|
|
if msgs and not err:
|
|
formatted = []
|
|
for m in msgs:
|
|
formatted.append({
|
|
"id": m.get("message_id") or f"msg-{m.get('seq')}",
|
|
"ts": None,
|
|
"from": {"name": m.get("role", "unknown")},
|
|
"text": m.get("text", "")
|
|
})
|
|
return True, formatted
|
|
except Exception:
|
|
pass
|
|
|
|
# Fallback to box-chat.py / CDP if gateway fails or returns empty
|
|
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 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 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-<agent>-<YYYYMMDD-HHMMSS> (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 <id>] 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 = [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
|
|
|
|
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)
|
|
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-<agent>-<YYYYMMDD-HHMMSS>]. 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))
|