From 86f0082ffc112820e77a2f0638ea2211299789a9 Mon Sep 17 00:00:00 2001 From: operator-main Date: Mon, 5 Oct 2026 00:48:38 +0000 Subject: [PATCH] =?UTF-8?q?feat(main-loop):=20digest=20response=20protocol?= =?UTF-8?q?=20=E2=80=94=20actionable=20digests,=20reply=20verbs,=20closure?= =?UTF-8?q?=20metrics?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Session: sidechat/main-loop-protocol --- bin/response-harvester.py | 85 ++++++++++++--- bin/self_main_loop.py | 185 +++++++++++++++++++++++++++++++-- docs/DIGEST-PROTOCOL.md | 213 ++++++++++++++++++++++++++++++++++++++ 3 files changed, 461 insertions(+), 22 deletions(-) create mode 100644 docs/DIGEST-PROTOCOL.md diff --git a/bin/response-harvester.py b/bin/response-harvester.py index b4288e3..096c6f8 100755 --- a/bin/response-harvester.py +++ b/bin/response-harvester.py @@ -4,8 +4,9 @@ response-harvester.py — Fleet agent readback and response harvesting daemon. Monitors Chromebox agents (muse, pip, 646, opm), harvests incoming messages from Main Chat and registered sidechats, maintains persistent watermarks, appends to -chat-history.jsonl, resolves pending follow-ups, and records [RESULT] completions -in job-log.jsonl. +chat-history.jsonl, resolves pending follow-ups (verb-aware: [ACK|CLAIM] -> +acknowledged, [RESULT|DECLINE|NO-ACTION] -> resolved, outcome recorded), and +records [RESULT] completions in job-log.jsonl. Features: - Direct CDP over host veth interfaces (fast, no sudo needed). @@ -83,6 +84,20 @@ def iter_result_markers(text): yield m.group(1).strip(), m.group(2).strip() +# Verb markers for the digest response protocol: +# [ACK|CLAIM|RESULT|DECLINE|NO-ACTION ]. +# ACK/CLAIM acknowledge a digest (nudge-suppressed, NOT closed); +# RESULT/DECLINE/NO-ACTION close the digest. Every verb match records +# outcome= on the followup record. +VERB_RE = re.compile(r"\[(ACK|CLAIM|RESULT|DECLINE|NO-ACTION)\s+([A-Za-z0-9_-]+)\]") + + +def iter_verb_markers(text): + """Yield (verb, job_id) for every [VERB ] marker in text.""" + for m in VERB_RE.finditer(text or ""): + yield m.group(1), m.group(2).strip() + + def utcnow(): return datetime.now(timezone.utc).isoformat() @@ -366,7 +381,8 @@ def harvest_main_feed(cdp, agent, watermarks, followups, dry_run=False): if author == "assistant": markers = list(iter_result_markers(text)) - if markers: + verbs = list(iter_verb_markers(text)) + if markers or verbs: for job_id, result_text in markers: is_fail = is_fail_result(result_text) job_results += 1 @@ -385,7 +401,12 @@ def harvest_main_feed(cdp, agent, watermarks, followups, dry_run=False): append_jsonl(JOB_LOG, job_record) trigger_chain_next(job_id, result_text, success=not is_fail) clear_matching_followups(followups, agent, "main", mid, text, - dry_run, job_id=job_id) + dry_run, job_id=job_id, verb="RESULT") + for verb, job_id in verbs: + if verb == "RESULT": + continue # resolved via the result-marker path above + clear_matching_followups(followups, agent, "main", mid, text, + dry_run, job_id=job_id, verb=verb) else: clear_matching_followups(followups, agent, "main", mid, text, dry_run) @@ -487,7 +508,8 @@ def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run if author == "assistant": markers = list(iter_result_markers(text)) - if markers: + verbs = list(iter_verb_markers(text)) + if markers or verbs: for job_id, result_text in markers: is_fail = is_fail_result(result_text) job_results += 1 @@ -507,7 +529,12 @@ def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run # Trigger pipeline chaining or next job if configured trigger_chain_next(job_id, result_text, success=not is_fail) clear_matching_followups(followups, agent, thread_id, mid, text, - dry_run, job_id=job_id) + dry_run, job_id=job_id, verb="RESULT") + for verb, job_id in verbs: + if verb == "RESULT": + continue # resolved via the result-marker path above + clear_matching_followups(followups, agent, thread_id, mid, text, + dry_run, job_id=job_id, verb=verb) else: clear_matching_followups(followups, agent, thread_id, mid, text, dry_run) @@ -611,7 +638,7 @@ def trigger_chain_next(job_id, result_text, success=True): def clear_matching_followups(followups, agent, thread_id, mid, text, dry_run=False, - job_id=None): + job_id=None, verb=None): """Resolve follow-up records if an assistant message is detected in the thread. Matches on thread identity (thread_uuid or target='main') OR on job_id @@ -619,13 +646,27 @@ def clear_matching_followups(followups, agent, thread_id, mid, text, dry_run=Fal thread_uuid or target, fixing ghost followups with null thread_uuid. Also matches a main-chat reply when the sweeper recorded final_nudge_target='main' (final nudge routed to main chat). + + When verb is given ([ACK|CLAIM|RESULT|DECLINE|NO-ACTION ]), + matching is job_id-scoped for non-RESULT verbs (a verb marker names the + digest it answers, so it must not touch unrelated pending followups that + merely share the thread); RESULT keeps the historical thread-or-job + matching. Every verb match records outcome=. ACK/CLAIM set status + 'acknowledged' (nudge-suppressed, NOT closed) instead of 'resolved', and + also match already-'acknowledged' records so an ACK -> RESULT lifecycle + closes correctly. """ if not followups: return modified = False for f_id, f_rec in followups.items(): - if f_rec.get("status") != "pending": + # Verb replies can follow an ACK (ACK -> RESULT lifecycle), so verbs + # also match 'acknowledged' records; plain replies only match pending. + if verb: + if f_rec.get("status") not in ("pending", "acknowledged"): + continue + elif f_rec.get("status") != "pending": continue if f_rec.get("recipient") != agent: continue @@ -642,17 +683,33 @@ def clear_matching_followups(followups, agent, thread_id, mid, text, dry_run=Fal # sidechat. match_thread = True - # Match by job_id (from [RESULT ]) -- works regardless of - # thread_uuid or target. This is an ADDITIONAL path, not a replacement. + # Match by job_id (from [RESULT ] or [VERB ]) -- + # works regardless of thread_uuid or target. This is an ADDITIONAL + # path, not a replacement. match_job = False if job_id and f_rec.get("job_id") and f_rec.get("job_id") == job_id: match_job = True + # Non-RESULT verbs are job-scoped: they must not acknowledge/resolve + # unrelated pending followups that merely share the thread. RESULT + # keeps the historical thread-or-job matching. + if verb and verb != "RESULT" and not match_job: + continue + if match_thread or match_job: - f_rec["status"] = "resolved" - f_rec["resolved_at"] = utcnow() - f_rec["resolved_by_mid"] = mid - f_rec["resolved_snippet"] = text[:150] + if verb: + f_rec["outcome"] = verb + if verb in ("ACK", "CLAIM"): + # Acknowledged: sweeper nudges stop (status != pending), but + # the digest is NOT closed until a closing verb arrives. + f_rec["status"] = "acknowledged" + f_rec["acknowledged_at"] = utcnow() + f_rec["acknowledged_by_mid"] = mid + else: + f_rec["status"] = "resolved" + f_rec["resolved_at"] = utcnow() + f_rec["resolved_by_mid"] = mid + f_rec["resolved_snippet"] = text[:150] modified = True if modified and not dry_run: diff --git a/bin/self_main_loop.py b/bin/self_main_loop.py index 04c39b9..607d9fd 100755 --- a/bin/self_main_loop.py +++ b/bin/self_main_loop.py @@ -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-- (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 ] 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--]. 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): diff --git a/docs/DIGEST-PROTOCOL.md b/docs/DIGEST-PROTOCOL.md new file mode 100644 index 0000000..969a8bf --- /dev/null +++ b/docs/DIGEST-PROTOCOL.md @@ -0,0 +1,213 @@ +# Digest Protocol Runbook + +How the main loop sends digests and how agents close them. + +**Status:** spec (2026-10-04). Implemented in `self_main_loop.py` + `response-harvester.py`; +no new machinery — the protocol wires digests into the existing `dm.py --expect-reply` +/ `followups.json` / harvester response path. + +--- + +## 1. Why this protocol exists + +Main-loop digests (`[main-loop] N new in main chat: ...`) were being sent +with no response contract: no stable ID, no declared timeout, no nudge/escalation +policy, and no way to distinguish "seen and closed" from "ignored". Across 43 +digests there were zero RESULT replies and zero measurable outcomes — not because +agents ignored them, but because nothing defined how to answer. + +This protocol closes that wiring gap without inventing a new system: +- digest composer → `dm.py send --expect-reply` (followup record, timeout/nudges + declared at send time) +- agent replies in-band with a verb: `[ACK|CLAIM|RESULT|DECLINE|NO-ACTION ]` +- `response-harvester.py` matches by `job_id` and records the outcome in + `followups.json` +- the sweeper nudges/escalates only unresolved followups, per the declared policy + +## 2. Digest classes + +Not every digest deserves a reply demand (notification fatigue; and the standing +rule that the human only sees real decisions and escalations). The composer +already marks `[?]` (question) and `[!]` (operator-needed / urgent). + +| Class | Marker | Response required? | Followup record? | Measured by | +|---|---|---|---|---| +| **ACTIONABLE** | carries `[?]` or `[!]` | **yes** | yes | delivered → acked → closed | +| **INFORMATIONAL** | neither | no | no | delivered only | + +Informational digests are sent exactly as today. They are **excluded from every +reply-rate metric** — an informational digest is measured as *delivered*, full stop. +Judging informational FYI traffic by answer rate is structurally misleading. + +## 3. Digest ID + +Every ACTIONABLE digest gets a stable ID: + +``` +ml-- +``` + +Example: `ml-646-20261004-213000`. + +The ID is embedded visibly in the digest as `[JOB ml-...]`. This marker survives +`extract_trailing_tags` (only known canonical keys are stripped) and is picked up +by `dm_send`'s job_id regex (`\[JOB\s+([A-Za-z0-9_-]+)\]`) anywhere in the message, +so it lands in `tags["job_id"]` → the followup record. **No new dm.py flags are +needed for the ID.** + +Replies match on this ID, so attribution games with `[from:X]` can never break +resolution. + +## 4. Reply verbs + +Printed in the contract footer of every ACTIONABLE digest (Section 7). + +| Verb | Syntax | Meaning | Followup effect | +|---|---|---|---| +| `ACK` | `[ACK ]` | Seen, noted | status `acknowledged` — stops nudges, **not closed** | +| `CLAIM` | `[CLAIM ]` | I'm handling it | status `acknowledged` — stops nudges, **not closed** | +| `RESULT` | `[RESULT ] ` | Done, result attached | status `resolved`, outcome `RESULT` — **closed** | +| `DECLINE` | `[DECLINE ]` | Won't act (reason optional) | status `resolved`, outcome `DECLINE` — **closed** | +| `NO-ACTION` | `[NO-ACTION ]` | Reviewed, nothing needed | status `resolved`, outcome `NO-ACTION` — **closed** | + +Semantics: +- `ACK` vs `CLAIM`: both acknowledge. `CLAIM` signals ownership; `ACK` is just + "seen". Use `CLAIM` when you're picking up the work so the human knows who's on it. +- `NO-ACTION` is a real answer, not a dodge. It means "I reviewed this and there + is genuinely nothing to do." It closes the digest so it stops being counted as open. +- `DECLINE` should carry a reason when one exists (`[DECLINE ml-646-...] blocked: + chromebox flapping`), but a bare DECLINE is still a valid close. +- `RESULT` closes only when the result text is present. An empty RESULT is a malformed + reply — the harvester treats it as unrecognized. + +Nudge interaction: `ACK`/`CLAIM` suppress further nudges but leave the followup open +until a closing verb arrives. This matches the real workflow — "I'm on it" shouldn't +trigger nags, but the loop isn't closed until there's a result. + +## 5. Timeout policy + +Declared at send time (standing rule: policy is part of the send, not decided later). + +| Digest class | Timeout | Nudges | +|---|---|---| +| normal (marked `[?]`) | 60 min | 2 | +| urgent (marked `[!]`) | 30 min | 2 | + +Implemented as `--reply-timeout {3600|1800} --reply-nudges 2` on the `dm.py send` +invocation for ACTIONABLE digests. The timeout and nudge count are stored in the +followup record; the sweeper honors them. + +## 6. Nudge / escalation policy + +- The sweeper nudges only followups with `status == "pending"`. Anything + `acknowledged` or `resolved` is never re-nudged (ACKed digests getting nagged + trains agents to ignore the protocol — verify this when wiring the sweeper). +- After the declared nudges are exhausted and the timeout passes, the followup is + **escalated**, not retried forever: + - **Phase 1:** the sweeper's existing escalation path (DM to the escalate target + declared at send time via `--reply-escalate`). + - **Phase 2:** route to the registered `main-loop brain` sidechat for triage. + The brain is registered in `job-sidechats.json` as `main-loop brain` and is + name-resolvable, so no hardcoded UUIDs are needed. +- `stale` = still `pending` past timeout + nudges → escalated. + +Digest routing must use sidechat **names** (`646 tasks`, `main-loop brain`) via the +existing name resolution — never hardcoded thread UUIDs (see the stale-UUID lesson +in AGENTS.md). + +Note: the harvester only monitors threads registered in `job-sidechats.json`. +`opm/heartbeat` is not registered (digests there resolve via fuzzy search), so +digest replies there would never resolve. **Route opm's digests to the registered +`main-loop brain` sidechat instead of `heartbeat`** — cleaner separation +(brain = triage workspace, heartbeat = heartbeats) and zero new registrations. + +## 7. In-band contract footer + +The digest budget is 600 chars, so the footer is compact. Actionable digests trim +message previews from 5 to 3 to fit it: + +``` +[main-loop] [JOB ml-646-20261004-213000] 3 new in 646 main chat — reply needed: +- human: [?] +Reply: [ACK id] seen | [CLAIM id] mine | [RESULT id] done | [DECLINE id] | [NO-ACTION id] +``` + +The footer is the primary contract. This doc is the reference. + +## 8. Example exchange + +``` +# opm → 646 tasks (ACTIONABLE, --expect-reply --reply-timeout 3600 --reply-nudges 2) + +[from:opm] [id:a6579a8d] [main-loop] [JOB ml-646-20261004-213000] 3 new in 646 main chat — reply needed: +- human: should we archive the stale pipe-7d2896 thread? [?] +- opm: propose yes, it's superseded by the brain channel [?] +Reply: [ACK id] seen | [CLAIM id] mine | [RESULT id] done | [DECLINE id] | [NO-ACTION id] + +# 646 picks it up + +[from:646] [CLAIM ml-646-20261004-213000] reviewing now + +→ followup status: acknowledged (nudges suppressed, loop still open) + +# 646 finishes + +[from:646] [RESULT ml-646-20261004-213000] archived pipe-7d2896 via muse-cli; brain channel is canonical + +→ followup status: resolved, outcome: RESULT (loop closed) +``` + +A shorter close: + +``` +# pip sees an informational-with-a-question, reviews, nothing to do + +[from:pip] [NO-ACTION ml-pip-20261004-214500] reviewed — nothing needed, timer stagger already handles it + +→ followup status: resolved, outcome: NO-ACTION (loop closed) +``` + +## 9. Metrics + +Source of truth: `followups.json` (`outcome` field, written by the harvester). + +| Metric | Definition | +|---|---| +| **delivered** | `dm.py` reported `verified:true` with placement confirmed for the digest message (see dm-log) | +| **acked** | outcome ∈ {`ACK`, `CLAIM`} | +| **closed** | outcome ∈ {`RESULT`, `DECLINE`, `NO-ACTION`} | +| **stale** | still `pending` past timeout + nudges → escalated | +| **closure_rate** | `closed / actionable_delivered` | + +Rules: +- Informational digests are excluded from every rate. They are measured as + delivered, full stop. +- Actionable digests are measured as closed. A digest that is ACKed but never + resolved counts as open until a closing verb arrives. +- `delivered` requires the placement confirmation, not just the `verified:true` + flag (see the AGENTS.md placement-blindness lesson). + +Reporting: `self_main_loop.py` `status` exposes `digest_health()` counts +(delivered/acked/closed/stale), surfaced through the existing Box +`main-loop/status` API. No new endpoint needed in phase 1. + +## 10. Implementation notes + +- `dm.py` — **no changes.** `--expect-reply`, job_id extraction from `[JOB …]`, + followup registration, and placement verification already exist and are tested. +- `response-harvester.py` — extend the reply regex to verb-aware matching: + `\[(ACK|CLAIM|RESULT|DECLINE|NO-ACTION)\s+([A-Za-z0-9_-]+)\]`; on match, resolve + by job_id and record `outcome`. `ACK`/`CLAIM` → `acknowledged`; + `RESULT`/`DECLINE`/`NO-ACTION` → `resolved` + outcome. +- `self_main_loop.py` — `compose_digest` classifies actionable vs informational; + actionable embeds `[JOB ml-…]` + footer; `send_prompt` passes + `--expect-reply --reply-timeout {3600|1800} --reply-nudges 2` for actionable only. +- Sweeper — verify it nudges only `status == "pending"` (no code change expected). + +All changes follow review-then-commit. No new daemons, no new state files beyond +the existing `followups.json`. + +--- + +*Companion docs: `DM-SPEC.md` (DM format), `THREAD-BOOKKEEPING.md` (pin/archive), +`WARP-EGRESS-FIX.md` (partition handling).*