#!/usr/bin/env python3 """ DM: Headless Direct Message API (muse.ai) — tagged, logged, verifiable. Wire format (every send, signed or not): [from:] [id:<8-hex>] Signed DMs (produced by bin/dm-sign.sh, namespace "dm") append the SSH signature block after a blank line; `verify-sig` checks it against the sender's key in dm-signers/.pub. An unsigned message carrying a [from:X] header is just a claim — only a GOOD verify-sig result is proof. Every send is appended to dm-log.jsonl. Delivery is confirmed by reading the RECIPIENT's chat for the message id (up to 3 attempts): SENT+VERIFIED means the id was seen in the recipient's chat; FAILED means it wasn't after 3 attempts. `send --raw` transmits verbatim (for pre-signed messages): no tagging, no truncation — the signed payload must survive byte-identical. Policy (2026-10-04): sidechat-first. Main-chat sends are refused unless --allow-main-chat is passed explicitly. Never saturate main threads by default; use a sidechat target instead. Usage: dm.py send --agent opm --to 646 --target "646 tasks" "message" dm.py send --agent opm --target main --allow-main-chat --raw "$(dm-sign.sh --from operator-main 'hi')" dm.py verify-sig --agent opm --target "646 tasks" # scan recent reads for signed DMs dm.py verify-sig "$(dm-sign.sh --from operator-main 'hi')" # verify text directly dm.py read --agent 646 --target "646 tasks" [n] dm.py log [--n 20] dm.py thread --from opm --to 646 --target "646 tasks" "message" """ import argparse import os import re import subprocess import time import sys import uuid import json import os import urllib.request import urllib.parse from datetime import datetime, timezone, timedelta # Shared per-agent rate limiter (jittered) — de-correlates fleet sends # so the 4 nodes don't hit the shared egress IP in lockstep. try: from rate_limiter import rate_limit_wait HAS_RATE_LIMITER = True except ImportError: HAS_RATE_LIMITER = False API = "/home/super/Projects/NetVM/bin/muse-chat-api.py" NETVM_EXEC = "/home/super/Projects/NetVM/bin/netvm-exec.sh" VALID_AGENTS = ["muse", "pip", "646", "opm", "dev", "def"] VALID_SENDERS = ["muse", "pip", "646", "opm", "dev", "def", "super"] VALID_RECIPIENTS = ["muse", "pip", "646", "opm", "dev", "def"] # ---- Canonical follow-up tags (2026-10-04, DEPLOY-DECISIONS.md) ---- # Tags declare follow-up policy at send time. They are metadata only: # stripped from the delivered text, recorded on send_start/sent log # events. The request-store sweeper (VM side) guarantees nudges, # retries, closure, and escalation. Untagged DMs behave exactly as before. TAG_TIMEOUT_MIN = 60 TAG_TIMEOUT_MAX = 604800 TAG_TIMEOUT_DEFAULT = 3600 # 1h, per deploy decisions TAG_NUDGES_MIN = 0 TAG_NUDGES_MAX = 10 TAG_NUDGES_DEFAULT = 2 TAG_ROUTE_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_-]{0,63}$") TAG_THREAD_RE = re.compile(r"^[A-Za-z0-9_-]{1,64}$") # Thread-UUID pattern, shared by nav parsing and the pre-send gate. UUID_RE = (r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-" r"[0-9a-f]{4}-[0-9a-f]{12}") # Trailing-token form: [reply:expected] [reply:timeout=7200] ... TAG_TOKEN_RE = re.compile(r"\[([A-Za-z][A-Za-z0-9:]*)" r"(?:=([^\]]*))?\]\s*$") def parse_canonical_tag(item): """Parse one canonical tag token. Returns (key, value). Bare 'reply:expected' -> ('reply:expected', True). 'route:X' / 'thread:X' use colon form; 'reply:timeout=N' etc use '='. Raises ValueError with a human message on any problem.""" item = item.strip() if "=" in item: k, v = item.split("=", 1) k, v = k.strip(), v.strip() elif item.startswith("route:"): k, v = "route", item[len("route:"):].strip() elif item.startswith("thread:"): k, v = "thread", item[len("thread:"):].strip() else: k, v = item, True if k == "reply:expected": # Bare 'reply:expected' or 'reply:expected=' (empty value from # --tag flag generation); anything else is an error. if v is not True and v != "": raise ValueError("reply:expected takes no value") return k, True if k == "reply:timeout": try: n = int(v) except (TypeError, ValueError): raise ValueError("reply:timeout must be integer seconds") if not TAG_TIMEOUT_MIN <= n <= TAG_TIMEOUT_MAX: raise ValueError("reply:timeout must be %d..%d" % (TAG_TIMEOUT_MIN, TAG_TIMEOUT_MAX)) return k, n if k == "reply:nudges": try: n = int(v) except (TypeError, ValueError): raise ValueError("reply:nudges must be an integer") if not TAG_NUDGES_MIN <= n <= TAG_NUDGES_MAX: raise ValueError("reply:nudges must be %d..%d" % (TAG_NUDGES_MIN, TAG_NUDGES_MAX)) return k, n if k == "reply:escalate": if not v or not TAG_ROUTE_RE.fullmatch(str(v)): raise ValueError("reply:escalate must be an identity") return k, str(v) if k == "route": if not TAG_ROUTE_RE.fullmatch(str(v)): raise ValueError("route must match ^[A-Za-z0-9][A-Za-z0-9_-]{0,63}$") return k, str(v) if k == "thread": if not TAG_THREAD_RE.fullmatch(str(v)): raise ValueError("thread must match ^[A-Za-z0-9_-]{1,64}$") return k, str(v) raise ValueError("unknown tag %r; known: reply:expected, reply:timeout=, " "reply:nudges=, reply:escalate=, route:, thread:" % (k,)) def parse_tags(tag_list): """Validate --tag items. Returns canonical dict. Raises ValueError.""" tags = {} for item in tag_list or []: k, v = parse_canonical_tag(item) if k in tags: raise ValueError("duplicate tag %r" % (k,)) tags[k] = v # Cross-field consistency: policy tags require reply:expected. if "reply:expected" not in tags: for k in ("reply:timeout", "reply:nudges", "reply:escalate"): if k in tags: raise ValueError("tag %s requires reply:expected" % (k,)) # Defaults when reply:expected without explicit policy. if "reply:expected" in tags: tags.setdefault("reply:timeout", TAG_TIMEOUT_DEFAULT) tags.setdefault("reply:nudges", TAG_NUDGES_DEFAULT) return tags def extract_trailing_tags(message): """Pull trailing [tag] tokens off message text. Returns (stripped_message, tags_dict). Unknown brackets are left alone.""" tags = {} text = message.rstrip() while True: m = TAG_TOKEN_RE.search(text) if not m: break key = m.group(1) # Only consume known canonical keys; anything else stays. probe = key + ("=" + m.group(2) if m.group(2) is not None else "") try: k, v = parse_canonical_tag(probe) except ValueError: break if k in tags: break # duplicate; leave it in the text tags[k] = v text = text[:m.start()].rstrip() return text, tags def merge_tags(flag_tags, cli_tags, text_tags): """Precedence: --tag > dedicated flags > trailing text tokens.""" merged = dict(text_tags) merged.update(flag_tags) merged.update(cli_tags) # Defaults when reply:expected without explicit policy. if "reply:expected" in merged: merged.setdefault("reply:timeout", TAG_TIMEOUT_DEFAULT) merged.setdefault("reply:nudges", TAG_NUDGES_DEFAULT) return merged def tags_from_flags(a): """Build canonical tags from argparse --expect-reply et al.""" tags = {} if getattr(a, "expect_reply", False): tags["reply:expected"] = True if getattr(a, "reply_timeout", None) is not None: k, v = parse_canonical_tag("reply:timeout=%s" % a.reply_timeout) tags[k] = v if getattr(a, "reply_nudges", None) is not None: k, v = parse_canonical_tag("reply:nudges=%s" % a.reply_nudges) tags[k] = v if getattr(a, "reply_escalate", None): k, v = parse_canonical_tag("reply:escalate=%s" % a.reply_escalate) tags[k] = v if getattr(a, "route", None): k, v = parse_canonical_tag("route:%s" % a.route) tags[k] = v if getattr(a, "thread", None): k, v = parse_canonical_tag("thread:%s" % a.thread) tags[k] = v return tags LOG_FILE = "/home/super/Projects/NetVM/dm-log.jsonl" # Well-known sidechat name -> thread UUID aliases. # These bypass fuzzy name matching for reliable placement. # Resolved 2026-10-04. Add new entries as sidechats are created. SIDCHAT_ALIASES = { # Intentionally empty. Thread UUIDs rotate (heartbeat-opm died twice in # one day: 5bd5b350 -> 0077e918 -> dead). Hardcoded UUIDs misdeliver or # drop. Targets fall through to job-sidechats.json dynamic mappings, # then name-based `sidechat use` with autoprovision, which self-heals. # Do NOT add hardcoded UUIDs here. } # NOTE 2026-10-04: "646 tasks", "heartbeat", and "646-opm-work" aliases # REMOVED -- all hardcoded UUIDs (1e75a740-..., 0077e918-..., d410b9ad-...) # went stale (threads rotate; the SPA redirects to another thread, causing # misdelivery or dropped digests). Targets now fall through to passthrough # -> job-sidechats.json dynamic mapping -> fuzzy name search -> # autoprovision, which registers the live UUID in job-sidechats.json. # Do NOT re-add hardcoded UUIDs here; use job-sidechats.json instead. def resolve_sidechat_target(target, agent=None): """Resolve target alias or name to UUID dynamically from job-sidechats.json. Supports agent-specific scoped targets (e.g. '646-pip' resolving to recipient's local thread).""" if not target or not str(target).strip(): raise ValueError( "DM target must not be empty: specify a sidechat name/UUID, " "or main (main requires --allow-main-chat)") if target == "main": return target if re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", target.lower()): return target.lower() if target in SIDCHAT_ALIASES: return SIDCHAT_ALIASES[target] sc_file = "/home/super/Projects/NetVM/job-sidechats.json" if os.path.exists(sc_file): try: with open(sc_file, "r", encoding="utf-8") as f: sc_data = json.load(f) # 1. Check agent-scoped alias first (e.g. target:646-pip, agent:646) if agent: scoped_key = f"{target}@{agent}" if scoped_key in sc_data: val = sc_data[scoped_key] if isinstance(val, dict): return val.get("thread_uuid") or val.get("uuid") or target elif isinstance(val, str): return val # 2. Check general key val = sc_data.get(target) if isinstance(val, dict): # If mapped entry defines an agent, but recipient differs, check for agent match if agent and val.get("agent") and val.get("agent") != agent: # Look for an entry explicitly matching this agent for k, item in sc_data.items(): if isinstance(item, dict) and item.get("agent") == agent and k.startswith(target): return item.get("thread_uuid") or target return val.get("thread_uuid") or val.get("uuid") or target elif isinstance(val, str): return val except Exception: pass return target def resolve_sidechat_source(target): """Return where a target resolved from, for placement auditing. Values: main / direct_uuid / static_alias / dynamic_mapping / passthrough. Mirrors resolve_sidechat_target() logic without changing its signature. """ if not target or not str(target).strip(): return "invalid" if target == "main": return "main" if re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", target.lower()): return "direct_uuid" if target in SIDCHAT_ALIASES: return "static_alias" sc_file = "/home/super/Projects/NetVM/job-sidechats.json" if os.path.exists(sc_file): try: with open(sc_file, "r", encoding="utf-8") as f: sc_data = json.load(f) if sc_data.get(target): return "dynamic_mapping" except Exception: pass return "passthrough" def log_event(event): """Append to JSONL log.""" event["ts"] = datetime.now(timezone.utc).isoformat() with open(LOG_FILE, "a") as f: f.write(json.dumps(event) + "\n") def run(cmd, timeout=60, priority=None): # priority: CDP queue priority for the subprocess (high/normal/low). # DM sends are user-facing -> high. Reads stay at default (normal). env = dict(os.environ, CDP_PRIORITY=priority) if priority else None result = subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=timeout, env=env) return result.stdout.strip() def run_full(cmd, timeout=60): """Variant returning (rc, stdout, stderr) for calls where a silent failure is worse than noise (observed 2026-10-03: opm's browser was dead and `dm.py read` printed nothing with exit 0, hiding the outage).""" 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). The send/verify loop only proves the message exists in whichever chat the browser was parked on at read time. A silent navigation drift (the SPA restoring main chat after a thread-URL navigation, a stale DOM read) used to produce verified:true for messages that actually landed in main. This re-navigates to the target by direct URL, asserts the post-nav URL matches the target, then reads the chat back and asserts the message ID is present IN THAT CHAT. Returns (ok, detail). Read-only retries only -- never resends: a misplaced send must fail loudly, not be duplicated. """ detail = {"target": target, "thread_uuid": thread_uuid} try: if target == "main": _rc, _out, _err = run_full( f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat main") if _rc != 0: detail.update({"reason": "nav_failed", "rc": _rc, "err": _err[:200]}) return False, detail time.sleep(2) _u_rc, _u_out, _u_err = run_full( f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} url") detail["actual_url"] = _u_out[:200] if "/thread/" in _u_out: detail["reason"] = "url_mismatch" detail["expected"] = "main chat (no /thread/ in URL)" return False, detail else: expected_uuid = (thread_uuid or "").lower() detail["expected_uuid"] = expected_uuid or None if not expected_uuid: detail["reason"] = "no_thread_uuid" return False, detail # Direct-URL navigation: sidechat use sets # window.location.href to /thread/ and confirms it. _rc, _out, _err = run_full( f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat use {expected_uuid}") if _rc != 0 or "NOTFOUND" in _out or expected_uuid not in _out.lower(): detail.update({"reason": "nav_failed", "rc": _rc, "out": _out[:200], "err": _err[:200]}) return False, detail time.sleep(2) # Independent URL assertion: the nav command's own confirm read # the same window.location.href, so sample it again here. _u_rc, _u_out, _u_err = run_full( f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} url") detail["actual_url"] = _u_out[:200] if expected_uuid not in _u_out.lower(): detail["reason"] = "url_mismatch" return False, detail # Read-back with read-only retries: the SPA may still be rendering # the thread after navigation. A miss here must NOT trigger a resend. for _r in range(3): try: check_msgs = run( f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} messages 5 200") except Exception as e: detail["read_error"] = str(e)[:100] time.sleep(2) continue if msg_id in check_msgs: return True, detail time.sleep(2) detail["reason"] = "msg_not_in_target_chat" return False, detail except Exception as e: detail["reason"] = "exception" detail["error"] = str(e)[:200] return False, detail def assert_pre_send_placement(recipient, target, thread_uuid, direct_nav_done=False): """Pre-send placement gate (2026-10-04; restored after 19:43Z clobber). Independently samples the recipient browser's URL and requires the expected thread UUID in it BEFORE any send. Fails closed: returns (False, detail) and the caller must abort the send. 'main' skips. """ detail = {"target": target} if target == "main": return True, detail if not thread_uuid: # Fresh autoprovision: SPA assigns UUID only on first send. # Strongest pre-send signal: browser must be on a /thread/ page # (the incident was a send parked on the muse.ai landing page). _u_rc, _u_out, _u_err = run_full( f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} url") actual = (_u_out or "").strip() detail["actual_url"] = actual[:200] if "/thread/" not in actual: detail["reason"] = "not_on_thread_page" detail["expected"] = "a /thread/ URL (fresh sidechat, UUID assigned on first send)" return False, detail return True, detail expected = thread_uuid.lower() detail["expected_uuid"] = expected if not direct_nav_done: # Title-search nav only: re-navigate by direct /thread/ URL. _rc, _out, _err = run_full( f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat use {expected}") if _rc != 0 or "NOTFOUND" in _out or expected not in _out.lower(): detail.update({"reason": "direct_nav_failed", "rc": _rc, "out": _out[:200], "err": _err[:200]}) return False, detail time.sleep(2) else: time.sleep(1) _u_rc, _u_out, _u_err = run_full( f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} url") actual = (_u_out or "").strip() detail["actual_url"] = actual[:200] if expected not in actual.lower(): detail["reason"] = "url_mismatch" return False, detail return True, detail def dm_send(agent, target, message, verify=True, raw=False, to_agent=None, tags=None, nudge_meta=None, allow_main_chat=False): """Send a DM. raw=True sends verbatim (for pre-signed messages from dm-sign.sh, which already carry [from:X] [id:Y]): no tagging, no truncation. Non-raw messages are tagged [from:] [id:].""" # Cross-operator: to_agent is the recipient (whose browser/chat to use). # agent is the sender (for attribution). If to_agent is None, send to own chat. recipient = to_agent if to_agent else agent if agent not in VALID_SENDERS: print(f"ERROR: Unknown agent {agent}", file=sys.stderr) sys.exit(1) if recipient not in VALID_RECIPIENTS: print(f"ERROR: Unknown recipient {recipient}", file=sys.stderr) sys.exit(1) if raw: tagged = message m = re.search(r'\[id:([^\]]+)\]', message) msg_id = m.group(1) if m else "raw" else: msg_id = str(uuid.uuid4())[:8] # Single unified attribution format (matches verify-sig's regex). tagged = f"[from:{agent}] [id:{msg_id}] {message}" # Empty target is an error, not a silent main redirect (2026-10-04). # Only explicit --allow-main-chat reaches main. if not target or not str(target).strip(): print(f"ERROR: DM target must not be empty (agent={agent} to={recipient}). " f"Specify a sidechat name/UUID, or main with --allow-main-chat.", file=sys.stderr) sys.exit(2) # Sidechat-first policy (2026-10-04): refuse main-chat sends unless the # caller explicitly opted in. Main threads must never be saturated by # default; use a sidechat target instead. if target == "main" and not allow_main_chat: log_event({"type": "main_chat_blocked", "id": msg_id, "agent": agent, "to": recipient, "target": target}) print(f"ERROR: Refusing main-chat send by policy (agent={agent} to={recipient}). " f"Use a sidechat target, or pass --allow-main-chat for explicit main-chat sends.", file=sys.stderr) sys.exit(2) # Canonical follow-up tags: trailing [bracket] tokens are metadata, # stripped from the delivered text, recorded on the log events. tags = tags or {} # Sidechat-first policy audit marker (2026-10-04): record explicit # main-chat opt-in so the main-chat watchdog can distinguish # authorized main sends from policy violations. if allow_main_chat and target == "main": tags["allow_main_chat"] = True # Extract job_id from [JOB ] marker for followup correlation. # This lets the harvester resolve followups by job_id when a # [RESULT ] reply arrives, even if thread_uuid is null. if "job_id" not in tags: _jm = re.search(r'\[JOB\s+([A-Za-z0-9_-]+)\]', message) if _jm: tags["job_id"] = _jm.group(1) if not raw: message, text_tags = extract_trailing_tags(message) if text_tags: # Re-tag with the stripped body so the wire text is clean. tagged = f"[from:{agent}] [id:{msg_id}] {message}" for k, v in text_tags.items(): tags.setdefault(k, v) _log = {"type": "send_start", "id": msg_id, "agent": agent, "to": recipient, "target": target, "msg": message[:100], "tags": tags} if nudge_meta: _log["nudge_meta"] = nudge_meta log_event(_log) # Navigate the RECIPIENT's browser to the target chat (the send and the # read-back both happen there; navigating the sender's browser was a bug # for cross-operator side-chat targets). # Use run_full: a silent navigation failure used to send the message to # whatever chat the browser happened to be parked on (2026-10-04). 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, recipient) if nav_target != target: log_event({"type": "alias_resolved", "id": msg_id, "target": target, "thread_uuid": nav_target, "source": resolve_sidechat_source(target)}) # Fast path: try fast headless gateway via muse_hybrid (isolated per-node WARP egress) is_uuid = bool(re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", nav_target.strip().lower())) if is_uuid and recipient in VALID_AGENTS: try: import muse_hybrid gw_res, gw_err = muse_hybrid.send_message(recipient, tagged, thread_id=nav_target, wait=0) if gw_res and not gw_err: thread_uuid = nav_target tags["thread"] = thread_uuid log_event({"type": "verified", "id": msg_id, "agent": agent, "to": recipient, "target": target, "thread_uuid": thread_uuid, "transport": "gateway", "placement": "confirmed"}) log_event({"type": "send_done", "id": msg_id, "agent": agent, "to": recipient, "target": target}) log_event({"type": "sent", "id": msg_id, "agent": agent, "to": recipient, "target": target, "verified": True, "transport": "gateway", "tags": tags}) print(f"DM {msg_id} from {agent} to {recipient}/{target}: SENT and VERIFIED thread={thread_uuid} (gateway)") # Follow-up record creation (same as the CDP path below): the # gateway fast-path must not skip followup registration, or # --expect-reply sends silently lose tracking. if tags.get("reply:expected"): _register_followup(msg_id, agent, recipient, target, tags) return msg_id except Exception as _e: log_event({"type": "gateway_fallback", "id": msg_id, "error": str(_e)[:100]}) if target == "main": _rc, _out, _err = run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat main") else: # Reset-at-Begin: If searching sidebar by title (not a direct UUID), reset to main first to expose the sidebar nav_is_uuid = is_uuid if not is_uuid: run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat main") time.sleep(1) _rc, _out, _err = run_full(f"{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} sidechat use {nav_target}") if _rc != 0: log_event({"type": "nav_failed", "id": msg_id, "agent": agent, "to": recipient, "target": target, "rc": _rc, "err": _err[:200]}) print(f"DM {msg_id} from {agent} to {recipient}/{target}: FAILED (navigation rc={_rc}: {_err[:120]})", file=sys.stderr) sys.exit(1) # cmd_sidechat_use exits 0 even on NOTFOUND (it prints "Navigated to: NOTFOUND"). # If sidechat is not found, auto-provision a new thread via `sidechat create`! if target != "main": 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, "target": target, "reason": "autoprovision_create_failed", "rc": _c_rc, "err": _c_err[:200], "out": _c_out[:200]}) print(f"DM {msg_id} from {agent} to {recipient}/{target}: FAILED (sidechat auto-creation failed: {_c_err[:120]})", file=sys.stderr) sys.exit(1) is_new_sidechat = True # If create returned a real thread URL (not /thread/new), # pin it now so the pre-send gate can assert it -- 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: _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}) else: m = re.search(r"/thread/([0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12})", _out, re.I) if m: thread_uuid = m.group(1).lower() thread_url = "https://muse.ai/thread/" + thread_uuid else: # Navigation "succeeded" but we can't confirm where we landed. # Don't trust it: fail instead of sending to an unknown chat. log_event({"type": "nav_failed", "id": msg_id, "agent": agent, "to": recipient, "target": target, "reason": "no_thread_url", "out": _out[:200]}) print(f"DM {msg_id} from {agent} to {recipient}/{target}: FAILED (could not confirm thread URL)", file=sys.stderr) sys.exit(1) # If we resolved via alias, assert the browser actually landed on # the expected thread. A mismatch means the navigation didn't take # (stale URL, blocked nav) -- fail instead of misdelivering. if nav_target != target and thread_uuid != nav_target.lower(): log_event({"type": "nav_failed", "id": msg_id, "agent": agent, "to": recipient, "target": target, "reason": "uuid_mismatch", "expected": nav_target.lower(), "got": thread_uuid}) print(f"DM {msg_id} from {agent} to {recipient}/{target}: FAILED (navigation landed on wrong thread)", file=sys.stderr) sys.exit(1) log_event({"type": "nav_ok", "id": msg_id, "agent": agent, "to": recipient, "target": target, "thread_uuid": thread_uuid, "browser_url": thread_url}) # Pre-send placement assertion: independently confirm the browser is on # the target thread BEFORE any send. Failure fails loudly (no ghost sends). _gate_uuid = thread_uuid or (nav_target.lower() if nav_is_uuid else None) _ps_ok, _ps_detail = assert_pre_send_placement( recipient, target, _gate_uuid, direct_nav_done=nav_is_uuid) if not _ps_ok: log_event({"type": "pre_send_assert_failed", "id": msg_id, "agent": agent, "to": recipient, **_ps_detail}) print(f"DM {msg_id} from {agent} to {recipient}/{target}: FAILED " f"(pre-send placement assertion failed: {_ps_detail.get('reason')}; " f"expected={_ps_detail.get('expected_uuid') or _ps_detail.get('expected')}; " f"actual_url={_ps_detail.get('actual_url')})", file=sys.stderr) sys.exit(1) # Per-agent rate limit (jittered): de-correlate fleet sends across the # shared egress IP. Near no-op at normal pacing (sends already sleep). if HAS_RATE_LIMITER: rate_limit_wait(agent) time.sleep(2) # Send (raw mode: no truncation — signatures must survive intact) safe = tagged.replace('"', '\\"').replace('$', '\\$').replace('`', '\\`') if not raw: safe = safe[:1000] # Send with verification retries # The underlying muse-chat-api.py send returns None/unreliable status, # so we verify by reading the recipient's chat for our message ID. max_retries = 3 delivered = False verified_attempt = None for attempt in range(max_retries): send_out = run(f'{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} send "{safe}"', priority="high") time.sleep(3) # Wait for message to propagate # The send command reports 'sent' / 'enter-sent' / 'NOINPUT'. # 'NOINPUT' (or empty output) means the compose box was never found: # retrying verification would be meaningless, so fail fast. if send_out.strip() not in ("sent", "enter-sent"): log_event({"type": "send_failed", "id": msg_id, "agent": agent, "to": recipient, "target": target, "attempt": attempt + 1, "send_out": send_out[:80]}) if attempt < max_retries - 1: time.sleep(2) continue # Assert the composer cleared: if our text (with the [id:...] tag) is # still sitting in the compose box, the send click never fired and any # later "verification" would be a false positive (2026-10-04 bug). compose = run(f'{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} compose_check') if compose and compose != "NOCOMPOSE" and msg_id in compose: log_event({"type": "compose_stuck", "id": msg_id, "agent": agent, "to": recipient, "target": target, "attempt": attempt + 1}) # Clear the stuck draft so a retry starts clean. run(f'{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} send ""', priority="high") if attempt < max_retries - 1: time.sleep(2) continue # Verify by reading recipient's chat (independent check, not local echo). # cmd_messages now excludes the composer subtree, so a hit here means # the message is actually in the chat history. try: check_msgs = run(f'{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} messages 5 200') if msg_id in check_msgs: delivered = True verified_attempt = attempt + 1 # Capture the actual browser URL at verify time. `messages` # reads whichever chat the browser is currently parked on; if # that isn't the thread we navigated to, the send went to the # wrong chat (placement-blindness, 2026-10-04). This is only # an early signal -- the authoritative placement check runs # after this loop, and verified:true is only logged there. _v_rc, _v_out, _v_err = run_full(f'{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} url') _v_m = re.search(r"/thread/([0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12})", _v_out, re.I) checked_uuid = _v_m.group(1).lower() if _v_m else None if target == "main": if checked_uuid: # Browser was parked in a sidechat thread while # sending to main: the message likely went to the # sidechat, not main. log_event({"type": "placement_mismatch", "id": msg_id, "agent": agent, "to": recipient, "target": target, "expected": "main", "actual_uuid": checked_uuid, "attempt": attempt + 1}) elif checked_uuid and thread_uuid and checked_uuid != thread_uuid: log_event({"type": "placement_mismatch", "id": msg_id, "agent": agent, "to": recipient, "target": target, "expected_uuid": thread_uuid, "actual_uuid": checked_uuid, "attempt": attempt + 1}) break else: log_event({"type": "retry", "id": msg_id, "agent": agent, "to": recipient, "attempt": attempt + 1}) except Exception as e: log_event({"type": "verify_error", "id": msg_id, "error": str(e)[:100]}) if attempt < max_retries - 1: time.sleep(2) # Brief pause before retry 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/(" + UUID_RE + r")", _u_out, re.I) if m_uuid: _captured = m_uuid.group(1).lower() break time.sleep(1) 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 or "")[:200]}) if delivered and thread_uuid and "thread" not in tags: tags["thread"] = thread_uuid if delivered: # Authoritative placement verification (2026-10-04): the loop above # only proved the message exists in whichever chat the browser was # parked on. Re-navigate to the target by direct URL, assert the # post-nav URL, and read the message back IN THAT CHAT before # logging verified:true. A failure here means the message landed in # the wrong chat: fail loudly (do NOT resend -- that would # duplicate the misplaced message). placed, pdetail = verify_placement(recipient, msg_id, target, thread_uuid) if not placed: log_event({"type": "placement_failed", "id": msg_id, "agent": agent, "to": recipient, "loop_attempt": verified_attempt, **pdetail}) log_event({"type": "failed", "id": msg_id, "agent": agent, "to": recipient, "target": target, "verified": False, "reason": "placement_failed"}) print(f"DM {msg_id} from {agent} to {recipient}/{target}: FAILED " f"(placement check failed: {pdetail.get('reason')}; " f"expected={pdetail.get('expected_uuid') or pdetail.get('expected')}; " f"actual_url={pdetail.get('actual_url')})", file=sys.stderr) sys.exit(1) log_event({"type": "verified", "id": msg_id, "agent": agent, "to": recipient, "target": target, "thread_uuid": thread_uuid, "placement": "confirmed", "attempt": verified_attempt}) # Note: Leave recipient browser parked in the sidechat thread to preserve Main Chat DOM log_event({"type": "send_done", "id": msg_id, "agent": agent, "to": recipient, "target": target}) if delivered: _sent = {"type": "sent", "id": msg_id, "agent": agent, "to": recipient, "target": target, "verified": True, "tags": tags} if nudge_meta: _sent["nudge_meta"] = nudge_meta log_event(_sent) # Thread UUID on the verification line: the board server's # _box_dm_send_run parses it into the auto-logged DM row. # Stays None for main-chat sends and UUID-capture failures. if thread_uuid: print(f"DM {msg_id} from {agent} to {recipient}/{target}: SENT and VERIFIED thread={thread_uuid}") else: print(f"DM {msg_id} from {agent} to {recipient}/{target}: SENT and VERIFIED") # Follow-up record creation: if tags request tracking, POST to the # VM's /api/box/followups endpoint. Fail gracefully -- the DM already # sent, so a record-creation failure is logged, not fatal. if tags.get("reply:expected"): _register_followup(msg_id, agent, recipient, target, tags) else: log_event({"type": "failed", "id": msg_id, "agent": agent, "to": recipient, "target": target, "verified": False}) print(f"DM {msg_id} from {agent} to {recipient}/{target}: FAILED (not found in recipient chat after {max_retries} attempts)", file=sys.stderr) sys.exit(1) return msg_id def dm_read(agent, target, n=5, quiet=False, width=200): """Read DMs via headless. width widens the per-paragraph slice (needed for multi-line signature blocks).""" if agent not in VALID_AGENTS: print(f"ERROR: Unknown agent {agent}", file=sys.stderr) sys.exit(1) nav_target = resolve_sidechat_target(target, agent) if target == "main": run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat main") else: is_uuid = bool(re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", nav_target.strip().lower())) if not is_uuid: run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat main") time.sleep(1) run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat use {nav_target}") time.sleep(2) rc, msgs, err = run_full(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} messages {n} {width}") if rc != 0 or not msgs: # Never fail silently: an empty read with exit 0 hid a dead browser # for hours (opm, 2026-10-03). Surface the last error line. detail = err.splitlines()[-1] if err.strip() else "no output" print(f"WARNING: read of {agent}/{target} failed (rc={rc}): {detail}", file=sys.stderr) if not quiet: print(msgs) return msgs SIGNERS_DIR = "/home/super/Projects/NetVM/dm-signers" def dm_select(agent, target=None): """Select active conversation for an agent. - With --target: switch directly to that conversation. - Without --target: list available conversations. Uses sidechat use/main under the hood; stores selection for future commands. """ import os, json, subprocess NETVM_EXEC = os.path.expanduser("~/Projects/NetVM/bin/netvm-exec.sh") API = os.path.expanduser("~/Projects/NetVM/bin/muse-chat-api.py") state_file = os.path.expanduser(f"~/.dm-select-{agent}.json") def run(cmd): r = subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=60) return r.stdout.strip() if target: nav_target = resolve_sidechat_target(target) # Switch to target conversation if target == "main": run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat main") else: run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat use {nav_target}") # Store selection with open(state_file, "w") as f: json.dump({"agent": agent, "target": target}, f) print(f"Selected: {agent}/{target}") return target # List available conversations out = run(f"{NETVM_EXEC} {agent} -- python3 {API} --account {agent} sidechat list") print(f"Conversations for {agent}:") print(" main") # Parse: filter out headers, timestamps, and status lines import re lines = [l.strip() for l in out.split("\n") if l.strip()] ts = re.compile(r"^\d+[mhd]$") skip = {"NOSIDEBAR", "Side chats", "Unread updates"} for line in lines: if line in skip: continue if ts.match(line): continue # Remaining lines are chat names print(f" {line}") print() print(f"Use: dm.py select --agent {agent} --target ") # Show current selection if os.path.exists(state_file): with open(state_file) as f: sel = json.load(f) print(f"Current: {sel.get('target', 'main')}") return None def dm_verify_sig(message=None, agent=None, target=None): """Verify an SSH-signed DM. Pass message text directly, or give --agent/--target to scan recent reads for signed messages.""" import tempfile, re texts = [] if message: texts = [message] elif agent and target: msgs = dm_read(agent, target, n=10, quiet=True, width=2000) # split read output into candidate blocks containing a signature chunks = re.split(r'\n---\n', msgs) texts = [c for c in chunks if '-----BEGIN SSH SIGNATURE-----' in c] if not texts: print("no signed messages found in recent reads", file=sys.stderr) sys.exit(1) else: print("ERROR: provide a message or --agent/--target", file=sys.stderr) sys.exit(1) ok_any = False for text in texts: text = text.strip() m = re.match(r'\[from:([^\]]+)\]\s*\[id:([^\]]+)\]', text) if not m: print("BAD: no [from:]/[id:] header", file=sys.stderr) continue sender, mid = m.group(1), m.group(2) sm = re.search(r'\n-----BEGIN SSH SIGNATURE-----\n(.*?)\n-----END SSH SIGNATURE-----', text, re.S) if not sm: print(f"BAD: [{mid}] no signature block", file=sys.stderr) continue 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") 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 tempfile.TemporaryDirectory() as td: allowed = os.path.join(td, "allowed") with open(allowed, "w") as f: f.write(f"{sender} {pubkey}\n") sigf = os.path.join(td, "sig") with open(sigf, "w") as f: f.write(sigblock) payf = os.path.join(td, "payload") with open(payf, "w") as f: f.write(payload) # verify reads the payload from stdin; feed it from the temp file # (redirect, not a pipe) per the file-based lesson from chat-400 with open(payf, "rb") as fin: r = subprocess.run( ["ssh-keygen", "-Y", "verify", "-f", allowed, "-I", sender, "-n", "dm", "-s", sigf], stdin=fin, capture_output=True, text=True, timeout=15) if r.returncode == 0: print(f"GOOD: [{mid}] signature valid — really from '{sender}'") ok_any = True else: print(f"BAD: [{mid}] signature FAILED for claimed sender '{sender}': " f"{r.stderr.strip()[:120]}", file=sys.stderr) sys.exit(0 if ok_any else 1) def dm_log(n=20): """Show recent log entries.""" if not os.path.exists(LOG_FILE): print("No log file yet") return with open(LOG_FILE) as f: lines = f.readlines() for line in lines[-n:]: e = json.loads(line) print(f"{e['ts'][:19]} {e.get('type','?'):12} {e.get('id','-'):8} {e.get('agent','-')}/{e.get('target','-')}") def dm_thread(from_agent, to_agent, target, message, allow_main_chat=False): """Thread from one agent to another. Attribution is applied exactly once, in the unified [from:X] [id:Y] format, by dm_send.""" return dm_send(from_agent, target, message, to_agent=to_agent, allow_main_chat=allow_main_chat) def main(): p = argparse.ArgumentParser(description="DM: Headless Direct Messages (tagged [from:X] [id:Y], logged; delivery confirmed by recipient read-back). Sidechat-first policy: --target main requires --allow-main-chat.") sub = p.add_subparsers(dest='cmd', required=True) ps = sub.add_parser('send', help='Send a DM (tagged; delivery confirmed by recipient read-back, up to 3 attempts)') ps.add_argument('--agent', required=True, choices=VALID_SENDERS) ps.add_argument('--to', required=False, choices=VALID_RECIPIENTS, default=None, help='Recipient operator (for cross-operator DMs). Uses recipient\'s browser/chat.') ps.add_argument('--target', required=True, help='Conversation: sidechat name or thread UUID. Sidechat-first policy: --target main requires --allow-main-chat.') ps.add_argument('--allow-main-chat', action='store_true', help='Explicit opt-in for main-chat sends (refused by default per sidechat-first policy)') ps.add_argument('--no-verify', action='store_true', help='Accepted for compatibility but ignored: recipient-side read-back verification always runs.') ps.add_argument('--raw', action='store_true', help='Send verbatim: no tagging, no truncation (for pre-signed messages from dm-sign.sh)') ps.add_argument('--expect-reply', action='store_true', help='Tag reply:expected: create a follow-up record (nudge/escalate per policy)') ps.add_argument('--reply-timeout', default=None, help='reply:timeout= seconds until nudge/escalation (60..604800, default 3600)') ps.add_argument('--reply-nudges', default=None, help='reply:nudges= max nudges before escalation (0..10, default 2)') ps.add_argument('--reply-escalate', default=None, help='reply:escalate= identity to alert on timeout') ps.add_argument('--route', default=None, help='route: attach follow-up to a conversation route') ps.add_argument('--thread', default=None, help='thread: associate DM with a thread UUID') ps.add_argument('--tag', action='append', default=[], help='Raw canonical tag (repeatable): reply:expected, reply:timeout=N, reply:nudges=N, reply:escalate=X, route:R, thread:U') ps.add_argument('--nudge-meta', default=None, help='JSON string of nudge metadata (followup_id, nudge_n, dm_id) for audit logging') ps.add_argument('message') def _send_with_tags(a): try: cli_tags = parse_tags(a.tag) flag_tags = tags_from_flags(a) tags = merge_tags(flag_tags, cli_tags, {}) except ValueError as e: print(f"TAG_ERROR: {e}", file=sys.stderr) sys.exit(2) nudge_meta = None if a.nudge_meta: try: import json as _json nudge_meta = _json.loads(a.nudge_meta) if not isinstance(nudge_meta, dict): raise ValueError("nudge-meta must be a JSON object") except Exception as e: print(f"NUDGE_META_ERROR: {e}", file=sys.stderr) sys.exit(2) return dm_send(a.agent, a.target, a.message, verify=not a.no_verify, raw=a.raw, to_agent=a.to, tags=tags, nudge_meta=nudge_meta, allow_main_chat=a.allow_main_chat) ps.set_defaults(func=_send_with_tags) psel = sub.add_parser('select', help='Select active conversation') psel.add_argument('--agent', required=True, choices=VALID_AGENTS) psel.add_argument('--target', required=False, default=None, help='Conversation: main or side chat name') psel.set_defaults(func=lambda a: dm_select(a.agent, a.target)) pr = sub.add_parser('read', help='Read DMs') pr.add_argument('--agent', required=True, choices=VALID_AGENTS) pr.add_argument('--target', required=True) pr.add_argument('--n', type=int, default=5) pr.set_defaults(func=lambda a: dm_read(a.agent, a.target, a.n)) pvs = sub.add_parser('verify-sig', help='Verify SSH signature on a signed DM (pass message text, or --agent/--target to scan recent reads)') pvs.add_argument('--agent', required=False, choices=VALID_AGENTS) pvs.add_argument('--target', required=False) pvs.add_argument('message', nargs='?') pvs.set_defaults(func=lambda a: dm_verify_sig(message=a.message, agent=a.agent, target=a.target)) pl = sub.add_parser('log', help='Show DM log') pl.add_argument('--n', type=int, default=20) pl.set_defaults(func=lambda a: dm_log(a.n)) pt = sub.add_parser('thread', help='Thread from one agent to another') pt.add_argument('--from', dest='from_agent', required=True, choices=VALID_SENDERS) pt.add_argument('--to', dest='to_agent', required=True, choices=VALID_RECIPIENTS) pt.add_argument('--target', required=True, help='Conversation: sidechat name or thread UUID. Sidechat-first policy: main requires --allow-main-chat.') pt.add_argument('--allow-main-chat', action='store_true', help='Explicit opt-in for main-chat sends (refused by default per sidechat-first policy)') pt.add_argument('message') pt.set_defaults(func=lambda a: dm_thread(a.from_agent, a.to_agent, a.target, a.message, allow_main_chat=a.allow_main_chat)) args = p.parse_args() args.func(args) BOX_API_BASE = os.environ.get("BOX_API_BASE", "https://box.muse-dev.online") BOX_SIGN_KEY = os.environ.get("BOX_SIGN_KEY", os.path.expanduser("~/.ssh/id_ed25519")) def _box_sign(identity, endpoint): """Sign a box API request. Returns (ts, sig_armored) or (None, None).""" ts = str(int(time.time())) payload = f"{ts}\n{endpoint}".encode() try: p = subprocess.run( ["ssh-keygen", "-Y", "sign", "-f", BOX_SIGN_KEY, "-n", "box"], input=payload, capture_output=True, timeout=15) if p.returncode != 0: return None, None return ts, p.stdout.decode() except Exception: return None, None def _register_followup(msg_id, agent, recipient, target, tags): """POST a follow-up record to the VM. Never raises -- logs and returns.""" identity_map = { "bl": "bl", "opm": "operator-main", "646": "operator-646", } identity = "bl" ts, sig = _box_sign(identity, "followups") if not ts or not sig: log_event({"type": "followup_register_failed", "id": msg_id, "error": "signing failed"}) return tags_obj = {} if tags.get("reply:expected"): tags_obj["reply:expected"] = True for key in ("reply:timeout", "reply:nudges", "reply:escalate", "route", "thread", "job_id"): if key in tags: tags_obj[key] = tags[key] body = { "dm_id": msg_id, "to": recipient, "from": agent, "target": target, "tags": tags_obj, "dm_sent_ts": datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ"), } query = urllib.parse.urlencode({ "identity": identity, "ts": ts, "sig": sig, }) url = f"{BOX_API_BASE}/api/box/followups?{query}" try: req = urllib.request.Request( url, data=json.dumps(body).encode(), headers={"Content-Type": "application/json", # Cloudflare 1010-blocks python-urllib's default UA; # use an identifiable custom UA instead. "User-Agent": "NetVM-dm.py/1.0 (bl)"}, method="POST") with urllib.request.urlopen(req, timeout=30) as r: resp = json.load(r) log_event({"type": "followup_registered", "id": msg_id, "request_id": resp.get("request_id"), "tracked": resp.get("tracked")}) except Exception as e: log_event({"type": "followup_register_failed", "id": msg_id, "error": str(e)[:200]}) # Also persist to local followups.json for bl autonomy try: f_path = "/home/super/Projects/NetVM/followups.json" followups = {} if os.path.exists(f_path): with open(f_path, "r", encoding="utf-8") as f: followups = json.load(f) now_dt = datetime.now(timezone.utc) timeout_s = int(tags.get("reply:timeout", 3600)) nudges_n = int(tags.get("reply:nudges", 2)) deadline_dt = now_dt + timedelta(seconds=timeout_s) followups[msg_id] = { "dm_id": msg_id, "sender": agent, "recipient": recipient, "target": target, "thread_uuid": tags.get("thread"), "job_id": tags.get("job_id"), "route": tags.get("route"), "sent_at": now_dt.isoformat(), "deadline": deadline_dt.isoformat(), "timeout_s": timeout_s, "nudges_allowed": nudges_n, "nudges_sent": 0, "escalate_to": tags.get("reply:escalate", "opm"), "status": "pending" } tmp = f"{f_path}.tmp.{os.getpid()}" with open(tmp, "w", encoding="utf-8") as f: json.dump(followups, f, indent=2) os.replace(tmp, f_path) except Exception as e: log_event({"type": "local_followup_save_failed", "id": msg_id, "error": str(e)[:200]}) if __name__ == '__main__': main()