diff --git a/bin/dm.py b/bin/dm.py old mode 100644 new mode 100755 index 26c7e48..dba2df6 --- a/bin/dm.py +++ b/bin/dm.py @@ -296,6 +296,51 @@ def run_full(cmd, timeout=60): result = subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=timeout) return result.returncode, result.stdout.strip(), result.stderr.strip() +SIDCHAT_MAP_FILE = "/home/super/Projects/NetVM/job-sidechats.json" + + +def load_sidechat_map(): + """Read the name->thread-UUID mapping. Returns {} on any error.""" + try: + if os.path.exists(SIDCHAT_MAP_FILE): + with open(SIDCHAT_MAP_FILE, "r", encoding="utf-8") as f: + data = json.load(f) + return data if isinstance(data, dict) else {} + except Exception: + pass + return {} + + +def autoprovision_adopt_ok(captured_uuid, pre_create_uuid, existing_map, target): + """Creation check for autoprovisioned thread UUIDs (2026-10-04 fix). + + Autoprovision must produce a NEW thread. Before this check, dm.py adopted + whatever UUID the post-send URL showed -- including the browser's parked + thread when the create produced nothing. On 2026-10-04 that corrupted the + "heartbeat" mapping with the parked pipe-demo thread's UUID, routing all + heartbeat DMs into the wrong thread. + + Returns (ok, reason). Refuses when: + - nothing was captured; + - the captured UUID matches the pre-create parked thread (no new + thread was created); or + - the UUID is already mapped under a DIFFERENT target name (adopting + it would corrupt that target's routing). + """ + cu = (captured_uuid or "").lower() + if not cu: + return False, "no_uuid_captured" + pc = (pre_create_uuid or "").lower() + if pc and cu == pc: + return False, "matches_pre_create_parked_thread" + for name, rec in (existing_map or {}).items(): + if name == target: + continue + if isinstance(rec, dict) and str(rec.get("thread_uuid", "")).lower() == cu: + return False, "already_mapped_under_" + str(name) + return True, "new_thread" + + def verify_placement(recipient, msg_id, target, thread_uuid): """Authoritative post-send placement verification (2026-10-04). @@ -493,6 +538,7 @@ def dm_send(agent, target, message, verify=True, raw=False, thread_uuid = None thread_url = None is_new_sidechat = False + pre_create_uuid = None # parked thread before autoprovision (creation check) nav_is_uuid = False # Resolve well-known aliases or dynamic thread mappings to UUIDs. nav_target = resolve_sidechat_target(target) @@ -521,6 +567,14 @@ def dm_send(agent, target, message, verify=True, raw=False, if "NOTFOUND" in _out or "Navigated to: None" in _out: log_event({"type": "sidechat_autoprovision_start", "id": msg_id, "agent": agent, "to": recipient, "target": target}) + # Creation check, part 1 (2026-10-04): record the parked thread + # UUID BEFORE create runs. If the post-send URL still shows this + # UUID, no new thread was created and the UUID must NOT be + # adopted into the mapping (heartbeat incident: the parked + # pipe-demo thread's UUID was adopted for "heartbeat"). + _pre_rc, _pre_out, _ = run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} url") + _pre_m = re.search(r"/thread/(" + UUID_RE + r")", _pre_out or "", re.I) + pre_create_uuid = _pre_m.group(1).lower() if _pre_m else None _c_rc, _c_out, _c_err = run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat create") if _c_rc != 0 or "Created:" not in _c_out: log_event({"type": "nav_failed", "id": msg_id, "agent": agent, "to": recipient, @@ -531,11 +585,23 @@ def dm_send(agent, target, message, verify=True, raw=False, 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. + # pin it now so the pre-send gate can assert it -- but only if + # it is genuinely new (creation check). A "create" that echoes + # the parked thread must fail loudly, not poison the gate. _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 + _cand = _cm.group(1).lower() + _ok, _why = autoprovision_adopt_ok(_cand, pre_create_uuid, load_sidechat_map(), target) + if _ok: + thread_uuid = _cand + thread_url = "https://muse.ai/thread/" + thread_uuid + else: + log_event({"type": "autoprovision_false_capture", "id": msg_id, "agent": agent, + "to": recipient, "target": target, "reason": _why, + "captured_uuid": _cand, "pre_create_uuid": pre_create_uuid}) + print(f"DM {msg_id} from {agent} to {recipient}/{target}: FAILED " + f"(autoprovision did not create a new thread: {_why})", file=sys.stderr) + sys.exit(1) log_event({"type": "nav_ok", "id": msg_id, "agent": agent, "to": recipient, "target": target, "status": "sidechat_created_pending_uuid", "browser_url": thread_url}) @@ -663,38 +729,47 @@ def dm_send(agent, target, message, verify=True, raw=False, if delivered and is_new_sidechat: # We sent to /thread/new; muse.ai now assigns a permanent UUID. # Capture the current URL BEFORE navigating away to main. + _captured = None for _ in range(10): _u_rc, _u_out, _u_err = run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} url") - m_uuid = re.search(r"/thread/([0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12})", _u_out, re.I) + m_uuid = re.search(r"/thread/(" + UUID_RE + r")", _u_out, re.I) if m_uuid: - thread_uuid = m_uuid.group(1).lower() - thread_url = "https://muse.ai/thread/" + thread_uuid + _captured = m_uuid.group(1).lower() break time.sleep(1) - if thread_uuid: - sc_file = "/home/super/Projects/NetVM/job-sidechats.json" - try: - sc_data = {} - if os.path.exists(sc_file): - with open(sc_file, "r", encoding="utf-8") as f: - sc_data = json.load(f) - sc_data[target] = { - "thread_uuid": thread_uuid, - "agent": recipient, - "created_at": datetime.now(timezone.utc).isoformat() - } - tmp_sc = f"{sc_file}.tmp.{os.getpid()}" - with open(tmp_sc, "w", encoding="utf-8") as f: - json.dump(sc_data, f, indent=2) - os.replace(tmp_sc, sc_file) - log_event({"type": "sidechat_autoprovisioned", "id": msg_id, "target": target, - "thread_uuid": thread_uuid, "agent": recipient}) - except Exception as e: - log_event({"type": "sidechat_persist_error", "id": msg_id, "error": str(e)[:100]}) - tags["thread"] = thread_uuid + if _captured: + # Creation check, part 2 (2026-10-04): only adopt the UUID if it + # is a genuinely new thread. On refusal, leave thread_uuid alone + # (a UUID pinned at create time stays) and do NOT write the + # mapping; verify_placement then fails closed (no_thread_uuid) + # instead of blessing a misplaced send. + _ok, _why = autoprovision_adopt_ok(_captured, pre_create_uuid, load_sidechat_map(), target) + if not _ok: + log_event({"type": "autoprovision_false_capture", "id": msg_id, "agent": agent, + "to": recipient, "target": target, "reason": _why, + "captured_uuid": _captured, "pre_create_uuid": pre_create_uuid}) + else: + thread_uuid = _captured + thread_url = "https://muse.ai/thread/" + thread_uuid + try: + sc_data = load_sidechat_map() + sc_data[target] = { + "thread_uuid": thread_uuid, + "agent": recipient, + "created_at": datetime.now(timezone.utc).isoformat() + } + tmp_sc = f"{SIDCHAT_MAP_FILE}.tmp.{os.getpid()}" + with open(tmp_sc, "w", encoding="utf-8") as f: + json.dump(sc_data, f, indent=2) + os.replace(tmp_sc, SIDCHAT_MAP_FILE) + log_event({"type": "sidechat_autoprovisioned", "id": msg_id, "target": target, + "thread_uuid": thread_uuid, "agent": recipient}) + except Exception as e: + log_event({"type": "sidechat_persist_error", "id": msg_id, "error": str(e)[:100]}) + tags["thread"] = thread_uuid else: - log_event({"type": "sidechat_uuid_capture_failed", "id": msg_id, "out": _u_out[:200]}) + log_event({"type": "sidechat_uuid_capture_failed", "id": msg_id, "out": (_u_out or "")[:200]}) if delivered and thread_uuid and "thread" not in tags: tags["thread"] = thread_uuid @@ -869,11 +944,32 @@ def dm_verify_sig(message=None, agent=None, target=None): payload = text[:sm.start()].strip() sigblock = "-----BEGIN SSH SIGNATURE-----\n" + sm.group(1).strip() + "\n-----END SSH SIGNATURE-----\n" pubpath = os.path.join(SIGNERS_DIR, sender + ".pub") - if not os.path.exists(pubpath): + pubkey = None + if os.path.exists(pubpath): + with open(pubpath) as f: + pubkey = f.read().strip() + else: + # Fallback to crypt.muse-dev.online + try: + import urllib.request + req = urllib.request.Request( + "https://crypt.muse-dev.online/keys/allowed_signers", + headers={"User-Agent": "dm-verify/1.0"} + ) + with urllib.request.urlopen(req, timeout=5) as resp: + if resp.status == 200: + lines = resp.read().decode("utf-8", errors="ignore").splitlines() + for line in lines: + parts = line.strip().split(None, 1) + if len(parts) == 2 and parts[0] == sender: + pubkey = parts[1] + break + except Exception as e: + sys.stderr.write(f"warning: failed to query crypt.muse-dev.online: {e}\n") + + if not pubkey: print(f"BAD: [{mid}] no public key registered for sender '{sender}'", file=sys.stderr) continue - with open(pubpath) as f: - pubkey = f.read().strip() with tempfile.TemporaryDirectory() as td: allowed = os.path.join(td, "allowed") with open(allowed, "w") as f: diff --git a/bin/tests/test_followup_fixes.py b/bin/tests/test_followup_fixes.py index ea70b91..10a8250 100644 --- a/bin/tests/test_followup_fixes.py +++ b/bin/tests/test_followup_fixes.py @@ -656,6 +656,68 @@ def test_wired_dm_send_asserts_post_nav_url_before_send(): "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) # --------------------------------------------------------------------------