feat(core-loop): implement dual-monitoring, fast gateway dispatch, dedicated pip tasks sidechat, and 2m cadence
This commit is contained in:
+184
-18
@@ -54,8 +54,8 @@ 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",
|
||||
"opm": "main-loop brain",
|
||||
"pip": "pip tasks",
|
||||
"muse": "muse tasks",
|
||||
}
|
||||
DEFAULT_SENDER = "opm" # neutral sender, mirrors `box notify`
|
||||
@@ -119,7 +119,14 @@ def get_config(st):
|
||||
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 {})
|
||||
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:
|
||||
@@ -206,8 +213,56 @@ def compose_digest(agent, new):
|
||||
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 send_prompt(sender, agent, sidechat, digest):
|
||||
"""Post the digest to the agent's prompting sidechat via dm.py."""
|
||||
"""Post the digest to the agent's prompting sidechat.
|
||||
Tries fast direct gateway send via muse_hybrid first, falling back to dm.py.
|
||||
"""
|
||||
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)"
|
||||
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:
|
||||
@@ -219,7 +274,108 @@ def send_prompt(sender, agent, sidechat, digest):
|
||||
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"
|
||||
return True, "sent (dm.py)"
|
||||
|
||||
|
||||
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 do_check(only_agent=None):
|
||||
@@ -235,11 +391,15 @@ def do_check(only_agent=None):
|
||||
prompted = 0
|
||||
errors = 0
|
||||
|
||||
import muse_hybrid
|
||||
|
||||
for agent in agents:
|
||||
if not enabled.get(agent, True):
|
||||
results[agent] = {"ok": True, "new": 0, "disabled": True}
|
||||
continue
|
||||
wm = agents_state.get(agent) or {}
|
||||
|
||||
# 1. Main chat check
|
||||
ok, payload = read_main_chat(agent, cfg["read_limit"])
|
||||
if not ok:
|
||||
log("%s: READ FAILED: %s" % (agent, payload))
|
||||
@@ -247,33 +407,39 @@ def do_check(only_agent=None):
|
||||
errors += 1
|
||||
continue
|
||||
messages = payload
|
||||
new = new_messages(messages, wm)
|
||||
main_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")}
|
||||
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)
|
||||
|
||||
agents_state[agent] = wm
|
||||
|
||||
total_new = len(main_new) + sum(len(a["messages"]) for a in sc_activity)
|
||||
if total_new == 0:
|
||||
results[agent] = {"ok": True, "new": 0}
|
||||
continue
|
||||
digest = compose_digest(agent, new)
|
||||
|
||||
digest = compose_dual_digest(agent, main_new, sc_activity)
|
||||
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}
|
||||
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:
|
||||
# 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)}
|
||||
results[agent] = {"ok": False, "error": detail, "new": total_new}
|
||||
errors += 1
|
||||
|
||||
# Save under the state lock with a fresh reload: an enable/disable may
|
||||
|
||||
Reference in New Issue
Block a user