Files
box/bin/muse-threads.py
T

251 lines
9.0 KiB
Python
Raw Normal View History

#!/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-<node>, per-node WARP egress identity) with its
per-node cookie jar (~/.config/muse-cli/<node>/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 <op> --agent <agent> [--thread <id>] [--title <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):
is_main = (
t.get("thread") is False or
t.get("is_main") is True or
(t.get("title") and t.get("title").lower() in ("main chat", "main", "start conversation with muse"))
)
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"),
"thread": t.get("thread", True),
"is_main": is_main,
}
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)]
mains = [t for t in threads if t.get("is_main")]
pinned = [t for t in threads if t.get("pinned") and not t.get("is_main")]
regular = [t for t in threads if not t.get("is_main") and not t.get("pinned")]
mains.sort(key=lambda t: t.get("updated") or "", reverse=True)
pinned.sort(key=lambda t: t.get("updated") or "", reverse=True)
regular.sort(key=lambda t: t.get("updated") or "", reverse=True)
sorted_threads = mains + pinned + regular
print(json.dumps({"ok": True, "agent": agent, "threads": sorted_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()