Files
box/bin/tests/test_followup_fixes.py
T

916 lines
40 KiB
Python

#!/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
_orig_sleep = dm.time.sleep
_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
# Settle sleeps (1s/2s per gate call) are production pacing, not
# asserted behavior: skip them like the browser subprocess above.
dm.time.sleep = lambda s: None
# 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
dm.time.sleep = _orig_sleep
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.")
# --------------------------------------------------------------------------
# wiring: dm.py autoprovision creation check (2026-10-04 fix)
# --------------------------------------------------------------------------
_ADOPT_PARKED = "bf7bf3e1-3ee8-4816-9b66-b77b34387986" # the 19:50 incident UUID
_ADOPT_NEW = "aaaaaaaa-1111-2222-3333-444444444444"
def _dm_mod():
return _load("dm_under_test", "dm.py")
def test_wired_adopt_genuine_new_uuid():
dm = _dm_mod()
ok, why = dm.autoprovision_adopt_ok(
_ADOPT_NEW, _ADOPT_PARKED,
{"pipe-1b4579": {"thread_uuid": _ADOPT_PARKED, "agent": "opm"}},
"heartbeat")
assert ok is True and why == "new_thread", (ok, why)
def test_wired_adopt_refuses_parked_thread():
# the 2026-10-04 19:50 UTC incident: captured == pre-create parked UUID
dm = _dm_mod()
ok, why = dm.autoprovision_adopt_ok(
_ADOPT_PARKED, _ADOPT_PARKED,
{"pipe-1b4579": {"thread_uuid": _ADOPT_PARKED, "agent": "opm"}},
"heartbeat")
assert ok is False and "matches_pre_create_parked_thread" in why, (ok, why)
def test_wired_adopt_refuses_already_mapped_elsewhere():
dm = _dm_mod()
ok, why = dm.autoprovision_adopt_ok(
_ADOPT_PARKED, None,
{"pipe-1b4579": {"thread_uuid": _ADOPT_PARKED, "agent": "opm"}},
"heartbeat")
assert ok is False and "already_mapped_under_pipe-1b4579" in why, (ok, why)
def test_wired_adopt_same_target_readopt_allowed():
dm = _dm_mod()
ok, why = dm.autoprovision_adopt_ok(
_ADOPT_PARKED, "99999999-1111-2222-3333-444444444444",
{"pipe-1b4579": {"thread_uuid": _ADOPT_PARKED, "agent": "opm"}},
"pipe-1b4579")
assert ok is True, (ok, why)
def test_wired_adopt_empty_capture_refused():
dm = _dm_mod()
for cap in (None, ""):
ok, why = dm.autoprovision_adopt_ok(cap, _ADOPT_PARKED, {}, "heartbeat")
assert ok is False and "no_uuid_captured" in why, (cap, ok, why)
def test_wired_adopt_fresh_uuid_empty_map():
dm = _dm_mod()
ok, why = dm.autoprovision_adopt_ok(_ADOPT_NEW, None, {}, "heartbeat")
assert ok is True and why == "new_thread", (ok, why)
# --------------------------------------------------------------------------
# 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)
return 1 if nfail else 0
import unittest
# --------------------------------------------------------------------------
# harvester resurrection -- dm_id markers, dry-run purity, scheduling
# --------------------------------------------------------------------------
def test_wired_harvester_dmid_marker_resolves():
"""Real clear_matching_followups: [RESULT <dm_id>] resolves a
DM-ordered followup whose job_id differs (live f4293153 pattern:
marker quoted the nudge's DM id, record job was ml-muse-*).
Unrelated thread isolates the dm_id path from thread matching."""
harv = _load("harvester_under_test", "response-harvester.py")
rec = _mk_rec(thread_uuid=UUID_B, job_id="ml-muse-20261007-013210")
fups = {"f4293153": rec}
harv.clear_matching_followups(fups, "646", "unrelated-thread", "mid-9",
"[RESULT f4293153] done", dry_run=True,
job_id="f4293153", verb="RESULT")
assert rec.get("status") == "resolved", (
"DEVIATION: [RESULT <dm_id>] does not resolve its followup -- "
"clear_matching_followups() matches marker ids against job_id "
f"only, never the followup key (status={rec.get('status')!r})")
# ... and a wrong id must not resolve.
rec2 = _mk_rec(thread_uuid=UUID_B, job_id="ml-muse-20261007-013210")
fups2 = {"f4293153": rec2}
harv.clear_matching_followups(fups2, "646", "unrelated-thread", "mid-9",
"[RESULT deadbeef] done", dry_run=True,
job_id="deadbeef", verb="RESULT")
assert rec2.get("status") == "pending", (
f"wrong marker id wrongly resolved (status={rec2.get('status')!r})")
def test_wired_harvester_dry_run_has_no_side_effects():
"""process_messages(dry_run=True) with an evidence-less RESULT must
still extract the marker but must not fire proof followups,
archive threads, or persist anything."""
harv = _load("harvester_under_test", "response-harvester.py")
calls = []
saved = {n: getattr(harv, n) for n in
("execute_agent_tool", "archive_ephemeral_thread",
"append_jsonl", "save_json_file")}
harv.execute_agent_tool = lambda *a, **k: calls.append("exec") or (True, {})
harv.archive_ephemeral_thread = (
lambda *a, **k: calls.append("archive"))
harv.append_jsonl = lambda *a, **k: calls.append("append")
harv.save_json_file = lambda *a, **k: calls.append("save")
try:
msgs = [{"id": "m1", "author": "assistant",
"text": "[RESULT j1] done",
"ts": "2026-10-07T00:00:00+00:00"}]
new, wm, nres = harv.process_messages(
msgs, "646", UUID_A, "t", None, {}, dry_run=True)
finally:
for n, fn in saved.items():
setattr(harv, n, fn)
assert nres == 1, "dry-run must still extract markers"
assert calls == [], f"dry-run leaked side effects: {calls}"
def test_wired_result_markers_bare_and_status_forms():
"""iter_result_markers handles the engine's 3-group shape: a bare
[RESULT <id>] <text> (status None) must not crash, and a status
token must survive into the result text for fail detection."""
harv = _load("harvester_under_test", "response-harvester.py")
assert list(harv.iter_result_markers("[RESULT f4293153] done")) == [
("f4293153", "done")], "bare RESULT marker must extract cleanly"
jid, text = list(harv.iter_result_markers("[RESULT j9] FAIL blew up"))[0]
assert jid == "j9" and "FAIL" in text and "blew up" in text, (
f"status token must survive into result text (got {jid!r} {text!r})")
def test_wired_poison_message_does_not_wedge_batch():
"""A marker-extraction crash degrades to plain-reply handling so
sibling messages still process and the watermark keeps advancing."""
harv = _load("harvester_under_test", "response-harvester.py")
real_iter = harv.iter_result_markers
real_nudge = harv.maybe_nudge_untagged_sidechat
def boom(text):
if "POISON" in (text or ""):
raise RuntimeError("boom")
return real_iter(text)
harv.iter_result_markers = boom
harv.maybe_nudge_untagged_sidechat = lambda *a, **k: None
try:
msgs = [
{"id": "m1", "author": "assistant",
"text": "POISON [RESULT x] y",
"ts": "2026-10-07T00:00:00+00:00"},
{"id": "m2", "author": "assistant",
"text": "[RESULT j2] ok",
"ts": "2026-10-07T00:01:00+00:00"},
]
new, wm, nres = harv.process_messages(
msgs, "646", UUID_A, "t", None, {}, dry_run=True)
finally:
harv.iter_result_markers = real_iter
harv.maybe_nudge_untagged_sidechat = real_nudge
assert nres == 1, "sibling marker must still extract"
assert [m["id"] for m in new] == ["m1", "m2"], \
"both messages must process past the poison one"
def test_harvester_timer_unit_wired():
"""The harvester must be scheduler-owned: unit files exist, the
service runs --once, and the timer fires on a short cadence.
Ingestion died silently for ~22h with no unit at all."""
root = BIN_DIR.parent
svc = (root / "systemd" / "response-harvester.service").read_text()
tmr = (root / "systemd" / "response-harvester.timer").read_text()
assert "response-harvester.py" in svc and "--once" in svc, \
"service must run the harvester --once"
assert "OnUnitActiveSec=" in tmr, "timer needs a repeat cadence"
assert "WantedBy=timers.target" in tmr, "timer must target timers.target"
def test_collection_adapter_is_single_and_pytest_opted_out():
"""Collection-shape guard (no 3x duplicates): exactly one TestCase
adapter is reachable from module globals (the adapter loop must not
leak a `_fn` alias that pytest collects as a second class), and the
adapter opts out of pytest (`__test__ = False`) so the module-level
functions are pytest's single source while unittest discovery still
runs the adapter."""
cases = [v for v in list(globals().values())
if inspect.isclass(v) and issubclass(v, unittest.TestCase)]
assert len(cases) == 1, (
f"expected exactly 1 TestCase adapter, found {len(cases)} "
f"(stray aliases reintroduce duplicate collection)")
assert TestFollowupFixes.__test__ is False, (
"TestFollowupFixes must set __test__ = False so pytest collects "
"each test once via the module-level functions")
class TestFollowupFixes(unittest.TestCase):
"""unittest discovery adapter for contract and wired test functions."""
# pytest collects the module-level functions; skip the adapter so each
# test runs once. (unittest discovery ignores __test__ and still runs
# the adapter, which is its only view of this file's tests.)
__test__ = False
for _name, _fn in list(globals().items()):
if _name.startswith("test_") and callable(_fn):
def _bind(f):
def _runner(self):
f()
return _runner
setattr(TestFollowupFixes, _name, _bind(_fn))
# Drop the loop temporaries: after the final iteration `_fn` aliases
# TestFollowupFixes, and pytest collects TestCase subclasses regardless of
# name -- that stray alias was the third copy (module fn + adapter + `_fn`).
del _name, _fn
if __name__ == "__main__":
sys.exit(main())