| 1 | #!/usr/bin/env python3 |
| 2 | import argparse |
| 3 | import base64 |
| 4 | import hashlib |
| 5 | import importlib |
| 6 | import json |
| 7 | from pathlib import Path |
| 8 | import subprocess |
| 9 | import sys |
| 10 | import time |
| 11 | import urllib.error |
| 12 | import urllib.parse |
| 13 | import urllib.request |
| 14 | import uuid |
| 15 | |
| 16 | |
| 17 | class NoRedirect(urllib.request.HTTPRedirectHandler): |
| 18 | def redirect_request(self, request, fp, code, message, headers, newurl): |
| 19 | return None |
| 20 | |
| 21 | |
| 22 | def 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 | |
| 234 | if __name__ == "__main__": |
| 235 | main() |