Files
box/bin/self_main_loop.py
T

802 lines
30 KiB
Python
Raw Normal View History

#!/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": "646 tasks",
"opm": "main-loop brain",
"pip": "pip tasks",
"muse": "muse tasks",
"dev": "onboarding-dev",
}
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
# 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)", 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
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, 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))