feat(main-loop): digest response protocol — actionable digests, reply verbs, closure metrics

Session: sidechat/main-loop-protocol
This commit is contained in:
operator-main
2026-10-05 00:48:38 +00:00
parent 664b4da952
commit 86f0082ffc
3 changed files with 461 additions and 22 deletions
+177 -8
View File
@@ -37,6 +37,7 @@ import os
import re
import subprocess
import sys
import time
from contextlib import contextmanager
from datetime import datetime, timezone
@@ -44,11 +45,20 @@ 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"]
# Prompting sidechat per agent — mirrors box-ctl.py NOTIFY_SIDECHATS.
@@ -63,6 +73,7 @@ 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
@@ -239,10 +250,69 @@ def get_monitored_sidechats(agent):
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.
Tries fast direct gateway send via muse_hybrid first, falling back to dm.py.
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
@@ -257,7 +327,7 @@ def send_prompt(sender, agent, sidechat, digest):
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)"
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))
@@ -268,13 +338,13 @@ def send_prompt(sender, agent, 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
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
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)
return True, "sent (dm.py)"
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):
@@ -472,7 +542,17 @@ def do_check(only_agent=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
@@ -519,9 +599,12 @@ def do_check(only_agent=None):
errors += 1
continue
sent, detail = send_prompt(cfg["sender"], agent, sidechat, digest)
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:
@@ -571,6 +654,91 @@ def _set_enabled(agent, value):
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)
@@ -580,7 +748,8 @@ def do_status():
"enabled": cfg.get("enabled") or {},
"watermark": agents_state,
"last_run": st.get("last_run"),
"last_result": st.get("last_result")}
"last_result": st.get("last_result"),
"digest_health": digest_health()}
def main(argv):