| 1 | #!/usr/bin/env python3 |
| 2 | import argparse |
| 3 | import base64 |
| 4 | from concurrent.futures import ThreadPoolExecutor |
| 5 | import hashlib |
| 6 | import importlib |
| 7 | import json |
| 8 | import os |
| 9 | from pathlib import Path |
| 10 | import queue |
| 11 | import re |
| 12 | import tempfile |
| 13 | import socket |
| 14 | import sqlite3 |
| 15 | import struct |
| 16 | import subprocess |
| 17 | import sys |
| 18 | import threading |
| 19 | import time |
| 20 | import urllib.error |
| 21 | import urllib.parse |
| 22 | import urllib.request |
| 23 | import uuid |
| 24 | |
| 25 | |
| 26 | class NoRedirect(urllib.request.HTTPRedirectHandler): |
| 27 | def redirect_request(self, request, fp, code, message, headers, newurl): |
| 28 | return None |
| 29 | |
| 30 | |
| 31 | class Agent: |
| 32 | def __init__(self, address, token, status=101, origin=None, host="globe.studio.test"): |
| 33 | target = urllib.parse.urlsplit(address) |
| 34 | self.socket = socket.create_connection((target.hostname, target.port), timeout=5) |
| 35 | self.socket.settimeout(45) |
| 36 | key = base64.b64encode(os.urandom(16)).decode() |
| 37 | fields = {"Host": host, "Upgrade": "websocket", "Connection": "Upgrade", "Sec-WebSocket-Version": "13", |
| 38 | "Sec-WebSocket-Key": key, "Authorization": "Bearer " + token} |
| 39 | if origin is not None: |
| 40 | fields["Origin"] = origin |
| 41 | self.socket.sendall(("GET /agent/connect HTTP/1.1\r\n" + "".join(k + ": " + v + "\r\n" for k, v in fields.items()) + "\r\n").encode()) |
| 42 | self.stream = self.socket.makefile("rb", buffering=0) |
| 43 | response = self.stream.readline().decode().split() |
| 44 | assert int(response[1]) == status, response |
| 45 | headers = {} |
| 46 | while line := self.stream.readline().strip(): |
| 47 | k, v = line.decode().split(":", 1) |
| 48 | headers[k.lower()] = v.strip() |
| 49 | if status != 101: |
| 50 | self.close() |
| 51 | return |
| 52 | expected = base64.b64encode(hashlib.sha1((key + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11").encode()).digest()).decode() |
| 53 | assert headers["sec-websocket-accept"] == expected |
| 54 | self.lock = threading.Lock() |
| 55 | self.frames = queue.Queue() |
| 56 | self.mode = "reply" |
| 57 | self.pings = 0 |
| 58 | self.closed = threading.Event() |
| 59 | self.machine = self.receive()["machine_id"] |
| 60 | self.thread = threading.Thread(target=self.run, daemon=True) |
| 61 | self.thread.start() |
| 62 | |
| 63 | def read(self, size): |
| 64 | result = b"" |
| 65 | while len(result) < size: |
| 66 | part = self.stream.read(size - len(result)) |
| 67 | if not part: |
| 68 | raise EOFError() |
| 69 | result += part |
| 70 | return result |
| 71 | |
| 72 | def receive(self): |
| 73 | while True: |
| 74 | first, size = self.read(2) |
| 75 | assert first & 128 and not size & 128 |
| 76 | size &= 127 |
| 77 | if size == 126: |
| 78 | size = struct.unpack("!H", self.read(2))[0] |
| 79 | elif size == 127: |
| 80 | size = struct.unpack("!Q", self.read(8))[0] |
| 81 | assert size <= 4 * 1024 * 1024 |
| 82 | data = self.read(size) |
| 83 | opcode = first & 15 |
| 84 | if opcode == 9: |
| 85 | self.pings += 1 |
| 86 | self.send(data, opcode=10) |
| 87 | elif opcode == 8: |
| 88 | raise EOFError() |
| 89 | else: |
| 90 | assert opcode == 1 |
| 91 | return json.loads(data) |
| 92 | |
| 93 | def send(self, value, opcode=1): |
| 94 | data = value if isinstance(value, bytes) else json.dumps(value).encode() |
| 95 | prefix = bytes([128 | opcode]) |
| 96 | length = len(data) |
| 97 | if length < 126: |
| 98 | prefix += bytes([128 | length]) |
| 99 | elif length < 65536: |
| 100 | prefix += bytes([128 | 126]) + struct.pack("!H", length) |
| 101 | else: |
| 102 | prefix += bytes([128 | 127]) + struct.pack("!Q", length) |
| 103 | mask = os.urandom(4) |
| 104 | packet = prefix + mask + bytes(byte ^ mask[i % 4] for i, byte in enumerate(data)) |
| 105 | with self.lock: |
| 106 | self.socket.sendall(packet) |
| 107 | |
| 108 | def run(self): |
| 109 | try: |
| 110 | while True: |
| 111 | frame = self.receive() |
| 112 | assert set(frame) == {"id", "method", "params"} |
| 113 | uuid.UUID(frame["id"]) |
| 114 | self.frames.put(frame) |
| 115 | if self.mode == "reply": |
| 116 | self.send({"id": frame["id"], "result": {"fixture": self.machine, "method": frame["method"], "params": frame["params"]}}) |
| 117 | elif self.mode == "drop": |
| 118 | self.socket.shutdown(socket.SHUT_RDWR) |
| 119 | break |
| 120 | except (EOFError, OSError): |
| 121 | pass |
| 122 | finally: |
| 123 | self.closed.set() |
| 124 | |
| 125 | def close(self): |
| 126 | try: |
| 127 | self.socket.shutdown(socket.SHUT_RDWR) |
| 128 | except OSError: |
| 129 | pass |
| 130 | self.stream.close() |
| 131 | self.socket.close() |
| 132 | if hasattr(self, "thread"): |
| 133 | self.thread.join(timeout=5) |
| 134 | |
| 135 | |
| 136 | def main(): |
| 137 | parser = argparse.ArgumentParser() |
| 138 | parser.add_argument("--url", required=True) |
| 139 | parser.add_argument("--proof-file", type=Path, required=True) |
| 140 | parser.add_argument("--data-dir", type=Path, required=True) |
| 141 | parser.add_argument("--restart-unit", required=True) |
| 142 | parser.add_argument("--output", type=Path) |
| 143 | parser.add_argument("--agent-dir", type=Path) |
| 144 | parser.add_argument("--agent-origin") |
| 145 | parser.add_argument("--agent-ca", type=Path) |
| 146 | args = parser.parse_args() |
| 147 | if args.output: |
| 148 | args.output.unlink(missing_ok=True) |
| 149 | origin = "https://globe.studio.test" |
| 150 | resource = origin + "/mcp/agents" |
| 151 | proof = args.proof_file.read_text().strip() |
| 152 | opener = urllib.request.build_opener(NoRedirect) |
| 153 | marker = "relay-fixture-" + uuid.uuid4().hex |
| 154 | names = [marker + "-one", marker + "-two"] |
| 155 | sys.path.insert(0, str(Path(__file__).resolve().parent.parent / "service/keycloak")) |
| 156 | from api import Keycloak |
| 157 | keycloak = Keycloak("keycloak.studio.test", importlib.import_module("dashboard-run").secret("get", "keycloak", "password"), attempts=1) |
| 158 | agents = [] |
| 159 | credentials = [] |
| 160 | |
| 161 | def http(path, method="GET", body=None, actor=None, token=None, status=200, form=False, headers=None): |
| 162 | fields = {"Host": "globe.studio.test", "Content-Type": "application/x-www-form-urlencoded" if form else "application/json"} |
| 163 | if actor: |
| 164 | fields.update({"Studio-Proxy-Token": proof, "User-Name": actor, "User-Groups": "", "Origin": origin}) |
| 165 | if token: |
| 166 | fields.update({"Authorization": "Bearer " + token, "Accept": "application/json, text/event-stream", "MCP-Protocol-Version": "2025-11-25"}) |
| 167 | fields.update(headers or {}) |
| 168 | data = urllib.parse.urlencode(body).encode() if form else json.dumps(body).encode() if body is not None else None |
| 169 | try: |
| 170 | response = opener.open(urllib.request.Request(args.url + path, method=method, headers=fields, data=data), timeout=65) |
| 171 | except urllib.error.HTTPError as error: |
| 172 | response = error |
| 173 | with response: |
| 174 | content = response.read() |
| 175 | assert response.status == status, (path, response.status, content[:300]) |
| 176 | value = json.loads(content) if content and response.headers.get("Content-Type", "").startswith("application/json") else content.decode() |
| 177 | return value, response.headers |
| 178 | |
| 179 | def pair(actor, label): |
| 180 | pending, _ = http("/pairing", "POST", {"name": label, "platform": "fixture"}, status=201) |
| 181 | credentials.append(pending["token"]) |
| 182 | http("/pairing", token=pending["token"], status=202) |
| 183 | machine, _ = http("/api/mcp/relay/pair", "POST", {"code": pending["code"]}, actor=actor) |
| 184 | http("/api/mcp/relay/pair", "POST", {"code": pending["code"]}, actor=actor, status=410) |
| 185 | assert "tokenHash" not in machine and "token" not in machine |
| 186 | result, _ = http("/pairing", token=pending["token"]) |
| 187 | assert result["machine_id"] == machine["id"] |
| 188 | return machine, pending["token"] |
| 189 | |
| 190 | def key(actor, machines, write=False): |
| 191 | result, _ = http("/api/mcp/relay/keys", "POST", {"name": marker, "resources": machines, "write": write}, actor=actor) |
| 192 | credentials.append(result["key"]) |
| 193 | return result |
| 194 | |
| 195 | def rpc(token, method, params=None, error=False): |
| 196 | response, _ = http("/mcp/agents", "POST", {"jsonrpc": "2.0", "id": 1, "method": method, **({"params": params} if params is not None else {})}, token=token) |
| 197 | assert "error" not in response, response |
| 198 | result = response["result"] |
| 199 | if method == "tools/call": |
| 200 | assert bool(result.get("isError")) == error, result |
| 201 | return result |
| 202 | |
| 203 | def call(token, name, fields=None, error=False): |
| 204 | return rpc(token, "tools/call", {"name": name, "arguments": fields or {}}, error=error) |
| 205 | |
| 206 | def command(token, machine, method="list_threads", params=None, status=200): |
| 207 | return http("/api/v1/machines/" + machine + "/commands", "POST", {"method": method, "params": params or {}}, token=token, status=status)[0] |
| 208 | |
| 209 | def connect(token): |
| 210 | agent = Agent(args.url, token) |
| 211 | agents.append(agent) |
| 212 | return agent |
| 213 | |
| 214 | try: |
| 215 | for name in names: |
| 216 | keycloak.request("/admin/realms/master/users", "POST", {"username": name, "enabled": True}) |
| 217 | first, token1 = pair(names[0], "First fixture") |
| 218 | second, token2 = pair(names[0], "Second fixture") |
| 219 | foreign, foreign_token = pair(names[1], "Foreign fixture") |
| 220 | http("/api/mcp/relay/machines/" + first["id"], "DELETE", actor=names[1], status=404) |
| 221 | http("/api/mcp/relay/machines/" + first["id"], "PATCH", {"name": "Renamed fixture"}, actor=names[0]) |
| 222 | http("/api/mcp/relay/pair", "POST", {"code": "random"}, actor=names[0], status=403, headers={"Origin": "https://other.invalid"}) |
| 223 | http("/api/mcp/relay/keys", "POST", {"name": marker, "resources": [foreign["id"]]}, actor=names[0], status=403) |
| 224 | http("/pairing", "POST", {}, status=421, headers={"Host": "other.invalid"}) |
| 225 | read = key(names[0], [first["id"], second["id"]]) |
| 226 | control = key(names[0], [first["id"], second["id"]], write=True) |
| 227 | command(control["key"], foreign["id"], status=403) |
| 228 | command(read["key"], first["id"], "start_thread", {"provider": "codex", "cwd": "/owned", "message": "fixture"}, status=403) |
| 229 | http("/api/v1/machines", status=401, headers={"User-Name": names[0], "User-Groups": "infra-admin", "Studio-Proxy-Token": proof}) |
| 230 | http("/api/v1/machines", token=read["key"], status=403, headers={"Origin": "https://other.invalid"}) |
| 231 | Agent(args.url, token1, status=401, origin=origin) |
| 232 | Agent(args.url, token1, status=421, host="other.invalid") |
| 233 | live = opener.open(urllib.request.Request(args.url + "/api/mcp/relay/live", headers={"Host": "globe.studio.test", "Studio-Proxy-Token": proof, "User-Name": names[0], "User-Groups": ""}), timeout=10) |
| 234 | def snapshot(): |
| 235 | while line := live.readline(): |
| 236 | if line.startswith(b"data:"): |
| 237 | value = json.loads(line.split(b":", 1)[1]) |
| 238 | assert {m["id"] for m in value} == {first["id"], second["id"]} |
| 239 | return value |
| 240 | raise AssertionError("Machine stream closed") |
| 241 | try: |
| 242 | assert not any(m["online"] for m in snapshot()) |
| 243 | a1 = connect(token1) |
| 244 | assert sum(m["online"] for m in snapshot()) == 1 |
| 245 | a2 = connect(token2) |
| 246 | assert all(m["online"] for m in snapshot()) |
| 247 | http("/api/mcp/relay/machines/" + first["id"], "PATCH", {"name": "Stream renamed fixture"}, actor=names[0]) |
| 248 | assert next(m for m in snapshot() if m["id"] == first["id"])["name"] == "Stream renamed fixture" |
| 249 | finally: |
| 250 | live.close() |
| 251 | Agent(args.url, token1, status=409) |
| 252 | assert a1.machine == first["id"] and a2.machine == second["id"] |
| 253 | init = rpc(control["key"], "initialize", {"protocolVersion": "2025-11-25", "capabilities": {}, "clientInfo": {"name": marker, "version": "1"}}) |
| 254 | assert init["capabilities"]["tools"] == {} |
| 255 | tools = rpc(control["key"], "tools/list")["tools"] |
| 256 | assert len(tools) == 7 |
| 257 | assert next(t for t in tools if t["name"] == "send_message")["annotations"]["readOnlyHint"] is False |
| 258 | call(control["key"], "send_message", {"provider": "codex", "thread_id": str(uuid.uuid4()), "message": "ambiguous fixture"}, error=True) |
| 259 | assert a1.frames.empty() and a2.frames.empty() |
| 260 | results = call(read["key"], "list_threads")["structuredContent"]["results"] |
| 261 | assert len(results) == 2 and all(r["result"]["params"] == {"limit": 30} for r in results) |
| 262 | a1.frames.get(timeout=2); a2.frames.get(timeout=2) |
| 263 | call(control["key"], "set_target_machines", {"machine_ids": [foreign["id"]]}, error=True) |
| 264 | call(control["key"], "set_target_machines", {"machine_ids": [first["id"]]}) |
| 265 | unchanged, _ = http("/api/v1/machines", token=read["key"]) |
| 266 | assert len(unchanged) == 2 and all(m["selected"] for m in unchanged) |
| 267 | write_params = {"provider": "claude", "thread_id": str(uuid.uuid4()), "message": "owned fixture", "expected_turn_id": str(uuid.uuid4())} |
| 268 | call(control["key"], "send_message", write_params) |
| 269 | assert a1.frames.get(timeout=2)["params"] == write_params and a2.frames.empty() |
| 270 | command(control["key"], first["id"], params={"limit": 101}, status=400) |
| 271 | assert a1.frames.empty() |
| 272 | call(control["key"], "send_message", {"provider": "codex", "thread_id": str(uuid.uuid4()), "message": "雪" * 100000}, error=True) |
| 273 | assert a1.frames.empty() |
| 274 | a1.mode = "hold" |
| 275 | with ThreadPoolExecutor(max_workers=9) as pool: |
| 276 | pending = [pool.submit(command, control["key"], first["id"]) for _ in range(8)] |
| 277 | frames = [a1.frames.get(timeout=10) for _ in range(8)] |
| 278 | command(control["key"], first["id"], status=429) |
| 279 | for frame in frames: |
| 280 | a1.send({"id": frame["id"], "result": {"bounded": True}}) |
| 281 | assert all(future.result(timeout=10)["result"] == {"bounded": True} for future in pending) |
| 282 | started = time.monotonic() |
| 283 | outcome = command(control["key"], first["id"], status=409) |
| 284 | elapsed = time.monotonic() - started |
| 285 | frame = a1.frames.get(timeout=2) |
| 286 | assert "unknown" in outcome["error"] and 29 <= elapsed < 40 and a1.pings >= 1 |
| 287 | a1.send({"id": frame["id"], "result": {"late": True}}) |
| 288 | a1.mode = "reply" |
| 289 | assert command(control["key"], first["id"])["result"]["params"] == {"limit": 30} |
| 290 | a1.frames.get(timeout=2) |
| 291 | a1.mode = "drop" |
| 292 | started = time.monotonic() |
| 293 | outcome = command(control["key"], first["id"], "send_message", write_params, status=409) |
| 294 | assert "unknown" in outcome["error"] and time.monotonic() - started < 5 |
| 295 | a1.frames.get(timeout=2) |
| 296 | a1.close() |
| 297 | a1 = connect(token1) |
| 298 | time.sleep(.2) |
| 299 | assert a1.frames.empty() |
| 300 | assert command(control["key"], first["id"])["result"]["method"] == "list_threads" |
| 301 | a1.frames.get(timeout=2) |
| 302 | metadata, _ = http("/.well-known/oauth-protected-resource/mcp/agents") |
| 303 | assert metadata["resource"] == resource and "sessions:write" in metadata["scopes_supported"] |
| 304 | client, _ = http("/oauth/register", "POST", {"client_name": marker, "redirect_uris": ["http://127.0.0.1:20001/callback"]}, status=201) |
| 305 | verifier = uuid.uuid4().hex + uuid.uuid4().hex |
| 306 | fields = {"response_type": "code", "client_id": client["client_id"], "redirect_uri": client["redirect_uris"][0], "resource": resource, "scope": "sessions:read sessions:write offline_access", |
| 307 | "code_challenge_method": "S256", "code_challenge": base64.urlsafe_b64encode(hashlib.sha256(verifier.encode()).digest()).decode().rstrip("=")} |
| 308 | _, headers = http("/oauth/authorize?" + urllib.parse.urlencode(fields), status=302) |
| 309 | request_id = urllib.parse.urlsplit(headers["Location"]).path.removeprefix("/connect/") |
| 310 | consent_path = "/api/mcp/consent/" + request_id |
| 311 | details, _ = http(consent_path, actor=names[0]) |
| 312 | assert {r["id"] for r in details["resources"]} == {first["id"], second["id"]} |
| 313 | http(consent_path, actor=names[1], status=403) |
| 314 | result, _ = http(consent_path, "POST", {"resources": [first["id"]]}, actor=names[0]) |
| 315 | code = urllib.parse.parse_qs(urllib.parse.urlsplit(result["redirect"]).query)["code"][0] |
| 316 | tokens, _ = http("/oauth/token", "POST", {"grant_type": "authorization_code", "client_id": client["client_id"], "redirect_uri": fields["redirect_uri"], "resource": resource, "code_verifier": verifier, "code": code}, form=True) |
| 317 | credentials.extend([tokens["access_token"], tokens["refresh_token"]]) |
| 318 | assert len(call(tokens["access_token"], "list_machines")["structuredContent"]["machines"]) == 1 |
| 319 | http("/mcp/observability", "POST", {}, token=tokens["access_token"], status=401) |
| 320 | http("/oauth/token", "POST", {"grant_type": "refresh_token", "client_id": client["client_id"], "refresh_token": tokens["refresh_token"], "resource": origin + "/mcp/observability"}, status=400, form=True) |
| 321 | rotated, _ = http("/oauth/token", "POST", {"grant_type": "refresh_token", "client_id": client["client_id"], "refresh_token": tokens["refresh_token"], "resource": resource}, form=True) |
| 322 | credentials.extend([rotated["access_token"], rotated["refresh_token"]]) |
| 323 | with sqlite3.connect(args.data_dir / "connections.sqlite") as db: |
| 324 | rows = db.execute("SELECT key,value FROM records").fetchall() |
| 325 | assert all(not any(secret in key + value for secret in credentials) for key, value in rows) |
| 326 | expired, _ = http("/pairing", "POST", {"name": marker, "platform": "test"}, status=201) |
| 327 | code_hash = hashlib.sha256(expired["code"].replace("-", "").encode()).hexdigest() |
| 328 | db.execute("UPDATE records SET expires=1 WHERE key=?", ("pair:" + code_hash,)) |
| 329 | http("/pairing", token=expired["token"], status=410) |
| 330 | subprocess.run(["systemctl", "restart", args.restart_unit], check=True, timeout=180, capture_output=True) |
| 331 | for agent in agents: |
| 332 | assert agent.closed.wait(5) |
| 333 | a1, a2 = connect(token1), connect(token2) |
| 334 | assert command(read["key"], first["id"])["result"]["method"] == "list_threads" |
| 335 | a1.frames.get(timeout=2) |
| 336 | assert len(call(rotated["access_token"], "list_machines")["structuredContent"]["machines"]) == 1 |
| 337 | http("/api/mcp/connections/" + control["id"], "DELETE", actor=names[0], status=204) |
| 338 | http("/api/v1/machines", token=control["key"], status=401) |
| 339 | assert a1.frames.empty() |
| 340 | http("/api/mcp/relay/machines/" + first["id"], "DELETE", actor=names[0], status=204) |
| 341 | assert a1.closed.wait(5) |
| 342 | Agent(args.url, token1, status=401) |
| 343 | http("/pairing", token=token1, status=410) |
| 344 | remaining = call(read["key"], "list_machines")["structuredContent"]["machines"] |
| 345 | assert len(remaining) == 1 and remaining[0]["id"] == second["id"] and remaining[0]["selected"] |
| 346 | http("/api/v1/targets", "PUT", {"machine_ids": [first["id"]]}, token=read["key"], status=403) |
| 347 | a2.send({"id": str(uuid.uuid4()), "result": {}, "error": "malformed"}) |
| 348 | assert a2.closed.wait(5) |
| 349 | original_agent_checks = None |
| 350 | if args.agent_dir: |
| 351 | node = next(Path("/nix/store").glob("*nodejs-24*/bin/node")) |
| 352 | with tempfile.TemporaryDirectory(prefix="studio-relay-cli-") as temporary: |
| 353 | local = Path(temporary) |
| 354 | home = local / "home" |
| 355 | home.mkdir() |
| 356 | (home / "codex").mkdir() |
| 357 | (home / "claude/projects").mkdir(parents=True) |
| 358 | with sqlite3.connect(home / "codex/state_5.sqlite") as catalog: |
| 359 | catalog.execute("CREATE TABLE threads(id TEXT,title TEXT,cwd TEXT,updated_at INTEGER,rollout_path TEXT,archived INTEGER)") |
| 360 | allowed = local / "workspace" |
| 361 | allowed.mkdir() |
| 362 | resolver = local / "lookup.cjs" |
| 363 | resolver.write_text("const dns=require('node:dns'); const original=dns.lookup; dns.lookup=function(host,opts,callback){if(host==='globe.studio.test'){if(typeof opts==='function'){callback=opts;opts={};} process.nextTick(()=>opts?.all?callback(null,[{address:'127.0.0.1',family:4}]):callback(null,'127.0.0.1',4));}else{return original.apply(this,arguments);}};\n") |
| 364 | executable = local / "model-fixture" |
| 365 | executable.write_text("#!/bin/sh\nexec " + str(node) + " " + str(args.agent_dir / "fake-cli.mjs") + ' "$@"\n') |
| 366 | executable.chmod(0o700) |
| 367 | log = local / "agent.log" |
| 368 | environment = {**os.environ, "HOME": str(home), "CODEX_HOME": str(home / "codex"), "CLAUDE_CONFIG_DIR": str(home / "claude"), |
| 369 | "PATH": str(node.parent) + ":" + os.environ.get("PATH", ""), "NODE_OPTIONS": "--require " + str(resolver), "NODE_EXTRA_CA_CERTS": str(args.agent_ca)} |
| 370 | output = queue.Queue() |
| 371 | with log.open("w") as stderr: |
| 372 | process = subprocess.Popen([str(node), str(args.agent_dir / "dist/agent.js"), "run", "--server", args.agent_origin, |
| 373 | "--name", marker, "--data-dir", str(local / "identity"), "--allow-root", str(allowed), |
| 374 | "--codex-bin", str(executable), "--claude-bin", str(executable)], env=environment, stdout=subprocess.PIPE, stderr=stderr, text=True) |
| 375 | reader = threading.Thread(target=lambda: [output.put(line) for line in process.stdout], daemon=True) |
| 376 | reader.start() |
| 377 | try: |
| 378 | deadline = time.monotonic() + 15 |
| 379 | code = None |
| 380 | while not code and time.monotonic() < deadline: |
| 381 | line = output.get(timeout=15) |
| 382 | match = re.search(r"code ([A-F0-9]{5}-[A-F0-9]{5})", line) |
| 383 | if match: |
| 384 | code = match[1] |
| 385 | assert code, "Original agent did not begin pairing" |
| 386 | machine, _ = http("/api/mcp/relay/pair", "POST", {"code": code}, actor=names[0]) |
| 387 | deadline = time.monotonic() + 15 |
| 388 | while time.monotonic() < deadline: |
| 389 | overview, _ = http("/api/mcp", actor=names[0]) |
| 390 | if any(m["id"] == machine["id"] and m["online"] for m in overview["machines"]): |
| 391 | break |
| 392 | time.sleep(.2) |
| 393 | else: |
| 394 | raise AssertionError("Original agent did not connect") |
| 395 | credential_file = local / "identity/agent.json" |
| 396 | assert credential_file.stat().st_mode & 0o777 == 0o600 |
| 397 | native_key = key(names[0], [machine["id"]], write=True) |
| 398 | assert command(native_key["key"], machine["id"])["result"] == {"threads": [], "errors": []} |
| 399 | for provider in ["codex", "claude"]: |
| 400 | started = command(native_key["key"], machine["id"], "start_thread", {"provider": provider, "cwd": str(allowed), "message": "owned fixture"})["result"] |
| 401 | thread = started["thread_id"] |
| 402 | time.sleep(.05) |
| 403 | read_fields = {"provider": provider, "thread_id": thread} |
| 404 | transcript = command(native_key["key"], machine["id"], "read_thread", read_fields)["result"] |
| 405 | assert "fixture reply" in json.dumps(transcript) |
| 406 | command(native_key["key"], machine["id"], "send_message", {**read_fields, "message": "slow"}) |
| 407 | if provider == "claude": |
| 408 | live = command(native_key["key"], machine["id"], "read_thread", read_fields)["result"] |
| 409 | command(native_key["key"], machine["id"], "interrupt_thread", {**read_fields, "expected_turn_id": live["active_turn_id"]}) |
| 410 | command(native_key["key"], machine["id"], "start_thread", {"provider": "codex", "cwd": str(home), "message": "outside allowed directory"}, status=400) |
| 411 | http("/api/mcp/relay/machines/" + machine["id"], "DELETE", actor=names[0], status=204) |
| 412 | assert process.wait(timeout=10) != 0 |
| 413 | original_agent_checks = {"tls_pairing_and_outbound_wss": True, "credential_mode_0600": True, "codex_and_claude_owned_session_roundtrip": True, |
| 414 | "session_control_and_expected_turn": True, "local_directory_allowlist": True, "unlink_rejects_reconnect": True} |
| 415 | finally: |
| 416 | if process.poll() is None: |
| 417 | process.terminate() |
| 418 | process.wait(timeout=10) |
| 419 | process.stdout.close() |
| 420 | reader.join(timeout=5) |
| 421 | result = {"native_pairing_single_use_and_expiry": True, "existing_agents_wire_protocol": True, "ordinary_realm_user_consent": True, |
| 422 | "machine_and_cross_catalog_isolation": True, "explicit_control_scope": True, "single_machine_writes": True, |
| 423 | "per_connection_target_selection": True, "native_mcp_tools": True, "owned_machine_stream": True, "eight_pending_limit": True, "unicode_command_frame_budget": True, |
| 424 | "heartbeat_and_timeout_unknown_outcome": True, "disconnect_unknown_outcome_no_replay": True, |
| 425 | "hash_only_credentials": True, "restart_preserves_pairing_keys_and_oauth": True, |
| 426 | "revoke_and_unlink_immediate": True, "origin_host_and_malformed_frames_refused": True, |
| 427 | "timeout_seconds": round(elapsed, 2)} |
| 428 | if original_agent_checks is not None: |
| 429 | result["original_node_agent"] = original_agent_checks |
| 430 | finally: |
| 431 | for agent in agents: |
| 432 | agent.close() |
| 433 | keycloak = Keycloak("keycloak.studio.test", importlib.import_module("dashboard-run").secret("get", "keycloak", "password"), attempts=1) |
| 434 | for name in names: |
| 435 | try: |
| 436 | overview, _ = http("/api/mcp", actor=name) |
| 437 | for connection in overview["connections"]: |
| 438 | http("/api/mcp/connections/" + connection["id"], "DELETE", actor=name, status=204) |
| 439 | for machine in overview["machines"]: |
| 440 | http("/api/mcp/relay/machines/" + machine["id"], "DELETE", actor=name, status=204) |
| 441 | finally: |
| 442 | for identity in keycloak.request("/admin/realms/master/users?username=" + name + "&exact=true"): |
| 443 | keycloak.request("/admin/realms/master/users/" + identity["id"], "DELETE") |
| 444 | assert not keycloak.request("/admin/realms/master/users?username=" + name + "&exact=true") |
| 445 | result["owned_users_and_machines_removed"] = True |
| 446 | if args.output: |
| 447 | args.output.write_text(json.dumps(result, indent=2) + "\n") |
| 448 | print(json.dumps(result)) |
| 449 | |
| 450 | |
| 451 | if __name__ == "__main__": |
| 452 | main() |