diff --git a/bin/muse-threads.py b/bin/muse-threads.py new file mode 100755 index 0000000..9f601cc --- /dev/null +++ b/bin/muse-threads.py @@ -0,0 +1,236 @@ +#!/usr/bin/env python3 +"""muse-threads.py — per-agent thread bookkeeping over the hybrid gateway. + +Thin JSON-contract wrapper around `bin/muse-cli-node` (the NetVM hybrid +gateway transport): runs `session-pin / session-unpin / session-archive / +session-unarchive / session-rename / threads` inside the agent's isolated +network namespace (warp-, per-node WARP egress identity) with its +per-node cookie jar (~/.config/muse-cli//cookies.txt, 0600) and +automatic CDP cookie refresh on auth failure. + +This helper adds what the Box API needs on top of the raw transport: + - one JSON object on stdout: {"ok": true, ...} or + {"ok": false, "code": "...", "error": "...", "detail": ...} + - exit 0 on success, nonzero on failure + - input validation (agent / thread-id / title) with spec error codes + - gateway failure mapping to THREAD_* codes server.py understands + - a small delay before gateway writes to respect rate limits + - never prints secrets, never uses a shell + +Usage: + bin/muse-threads.py --agent [--thread ] [--title ] + ops: list | pin | unpin | archive | unarchive | rename + agent: one of muse, pip, 646, opm (node == agent == account) + +Error codes (server.py maps these to HTTP): + AGENT_NOT_FOUND unknown agent name + THREAD_NOT_CONFIGURED no cookie jar for this agent (fail closed) + INVALID_NAME agent or thread id fails the id regex + INVALID_TITLE rename title missing/empty/too long/has controls + THREAD_AUTH_FAILED auth failed even after cookie refresh + THREAD_NOT_FOUND gateway says no such session + GATEWAY_UNAVAILABLE timeout, or gateway vm_resolution_issue + GATEWAY_ERROR other gateway failures (server maps to 502) + +Session: sidechat/muse-cli-threads-helper +""" + +import argparse +import json +import os +import re +import subprocess +import sys +import time + +KNOWN_AGENTS = ["muse", "pip", "646", "opm"] +AGENT_RE = re.compile(r"^[a-z0-9-]{1,64}$") +THREAD_RE = re.compile(r"^[A-Za-z0-9-]{1,64}$") +MUTATE_DELAY = float(os.environ.get("MUSE_THREADS_DELAY", "2.0")) +SUBPROCESS_TIMEOUT = 120 + +EXIT_USAGE = 1 +EXIT_AUTH = 2 +EXIT_GATEWAY = 3 +EXIT_NOT_CONFIGURED = 5 + +NETVM_BIN = os.path.dirname(os.path.abspath(__file__)) +MUSE_CLI_NODE = os.path.join(NETVM_BIN, "muse-cli-node") + + +def fail(code, error, detail=None, exit_code=EXIT_GATEWAY): + payload = {"ok": False, "code": code, "error": error} + if detail: + payload["detail"] = detail + print(json.dumps(payload)) + sys.exit(exit_code) + + +def cookie_jar(agent): + return os.path.join(os.path.expanduser("~"), ".config", "muse-cli", + agent, "cookies.txt") + + +def run_node_cli(agent, argv): + """Run muse-cli-node <agent> <argv>. Returns (rc, stdout, stderr_tail).""" + if not os.path.isfile(MUSE_CLI_NODE): + fail("GATEWAY_ERROR", + f"muse-cli-node not found at {MUSE_CLI_NODE}", + exit_code=EXIT_GATEWAY) + try: + proc = subprocess.run( + [MUSE_CLI_NODE, agent] + argv, + capture_output=True, + text=True, + timeout=SUBPROCESS_TIMEOUT, + ) + except subprocess.TimeoutExpired: + fail("GATEWAY_UNAVAILABLE", "muse-cli-node timed out", + exit_code=EXIT_GATEWAY) + except OSError as e: + fail("GATEWAY_ERROR", f"failed to exec muse-cli-node: {e}", + exit_code=EXIT_GATEWAY) + err_lines = proc.stderr.strip().splitlines() + err_tail = "\n".join(err_lines[-5:]) if err_lines else "" + return proc.returncode, proc.stdout, err_tail + + +def map_failure(rc, err_tail): + """Map a transport failure to (code, error, exit_code). Never leaks cookies.""" + low = err_tail.lower() + if rc == 2 or "autherror" in low or "auth error" in low: + return ("THREAD_AUTH_FAILED", + "gateway auth failed even after cookie refresh — " + "re-export cookies via refresh-node-cookies.py", + err_tail or None, EXIT_AUTH) + if "vm_resolution_issue" in low: + return ("GATEWAY_UNAVAILABLE", "gateway vm_resolution_issue", + err_tail or None, EXIT_GATEWAY) + if rc == 4: + return ("GATEWAY_UNAVAILABLE", "gateway timeout", + err_tail or None, EXIT_GATEWAY) + if ("not found" in low or "no such session" in low or " 404" in low + or "does not exist" in low): + return ("THREAD_NOT_FOUND", "no such session on that account", + err_tail or None, EXIT_GATEWAY) + return ("GATEWAY_ERROR", "gateway error", err_tail or None, EXIT_GATEWAY) + + +def parse_threads_blob(blob): + try: + d = json.loads(blob) + except json.JSONDecodeError as e: + fail("BL_BAD_DATA", "threads output was unparseable", + detail=str(e), exit_code=EXIT_GATEWAY) + if isinstance(d, dict) and "threads" in d: + return d["threads"] + if isinstance(d, list): + return d + fail("BL_BAD_DATA", "threads output had unexpected shape", + exit_code=EXIT_GATEWAY) + + +def normalize_thread(t): + return { + "thread_id": t.get("session_id") or t.get("thread_id") or t.get("id"), + "title": t.get("title"), + "pinned": bool(t.get("pinned")), + "archived": bool(t.get("archived")), + "updated": t.get("updated"), + } + + +def cmd_list(agent): + rc, out, err = run_node_cli(agent, ["threads", "--all"]) + if rc != 0: + code, error, detail, ec = map_failure(rc, err) + fail(code, error, detail=detail, exit_code=ec) + threads = [normalize_thread(t) for t in parse_threads_blob(out)] + print(json.dumps({"ok": True, "agent": agent, "threads": threads})) + + +def cmd_mutate(agent, op, thread_id, title=None): + if op == "rename": + argv = ["session-rename", thread_id, title] + else: + argv = [f"session-{op}", thread_id] + time.sleep(MUTATE_DELAY) + rc, out, err = run_node_cli(agent, argv) + if rc != 0: + code, error, detail, ec = map_failure(rc, err) + fail(code, error, detail=detail, exit_code=ec) + try: + result = json.loads(out) if out.strip() else {} + except json.JSONDecodeError: + result = {"raw": out.strip()[:500]} + payload = {"ok": True, "agent": agent, "op": op, "thread_id": thread_id, + "result": result} + if op == "pin": + payload["pinned"] = True + elif op == "unpin": + payload["pinned"] = False + elif op == "archive": + payload["archived"] = True + elif op == "unarchive": + payload["archived"] = False + elif op == "rename": + payload["title"] = title + print(json.dumps(payload)) + + +def valid_title(title): + if title is None: + return False + s = title.strip() + if not (1 <= len(s) <= 128): + return False + if any(ord(c) < 32 or ord(c) == 127 for c in s): + return False + return True + + +def main(): + ap = argparse.ArgumentParser( + description="per-agent thread bookkeeping via the hybrid gateway") + ap.add_argument("op", choices=["list", "pin", "unpin", "archive", + "unarchive", "rename"]) + ap.add_argument("--agent", required=True, + help="one of: " + ", ".join(KNOWN_AGENTS)) + ap.add_argument("--thread", default=None, help="session/thread id") + ap.add_argument("--title", default=None, help="new title (rename only)") + args = ap.parse_args() + + agent = args.agent + if not AGENT_RE.match(agent) or agent not in KNOWN_AGENTS: + fail("AGENT_NOT_FOUND", f"unknown agent '{agent}'", + exit_code=EXIT_USAGE) + + if not os.path.isfile(cookie_jar(agent)): + fail("THREAD_NOT_CONFIGURED", + f"no cookie jar for agent '{agent}' " + f"(expected {cookie_jar(agent)}); " + f"run bin/refresh-node-cookies.py {agent}", + exit_code=EXIT_NOT_CONFIGURED) + + if args.op == "list": + cmd_list(agent) + return + + thread_id = args.thread + if not thread_id or not THREAD_RE.match(thread_id): + fail("INVALID_NAME", "thread id must match ^[A-Za-z0-9-]{1,64}$", + exit_code=EXIT_USAGE) + + if args.op == "rename": + if not valid_title(args.title): + fail("INVALID_TITLE", + "title must be 1-128 chars after stripping, " + "no control characters", + exit_code=EXIT_USAGE) + cmd_mutate(agent, "rename", thread_id, title=args.title.strip()) + else: + cmd_mutate(agent, args.op, thread_id) + + +if __name__ == "__main__": + main()