dm.py: autoprovision creation check — refuse to adopt parked/already-mapped thread UUIDs (fixes 2026-10-04 heartbeat misroute)
This commit is contained in:
@@ -296,6 +296,51 @@ def run_full(cmd, timeout=60):
|
|||||||
result = subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=timeout)
|
result = subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=timeout)
|
||||||
return result.returncode, result.stdout.strip(), result.stderr.strip()
|
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):
|
def verify_placement(recipient, msg_id, target, thread_uuid):
|
||||||
"""Authoritative post-send placement verification (2026-10-04).
|
"""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_uuid = None
|
||||||
thread_url = None
|
thread_url = None
|
||||||
is_new_sidechat = False
|
is_new_sidechat = False
|
||||||
|
pre_create_uuid = None # parked thread before autoprovision (creation check)
|
||||||
nav_is_uuid = False
|
nav_is_uuid = False
|
||||||
# Resolve well-known aliases or dynamic thread mappings to UUIDs.
|
# Resolve well-known aliases or dynamic thread mappings to UUIDs.
|
||||||
nav_target = resolve_sidechat_target(target)
|
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:
|
if "NOTFOUND" in _out or "Navigated to: None" in _out:
|
||||||
log_event({"type": "sidechat_autoprovision_start", "id": msg_id, "agent": agent,
|
log_event({"type": "sidechat_autoprovision_start", "id": msg_id, "agent": agent,
|
||||||
"to": recipient, "target": target})
|
"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")
|
_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:
|
if _c_rc != 0 or "Created:" not in _c_out:
|
||||||
log_event({"type": "nav_failed", "id": msg_id, "agent": agent, "to": recipient,
|
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)
|
sys.exit(1)
|
||||||
is_new_sidechat = True
|
is_new_sidechat = True
|
||||||
# If create returned a real thread URL (not /thread/new),
|
# 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)
|
_cm = re.search(r"/thread/(" + UUID_RE + r")", _c_out, re.I)
|
||||||
if _cm:
|
if _cm:
|
||||||
thread_uuid = _cm.group(1).lower()
|
_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
|
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,
|
log_event({"type": "nav_ok", "id": msg_id, "agent": agent, "to": recipient,
|
||||||
"target": target, "status": "sidechat_created_pending_uuid",
|
"target": target, "status": "sidechat_created_pending_uuid",
|
||||||
"browser_url": thread_url})
|
"browser_url": thread_url})
|
||||||
@@ -663,38 +729,47 @@ def dm_send(agent, target, message, verify=True, raw=False,
|
|||||||
if delivered and is_new_sidechat:
|
if delivered and is_new_sidechat:
|
||||||
# We sent to /thread/new; muse.ai now assigns a permanent UUID.
|
# We sent to /thread/new; muse.ai now assigns a permanent UUID.
|
||||||
# Capture the current URL BEFORE navigating away to main.
|
# Capture the current URL BEFORE navigating away to main.
|
||||||
|
_captured = None
|
||||||
for _ in range(10):
|
for _ in range(10):
|
||||||
_u_rc, _u_out, _u_err = run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} url")
|
_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:
|
if m_uuid:
|
||||||
thread_uuid = m_uuid.group(1).lower()
|
_captured = m_uuid.group(1).lower()
|
||||||
thread_url = "https://muse.ai/thread/" + thread_uuid
|
|
||||||
break
|
break
|
||||||
time.sleep(1)
|
time.sleep(1)
|
||||||
|
|
||||||
if thread_uuid:
|
if _captured:
|
||||||
sc_file = "/home/super/Projects/NetVM/job-sidechats.json"
|
# 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:
|
try:
|
||||||
sc_data = {}
|
sc_data = load_sidechat_map()
|
||||||
if os.path.exists(sc_file):
|
|
||||||
with open(sc_file, "r", encoding="utf-8") as f:
|
|
||||||
sc_data = json.load(f)
|
|
||||||
sc_data[target] = {
|
sc_data[target] = {
|
||||||
"thread_uuid": thread_uuid,
|
"thread_uuid": thread_uuid,
|
||||||
"agent": recipient,
|
"agent": recipient,
|
||||||
"created_at": datetime.now(timezone.utc).isoformat()
|
"created_at": datetime.now(timezone.utc).isoformat()
|
||||||
}
|
}
|
||||||
tmp_sc = f"{sc_file}.tmp.{os.getpid()}"
|
tmp_sc = f"{SIDCHAT_MAP_FILE}.tmp.{os.getpid()}"
|
||||||
with open(tmp_sc, "w", encoding="utf-8") as f:
|
with open(tmp_sc, "w", encoding="utf-8") as f:
|
||||||
json.dump(sc_data, f, indent=2)
|
json.dump(sc_data, f, indent=2)
|
||||||
os.replace(tmp_sc, sc_file)
|
os.replace(tmp_sc, SIDCHAT_MAP_FILE)
|
||||||
log_event({"type": "sidechat_autoprovisioned", "id": msg_id, "target": target,
|
log_event({"type": "sidechat_autoprovisioned", "id": msg_id, "target": target,
|
||||||
"thread_uuid": thread_uuid, "agent": recipient})
|
"thread_uuid": thread_uuid, "agent": recipient})
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
log_event({"type": "sidechat_persist_error", "id": msg_id, "error": str(e)[:100]})
|
log_event({"type": "sidechat_persist_error", "id": msg_id, "error": str(e)[:100]})
|
||||||
tags["thread"] = thread_uuid
|
tags["thread"] = thread_uuid
|
||||||
else:
|
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:
|
if delivered and thread_uuid and "thread" not in tags:
|
||||||
tags["thread"] = thread_uuid
|
tags["thread"] = thread_uuid
|
||||||
@@ -869,11 +944,32 @@ def dm_verify_sig(message=None, agent=None, target=None):
|
|||||||
payload = text[:sm.start()].strip()
|
payload = text[:sm.start()].strip()
|
||||||
sigblock = "-----BEGIN SSH SIGNATURE-----\n" + sm.group(1).strip() + "\n-----END SSH SIGNATURE-----\n"
|
sigblock = "-----BEGIN SSH SIGNATURE-----\n" + sm.group(1).strip() + "\n-----END SSH SIGNATURE-----\n"
|
||||||
pubpath = os.path.join(SIGNERS_DIR, sender + ".pub")
|
pubpath = os.path.join(SIGNERS_DIR, sender + ".pub")
|
||||||
if not os.path.exists(pubpath):
|
pubkey = None
|
||||||
print(f"BAD: [{mid}] no public key registered for sender '{sender}'", file=sys.stderr)
|
if os.path.exists(pubpath):
|
||||||
continue
|
|
||||||
with open(pubpath) as f:
|
with open(pubpath) as f:
|
||||||
pubkey = f.read().strip()
|
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 tempfile.TemporaryDirectory() as td:
|
with tempfile.TemporaryDirectory() as td:
|
||||||
allowed = os.path.join(td, "allowed")
|
allowed = os.path.join(td, "allowed")
|
||||||
with open(allowed, "w") as f:
|
with open(allowed, "w") as f:
|
||||||
|
|||||||
@@ -656,6 +656,68 @@ def test_wired_dm_send_asserts_post_nav_url_before_send():
|
|||||||
"post-send, which is a different (later) check.")
|
"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)
|
# runner (also pytest-compatible: plain test_* functions, no args)
|
||||||
# --------------------------------------------------------------------------
|
# --------------------------------------------------------------------------
|
||||||
|
|||||||
Reference in New Issue
Block a user