Add self_main_loop.py main-chat self-monitor + box main-loop action
Reads each fleet agent's muse.ai Main chat on a 5-min systemd timer (self-main-loop.timer). On new messages since the per-agent watermark, posts a concise digest to that agent's prompting sidechat (dm.py, opm as neutral sender - mirrors box notify), prompting the operator to check main chat via DM/box. Escapes the main-chat-goes-unread failure mode. No backfill on first run; silent when nothing new; read/send failures logged without killing the timer; flock overlap guard. box-ctl.py: main-loop check | status (fleet section). Session: sidechat/chromebox-ops
This commit is contained in:
@@ -650,6 +650,30 @@ def act_notify(agent, message, sidechat=None, allow_main_chat=False):
|
|||||||
out(True, agent=agent, target=target, sent=True)
|
out(True, agent=agent, target=target, sent=True)
|
||||||
|
|
||||||
|
|
||||||
|
def act_main_loop(sub):
|
||||||
|
"""Run the main-chat self-monitor loop (self_main_loop.py).
|
||||||
|
|
||||||
|
check: one loop iteration -> JSON {ok, prompted, errors, new_total, agents}
|
||||||
|
status: watermark / last run / prompt-sidechat config, no chat reads.
|
||||||
|
"""
|
||||||
|
if sub not in ("check", "status"):
|
||||||
|
fail("BAD_NAME", "usage: main-loop check|status")
|
||||||
|
audit("main-loop-" + sub)
|
||||||
|
script = BIN / "self_main_loop.py"
|
||||||
|
try:
|
||||||
|
r = subprocess.run([sys.executable, str(script), sub],
|
||||||
|
capture_output=True, text=True, timeout=600)
|
||||||
|
except subprocess.TimeoutExpired:
|
||||||
|
fail("LOOP_TIMEOUT", "self_main_loop.py timed out")
|
||||||
|
try:
|
||||||
|
data = json.loads(r.stdout.strip())
|
||||||
|
except Exception:
|
||||||
|
fail("LOOP_ERROR", "self_main_loop.py returned non-JSON",
|
||||||
|
{"stdout": (r.stdout or "")[-300:], "stderr": (r.stderr or "")[-300:]})
|
||||||
|
ok = bool(data.pop("ok", False))
|
||||||
|
out(ok, **data)
|
||||||
|
|
||||||
|
|
||||||
def act_watchdog_alerts():
|
def act_watchdog_alerts():
|
||||||
"""Check for new failed browser relaunches since the watermark.
|
"""Check for new failed browser relaunches since the watermark.
|
||||||
|
|
||||||
@@ -1147,6 +1171,7 @@ fleet:
|
|||||||
relay-health
|
relay-health
|
||||||
cdp-latency
|
cdp-latency
|
||||||
identity-audit VM identity audit drift check
|
identity-audit VM identity audit drift check
|
||||||
|
main-loop check|status main-chat self-monitor: prompt sidechat on new main msgs
|
||||||
|
|
||||||
timer actions:
|
timer actions:
|
||||||
timer-list
|
timer-list
|
||||||
@@ -1389,6 +1414,10 @@ def main(argv):
|
|||||||
if rest:
|
if rest:
|
||||||
fail("BAD_ARGS", "usage: chrome-errors")
|
fail("BAD_ARGS", "usage: chrome-errors")
|
||||||
act_chrome_errors()
|
act_chrome_errors()
|
||||||
|
elif action == "main-loop":
|
||||||
|
if len(rest) != 1 or rest[0] not in ("check", "status"):
|
||||||
|
fail("BAD_NAME", "usage: main-loop check|status")
|
||||||
|
act_main_loop(rest[0])
|
||||||
elif action == "dm-log":
|
elif action == "dm-log":
|
||||||
limit = 50
|
limit = 50
|
||||||
if rest:
|
if rest:
|
||||||
|
|||||||
Executable
+289
@@ -0,0 +1,289 @@
|
|||||||
|
#!/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 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 {})
|
||||||
|
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),
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
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 {}
|
||||||
|
results = {}
|
||||||
|
prompted = 0
|
||||||
|
errors = 0
|
||||||
|
|
||||||
|
for agent in agents:
|
||||||
|
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"]}
|
||||||
|
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 do_status():
|
||||||
|
st = load_state()
|
||||||
|
cfg = get_config(st)
|
||||||
|
agents_state = st.get("agents") or {}
|
||||||
|
return {"ok": True,
|
||||||
|
"config": cfg,
|
||||||
|
"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 != "check":
|
||||||
|
print(json.dumps({"ok": False, "error": "usage: check|status [--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))
|
||||||
Reference in New Issue
Block a user