| 1 | #!/usr/bin/env python3 |
| 2 | import argparse |
| 3 | from http.server import BaseHTTPRequestHandler, HTTPServer |
| 4 | import json |
| 5 | import os |
| 6 | from pathlib import Path |
| 7 | import socket |
| 8 | import subprocess |
| 9 | import tempfile |
| 10 | import threading |
| 11 | import time |
| 12 | import urllib.error |
| 13 | import urllib.request |
| 14 | import uuid |
| 15 | |
| 16 | import router |
| 17 | |
| 18 | |
| 19 | def main(): |
| 20 | parser = argparse.ArgumentParser() |
| 21 | parser.add_argument("--image", required=True) |
| 22 | parser.add_argument("--output", type=Path) |
| 23 | args = parser.parse_args() |
| 24 | proof = uuid.uuid4().hex + uuid.uuid4().hex |
| 25 | container = "studio-dashboard-routing-test-" + uuid.uuid4().hex[:12] |
| 26 | fixture_id = "studio-routing-fixture-" + uuid.uuid4().hex[:12] |
| 27 | trace_id = uuid.uuid4().hex |
| 28 | redirected = [] |
| 29 | |
| 30 | def shell(*argv): |
| 31 | result = subprocess.run(argv, capture_output=True, text=True, timeout=60) |
| 32 | if result.returncode: |
| 33 | raise AssertionError(result.stderr) |
| 34 | return result.stdout |
| 35 | |
| 36 | class Fixture(BaseHTTPRequestHandler): |
| 37 | def do_GET(self): |
| 38 | if self.path.startswith("/api/v2/"): |
| 39 | self.send_response(302) |
| 40 | self.send_header("Location", f"http://10.88.0.1:{self.server.server_port}/trap") |
| 41 | self.end_headers() |
| 42 | return |
| 43 | if self.path == "/trap": |
| 44 | redirected.append(dict(self.headers)) |
| 45 | self.send_response(200) |
| 46 | self.send_header("Content-Type", "application/json") |
| 47 | self.end_headers() |
| 48 | self.wfile.write(json.dumps({"path": self.path, "host": self.headers.get("Host"), |
| 49 | "proof": self.headers.get("Studio-Proxy-Token"), |
| 50 | "user": self.headers.get("User-Name")}).encode()) |
| 51 | |
| 52 | def log_message(self, *args): |
| 53 | pass |
| 54 | |
| 55 | with tempfile.TemporaryDirectory(prefix="studio-dashboard-routing-", dir="/run") as temporary, HTTPServer(("127.0.0.1", 0), Fixture) as fixture: |
| 56 | root = Path(temporary) |
| 57 | threading.Thread(target=fixture.serve_forever, daemon=True).start() |
| 58 | token = root / "proxy.token" |
| 59 | token.write_text(proof) |
| 60 | readonly = root / "nomad.token" |
| 61 | readonly.write_bytes(Path("/var/lib/studio/dashboard.token").read_bytes()) |
| 62 | for file in [token, readonly]: |
| 63 | os.chown(file, 0, 65534) |
| 64 | file.chmod(0o440) |
| 65 | with socket.socket() as reservation: |
| 66 | for port in range(20000, 32001): |
| 67 | try: |
| 68 | reservation.bind(("0.0.0.0", port)) |
| 69 | break |
| 70 | except OSError: |
| 71 | continue |
| 72 | else: |
| 73 | raise AssertionError("no fixture listener available") |
| 74 | os.environ.update(STUDIO_DOMAIN="studio.test", STUDIO_INTERNAL_PORT=str(port), |
| 75 | STUDIO_DASHBOARD_PORT="7072", STUDIO_PROXY_TOKEN_FILE=str(token)) |
| 76 | real_nomad = router.nomad |
| 77 | |
| 78 | def discovery(path, supplied): |
| 79 | if path in ["/v1/service/routing-fixture", "/v1/service/qbittorrent"]: |
| 80 | return [{"ServiceName": path.rsplit("/", 1)[1], "AllocID": "routing-fixture", "Address": "127.0.0.1", |
| 81 | "Port": fixture.server_port, "Tags": ["caddy-host=routing-fixture.studio.test"]}] |
| 82 | if path == "/v1/allocation/routing-fixture/checks": |
| 83 | return {"ready": {"Status": "failure"}} |
| 84 | value = real_nomad(path, supplied) |
| 85 | if path == "/v1/services": |
| 86 | next(item for item in value if item["Namespace"] == "default")["Services"].append({"ServiceName": "routing-fixture"}) |
| 87 | return value |
| 88 | |
| 89 | router.nomad = discovery |
| 90 | host = "dashboard.internal.studio.test" |
| 91 | content = router.render(Path(router.TOKEN).read_text().strip()) |
| 92 | gateway = content[content.index(host + ":" + str(port) + " {"):] |
| 93 | config = root / "Caddyfile" |
| 94 | config.write_text("{\n admin off\n auto_https disable_redirects\n skip_install_trust\n}\n" + gateway) |
| 95 | config.chmod(0o600) |
| 96 | binary = (Path("/proc") / shell("systemctl", "show", "-P", "MainPID", "caddy").strip() / "exe").resolve(strict=True) |
| 97 | caddy = None |
| 98 | try: |
| 99 | with (root / "caddy.log").open("w") as log: |
| 100 | caddy = subprocess.Popen([str(binary), "run", "--config", str(config), "--adapter", "caddyfile"], |
| 101 | stdout=log, stderr=log, env={**os.environ, "XDG_DATA_HOME": temporary, "XDG_CONFIG_HOME": temporary}) |
| 102 | ca = root / "caddy/pki/authorities/local/root.crt" |
| 103 | deadline = time.monotonic() + 15 |
| 104 | while not ca.exists(): |
| 105 | if caddy.poll() is not None or time.monotonic() >= deadline: |
| 106 | raise AssertionError((root / "caddy.log").read_text()) |
| 107 | time.sleep(.1) |
| 108 | ca.chmod(0o444) |
| 109 | image = json.loads(shell("podman", "image", "inspect", args.image))[0] |
| 110 | python = next(value.split("=", 1)[1] for value in image["Config"]["Env"] if value.startswith("STUDIO_YT_PYTHON=")) |
| 111 | client = f"""import json, ssl, time, urllib.request, urllib.error |
| 112 | base = 'https://{host}:{port}' |
| 113 | proof = open('/proxy.token').read().strip() |
| 114 | nomad = open('/nomad.token').read().strip() |
| 115 | context = ssl.create_default_context(cafile='/ca.crt') |
| 116 | def request(path, supplied=proof, method='GET', body=None, content_type='application/json'): |
| 117 | headers = {{'User-Name': 'forged', 'X-Nomad-Token': nomad, 'Content-Type': content_type}} |
| 118 | if supplied is not None: |
| 119 | headers['Studio-Proxy-Token'] = supplied |
| 120 | query = urllib.request.Request(base + path, data=body, headers=headers, method=method) |
| 121 | try: |
| 122 | with urllib.request.urlopen(query, context=context, timeout=15) as response: |
| 123 | return response.status, response.read() |
| 124 | except urllib.error.HTTPError as error: |
| 125 | return error.code, error.read() |
| 126 | for supplied in [None, '0' * 64, proof[:-1], proof + '0']: |
| 127 | assert request('/nomad/v1/services', supplied)[0] == 403 |
| 128 | status, body = request('/nomad/v1/services') |
| 129 | assert status == 200 and isinstance(json.loads(body), list), (status, body) |
| 130 | assert request('/nomad/v1/var/studio-routing-fixture')[0] == 403 |
| 131 | status, body = request('/nomad/v1/jobs', method='POST', body=json.dumps({{'Job': {{'ID': 'studio-routing-fixture', 'Name': 'studio-routing-fixture', 'Type': 'service', 'Datacenters': ['clover'], 'TaskGroups': [{{'Name': 'fixture', 'Count': 0, 'Tasks': [{{'Name': 'probe', 'Driver': 'podman', 'Config': {{'image': 'alpine:3.22'}}, 'Resources': {{'CPU': 100, 'MemoryMB': 32}}}}]}}]}}}}).encode()) |
| 132 | assert status == 403, (status, body) |
| 133 | assert request('/services/unknown-fixture/health')[0] == 404 |
| 134 | for service in ['victoria-metrics', 'victoria-logs', 'victoria-traces']: |
| 135 | status, body = request('/services/' + service + '/health') |
| 136 | assert status == 200, (service, status, body) |
| 137 | status, body = request('/services/routing-fixture/probe?fixture=1') |
| 138 | result = json.loads(body) |
| 139 | assert status == 200 and result == {{'path': '/probe?fixture=1', 'host': '127.0.0.1:{fixture.server_port}', 'proof': None, 'user': None}}, result |
| 140 | # Keep fixture samples outside VictoriaMetrics' query latency window. |
| 141 | at = time.time() - 60 |
| 142 | metric = 'studio_service_cpu_cores{{service="{fixture_id}"}} 0.125 ' + str(int(at * 1000)) + '\\n' |
| 143 | status, body = request('/services/victoria-metrics/api/v1/import/prometheus', method='POST', body=metric.encode(), content_type='text/plain') |
| 144 | assert status == 204, (status, body) |
| 145 | row = {{'_time': str(at), '_msg': 'Routing fixture', 'source': 'nomad', 'job': '{fixture_id}', 'task': 'probe', 'stream': 'stdout'}} |
| 146 | status, body = request('/services/victoria-logs/insert/jsonline', method='POST', body=(json.dumps(row) + '\\n').encode(), content_type='application/stream+json') |
| 147 | assert status == 200, (status, body) |
| 148 | trace = {{'resourceSpans': [{{'resource': {{'attributes': [{{'key': 'service.name', 'value': {{'stringValue': '{fixture_id}'}}}}]}}, |
| 149 | 'scopeSpans': [{{'spans': [{{'traceId': '{trace_id}', 'spanId': '1234567890abcdef', 'name': 'Routing fixture', 'kind': 1, |
| 150 | 'startTimeUnixNano': str(int(at * 1e9)), 'endTimeUnixNano': str(int((at + .01) * 1e9))}}]}}]}}]}} |
| 151 | status, body = request('/services/victoria-traces/insert/opentelemetry/v1/traces', method='POST', body=json.dumps(trace).encode()) |
| 152 | assert status == 200, (status, body) |
| 153 | print(json.dumps({{'private_bridge': True, 'nonroot': __import__('os').getuid() == 65534, |
| 154 | 'tls_verified': True, 'forged_gateway_refused': 4, 'nomad_metadata': True, |
| 155 | 'nomad_variables_refused': True, 'nomad_writes_refused': True, |
| 156 | 'unknown_service_refused': True, 'actual_telemetry_health': 3, |
| 157 | 'unhealthy_service_reachable': True, 'upstream_headers_scrubbed': True, |
| 158 | 'upstream_path_and_host': True}})) |
| 159 | """ |
| 160 | output = shell("podman", "run", "--rm", "--name", container, "--read-only", "--cap-drop=ALL", |
| 161 | "--security-opt=no-new-privileges", "--memory=256m", "--cpus=1", "--pids-limit=32", |
| 162 | "--add-host=" + host + ":host-gateway", "--entrypoint=" + python, |
| 163 | "--volume=" + str(ca) + ":/ca.crt:ro", "--volume=" + str(token) + ":/proxy.token:ro", |
| 164 | "--volume=" + str(readonly) + ":/nomad.token:ro", args.image, "-c", client) |
| 165 | result = json.loads(output) |
| 166 | assert result["nonroot"] |
| 167 | data = root / "data" |
| 168 | data.mkdir() |
| 169 | os.chown(data, 65534, 65534) |
| 170 | shell("podman", "run", "--detach", "--name", container, "--read-only", "--cap-drop=ALL", |
| 171 | "--security-opt=no-new-privileges", "--memory=512m", "--cpus=2", "--pids-limit=128", |
| 172 | "--add-host=" + host + ":host-gateway", "--publish=127.0.0.1::7072", |
| 173 | "--volume=" + str(ca) + ":/ca.crt:ro", "--volume=" + str(token) + ":/proxy.token:ro", |
| 174 | "--volume=" + str(readonly) + ":/nomad.token:ro", "--volume=" + str(data) + ":/data:rw", |
| 175 | "--env=STUDIO_PROXY_TOKEN_FILE=/proxy.token", "--env=STUDIO_NOMAD_TOKEN_FILE=/nomad.token", |
| 176 | "--env=STUDIO_CA_BUNDLE=/ca.crt", f"--env=STUDIO_INTERNAL_URL=https://{host}:{port}", |
| 177 | "--env=STUDIO_DATA_DIR=/data", "--env=STUDIO_DOMAIN=studio.test", args.image) |
| 178 | info = json.loads(shell("podman", "inspect", container))[0] |
| 179 | app_port = info["NetworkSettings"]["Ports"]["7072/tcp"][0]["HostPort"] |
| 180 | |
| 181 | def dashboard(path): |
| 182 | request = urllib.request.Request("http://127.0.0.1:" + app_port + path, headers={ |
| 183 | "Studio-Proxy-Token": proof, "User-Name": "fixture", "User-Groups": "infra-admin", |
| 184 | }) |
| 185 | with urllib.request.urlopen(request, timeout=20) as response: |
| 186 | return json.load(response) |
| 187 | |
| 188 | deadline = time.monotonic() + 30 |
| 189 | while True: |
| 190 | try: |
| 191 | assert dashboard("/api/me")["name"] == "fixture" |
| 192 | break |
| 193 | except OSError: |
| 194 | if time.monotonic() >= deadline: |
| 195 | raise |
| 196 | time.sleep(.1) |
| 197 | assert isinstance(dashboard("/api/launcher"), list) |
| 198 | deadline = time.monotonic() + 20 |
| 199 | while True: |
| 200 | try: |
| 201 | metrics = dashboard("/api/metrics/service.cpu?range=300&service=" + fixture_id) |
| 202 | assert any(.125 in item["v"] for item in metrics), metrics |
| 203 | logs = dashboard("/api/services/" + fixture_id + "/logs?limit=1") |
| 204 | assert logs and logs[0]["text"] == "Routing fixture", logs |
| 205 | trace = dashboard("/api/traces/" + trace_id) |
| 206 | assert trace["id"] == trace_id and trace["spans"][0]["name"] == "Routing fixture", trace |
| 207 | break |
| 208 | except (AssertionError, urllib.error.HTTPError): |
| 209 | if time.monotonic() >= deadline: |
| 210 | raise |
| 211 | time.sleep(2) |
| 212 | try: |
| 213 | dashboard("/api/seedbox") |
| 214 | raise AssertionError("backend redirect followed") |
| 215 | except urllib.error.HTTPError as error: |
| 216 | assert error.code == 502 and b"302" in error.read() |
| 217 | assert not redirected, redirected |
| 218 | result.update(rust_nomad_metadata=True, rust_actual_metric_ingestion=True, |
| 219 | rust_actual_log_ingestion=True, rust_actual_trace_ingestion=True, |
| 220 | rust_backend_redirect_refused=True) |
| 221 | if args.output: |
| 222 | args.output.parent.mkdir(parents=True, exist_ok=True) |
| 223 | args.output.write_text(json.dumps(result, indent=2) + "\n") |
| 224 | print(json.dumps(result)) |
| 225 | finally: |
| 226 | subprocess.run(["podman", "rm", "--force", container], capture_output=True, timeout=15) |
| 227 | if caddy is not None: |
| 228 | caddy.terminate() |
| 229 | caddy.wait(timeout=10) |
| 230 | fixture.shutdown() |
| 231 | |
| 232 | |
| 233 | if __name__ == "__main__": |
| 234 | main() |