Restore followup/DM reliability fixes wiped by 19:43Z tree-clean

Re-applies three workstreams lost when 793d3d7 committed over uncommitted
edits, reconciled against the parallel track's committed dm.py changes:
- followup-sweeper.py: backfill thread_uuid after successful nudge sends;
  record final_nudge_target=main on final-nudge routing (C1/C2)
- response-harvester.py: resolve followups on main-chat replies when
  final_nudge_target=main (C3); harvest ALL [RESULT] markers per message
- dm.py: pre-send placement gate (fail closed when post-nav URL lacks the
  target thread UUID; skips main) — purely additive over 793d3d7+f268d3d
- sidechat_manager.py: wait_for_chat_list() settle-poll for list population
  race (sidebar button renders before titles load)
- new: bin/tests/test_followup_fixes.py (25 tests), bin/placement-audit.py,
  bin/dm-log-taxonomy.py, bin/session-probe.py,
  docs/SIDECHAT-RELIABILITY.md, docs/UUID-ROTATION.md

Verified: 25/25 tests pass, py_compile clean, sweeper/harvester dry-runs clean.
Known limitation: gate catches wrong-placement, not wrong-mapping (false
autoprovision adopting the parked thread needs a creation check).
This commit is contained in:
operator-main
2026-10-04 20:06:58 +00:00
parent b5c5e2ff4c
commit ad9dbca7eb
10 changed files with 1984 additions and 51 deletions
+267
View File
@@ -0,0 +1,267 @@
#!/usr/bin/env python3
"""dm-log-taxonomy.py — READ-ONLY failure taxonomy for fleet DM sidechat reliability.
Reads /home/super/Projects/NetVM/dm-log.jsonl, prints:
1. Event-type counts and send outcome rates
2. Sidechat nav failure taxonomy (per-target, per-agent-pair)
3. Failure timeline (hourly buckets, worst 10-min windows, by node)
4. "Ghost" rate: verified:true sidechat sends with no UUID anywhere
5. Hypothesis evidence tables (nav_failed reasons, placement pairs,
alias sources, autoprovision success, retry distribution)
No writes to any state files. Runs in <1s on the current log size.
"""
import json
import re
import sys
from collections import Counter, defaultdict
from datetime import datetime, timezone, timedelta
LOG = "/home/super/Projects/NetVM/dm-log.jsonl"
WINDOW_H = 48
def parse_ts(s):
if not s:
return None
try:
dt = datetime.fromisoformat(s.replace("Z", "+00:00"))
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
return dt
except Exception:
return None
def main():
now = datetime.now(timezone.utc)
cutoff = now - timedelta(hours=WINDOW_H)
events = []
parse_err = 0
with open(LOG, encoding="utf-8") as f:
for line in f:
line = line.strip()
if not line:
continue
try:
e = json.loads(line)
except Exception:
parse_err += 1
continue
e["_dt"] = parse_ts(e.get("ts"))
events.append(e)
in_win = [e for e in events if e["_dt"] and e["_dt"] >= cutoff]
print(f"dm-log.jsonl: {len(events)} total lines ({parse_err} parse errors)")
print(f"window: last {WINDOW_H}h -> {len(in_win)} events")
if in_win:
print(f" range: {in_win[0]['ts']} .. {in_win[-1]['ts']}")
print("=" * 78)
# ---- 1. event-type counts ------------------------------------------------
types = Counter(e.get("type", "?") for e in in_win)
print("\n[1] EVENT-TYPE COUNTS")
for t, n in types.most_common():
print(f" {n:6d} {t}")
# ---- per-send assembly ---------------------------------------------------
sends = {}
order = []
def S(e):
sid = e.get("id")
if not sid:
return None
if sid not in sends:
sends[sid] = {"id": sid, "events": [], "first_ts": e["_dt"]}
order.append(sid)
s = sends[sid]
s["events"].append(e)
for k in ("agent", "to", "target"):
if k not in s and e.get(k) is not None:
s[k] = e.get(k)
t = e.get("type")
if t == "verified":
s["verified_seen"] = True
if e.get("thread_uuid"):
s["uuids"] = s.get("uuids", set()) | {e["thread_uuid"]}
if t == "sent":
s["sent_seen"] = True
s["sent_verified"] = bool(e.get("verified"))
if e.get("thread_uuid"):
s["uuids"] = s.get("uuids", set()) | {e["thread_uuid"]}
tags = e.get("tags") or {}
if isinstance(tags, dict) and tags.get("thread"):
s["uuids"] = s.get("uuids", set()) | {tags["thread"]}
if t == "send_done":
s["done"] = True
if t == "sidechat_uuid_capture_failed":
s["capture_failed"] = True
s["capture_out"] = (e.get("out") or "")[:80]
if t == "sidechat_autoprovisioned" and e.get("thread_uuid"):
s["uuids"] = s.get("uuids", set()) | {e["thread_uuid"]}
if t == "alias_resolved" and e.get("thread_uuid"):
s["uuids"] = s.get("uuids", set()) | {e["thread_uuid"]}
if t == "placement_mismatch":
s["placement_mismatch"] = True
if t == "pre_send_assert_failed":
s["gate_failed"] = True
s["gate_reason"] = e.get("reason")
if t == "nav_failed":
s["nav_failed"] = True
s["nav_reason"] = e.get("reason") or "transport"
return s
for e in in_win:
S(e)
started = [sends[i] for i in order
if any(e.get("type") == "send_start" for e in sends[i]["events"])]
print(f"\n[1b] SEND OUTCOMES: {len(started)} sends attempted")
oc = Counter()
for s in started:
if s.get("gate_failed"):
oc["loud-fail: pre_send_assert_failed"] += 1
elif s.get("nav_failed"):
oc["loud-fail: nav_failed"] += 1
elif s.get("sent_seen") and s.get("sent_verified"):
oc["sent verified:true"] += 1
elif s.get("sent_seen"):
oc["sent verified:false"] += 1
elif s.get("done"):
oc["send_done, no sent event"] += 1
else:
oc["abandoned (no terminal event)"] += 1
for k, n in oc.most_common():
print(f" {n:5d} {k}")
print(" (note: main_chat_blocked events are policy blocks, not failures;")
print(" followup_register_failed=signing issues, tangential)")
# ---- 2. sidechat taxonomy ------------------------------------------------
def is_sc(s):
return (s.get("target") or "") not in ("main", "", None)
sc = [s for s in started if is_sc(s)]
print(f"\n[2] SIDECHAT TAXONOMY ({len(sc)} sidechat-targeted sends)")
tax = Counter()
per_target = defaultdict(Counter)
per_pair = defaultdict(Counter)
for s in sc:
if s.get("gate_failed"):
cls = "gate_failed:" + str(s.get("gate_reason"))
elif s.get("nav_failed"):
cls = "nav_failed:" + str(s.get("nav_reason"))
elif s.get("capture_failed"):
cls = "uuid_capture_failed"
elif s.get("placement_mismatch"):
cls = "placement_mismatch"
elif s.get("sent_seen") and s.get("sent_verified"):
cls = ("verified:true, NO uuid anywhere (ghost)"
if not s.get("uuids") else "clean verified")
elif s.get("sent_seen"):
cls = "sent verified:false"
elif s.get("done"):
cls = "done, no sent event"
else:
cls = "abandoned"
tax[cls] += 1
per_target[s.get("target") or "?"][cls] += 1
per_pair[f"{s.get('agent') or '?'}->{s.get('to') or '?'}"][cls] += 1
for k, n in tax.most_common():
print(f" {n:5d} {k}")
print("\n per-target (sends, non-clean, top classes):")
for tgt, c in sorted(per_target.items(), key=lambda x: -sum(x[1].values())):
tot = sum(c.values())
bad = tot - c.get("clean verified", 0)
print(f" {tgt}: {tot} sends, {bad} non-clean {dict(c.most_common(4))}")
print("\n per agent-pair:")
for pair, c in sorted(per_pair.items(), key=lambda x: -sum(x[1].values())):
tot = sum(c.values())
bad = tot - c.get("clean verified", 0)
print(f" {pair}: {tot} sends, {bad} non-clean {dict(c.most_common(4))}")
# ---- 4. ghosts ------------------------------------------------------------
ghosts = [s for s in sc if s.get("sent_seen") and s.get("sent_verified")
and not s.get("uuids")]
print(f"\n[4] GHOST RATE: {len(ghosts)}/{len(sc)} "
f"({100.0 * len(ghosts) / len(sc) if sc else 0:.0f}%) verified:true "
f"sidechat sends with no UUID in any event")
gh = Counter(s["first_ts"].strftime("%m-%d %H") for s in ghosts if s["first_ts"])
print(" ghost hours:", dict(sorted(gh.items())))
print(" proxies: uuid_capture_failed="
f"{sum(1 for s in sc if s.get('capture_failed'))}, "
f"placement_mismatch={sum(1 for s in sc if s.get('placement_mismatch'))}, "
f"pre_send_assert_failed={sum(1 for s in sc if s.get('gate_failed'))}")
# ---- 3. timeline ------------------------------------------------------------
print("\n[3] TIMELINE (hourly, sidechat sends; #=bad, .=ok)")
buckets = defaultdict(Counter)
for s in sc:
if not s.get("first_ts"):
continue
hr = s["first_ts"].strftime("%m-%d %H:00")
buckets[hr]["total"] += 1
bad = not (s.get("sent_seen") and s.get("sent_verified")
and not s.get("capture_failed")
and not s.get("placement_mismatch")
and not s.get("gate_failed") and not s.get("nav_failed"))
if bad:
buckets[hr]["bad"] += 1
for hr in sorted(buckets):
t, b = buckets[hr]["total"], buckets[hr]["bad"]
print(f" {hr} total={t:3d} bad={b:3d} {'#' * b}{'.' * (t - b)}")
wins = defaultdict(Counter)
for s in sc:
if not s.get("first_ts"):
continue
w = s["first_ts"].strftime("%m-%d %H:%M")[:-1] + "0"
wins[w]["total"] += 1
if not (s.get("sent_seen") and s.get("sent_verified")):
wins[w]["bad"] += 1
print(" worst 10-min windows (>=3 sends):")
shown = 0
for w, c in sorted(wins.items(), key=lambda x: -x[1]["bad"]):
if c["total"] >= 3 and shown < 8:
print(f" {w} total={c['total']} bad={c['bad']}")
shown += 1
print(" bad rate by sending node:")
by_node = defaultdict(Counter)
for s in sc:
by_node[s.get("agent") or "?"]["total"] += 1
if not (s.get("sent_seen") and s.get("sent_verified")):
by_node[s.get("agent") or "?"]["bad"] += 1
for node, c in sorted(by_node.items(), key=lambda x: -x[1]["total"]):
r = 100.0 * c["bad"] / c["total"] if c["total"] else 0
print(f" {node}: {c['bad']}/{c['total']} bad ({r:.0f}%)")
# ---- 5. hypothesis evidence ---------------------------------------------------
print("\n[5] HYPOTHESIS EVIDENCE")
nfr = Counter(s.get("nav_reason") for s in sc if s.get("nav_failed"))
print(f" nav_failed reasons: {dict(nfr)}")
land = [s for s in sc if s.get("capture_failed")]
print(f" uuid_capture_failed: {len(land)}, "
f"landing-page outs: {sum(1 for s in land if 'muse.ai/' in (s.get('capture_out') or ''))}")
pairs = Counter()
for e in in_win:
if e.get("type") == "placement_mismatch":
pairs[((e.get("expected_uuid") or "?")[:8],
(e.get("actual_uuid") or "?")[:8], e.get("target"))] += 1
print(" placement_mismatch expected->actual:")
for (a, b, t), n in pairs.most_common(6):
print(f" {n:3d} exp={a}.. act={b}.. target={t}")
ar = Counter(e.get("source") for e in in_win if e.get("type") == "alias_resolved")
print(f" alias_resolved sources: {dict(ar)} (None = field absent, older events)")
ap_try = sum(1 for e in in_win if e.get("type") == "sidechat_autoprovision_start")
ap_ok = sum(1 for e in in_win if e.get("type") == "sidechat_autoprovisioned"
and e.get("thread_uuid"))
print(f" autoprovision: {ap_try} attempts -> {ap_ok} with uuid ({100.0 * ap_ok / ap_try if ap_try else 0:.0f}%)")
rt = Counter()
for s in started:
n = sum(1 for e in s["events"] if e.get("type") == "retry")
if n:
rt[n] += 1
print(f" retry distribution (sends with >=1 retry): {dict(sorted(rt.items()))}")
print("\n DONE.")
if __name__ == "__main__":
sys.exit(main())
Executable → Regular
+71
View File
@@ -62,6 +62,9 @@ TAG_NUDGES_MAX = 10
TAG_NUDGES_DEFAULT = 2
TAG_ROUTE_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_-]{0,63}$")
TAG_THREAD_RE = re.compile(r"^[A-Za-z0-9_-]{1,64}$")
# Thread-UUID pattern, shared by nav parsing and the pre-send gate.
UUID_RE = (r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-"
r"[0-9a-f]{4}-[0-9a-f]{12}")
# Trailing-token form: [reply:expected] [reply:timeout=7200] ...
TAG_TOKEN_RE = re.compile(r"\[([A-Za-z][A-Za-z0-9:]*)"
r"(?:=([^\]]*))?\]\s*$")
@@ -365,6 +368,52 @@ def verify_placement(recipient, msg_id, target, thread_uuid):
detail["error"] = str(e)[:200]
return False, detail
def assert_pre_send_placement(recipient, target, thread_uuid, direct_nav_done=False):
"""Pre-send placement gate (2026-10-04; restored after 19:43Z clobber).
Independently samples the recipient browser's URL and requires the
expected thread UUID in it BEFORE any send. Fails closed: returns
(False, detail) and the caller must abort the send. 'main' skips.
"""
detail = {"target": target}
if target == "main":
return True, detail
if not thread_uuid:
# Fresh autoprovision: SPA assigns UUID only on first send.
# Strongest pre-send signal: browser must be on a /thread/ page
# (the incident was a send parked on the muse.ai landing page).
_u_rc, _u_out, _u_err = run_full(
f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} url")
actual = (_u_out or "").strip()
detail["actual_url"] = actual[:200]
if "/thread/" not in actual:
detail["reason"] = "not_on_thread_page"
detail["expected"] = "a /thread/ URL (fresh sidechat, UUID assigned on first send)"
return False, detail
return True, detail
expected = thread_uuid.lower()
detail["expected_uuid"] = expected
if not direct_nav_done:
# Title-search nav only: re-navigate by direct /thread/<uuid> URL.
_rc, _out, _err = run_full(
f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat use {expected}")
if _rc != 0 or "NOTFOUND" in _out or expected not in _out.lower():
detail.update({"reason": "direct_nav_failed", "rc": _rc,
"out": _out[:200], "err": _err[:200]})
return False, detail
time.sleep(2)
else:
time.sleep(1)
_u_rc, _u_out, _u_err = run_full(
f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} url")
actual = (_u_out or "").strip()
detail["actual_url"] = actual[:200]
if expected not in actual.lower():
detail["reason"] = "url_mismatch"
return False, detail
return True, detail
def dm_send(agent, target, message, verify=True, raw=False,
to_agent=None, tags=None, nudge_meta=None,
allow_main_chat=False):
@@ -444,6 +493,7 @@ def dm_send(agent, target, message, verify=True, raw=False,
thread_uuid = None
thread_url = None
is_new_sidechat = False
nav_is_uuid = False
# Resolve well-known aliases or dynamic thread mappings to UUIDs.
nav_target = resolve_sidechat_target(target)
if nav_target != target:
@@ -454,6 +504,7 @@ def dm_send(agent, target, message, verify=True, raw=False,
else:
# Reset-at-Begin: If searching sidebar by title (not a direct UUID), reset to main first to expose the sidebar
is_uuid = bool(re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", nav_target.strip().lower()))
nav_is_uuid = is_uuid
if not is_uuid:
run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat main")
time.sleep(1)
@@ -479,6 +530,12 @@ def dm_send(agent, target, message, verify=True, raw=False,
file=sys.stderr)
sys.exit(1)
is_new_sidechat = True
# If create returned a real thread URL (not /thread/new),
# pin it now so the pre-send gate can assert it.
_cm = re.search(r"/thread/(" + UUID_RE + r")", _c_out, re.I)
if _cm:
thread_uuid = _cm.group(1).lower()
thread_url = "https://muse.ai/thread/" + thread_uuid
log_event({"type": "nav_ok", "id": msg_id, "agent": agent, "to": recipient,
"target": target, "status": "sidechat_created_pending_uuid",
"browser_url": thread_url})
@@ -509,6 +566,20 @@ def dm_send(agent, target, message, verify=True, raw=False,
"target": target, "thread_uuid": thread_uuid,
"browser_url": thread_url})
# Pre-send placement assertion: independently confirm the browser is on
# the target thread BEFORE any send. Failure fails loudly (no ghost sends).
_gate_uuid = thread_uuid or (nav_target.lower() if nav_is_uuid else None)
_ps_ok, _ps_detail = assert_pre_send_placement(
recipient, target, _gate_uuid, direct_nav_done=nav_is_uuid)
if not _ps_ok:
log_event({"type": "pre_send_assert_failed", "id": msg_id, "agent": agent,
"to": recipient, **_ps_detail})
print(f"DM {msg_id} from {agent} to {recipient}/{target}: FAILED "
f"(pre-send placement assertion failed: {_ps_detail.get('reason')}; "
f"expected={_ps_detail.get('expected_uuid') or _ps_detail.get('expected')}; "
f"actual_url={_ps_detail.get('actual_url')})", file=sys.stderr)
sys.exit(1)
time.sleep(2)
# Send (raw mode: no truncation — signatures must survive intact)
Executable → Regular
+52
View File
@@ -14,6 +14,7 @@ Usage:
import argparse
import json
import os
import re
import subprocess
import sys
import time
@@ -26,6 +27,7 @@ BIN_DIR = NETVM_ROOT / "bin"
JOBS_DIR = NETVM_ROOT / "jobs"
FOLLOWUPS_FILE = NETVM_ROOT / "followups.json"
JOB_LOG = NETVM_ROOT / "job-log.jsonl"
DM_LOG = NETVM_ROOT / "dm-log.jsonl"
DM_PY = BIN_DIR / "dm.py"
DISPATCH_PY = BIN_DIR / "job-dispatch.py"
@@ -101,6 +103,46 @@ def send_dm(sender, recipient, target, text):
except Exception as e:
return False, str(e)
_DM_ID_RE = re.compile(r"\bDM ([0-9a-f]{8})\b")
_THREAD_UUID_RE = re.compile(
r"^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$",
re.IGNORECASE,
)
def _resolve_nudge_thread_uuid(nudge_output):
"""Parse the DM id from dm.py's stdout, scan the tail (last 5000 lines)
of dm-log.jsonl for that id's sidechat_autoprovisioned event, falling
back to the sent event's tags.thread. Returns the UUID or None."""
m = _DM_ID_RE.search(nudge_output or "")
if not m:
return None
nudge_id = m.group(1)
fallback = None
try:
with open(DM_LOG, "r", encoding="utf-8") as f:
lines = f.readlines()
except FileNotFoundError:
return None
for line in lines[-5000:]:
line = line.strip()
if not line:
continue
try:
e = json.loads(line)
except Exception:
continue
if e.get("id") != nudge_id:
continue
if e.get("type") == "sidechat_autoprovisioned":
uuid = e.get("thread_uuid") or ""
if _THREAD_UUID_RE.fullmatch(uuid):
return uuid
elif e.get("type") == "sent":
cand = ((e.get("tags") or {}).get("thread")) or ""
if _THREAD_UUID_RE.fullmatch(cand):
fallback = cand
return fallback
def sweep_cycle(dry_run=False):
followups = load_followups()
@@ -181,6 +223,15 @@ def sweep_cycle(dry_run=False):
timeout_s = rec.get("timeout_s", 1800)
step_s = max(300, timeout_s // (nudges_allowed + 1))
rec["deadline"] = (now + timedelta(seconds=step_s)).isoformat()
# C1: backfill the thread the nudge actually landed in.
landed_uuid = _resolve_nudge_thread_uuid(out)
if landed_uuid and rec.get("thread_uuid") != landed_uuid:
rec["thread_uuid"] = landed_uuid
print(f"Sweeper: backfilled thread_uuid={landed_uuid} "
f"for followup {dm_id}", file=sys.stderr)
# C2: record final-nudge routing for the harvester.
if is_final:
rec["final_nudge_target"] = "main"
modified = True
append_job_log({
"ts": utcnow_str(),
@@ -189,6 +240,7 @@ def sweep_cycle(dry_run=False):
"nudge_num": nudge_num,
"recipient": recipient,
"target": delivery_target,
"thread_uuid": rec.get("thread_uuid"),
})
else:
print(f"Sweeper WARNING: nudge send failed: {out}", file=sys.stderr)
+253
View File
@@ -0,0 +1,253 @@
#!/usr/bin/env python3
"""placement-audit.py -- readable audit of sidechat placement / pre-send gate
activity in dm-log.jsonl.
The gate and navigation emit machine-readable events (pre_send_assert_failed,
placement_failed, placement_mismatch, sidechat_uuid_capture_failed, nav_failed,
...). This CLI turns them into human-readable summaries so the logs stay
useful for troubleshooting sidechat reliability.
Usage:
placement-audit.py --since 24h # summary over the last 24 hours
placement-audit.py --since 7d # ... 7 days
placement-audit.py --since 60m # ... 60 minutes
placement-audit.py --watch # live tail of placement-relevant events
Read-only: never writes, no network, no browser.
"""
import argparse
import json
import re
import sys
import time
from collections import Counter, defaultdict
from datetime import datetime, timedelta, timezone
DM_LOG = "/home/super/Projects/NetVM/dm-log.jsonl"
JOB_LOG = "/home/super/Projects/NetVM/job-log.jsonl"
# Event types that count as placement/navigation failures.
FAILURE_TYPES = {
"sidechat_uuid_capture_failed",
"placement_mismatch",
"placement_failed",
"pre_send_assert_failed",
"nav_failed",
"send_failed",
"failed",
}
# Event types worth streaming in --watch (dm-log).
WATCH_TYPES_DM = FAILURE_TYPES | {
"verified",
"sent",
"nav_ok",
"sidechat_autoprovisioned",
}
def parse_since(s):
m = re.fullmatch(r"(\d+)(m|h|d)", (s or "").strip().lower())
if not m:
raise ValueError(f"bad --since value {s!r}; use like 60m, 24h, 7d")
n, unit = int(m.group(1)), m.group(2)
return timedelta(minutes=n) if unit == "m" else timedelta(hours=n) if unit == "h" else timedelta(days=n)
def parse_ts(ts):
if not ts:
return None
try:
dt = datetime.fromisoformat(str(ts).replace("Z", "+00:00"))
except ValueError:
return None
if dt.tzinfo is None:
dt = dt.replace(tzinfo=timezone.utc)
return dt
def short_uuid(u):
u = str(u or "")
return u[:8] + "..." if len(u) > 12 else u
def fmt_ts(e):
dt = parse_ts(e.get("ts"))
return dt.strftime("%m-%d %H:%M:%S") if dt else "?"
def expected_actual(e):
"""Return (expected, actual) strings for a failure event."""
t = e.get("type")
if t in ("pre_send_assert_failed", "placement_failed"):
return e.get("expected_uuid"), e.get("actual_url")
if t == "placement_mismatch":
return e.get("expected_uuid"), e.get("actual_uuid")
if t == "sidechat_uuid_capture_failed":
return "(uuid capture)", e.get("out")
if t == "nav_failed":
return e.get("expected") or e.get("reason"), e.get("got") or e.get("out")
if t in ("send_failed", "failed"):
return e.get("reason"), (e.get("send_out") or e.get("err") or "")[:120]
return None, None
def one_line(e):
"""Compact one-line summary of an event."""
t = e.get("type", "?")
who = f"{e.get('agent', '?')}->{e.get('to', '?')}/{e.get('target', '?')}"
base = f"{fmt_ts(e)} {t:28s} {str(e.get('id', ''))[:8]:8s} {who}"
if t in FAILURE_TYPES:
reason = e.get("reason") or ""
exp, act = expected_actual(e)
extra = f" reason={reason}" if reason else ""
if exp or act:
extra += f" expected={short_uuid(exp) if exp and len(str(exp)) > 20 else exp} actual={act}"
return base + extra
if t == "sent":
return base + f" verified={e.get('verified')}"
if t == "verified":
return base + f" placement={e.get('placement')} thread={short_uuid(e.get('thread_uuid'))}"
if t == "nav_ok":
return base + f" status={e.get('status')} url={e.get('browser_url')}"
if t == "sidechat_autoprovisioned":
return base + f" thread={short_uuid(e.get('thread_uuid'))}"
return base
def iter_log(path, cutoff=None):
"""Yield parsed events from a jsonl file, optionally filtered by ts."""
try:
f = open(path, encoding="utf-8")
except OSError as ex:
print(f"warning: cannot open {path}: {ex}", file=sys.stderr)
return
with f:
for line in f:
line = line.strip()
if not line:
continue
try:
e = json.loads(line)
except json.JSONDecodeError:
continue
if cutoff is not None:
dt = parse_ts(e.get("ts"))
if dt is None or dt < cutoff:
continue
yield e
def cmd_since(args):
try:
delta = parse_since(args.since)
except ValueError as ex:
print(str(ex), file=sys.stderr)
return 1
cutoff = datetime.now(timezone.utc) - delta
sends = Counter() # (agent, target) -> sent events
sent_ok = Counter() # (agent, target) -> sent with verified=True
taxonomy = Counter() # failure label -> count
failures = [] # failure events, for the recent list
total = 0
for e in iter_log(DM_LOG, cutoff):
total += 1
t = e.get("type")
key = (e.get("agent") or "?", e.get("target") or "?")
if t == "sent":
sends[key] += 1
if e.get("verified") is True:
sent_ok[key] += 1
if t in FAILURE_TYPES:
reason = e.get("reason") or ""
label = f"{t}" + (f":{reason}" if reason else "")
taxonomy[label] += 1
failures.append(e)
# followup nudge failures live in job-log.jsonl
for e in iter_log(JOB_LOG, cutoff):
if e.get("type") == "followup_nudge_failed":
taxonomy["followup_nudge_failed"] += 1
failures.append({"type": "followup_nudge_failed",
"ts": e.get("ts"), "id": e.get("dm_id"),
"agent": "sweeper", "to": e.get("recipient"),
"target": e.get("target"),
"reason": (e.get("error") or "")[:100]})
print(f"== placement audit: last {args.since} (since {cutoff.strftime('%Y-%m-%d %H:%M UTC')}) ==")
print(f"dm-log events scanned: {total}")
print()
print("-- per-target sends --")
print(f"{'agent':10s} {'target':28s} {'sent':>5s} {'verified':>8s} {'unver':>6s}")
for (agent, target), n in sorted(sends.items(), key=lambda kv: -kv[1]):
ok = sent_ok.get((agent, target), 0)
print(f"{agent:10s} {target:28s} {n:5d} {ok:8d} {n - ok:6d}")
if not sends:
print("(no sends in window)")
print()
print("-- failure taxonomy --")
if taxonomy:
for label, n in taxonomy.most_common():
print(f"{n:5d} {label}")
else:
print("(no placement failures in window)")
print()
print("-- 10 most recent failures --")
failures.sort(key=lambda e: parse_ts(e.get("ts")) or datetime.min.replace(tzinfo=timezone.utc),
reverse=True)
for e in failures[:10]:
print(one_line(e))
if not failures:
print("(none)")
return 0
def cmd_watch(_args):
print("watching dm-log.jsonl for placement events (Ctrl-C to stop)...", file=sys.stderr)
try:
f = open(DM_LOG, encoding="utf-8")
except OSError as ex:
print(f"cannot open {DM_LOG}: {ex}", file=sys.stderr)
return 1
with f:
f.seek(0, 2) # start at end: live view only
try:
while True:
line = f.readline()
if not line:
time.sleep(2)
continue
line = line.strip()
if not line:
continue
try:
e = json.loads(line)
except json.JSONDecodeError:
continue
if e.get("type") in WATCH_TYPES_DM:
print(one_line(e), flush=True)
except KeyboardInterrupt:
print("\nstopped.", file=sys.stderr)
return 0
def main(argv=None):
ap = argparse.ArgumentParser(description="Audit sidechat placement / gate events in dm-log.jsonl")
ap.add_argument("--since", metavar="60m|24h|7d",
help="summarize events newer than this (e.g. 60m, 24h, 7d)")
ap.add_argument("--watch", action="store_true",
help="live tail of placement-relevant events")
args = ap.parse_args(argv)
if args.watch:
return cmd_watch(args)
if args.since:
return cmd_since(args)
ap.print_help()
return 2
if __name__ == "__main__":
sys.exit(main())
Executable → Regular
+64 -50
View File
@@ -71,6 +71,18 @@ except ImportError:
VALID_AGENTS = ["muse", "pip", "646", "opm"]
DEFAULT_PORTS = {"muse": 9410, "pip": 9420, "646": 9430, "opm": 9440}
# Matches EVERY [RESULT <job_id>] marker in a message (use with finditer, not
# search). The result text is lazy and stops before the next marker (or end of
# text), so a message closing two jobs records each with its own text instead
# of the first marker greedily swallowing the second.
RESULT_RE = re.compile(r"\[RESULT\s+([A-Za-z0-9_-]+)\]\s*(.*?)(?=\[RESULT\s|\Z)", re.S)
def iter_result_markers(text):
"""Yield (job_id, result_text) for every [RESULT <job_id>] marker in text."""
for m in RESULT_RE.finditer(text or ""):
yield m.group(1).strip(), m.group(2).strip()
def utcnow():
return datetime.now(timezone.utc).isoformat()
@@ -353,30 +365,29 @@ def harvest_main_feed(cdp, agent, watermarks, followups, dry_run=False):
append_jsonl(CHAT_HISTORY_LOG, record)
if author == "assistant":
m_res = re.search(r"\[RESULT\s+([A-Za-z0-9_-]+)\]\s*(.*)", text, re.S)
result_job_id = None
if m_res:
job_id = m_res.group(1).strip()
result_job_id = job_id
result_text = m_res.group(2).strip()
is_fail = is_fail_result(result_text)
job_results += 1
job_record = {
"ts": utcnow(),
"type": "job_result",
"job_id": job_id,
"agent": agent,
"success": not is_fail,
"result_snippet": result_text[:300],
"thread_id": "main",
"msg_id": mid,
}
if not dry_run:
append_jsonl(JOB_LOG, job_record)
trigger_chain_next(job_id, result_text, success=not is_fail)
markers = list(iter_result_markers(text))
if markers:
for job_id, result_text in markers:
is_fail = is_fail_result(result_text)
job_results += 1
clear_matching_followups(followups, agent, "main", mid, text, dry_run,
job_id=result_job_id)
job_record = {
"ts": utcnow(),
"type": "job_result",
"job_id": job_id,
"agent": agent,
"success": not is_fail,
"result_snippet": result_text[:300],
"thread_id": "main",
"msg_id": mid,
}
if not dry_run:
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)
else:
clear_matching_followups(followups, agent, "main", mid, text, dry_run)
return new_messages, new_wm, job_results
@@ -474,35 +485,31 @@ def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run
if not dry_run:
append_jsonl(CHAT_HISTORY_LOG, record)
# Check for [RESULT <job_id>] in assistant messages
if author == "assistant":
m_res = re.search(r"\[RESULT\s+([A-Za-z0-9_-]+)\]\s*(.*)", text, re.S)
result_job_id = None
if m_res:
job_id = m_res.group(1).strip()
result_job_id = job_id
result_text = m_res.group(2).strip()
is_fail = is_fail_result(result_text)
job_results += 1
markers = list(iter_result_markers(text))
if markers:
for job_id, result_text in markers:
is_fail = is_fail_result(result_text)
job_results += 1
job_record = {
"ts": utcnow(),
"type": "job_result",
"job_id": job_id,
"agent": agent,
"success": not is_fail,
"result_snippet": result_text[:300],
"thread_id": thread_id,
"msg_id": mid,
}
if not dry_run:
append_jsonl(JOB_LOG, job_record)
# Trigger pipeline chaining or next job if configured
trigger_chain_next(job_id, result_text, success=not is_fail)
# Check and clear pending follow-ups
clear_matching_followups(followups, agent, thread_id, mid, text, dry_run,
job_id=result_job_id)
job_record = {
"ts": utcnow(),
"type": "job_result",
"job_id": job_id,
"agent": agent,
"success": not is_fail,
"result_snippet": result_text[:300],
"thread_id": thread_id,
"msg_id": mid,
}
if not dry_run:
append_jsonl(JOB_LOG, job_record)
# 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)
else:
clear_matching_followups(followups, agent, thread_id, mid, text, dry_run)
return new_messages, new_wm, job_results
@@ -610,6 +617,8 @@ def clear_matching_followups(followups, agent, thread_id, mid, text, dry_run=Fal
Matches on thread identity (thread_uuid or target='main') OR on job_id
(from a [RESULT <job_id>] reply). The job_id path works regardless of
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).
"""
if not followups:
return
@@ -627,6 +636,11 @@ def clear_matching_followups(followups, agent, thread_id, mid, text, dry_run=Fal
match_thread = True
elif f_rec.get("thread_uuid") and f_rec.get("thread_uuid") == thread_id:
match_thread = True
elif f_rec.get("final_nudge_target") == "main" and thread_id == "main":
# C3: the sweeper routed the final nudge to main chat, so a
# main-chat reply resolves even when the followup target is a
# sidechat.
match_thread = True
# Match by job_id (from [RESULT <job_id>]) -- works regardless of
# thread_uuid or target. This is an ADDITIONAL path, not a replacement.
+183
View File
@@ -0,0 +1,183 @@
#!/usr/bin/env python3
"""
session-probe.py — classify a fleet agent's muse.ai login state via CDP.
READ-ONLY: performs a single Runtime.evaluate reading document.title,
location.href, and body markers. Never clicks, navigates, types, or mutates
session state in any way. Uses the shared CDP queue at PRIORITY_LOW so it
never blocks operator sends or the harvester.
Login states (per operator AGENTS.md, 2026-10-03):
LOGGED_IN title contains "Chat —"/"Muse —" (SPA booted authenticated)
LANDING title == "muse.ai" (landing page, not logged in)
LOGGED_OUT body shows a "Log in" affordance
OTP_PROMPT body contains "To log in, enter the code"
ACCOUNT_SELECTION body contains "Your email matches multiple accounts"
UNKNOWN none of the above matched
CDP_UNREACHABLE browser/CDP could not be reached at all
parked_on_landing (bool, separate field): LOGGED_IN but the current URL is
https://muse.ai/ — the watchdog relaunches browsers with the landing page
as start URL, so a fresh browser is parked there until first nav. Session
is fine; nav may proceed (but allow for SPA boot race).
Exit code: 0 when LOGGED_IN, 1 otherwise (suitable for pre-nav gating).
Stdout: one JSON line: {agent, state, title, url, has_input, checked_at}.
INVOCATION (important): CDP is only reachable from inside the node's netns
(Chromium binds DevTools to loopback). Always run via:
netvm-exec.sh <agent> -- python3 /home/super/Projects/NetVM/bin/session-probe.py --agent <agent>
Running it on the bl host directly will report CDP_UNREACHABLE even when the
browser is healthy.
"""
import argparse
import contextlib
import importlib.util
import json
import os
import sys
import urllib.request
from datetime import datetime, timezone
BIN_DIR = os.path.dirname(os.path.abspath(__file__))
sys.path.insert(0, BIN_DIR)
try:
from cdp_queue import cdp_slot, PRIORITY_LOW
HAS_CDP_QUEUE = True
except ImportError:
HAS_CDP_QUEUE = False
import websocket # noqa: E402 (after sys.path tweak, mirrors muse-chat-api.py)
def _load_accounts():
path = os.path.join(BIN_DIR, "netvm-registry.py")
spec = importlib.util.spec_from_file_location("netvm_registry", path)
mod = importlib.util.module_from_spec(spec)
spec.loader.exec_module(mod)
accounts = {}
for node, rec in mod.load().items():
accounts[node] = (node, "http://127.0.0.1:%d/json/list" % rec["cdp_port"])
return accounts
def _ev(ws, expr):
"""Runtime.evaluate with event draining (copied pattern from muse-chat-api.py)."""
ws.send(json.dumps({
"id": 1, "method": "Runtime.evaluate",
"params": {"expression": expr, "returnByValue": True},
}))
for _ in range(50):
resp = json.loads(ws.recv())
if resp.get("id") == 1:
break
else:
return None
return resp.get("result", {}).get("result", {}).get("value")
_STATE_JS = """(() => {
const body = document.body ? document.body.innerText.slice(0, 4000) : "";
const has_input = !!document.querySelector(
'[contenteditable="true"], textarea[placeholder*="Message"], div[role="textbox"]');
return JSON.stringify({
title: document.title || "",
url: location.href || "",
body: body,
has_input: has_input,
});
})()"""
def classify(title, url, body, has_input):
t = (title or "").strip()
b = (body or "")
if "Your email matches multiple accounts" in b:
return "ACCOUNT_SELECTION"
if "To log in, enter the code" in b:
return "OTP_PROMPT"
if t == "muse.ai":
return "LANDING"
# A "Chat —"/"Muse —" title means the SPA booted with an authenticated
# session (the landing page title is exactly "muse.ai"). The browser may
# still be parked on "/" (fresh relaunch start URL) — reported separately
# via parked_on_landing, not as a session failure.
if ("Chat \u2014" in t) or ("Muse \u2014" in t) or ("Chat -" in t) or ("Muse -" in t):
return "LOGGED_IN"
if "Log in" in b:
return "LOGGED_OUT"
return "UNKNOWN"
def probe(agent, timeout=15):
accounts = _load_accounts()
if agent not in accounts:
return {"agent": agent, "state": "UNKNOWN",
"error": "no such agent in registry",
"checked_at": _now()}
node, cdp_url = accounts[agent]
try:
with urllib.request.urlopen(cdp_url, timeout=5) as r:
targets = json.load(r)
except Exception as e:
return {"agent": agent, "state": "CDP_UNREACHABLE",
"error": "cdp list failed: %s" % str(e)[:120],
"checked_at": _now()}
pages = [t for t in targets if t.get("type") == "page"]
if not pages:
return {"agent": agent, "state": "CDP_UNREACHABLE",
"error": "no page target", "checked_at": _now()}
slot = cdp_slot(node, priority=PRIORITY_LOW) if HAS_CDP_QUEUE \
else contextlib.nullcontext()
try:
with slot:
ws = websocket.create_connection(
pages[0]["webSocketDebuggerUrl"], timeout=timeout)
try:
raw = _ev(ws, _STATE_JS)
finally:
ws.close()
except Exception as e:
return {"agent": agent, "state": "CDP_UNREACHABLE",
"error": "cdp session failed: %s" % str(e)[:120],
"checked_at": _now()}
if not raw:
return {"agent": agent, "state": "UNKNOWN",
"error": "empty evaluate result", "checked_at": _now()}
try:
snap = json.loads(raw)
except Exception:
return {"agent": agent, "state": "UNKNOWN",
"error": "unparseable evaluate result", "checked_at": _now()}
state = classify(snap.get("title"), snap.get("url"),
snap.get("body"), snap.get("has_input"))
url = snap.get("url") or ""
parked = url.rstrip("/") in ("https://muse.ai", "https://muse.ai/")
return {"agent": agent, "state": state, "title": snap.get("title"),
"url": url[:120], "parked_on_landing": parked,
"has_input": snap.get("has_input"),
"checked_at": _now()}
def _now():
return datetime.now(timezone.utc).isoformat()
def main():
accounts = _load_accounts()
p = argparse.ArgumentParser(
description="Classify a fleet agent's muse.ai login state (read-only).")
p.add_argument("--agent", required=True, choices=sorted(accounts.keys()))
p.add_argument("--timeout", type=int, default=15)
args = p.parse_args()
result = probe(args.agent, timeout=args.timeout)
print(json.dumps(result))
sys.stdout.flush()
sys.exit(0 if result.get("state") == "LOGGED_IN" else 1)
if __name__ == "__main__":
main()
+40 -1
View File
@@ -4,6 +4,7 @@ Robust sidechat management for muse-chat-api.py.
Provides:
- ensure_sidebar(ws): Opens sidebar if closed, with retry
- wait_for_chat_list(ws, name=None): Polls until sidebar list content loads
- list_sidechats(ws): Returns list of side chat names, with retry
- All operations logged for audit
@@ -56,13 +57,51 @@ def ensure_sidebar(ws, max_retries=3):
return !!document.querySelector('[data-testid="hatch-chat-compose"]');
})()""")
def wait_for_chat_list(ws, name=None, timeout_s=15):
"""Poll until the sidebar chat list content has loaded (and optionally
contains `name`, case-insensitive substring match).
Readiness requires at least one plausible chat-title line under the
"Side chats" header -- the header itself renders before items populate,
and the compose button that ensure_sidebar() keys on renders earlier
still. Clicking a chat title before the list is ready is a known
nav-miss contributor (the click lands nowhere and the SPA stays on /).
Per-poll CDP errors are treated as not-ready (browser may be mid-render
or mid-flap). Returns True when ready, False on timeout. Note: an
account with genuinely zero sidechats will always time out here.
"""
import re as _re
want = (name or "").strip().lower()
deadline = time.time() + timeout_s
while time.time() < deadline:
try:
text = _ev(ws, "document.body.innerText") or ""
except Exception:
text = ""
idx = text.find("Side chats")
if idx != -1:
section = text[idx:idx + 2000].lower()
lines = [l.strip() for l in section.split("\n")]
titles = [l for l in lines[1:21]
if 5 < len(l) < 80
and not _re.fullmatch(r"\d+[mh]", l)
and l != "unread updates"]
if titles and (not want or want in section):
return True
time.sleep(1)
return False
def list_sidechats(ws):
"""
List side chat names. Returns list of strings.
Ensures sidebar is open first.
Ensures sidebar is open first, then waits for the list content to
populate before scraping (the open signal fires before content loads).
"""
if not ensure_sidebar(ws):
return []
if not wait_for_chat_list(ws):
return []
result = _ev(ws, """(() => {
const text = document.body.innerText;
+694
View File
@@ -0,0 +1,694 @@
#!/usr/bin/env python3
"""
test_followup_fixes.py -- verification harness for the fleet DM machinery fixes.
Workstream 4 of 5 (verification). Covers the contracts of three sibling
workstreams editing code in /home/super/Projects/NetVM/bin/ :
ws-1 followup-sweeper.py
- backfill rec["thread_uuid"] from dm-log sidechat_autoprovisioned
events (so a followup whose thread provisioned late becomes
resolvable instead of a ghost);
- record rec["final_nudge_target"] = "main" when the final nudge is
routed to main chat.
ws-2 response-harvester.py
- resolve followups on main-chat assistant replies when the followup
carries final_nudge_target == "main";
- extract ALL [RESULT <job_id>] markers per message (finditer),
not just the first (re.search).
ws-3 dm.py
- assert the post-nav browser URL contains the target thread UUID
BEFORE sending; fail loudly otherwise. target == "main" skips the
assertion.
Test layout
-----------
test_contract_* Executable specs of the intended behavior, written against
small local reference predicates. These run GREEN now and
pin the exact semantics the siblings must satisfy.
test_wired_* The same behaviors exercised against the REAL modules
(importlib-loaded from bin/; stdlib-only, side-effect-free
imports; tmp fixtures). Where a sibling has not landed the
change yet, these FAIL with an exact deviation report --
that is the intended signal, not a bug in the harness.
Pure unit tests: NO browser, NO network, NO live DMs, NO writes to live state
files (followups.json, job-log.jsonl, dm-log.jsonl). Tmp copies/fixtures only.
Run: python3 test_followup_fixes.py (built-in runner below)
pytest test_followup_fixes.py (also compatible)
Snapshot note (2026-10-04 ~19:45Z): ws-2 has fully landed --
clear_matching_followups() grew final_nudge_target AND job_id matching paths,
and iter_result_markers() extracts all markers via finditer. ws-1 has
partially landed -- DM_LOG constant, _resolve_nudge_thread_uuid() helper, and
final_nudge_target marking are in; the ghost fail-fast for null-thread_uuid
recs is unchanged (no pre-check backfill). ws-3 (dm.py pre-send URL
assertion) had not landed at the time of writing.
"""
import importlib.util
import inspect
import json
import os
import re
import sys
import tempfile
from pathlib import Path
BIN_DIR = Path(__file__).resolve().parent.parent # .../bin/tests -> .../bin
# --------------------------------------------------------------------------
# module loading (side-effect-free: all three modules are stdlib-only at
# import time and do real work only inside functions / __main__)
# --------------------------------------------------------------------------
def _load(mod_name, filename):
path = BIN_DIR / filename
spec = importlib.util.spec_from_file_location(mod_name, str(path))
mod = importlib.util.module_from_spec(spec)
spec.loader.exec_module(mod)
return mod
class _Skip(Exception):
pass
def skip(reason):
raise _Skip(reason)
# --------------------------------------------------------------------------
# CONTRACT REFERENCE IMPLEMENTATIONS (executable specs -- green now)
# --------------------------------------------------------------------------
def contract_backfill_thread_uuid(rec, dm_log_events):
"""ws-1 spec: adopt thread_uuid from a matching sidechat_autoprovisioned
dm-log event. Never overwrite an existing UUID; never write null."""
if rec.get("thread_uuid"):
return False
for ev in dm_log_events or []:
if ev.get("type") != "sidechat_autoprovisioned":
continue
uuid = ev.get("thread_uuid")
if not uuid:
continue
if ev.get("id") == rec.get("dm_id") or ev.get("target") == rec.get("target"):
rec["thread_uuid"] = uuid
return True
return False
def contract_route_nudge(rec):
"""ws-1 spec: final-nudge routing + marking. Returns (delivery_target,
is_final); marks rec['final_nudge_target']='main' on the final nudge."""
nudge_num = rec.get("nudges_sent", 0) + 1
nudges_allowed = rec.get("nudges_allowed", 2)
is_final = (nudge_num == nudges_allowed)
delivery_target = "main" if is_final else rec.get("target", "main")
if is_final:
rec["final_nudge_target"] = "main"
return delivery_target, is_final
def contract_followup_matches(rec, agent, thread_id):
"""ws-2 spec: extended matching predicate for clear_matching_followups."""
if rec.get("status") != "pending":
return False
if rec.get("recipient") != agent:
return False
if rec.get("target") == "main" and thread_id == "main":
return True
if rec.get("thread_uuid") and rec.get("thread_uuid") == thread_id:
return True
# NEW (ws-2 contract): the final nudge went to main, so a main-chat
# assistant reply resolves the followup even though target != "main".
if rec.get("final_nudge_target") == "main" and thread_id == "main":
return True
return False
_RESULT_MARKER_RE = re.compile(r"\[RESULT\s+([A-Za-z0-9_-]+)\]")
def contract_extract_results(text):
"""ws-2 spec: extract EVERY [RESULT <job_id>] marker; each result text
runs from its marker to the next marker (or end of text)."""
ms = list(_RESULT_MARKER_RE.finditer(text or ""))
out = []
for i, m in enumerate(ms):
seg_end = ms[i + 1].start() if i + 1 < len(ms) else len(text)
out.append((m.group(1), text[m.end():seg_end].strip()))
return out
_UUID_RE = re.compile(
r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", re.I)
def contract_assert_post_nav_url(actual_url, expected_uuid, target):
"""ws-3 spec: (ok, reason). target == 'main' skips the assertion."""
if target == "main":
return True, "main_skipped"
if not expected_uuid:
return False, "no_expected_uuid"
if expected_uuid.lower() in (actual_url or "").lower():
return True, "uuid_present"
return False, "url_mismatch"
# --------------------------------------------------------------------------
# fixtures
# --------------------------------------------------------------------------
# Exact text of opm's 2026-10-04 16:54:17Z message (chat-history.jsonl,
# msg assistant-msg-a5afa9c6-...): carries TWO [RESULT] markers; the old
# re.search path recorded only the first, swallowing the 163439 marker.
TWO_MARKER_TEXT = 'All 10 public services are up, bl is up, tunnels are holding except muse\'s long-standing 7682 terminal gap, and the DM queue is drained — fleet operational.\n\n[RESULT pipe-demo-step1-20261004-165058-b08f5fb5] OK: Fleet healthy and operational\n\nOn the nudge: I never received the 16:35:04 request. I traced it in bl\'s dm-log — job `3ef91457` (pipe-demo-step1 from 16:34:39) targeted side chat "pipe-1b4579", but `sidechat_uuid_capture_failed` fired at 16:34:59 (nav landed on the muse.ai landing page, thread_uuid null), and the "verified: true" at 16:35:04 is the known placement-blind verification. The message went somewhere I can\'t see, so the followup tracker correctly saw no reply and nudged. Two things to fix on the pipeline side: assert the post-nav URL holds the target UUID before sending (never trust the verified flag alone), and dedupe nudges — a second identical nudge (`5e078808`) fired at 16:52:01.\n\nClosing the original loop anyway:\n\n[RESULT pipe-demo-step1-20261004-163439-8a408441] OK: Fleet healthy and operational'
UUID_A = "bf7bf3e1-3ee8-4816-9b66-b77b34387986"
UUID_B = "1e75a740-d08f-443d-a0f9-793db196e24f"
def _mk_rec(**kw):
rec = {
"dm_id": "testdm01",
"sender": "opm",
"recipient": "646",
"target": "646 tasks",
"thread_uuid": UUID_B,
"sent_at": "2026-10-04T17:00:00+00:00",
"deadline": "2020-01-01T00:00:00+00:00", # long past -> sweep acts now
"timeout_s": 900,
"nudges_allowed": 2,
"nudges_sent": 0,
"escalate_to": "opm",
"status": "pending",
}
rec.update(kw)
return rec
# --------------------------------------------------------------------------
# 1. sweeper backfill -- contract
# --------------------------------------------------------------------------
def test_contract_backfill_sets_uuid_on_matching_event():
rec = _mk_rec(thread_uuid=None)
evs = [{"type": "sidechat_autoprovisioned", "id": "testdm01",
"target": "646 tasks", "thread_uuid": UUID_A}]
assert contract_backfill_thread_uuid(rec, evs) is True
assert rec["thread_uuid"] == UUID_A, rec
def test_contract_backfill_no_event_stays_null():
rec = _mk_rec(thread_uuid=None)
assert contract_backfill_thread_uuid(rec, []) is False
assert rec["thread_uuid"] is None, rec
evs = [{"type": "nav_ok", "id": "testdm01", "target": "646 tasks"}]
assert contract_backfill_thread_uuid(rec, evs) is False
assert rec["thread_uuid"] is None, rec
def test_contract_backfill_never_overwrites_existing_uuid():
rec = _mk_rec(thread_uuid=UUID_B)
evs = [{"type": "sidechat_autoprovisioned", "id": "testdm01",
"target": "646 tasks", "thread_uuid": UUID_A}]
assert contract_backfill_thread_uuid(rec, evs) is False
assert rec["thread_uuid"] == UUID_B, "existing UUID must never be overwritten"
# ... and a null-UUID event must never blank it either
evs2 = [{"type": "sidechat_autoprovisioned", "id": "testdm01",
"target": "646 tasks", "thread_uuid": None}]
assert contract_backfill_thread_uuid(rec, evs2) is False
assert rec["thread_uuid"] == UUID_B, rec
def test_contract_backfill_ignores_unrelated_events():
rec = _mk_rec(thread_uuid=None)
evs = [{"type": "sidechat_autoprovisioned", "id": "otherdm99",
"target": "other-target", "thread_uuid": UUID_A}]
assert contract_backfill_thread_uuid(rec, evs) is False
assert rec["thread_uuid"] is None, rec
# --------------------------------------------------------------------------
# 2. sweeper final-nudge marking -- contract
# --------------------------------------------------------------------------
def test_contract_final_nudge_marks_and_routes_main():
rec = _mk_rec(nudges_sent=1, nudges_allowed=2) # about to send 2/2
target, is_final = contract_route_nudge(rec)
assert is_final is True
assert target == "main", target
assert rec.get("final_nudge_target") == "main", rec
def test_contract_nonfinal_nudge_keeps_target_unmarked():
rec = _mk_rec(nudges_sent=0, nudges_allowed=2) # about to send 1/2
target, is_final = contract_route_nudge(rec)
assert is_final is False
assert target == "646 tasks", target
assert "final_nudge_target" not in rec, rec
# --------------------------------------------------------------------------
# 3. harvester matching -- contract
# --------------------------------------------------------------------------
def test_contract_harvester_main_reply_resolves_via_final_nudge_target():
rec = _mk_rec(target="646 tasks", thread_uuid=UUID_B,
final_nudge_target="main")
assert contract_followup_matches(rec, "646", "main") is True
def test_contract_harvester_unrelated_thread_does_not_resolve():
rec = _mk_rec(target="646 tasks", thread_uuid=UUID_B,
final_nudge_target="main")
assert contract_followup_matches(rec, "646", "deadbeef-thread") is False
assert contract_followup_matches(rec, "pip", "main") is False
def test_contract_harvester_existing_rules_still_work():
# exact thread_uuid match
assert contract_followup_matches(_mk_rec(), "646", UUID_B) is True
# target == "main" + main thread
assert contract_followup_matches(_mk_rec(target="main"), "646", "main") is True
# non-pending never resolves
r = _mk_rec(status="resolved", final_nudge_target="main")
assert contract_followup_matches(r, "646", "main") is False
# --------------------------------------------------------------------------
# 4. multi-RESULT extraction -- contract
# --------------------------------------------------------------------------
def test_contract_multi_result_extracts_both():
got = contract_extract_results(TWO_MARKER_TEXT)
ids = [j for j, _ in got]
assert ids == ["pipe-demo-step1-20261004-165058-b08f5fb5",
"pipe-demo-step1-20261004-163439-8a408441"], ids
texts = dict(got)
assert texts["pipe-demo-step1-20261004-163439-8a408441"] == \
"OK: Fleet healthy and operational", texts
first = texts["pipe-demo-step1-20261004-165058-b08f5fb5"]
assert first.startswith("OK: Fleet healthy and operational"), first
assert "[RESULT" not in first, "first result text must stop at the 2nd marker"
def test_contract_single_result():
got = contract_extract_results("[RESULT abc-123] OK: done")
assert got == [("abc-123", "OK: done")], got
def test_contract_no_result():
assert contract_extract_results("just a normal message") == []
assert contract_extract_results("") == []
# --------------------------------------------------------------------------
# 5. dm.py URL assertion -- contract
# --------------------------------------------------------------------------
def test_contract_dm_url_contains_uuid_passes():
ok, reason = contract_assert_post_nav_url(
"https://muse.ai/thread/" + UUID_A, UUID_A, "pipe-1b4579")
assert ok is True and reason == "uuid_present", (ok, reason)
def test_contract_dm_url_landing_page_fails():
ok, reason = contract_assert_post_nav_url(
"https://muse.ai/", UUID_A, "pipe-1b4579")
assert ok is False and reason == "url_mismatch", (ok, reason)
def test_contract_dm_url_wrong_thread_fails():
ok, reason = contract_assert_post_nav_url(
"https://muse.ai/thread/" + UUID_B, UUID_A, "pipe-1b4579")
assert ok is False and reason == "url_mismatch", (ok, reason)
def test_contract_dm_url_main_skips_assertion():
ok, reason = contract_assert_post_nav_url("https://muse.ai/", None, "main")
assert ok is True and reason == "main_skipped", (ok, reason)
def test_contract_dm_url_empty_actual_fails():
ok, reason = contract_assert_post_nav_url("", UUID_A, "pipe-1b4579")
assert ok is False, (ok, reason)
# --------------------------------------------------------------------------
# wiring: real modules
# --------------------------------------------------------------------------
def _tmp_sweeper_env(sweeper):
"""Point the sweeper's file IO at a tmp dir; return (tmpdir, saved).
NOTE: the module uses pathlib.Path objects (FOLLOWUPS_FILE.exists()),
so the monkeypatched values must be Paths, not strs."""
tmp = tempfile.mkdtemp(prefix="w4sweep")
saved = {}
for attr, fname in (("FOLLOWUPS_FILE", "followups.json"),
("JOB_LOG", "job-log.jsonl")):
saved[attr] = getattr(sweeper, attr)
setattr(sweeper, attr, Path(tmp) / fname)
return tmp, saved
def _restore(mod, saved):
for attr, val in saved.items():
setattr(mod, attr, val)
def test_wired_sweeper_final_nudge_routes_main_and_marks():
"""Real sweep_cycle on a followup due its final nudge: delivery must go
to main (existing behavior) AND rec must gain final_nudge_target='main'
(ws-1 contract)."""
sweeper = _load("sweeper_under_test", "followup-sweeper.py")
tmp, saved = _tmp_sweeper_env(sweeper)
saved["send_dm"] = sweeper.send_dm
seen = {}
def fake_send_dm(sender, recipient, target, text):
seen.update(sender=sender, recipient=recipient, target=target)
return True, "SENT id=fake01"
sweeper.send_dm = fake_send_dm
try:
rec = _mk_rec(nudges_sent=1, nudges_allowed=2) # final nudge due
with open(os.path.join(tmp, "followups.json"), "w") as f:
json.dump({"w2nudge": rec}, f)
stats = sweeper.sweep_cycle(dry_run=False)
assert stats["nudges_sent"] == 1, stats
assert seen.get("target") == "main", \
f"final nudge must route to main, went to {seen.get('target')!r}"
with open(os.path.join(tmp, "followups.json")) as f:
rec2 = json.load(f)["w2nudge"]
assert rec2.get("final_nudge_target") == "main", (
"DEVIATION: ws-1 has not landed final_nudge_target marking -- "
f"sweep_cycle routed the final nudge to main but did not record "
f"final_nudge_target on the followup rec (rec keys: "
f"{sorted(rec2.keys())})")
finally:
_restore(sweeper, saved)
def test_wired_sweeper_backfills_thread_uuid():
"""Real _resolve_nudge_thread_uuid (ws-1): parses the nudge DM id from
send output, adopts the thread UUID from that id's
sidechat_autoprovisioned dm-log event (falling back to the sent event's
tags.thread); returns None when nothing usable exists (never fabricates
a UUID). Fixture dm-log via the module's DM_LOG hook -- no live state."""
sweeper = _load("sweeper_under_test", "followup-sweeper.py")
resolver = getattr(sweeper, "_resolve_nudge_thread_uuid", None)
assert resolver is not None, (
"DEVIATION: ws-1 backfill not landed -- followup-sweeper.py has no "
"_resolve_nudge_thread_uuid helper")
tmp = tempfile.mkdtemp(prefix="w4blog")
dmpath = Path(tmp) / "dm-log.jsonl"
saved = sweeper.DM_LOG
sweeper.DM_LOG = dmpath
nid = "f00dbabe"
out = (f"DM {nid} from opm to 646/pipe-x: SENT and VERIFIED "
f"thread={UUID_A}")
try:
# 1. autoprovisioned event -> UUID adopted
with open(dmpath, "w") as f:
f.write(json.dumps({"type": "sidechat_autoprovisioned", "id": nid,
"target": "pipe-x", "thread_uuid": UUID_A}) + "\n")
f.write(json.dumps({"type": "sent", "id": nid,
"tags": {"thread": UUID_A}}) + "\n")
assert resolver(out) == UUID_A, "autoprovisioned UUID not adopted"
# 2. no autoprovisioned event -> falls back to sent.tags.thread
with open(dmpath, "w") as f:
f.write(json.dumps({"type": "sent", "id": nid,
"tags": {"thread": UUID_A}}) + "\n")
assert resolver(out) == UUID_A, "sent.tags.thread fallback broken"
# 3. no usable events -> None (never fabricates / never null-writes)
with open(dmpath, "w") as f:
f.write(json.dumps({"type": "sent", "id": "other12",
"tags": {}}) + "\n")
assert resolver(out) is None, "resolver fabricated a UUID"
# 4. malformed UUID in event -> skipped
with open(dmpath, "w") as f:
f.write(json.dumps({"type": "sidechat_autoprovisioned", "id": nid,
"thread_uuid": "not-a-uuid"}) + "\n")
assert resolver(out) is None, "malformed UUID accepted"
# 5. unparseable nudge output -> None
assert resolver("some garbage without an id") is None
finally:
sweeper.DM_LOG = saved
def test_wired_sweeper_never_overwrites_uuid_with_null():
"""Real sweep_cycle: after a successful nudge whose thread cannot be
determined, an existing rec thread_uuid is left untouched (ws-1
contract: only a real UUID is ever written)."""
sweeper = _load("sweeper_under_test", "followup-sweeper.py")
tmp, saved = _tmp_sweeper_env(sweeper)
saved["send_dm"] = sweeper.send_dm
saved_dm_log = sweeper.DM_LOG
sweeper.DM_LOG = Path(tmp) / "dm-log.jsonl" # empty: no events
Path(tmp, "dm-log.jsonl").write_text("")
nid = "b00bf00d"
sweeper.send_dm = lambda *a: (
True, f"DM {nid} from opm to 646/pipe-x: SENT and VERIFIED")
try:
rec = _mk_rec(thread_uuid=UUID_B, nudges_sent=0, nudges_allowed=2)
with open(os.path.join(tmp, "followups.json"), "w") as f:
json.dump({"w2null": rec}, f)
sweeper.sweep_cycle(dry_run=False)
with open(os.path.join(tmp, "followups.json")) as f:
rec2 = json.load(f)["w2null"]
assert rec2["thread_uuid"] == UUID_B, (
f"DEVIATION: existing thread_uuid was overwritten "
f"(now {rec2['thread_uuid']!r}) despite no resolvable nudge thread")
assert rec2["nudges_sent"] == 1
finally:
sweeper.DM_LOG = saved_dm_log
_restore(sweeper, saved)
def test_wired_sweeper_backfill_reprovision_updates_uuid():
"""Real sweep_cycle: when a nudge lands in a newly provisioned thread,
the rec's stale UUID is replaced by the real new one (ws-1 reprovision
case)."""
sweeper = _load("sweeper_under_test", "followup-sweeper.py")
tmp, saved = _tmp_sweeper_env(sweeper)
saved["send_dm"] = sweeper.send_dm
saved_dm_log = sweeper.DM_LOG
sweeper.DM_LOG = Path(tmp) / "dm-log.jsonl"
nid = "c0ffee42"
with open(os.path.join(tmp, "dm-log.jsonl"), "w") as f:
f.write(json.dumps({"type": "sidechat_autoprovisioned", "id": nid,
"target": "pipe-x", "thread_uuid": UUID_A}) + "\n")
sweeper.send_dm = lambda *a: (
True, f"DM {nid} from opm to 646/pipe-x: SENT and VERIFIED "
f"thread={UUID_A}")
try:
rec = _mk_rec(thread_uuid=UUID_B, nudges_sent=0, nudges_allowed=2)
with open(os.path.join(tmp, "followups.json"), "w") as f:
json.dump({"w2re": rec}, f)
sweeper.sweep_cycle(dry_run=False)
with open(os.path.join(tmp, "followups.json")) as f:
rec2 = json.load(f)["w2re"]
assert rec2["thread_uuid"] == UUID_A, (
f"DEVIATION: reprovisioned thread UUID not adopted "
f"(still {rec2['thread_uuid']!r})")
finally:
sweeper.DM_LOG = saved_dm_log
_restore(sweeper, saved)
def test_wired_harvester_final_nudge_target_main_resolves():
"""Real clear_matching_followups: followup target='646 tasks',
final_nudge_target='main' must resolve on an assistant message in
thread 'main' (ws-2 contract)."""
harv = _load("harvester_under_test", "response-harvester.py")
rec = _mk_rec(target="646 tasks", thread_uuid=UUID_B,
final_nudge_target="main")
fups = {"hx1": rec}
harv.clear_matching_followups(fups, "646", "main", "mid-1",
"some assistant reply", dry_run=True)
assert rec.get("status") == "resolved", (
"DEVIATION: ws-2 final_nudge_target path not landed -- "
"clear_matching_followups() does not resolve a followup with "
"final_nudge_target='main' on a main-thread assistant reply "
f"(status={rec.get('status')!r}; fn signature: "
f"{inspect.signature(harv.clear_matching_followups)})")
# ... and must NOT resolve on an unrelated thread
rec2 = _mk_rec(target="646 tasks", thread_uuid=UUID_B,
final_nudge_target="main")
fups2 = {"hx2": rec2}
harv.clear_matching_followups(fups2, "646", "unrelated-thread", "mid-2",
"some assistant reply", dry_run=True)
assert rec2.get("status") == "pending", \
f"unrelated thread must not resolve (status={rec2.get('status')!r})"
def test_wired_harvester_existing_rules_still_hold():
"""Real clear_matching_followups: pre-existing rules keep working."""
harv = _load("harvester_under_test", "response-harvester.py")
r1 = _mk_rec() # thread_uuid == UUID_B
fups = {"e1": r1}
harv.clear_matching_followups(fups, "646", UUID_B, "m", "t", dry_run=True)
assert r1["status"] == "resolved", "exact thread_uuid match broke"
r2 = _mk_rec(target="main", thread_uuid=None)
fups = {"e2": r2}
harv.clear_matching_followups(fups, "646", "main", "m", "t", dry_run=True)
assert r2["status"] == "resolved", "target=='main' rule broke"
r3 = _mk_rec()
fups = {"e3": r3}
harv.clear_matching_followups(fups, "646", "nope", "m", "t", dry_run=True)
assert r3["status"] == "pending", "unrelated thread wrongly resolved"
def test_wired_harvester_extracts_all_result_markers():
"""The real extraction path must yield BOTH [RESULT] markers (ws-2:
finditer instead of first-only re.search)."""
harv = _load("harvester_under_test", "response-harvester.py")
extractor = next(
(getattr(harv, n) for n in dir(harv)
if "result" in n.lower()
and any(k in n.lower() for k in ("extract", "iter", "marker"))
and callable(getattr(harv, n))),
None)
src = inspect.getsource(harv)
if extractor is not None:
got = extractor(TWO_MARKER_TEXT)
ids = [j for j, _ in got]
assert ids == ["pipe-demo-step1-20261004-165058-b08f5fb5",
"pipe-demo-step1-20261004-163439-8a408441"], \
f"extractor {extractor.__name__} missed markers: {ids}"
return
if "finditer" in src and "RESULT" in src:
return # inline finditer implementation detected; contract assumed met
# Deviation evidence: the current inline path uses re.search (first only).
cur = re.search(r"\[RESULT\s+([A-Za-z0-9_-]+)\]\s*(.*)",
TWO_MARKER_TEXT, re.S)
raise AssertionError(
"DEVIATION: ws-2 finditer change not landed -- response-harvester.py "
"exposes no result extractor and its inline path is re.search "
f"(first marker only): it captures {cur.group(1)!r} and swallows the "
"second marker [RESULT pipe-demo-step1-20261004-163439-8a408441] "
"inside group(2)")
def test_wired_dm_send_asserts_post_nav_url_before_send():
"""dm_send must independently assert the post-nav browser URL contains
the target thread UUID BEFORE sending (ws-3 contract); 'main' skips."""
path = BIN_DIR / "dm.py"
dm = None
via = "file-text fallback"
try:
dm = _load("dm_under_test", "dm.py")
src = inspect.getsource(dm.dm_send)
via = "inspect(dm.dm_send)"
except Exception:
src = path.read_text()
m = re.search(r"def dm_send\(.*?(?=\ndef |\Z)", src, re.S)
src = m.group(0) if m else src
# Exposed predicate? test it directly.
pred = None
if dm is not None:
pred = next(
(getattr(dm, n) for n in dir(dm)
if callable(getattr(dm, n))
and "nav" in n.lower() and "url" in n.lower()),
None)
if pred is not None:
assert pred("https://muse.ai/thread/" + UUID_A, UUID_A, "pipe-x")[0] is True
assert pred("https://muse.ai/", UUID_A, "pipe-x")[0] is False
return
# ws-3 landed shape: module-level assert_pre_send_placement() called from
# dm_send before the send. Verify the call site and exercise the helper.
gate = getattr(dm, "assert_pre_send_placement", None) if dm is not None else None
if callable(gate):
assert "assert_pre_send_placement" in src, \
"gate exists but dm_send never calls it"
_orig_run_full = dm.run_full
_calls = []
def _stub(cmd, timeout=60):
_calls.append(cmd)
if "sidechat use" in cmd:
return 0, "navigated https://muse.ai/thread/" + UUID_A, ""
return 0, _stub.url, ""
try:
dm.run_full = _stub
# 1. UUID-known thread, correct placement -> pass
_stub.url = "https://muse.ai/thread/" + UUID_A
ok, detail = gate("opm", "pipe-x", UUID_A, direct_nav_done=True)
assert ok is True, f"expected pass on matching URL: {detail}"
# 2. landing page -> loud fail (the 2026-10-04 incident mode)
_stub.url = "https://muse.ai/"
ok, detail = gate("opm", "pipe-x", UUID_A, direct_nav_done=True)
assert ok is False and detail.get("reason") == "url_mismatch", detail
# 3. wrong thread -> loud fail
_stub.url = "https://muse.ai/thread/" + UUID_B
ok, detail = gate("opm", "pipe-x", UUID_A, direct_nav_done=True)
assert ok is False and detail.get("reason") == "url_mismatch", detail
# 4. main target skips assertion with zero subprocess calls
_calls.clear()
ok, _ = gate("opm", "main", None, direct_nav_done=False)
assert ok is True and not _calls, "main must skip without subprocess"
# 5. re-nav path (direct_nav_done=False) -> pass
_stub.url = "https://muse.ai/thread/" + UUID_A
ok, detail = gate("opm", "pipe-x", UUID_A, direct_nav_done=False)
assert ok is True, f"re-nav path should pass: {detail}"
finally:
dm.run_full = _orig_run_full
return
nav_i = src.find("sidechat use")
send_i = src.find("Send with verification retries")
segment = src[nav_i:send_i] if 0 <= nav_i < send_i else ""
has_url_fetch = re.search(r"""['"]\s*url['"]|account\s+\S+\s+url\b""", segment)
has_abort = ("sys.exit" in segment) or ("raise " in segment)
assert has_url_fetch and has_abort, (
f"DEVIATION: ws-3 pre-send URL assertion not landed ({via}) -- "
"dm_send's nav->send path performs no independent post-nav URL "
"fetch+containment check that aborts the send; it trusts the "
"`sidechat use` command output (no_thread_url/uuid_mismatch on the "
"nav output only). Note: verify_placement() asserts the URL "
"post-send, which is a different (later) check.")
# --------------------------------------------------------------------------
# runner (also pytest-compatible: plain test_* functions, no args)
# --------------------------------------------------------------------------
def main():
fns = [(n, f) for n, f in sorted(globals().items())
if n.startswith("test_") and callable(f)]
results = []
for name, fn in fns:
try:
fn()
results.append((name, "PASS", ""))
except _Skip as e:
results.append((name, "SKIP", str(e)))
except AssertionError as e:
results.append((name, "FAIL", str(e) or "assertion failed"))
except Exception as e: # noqa: BLE001 - harness must not crash
results.append((name, "ERROR",
f"{type(e).__name__}: {e}"))
npass = sum(1 for _, s, _ in results if s == "PASS")
nfail = sum(1 for _, s, _ in results if s in ("FAIL", "ERROR"))
nskip = sum(1 for _, s, _ in results if s == "SKIP")
print(f"\n{'test':58} result")
print("-" * 80)
for name, status, detail in results:
print(f"{name:58} {status}")
if detail and status in ("FAIL", "ERROR"):
for line in detail.splitlines():
print(f" {line}")
print("-" * 80)
print(f"{len(results)} tests: {npass} pass, {nfail} fail/error, {nskip} skip")
return 1 if nfail else 0
if __name__ == "__main__":
sys.exit(main())