1#!/usr/bin/env python3
2import argparse
3import base64
4import hashlib
5import importlib
6import json
7from pathlib import Path
8import subprocess
9import sys
10import time
11import urllib.error
12import urllib.parse
13import urllib.request
14import uuid
15
16
17class NoRedirect(urllib.request.HTTPRedirectHandler):
18 def redirect_request(self, request, fp, code, message, headers, newurl):
19 return None
20
21
22def main():
23 parser = argparse.ArgumentParser()
24 parser.add_argument("--url", required=True)
25 parser.add_argument("--proof-file", type=Path, required=True)
26 parser.add_argument("--restart-unit")
27 parser.add_argument("--output", type=Path)
28 args = parser.parse_args()
29 if args.output:
30 args.output.unlink(missing_ok=True)
31 origin = "https://globe.studio.test"
32 resource = origin + "/mcp/observability"
33 opener = urllib.request.build_opener(NoRedirect)
34 proof = args.proof_file.read_text().strip()
35 sys.path.insert(0, str(Path(__file__).resolve().parent.parent / "service/keycloak"))
36 from api import Keycloak
37 keycloak = Keycloak("keycloak.studio.test", importlib.import_module("dashboard-run").secret("get", "keycloak", "password"), attempts=1)
38 marker = "mcp-fixture-" + uuid.uuid4().hex
39 names = [marker + "-one", marker + "-two"]
40 identities = {}
41
42 def http(path, method="GET", body=None, status=200, actor=None, groups="infra-admin", token=None, form=False, headers=None):
43 fields = {"Host": "globe.studio.test", "Content-Type": "application/x-www-form-urlencoded" if form else "application/json"}
44 if actor is not None:
45 fields.update({"Studio-Proxy-Token": proof, "User-Name": actor, "User-Groups": groups, "Origin": origin})
46 if token is not None:
47 fields.update({"Authorization": "Bearer " + token, "Accept": "application/json, text/event-stream", "MCP-Protocol-Version": "2025-11-25"})
48 fields.update(headers or {})
49 encoded = urllib.parse.urlencode(body).encode() if form else json.dumps(body).encode() if body is not None else None
50 request = urllib.request.Request(args.url + path, method=method, data=encoded, headers=fields)
51 try:
52 response = opener.open(request, timeout=70)
53 except urllib.error.HTTPError as error:
54 response = error
55 with response:
56 content = response.read()
57 assert response.status == status, (path, response.status, content[:300])
58 if content and response.headers.get("Content-Type", "").startswith("application/json"):
59 value = json.loads(content)
60 else:
61 value = content.decode()
62 return value, response.headers
63
64 def consent(client, actor, chosen, scope="observability:read offline_access"):
65 verifier = uuid.uuid4().hex + uuid.uuid4().hex
66 challenge = base64.urlsafe_b64encode(hashlib.sha256(verifier.encode()).digest()).decode().rstrip("=")
67 request = {"response_type": "code", "client_id": client["client_id"], "redirect_uri": client["redirect_uris"][0],
68 "code_challenge_method": "S256", "code_challenge": challenge, "resource": resource, "scope": scope, "state": marker}
69 _, headers = http("/oauth/authorize?" + urllib.parse.urlencode(request), status=302)
70 pending = urllib.parse.urlsplit(headers["Location"]).path.removeprefix("/connect/")
71 path = "/api/mcp/consent/" + pending
72 details, _ = http(path, actor=actor)
73 assert details["client"] == client["client_name"]
74 http(path, actor=names[1] if actor == names[0] else names[0], status=403)
75 http(path, "POST", {"resources": ["outside-grant"]}, actor=actor, status=403)
76 http(path, "POST", {"resources": chosen}, actor=actor, headers={"Origin": "https://other.invalid"}, status=403)
77 result, _ = http(path, "POST", {"resources": chosen}, actor=actor)
78 query = urllib.parse.parse_qs(urllib.parse.urlsplit(result["redirect"]).query)
79 assert query["state"] == [marker] and query["iss"] == [origin + "/"]
80 http(path, actor=actor, status=404)
81 return {"grant_type": "authorization_code", "client_id": client["client_id"], "redirect_uri": client["redirect_uris"][0],
82 "code_verifier": verifier, "code": query["code"][0], "resource": resource}
83
84 def rpc(token, method, params=None):
85 result, _ = http("/mcp/observability", "POST", {"jsonrpc": "2.0", "id": 1, "method": method,
86 **({"params": params} if params is not None else {})}, token=token)
87 assert "error" not in result, (method, result)
88 return result["result"]
89
90 def ingest(service, path, body, content_type="application/json"):
91 readonly = Path("/var/lib/studio/dashboard.token").read_text().strip()
92 request = urllib.request.Request("http://127.0.0.1:4646/v1/service/" + service, headers={"X-Nomad-Token": readonly})
93 with urllib.request.urlopen(request, timeout=10) as response:
94 upstream = json.load(response)[0]
95 request = urllib.request.Request(f"http://{upstream['Address']}:{upstream['Port']}" + path,
96 data=body, headers={"Content-Type": content_type})
97 with urllib.request.urlopen(request, timeout=10) as response:
98 assert response.status == 200
99
100 try:
101 role = keycloak.request("/admin/realms/master/roles/infra-admin")
102 for name in names:
103 result = keycloak.request("/admin/realms/master/users", "POST", {"username": name, "enabled": True}, full=True)
104 identities[name] = result["id"]
105 keycloak.request("/admin/realms/master/users/" + result["id"] + "/role-mappings/realm", "POST", [role])
106 metadata, _ = http("/.well-known/oauth-protected-resource/mcp/observability")
107 assert metadata["resource"] == resource and metadata["authorization_servers"] == [origin + "/"]
108 metadata, _ = http("/.well-known/oauth-authorization-server")
109 assert metadata["issuer"] == origin + "/" and metadata["code_challenge_methods_supported"] == ["S256"]
110 _, headers = http("/mcp/observability", status=401, headers={"User-Name": names[0], "User-Groups": "infra-admin"})
111 assert headers["WWW-Authenticate"].startswith("Bearer resource_metadata=")
112 http("/api/mcp", status=403, headers={"User-Name": names[0], "User-Groups": "infra-admin"})
113 overview, _ = http("/api/mcp", actor=names[0])
114 assert not overview["connections"] and overview["catalogs"][0]["endpoint"] == resource
115 readonly = Path("/var/lib/studio/dashboard.token").read_text().strip()
116 request = urllib.request.Request("http://127.0.0.1:4646/v1/jobs", headers={"X-Nomad-Token": readonly})
117 with urllib.request.urlopen(request, timeout=10) as response:
118 jobs = json.load(response)
119 aliases = {}
120 for job in jobs:
121 request = urllib.request.Request("http://127.0.0.1:4646/v1/job/" + urllib.parse.quote(job["ID"], safe=""), headers={"X-Nomad-Token": readonly})
122 with urllib.request.urlopen(request, timeout=10) as response:
123 aliases[job["ID"]] = json.load(response).get("Meta", {}).get("studio_trace_service")
124 services = [job for job, alias in aliases.items() if alias and list(aliases.values()).count(alias) == 1]
125 assert len(services) >= 2, "need two trace-enabled services for scoped native fixture"
126 allowed, denied = services[:2]
127 client, _ = http("/oauth/register", "POST", {"client_name": marker, "redirect_uris": ["http://127.0.0.1:29999/callback"],
128 "token_endpoint_auth_method": "none"}, status=201)
129 form = consent(client, names[0], [allowed])
130 wrong = {**form, "code_verifier": "wrong"}
131 http("/oauth/token", "POST", wrong, form=True, status=400)
132 tokens, _ = http("/oauth/token", "POST", form, form=True)
133 http("/oauth/token", "POST", form, form=True, status=400)
134 token = tokens["access_token"]
135 assert rpc(token, "initialize", {"protocolVersion": "2025-11-25", "capabilities": {}, "clientInfo": {"name": marker, "version": "1"}})["capabilities"]["tools"] == {}
136 tools = rpc(token, "tools/list")["tools"]
137 assert {tool["name"] for tool in tools} == {"get_logs", "get_traces", "get_trace"}
138 assert all(tool["annotations"]["readOnlyHint"] for tool in tools)
139 refused = rpc(token, "tools/call", {"name": "get_logs", "arguments": {"service": denied}})
140 assert refused["isError"]
141 http("/mcp/observability", "POST", {"jsonrpc": "2.0", "id": 1, "method": "tools/list"}, token=token, headers={"Origin": "https://other.invalid"}, status=403)
142 at = time.time() - 60
143 row = {"_time": str(at), "_msg": marker, "source": "nomad", "job": allowed, "task": "fixture", "stream": "stdout"}
144 ingest("victoria-logs", "/insert/jsonline", (json.dumps(row) + "\n").encode(), "application/stream+json")
145 trace_id = uuid.uuid4().hex
146 spans = []
147 for service, span_id, parent in [(denied, "1234567890abcdef", ""), (allowed, "abcdef1234567890", "1234567890abcdef")]:
148 spans.append({"resource": {"attributes": [{"key": "service.name", "value": {"stringValue": aliases[service]}}]},
149 "scopeSpans": [{"spans": [{"traceId": trace_id, "spanId": span_id, "parentSpanId": parent,
150 "name": marker + ("-foreign-secret" if service == denied else "-allowed"), "kind": 1,
151 "startTimeUnixNano": str(int(at * 1e9)), "endTimeUnixNano": str(int((at + .01) * 1e9)),
152 "status": {"code": 2 if service == denied else 1},
153 "attributes": [{"key": "studio.service", "value": {"stringValue": service}}]}]}]})
154 ingest("victoria-traces", "/insert/opentelemetry/v1/traces", json.dumps({"resourceSpans": spans}).encode())
155 deadline = time.monotonic() + 30
156 while True:
157 try:
158 logs = rpc(token, "tools/call", {"name": "get_logs", "arguments": {"service": allowed, "q": marker}})
159 assert not logs.get("isError") and any(row["text"] == marker for row in logs["structuredContent"]["logs"]), logs
160 trace = rpc(token, "tools/call", {"name": "get_trace", "arguments": {"service": allowed, "trace_id": trace_id}})
161 assert not trace.get("isError") and len(trace["structuredContent"]["trace"]["spans"]) == 1, trace
162 traces = rpc(token, "tools/call", {"name": "get_traces", "arguments": {"service": allowed, "q": marker}})
163 assert not traces.get("isError") and any(row["id"] == trace_id and row["spans"] == 1 for row in traces["structuredContent"]["traces"]), traces
164 assert "foreign-secret" not in json.dumps(trace) + json.dumps(traces)
165 break
166 except (AssertionError, urllib.error.HTTPError):
167 if time.monotonic() > deadline:
168 raise
169 time.sleep(1)
170 overview, _ = http("/api/mcp", actor=names[0])
171 grant_id = overview["connections"][0]["id"]
172 other, _ = http("/api/mcp", actor=names[1])
173 assert not other["connections"]
174 http("/api/mcp/connections/" + grant_id, "DELETE", actor=names[1], status=404)
175 if args.restart_unit:
176 subprocess.run(["systemctl", "restart", args.restart_unit], check=True, capture_output=True, timeout=60)
177 rpc(token, "tools/list")
178 refresh = {"grant_type": "refresh_token", "client_id": client["client_id"], "refresh_token": tokens["refresh_token"], "resource": resource}
179 rotated, _ = http("/oauth/token", "POST", refresh, form=True)
180 assert rotated["refresh_token"] != tokens["refresh_token"]
181 http("/oauth/token", "POST", refresh, form=True, status=400)
182 http("/mcp/observability", token=token, status=401)
183 http("/mcp/observability", token=rotated["access_token"], status=401)
184 second = consent(client, names[0], [allowed], "observability:read")
185 tokens, _ = http("/oauth/token", "POST", second, form=True)
186 assert "refresh_token" not in tokens
187 overview, _ = http("/api/mcp", actor=names[0])
188 http("/api/mcp/connections/" + overview["connections"][0]["id"], "DELETE", actor=names[0], status=204)
189 http("/mcp/observability", token=tokens["access_token"], status=401)
190 for change in ["disable", "remove-role", "delete-user"]:
191 exchange = consent(client, names[0], [allowed])
192 tokens, _ = http("/oauth/token", "POST", exchange, form=True)
193 path = "/admin/realms/master/users/" + identities[names[0]]
194 if change == "disable":
195 keycloak.request(path, "PUT", {"enabled": False})
196 elif change == "remove-role":
197 keycloak.request(path + "/role-mappings/realm", "DELETE", [role])
198 else:
199 keycloak.request(path, "DELETE")
200 refresh = {"grant_type": "refresh_token", "client_id": client["client_id"], "refresh_token": tokens["refresh_token"]}
201 rejected, _ = http("/oauth/token", "POST", refresh, form=True, status=400)
202 assert rejected["error"] == "invalid_grant"
203 http("/mcp/observability", token=tokens["access_token"], status=401)
204 if change == "disable":
205 keycloak.request(path, "PUT", {"enabled": True})
206 elif change == "remove-role":
207 keycloak.request(path + "/role-mappings/realm", "POST", [role])
208 result = {"resource_metadata": True, "pkce_code_single_use": True, "shared_realm_identity": True,
209 "consent_owner_and_origin": True, "sdk_initialization_and_tools": True, "explicit_service_grants": True,
210 "native_log_retrieval": True, "native_trace_retrieval": True, "foreign_spans_and_summary_filtered": True,
211 "cross_user_connections_refused": True, "refresh_rotation_and_replay_revocation": True,
212 "dashboard_revocation_immediate": True, "refresh_requires_consent": True,
213 "refresh_honors_account_disable_role_removal_and_deletion": True,
214 "restart_preserves_tokens": bool(args.restart_unit)}
215 finally:
216 for name in names:
217 try:
218 overview, _ = http("/api/mcp", actor=name)
219 for connection in overview["connections"]:
220 http("/api/mcp/connections/" + connection["id"], "DELETE", actor=name, status=204)
221 except (AssertionError, urllib.error.URLError):
222 pass
223 for name in names:
224 for identity in keycloak.request("/admin/realms/master/users?username=" + name + "&exact=true"):
225 keycloak.request("/admin/realms/master/users/" + identity["id"], "DELETE")
226 assert not keycloak.request("/admin/realms/master/users?username=" + name + "&exact=true")
227 result["owned_users_removed"] = True
228 if args.output:
229 args.output.parent.mkdir(parents=True, exist_ok=True)
230 args.output.write_text(json.dumps(result, indent=2) + "\n")
231 print(json.dumps(result))
232
233
234if __name__ == "__main__":
235 main()