916 lines
40 KiB
Python
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())
|
|
|