| 1 | #!/usr/bin/env python3 |
| 2 | import argparse |
| 3 | from concurrent.futures import ThreadPoolExecutor |
| 4 | import json |
| 5 | import os |
| 6 | from pathlib import Path |
| 7 | import shutil |
| 8 | import subprocess |
| 9 | import sys |
| 10 | import tempfile |
| 11 | import time |
| 12 | import urllib.error |
| 13 | import urllib.request |
| 14 | import uuid |
| 15 | |
| 16 | import release |
| 17 | |
| 18 | |
| 19 | def main(): |
| 20 | parser = argparse.ArgumentParser() |
| 21 | parser.add_argument("socket", type=Path) |
| 22 | parser.add_argument("--image", required=True) |
| 23 | parser.add_argument("--output", type=Path) |
| 24 | args = parser.parse_args() |
| 25 | suffix = uuid.uuid4().hex[:12] |
| 26 | proof = uuid.uuid4().hex + uuid.uuid4().hex |
| 27 | stage_id = "boundary-preview-" + suffix |
| 28 | stage_file = release.STATE / "stages" / (stage_id + ".json") |
| 29 | service_id = "boundary-service-" + suffix |
| 30 | managed_file = release.STATE / "managed-jobs.json" |
| 31 | managed_bytes = managed_file.read_bytes() |
| 32 | container = "studio-dashboard-deploy-test-" + suffix |
| 33 | current = release.current_release() |
| 34 | history_bytes = release.HISTORY.read_bytes() |
| 35 | stages_before = {path.name for path in (release.STATE / "stages").glob("*.json")} |
| 36 | client = Path(__file__).with_name("dashboard-host-vm-test.py") |
| 37 | |
| 38 | def shell(*argv): |
| 39 | process = subprocess.run(argv, capture_output=True, text=True) |
| 40 | if process.returncode: |
| 41 | raise RuntimeError(process.stderr) |
| 42 | return process.stdout |
| 43 | |
| 44 | def call(operation, status=None, **fields): |
| 45 | payload = json.dumps({"operation": operation, **fields}).encode() + b"\n" |
| 46 | output = subprocess.run([sys.executable, str(client), str(args.socket), "--client"], input=payload, |
| 47 | capture_output=True, check=True) |
| 48 | value = json.loads(output.stdout) |
| 49 | if status: |
| 50 | assert value["status"] == status, value |
| 51 | return value |
| 52 | assert "value" in value, value |
| 53 | return value["value"] |
| 54 | |
| 55 | def http(path, body=None, status=200, section="deploys", groups="infra-admin", method=None): |
| 56 | request = urllib.request.Request("http://127.0.0.1:7076/api/" + section + path, |
| 57 | data=json.dumps(body).encode() if body is not None else None, method=method, |
| 58 | headers={"User-Name": "fixture", "User-Groups": groups, "Studio-Proxy-Token": proof, "Content-Type": "application/json"}) |
| 59 | try: |
| 60 | response = urllib.request.urlopen(request, timeout=40) |
| 61 | except urllib.error.HTTPError as error: |
| 62 | response = error |
| 63 | with response: |
| 64 | payload = response.read() |
| 65 | assert response.status == status, (path, response.status, payload[:300]) |
| 66 | return json.loads(payload) if response.status == 200 and not path.startswith("/runs/") else payload |
| 67 | |
| 68 | def nomad(path, body=None, method=None): |
| 69 | request = urllib.request.Request("http://127.0.0.1:4646/v1/" + path, |
| 70 | data=json.dumps(body).encode() if body is not None else None, method=method, |
| 71 | headers={"X-Nomad-Token": (release.STATE / "nomad.token").read_text().strip(), "Content-Type": "application/json"}) |
| 72 | try: |
| 73 | with urllib.request.urlopen(request, timeout=10) as response: |
| 74 | return None if method == "DELETE" else json.load(response) |
| 75 | except urllib.error.HTTPError as error: |
| 76 | if method == "DELETE" and error.code == 404: |
| 77 | return None |
| 78 | raise |
| 79 | |
| 80 | def done(identity): |
| 81 | deadline = time.monotonic() + 30 |
| 82 | while time.monotonic() < deadline: |
| 83 | result = call("deploy.run", id=identity) |
| 84 | if result["code"] is not None: |
| 85 | return result |
| 86 | time.sleep(.1) |
| 87 | raise AssertionError("deployment worker did not finish") |
| 88 | |
| 89 | def ready(): |
| 90 | deadline = time.monotonic() + 40 |
| 91 | while time.monotonic() < deadline: |
| 92 | try: |
| 93 | return http("") |
| 94 | except OSError: |
| 95 | time.sleep(.1) |
| 96 | raise AssertionError("container deployment view did not start") |
| 97 | |
| 98 | with tempfile.TemporaryDirectory(prefix="studio-deploy-boundary-", dir="/run") as temporary: |
| 99 | data = Path(temporary) / "data" |
| 100 | data.mkdir() |
| 101 | os.chown(data, 65534, 65534) |
| 102 | proxy_token = Path(temporary) / "proxy.token" |
| 103 | proxy_token.write_text(proof) |
| 104 | os.chown(proxy_token, 0, 65534) |
| 105 | proxy_token.chmod(0o440) |
| 106 | fixture_release = None |
| 107 | owned_runs = [] |
| 108 | service_created = False |
| 109 | mount = Path("/srv/staging") / stage_id |
| 110 | try: |
| 111 | assert call("deploy.current") == current |
| 112 | assert call("deploy.history") == json.loads(history_bytes) |
| 113 | assert {stage["id"] + ".json" for stage in call("deploy.stages")} == stages_before |
| 114 | assert call("deploy.release", release=current) |
| 115 | call("deploy.start", action="sh", target=stage_id, status=400) |
| 116 | call("deploy.start", action="destroy", target="../escape", status=400) |
| 117 | call("deploy.start", action="rollback", target=current, status=400) |
| 118 | call("deploy.start", action="destroy", target=stage_id, env={"PATH": "/tmp"}, status=400) |
| 119 | call("deploy.run", id="../../nomad.token", status=400) |
| 120 | stage_file.write_text(json.dumps({"sourceId": "navidrome", "release": current, "ready": True, |
| 121 | "overrides": [], "clone": None, "mount": str(mount), "inputs": {}, "tasks": []})) |
| 122 | candidate = call("deploy.main") |
| 123 | call("deploy.start", action="promote", target=stage_id, status=400) |
| 124 | call("deploy.start", action="deploy", target=stage_id, status=409) |
| 125 | if candidate and candidate["release"] == current: |
| 126 | deployed = call("deploy.start", action="deploy", target=current) |
| 127 | owned_runs.append(deployed["id"]) |
| 128 | result = done(deployed["id"]) |
| 129 | assert result["code"] == 0 and "already running " + current in result["lines"], result |
| 130 | rolled = call("deploy.start", action="rollback", target=str(len(json.loads(history_bytes)))) |
| 131 | owned_runs.append(rolled["id"]) |
| 132 | assert done(rolled["id"])["code"] == 0 |
| 133 | |
| 134 | fixture_release = release.RELEASES / ("boundary-" + suffix) |
| 135 | fixture_release.mkdir() |
| 136 | for name in release.SOURCES: |
| 137 | path = fixture_release / name |
| 138 | if "." in name: |
| 139 | path.write_text("fixture") |
| 140 | else: |
| 141 | path.mkdir() |
| 142 | script = fixture_release / "tools/studio.py" |
| 143 | script.write_text("import time\ntime.sleep(3)\nprint('approved fixture worker')\n") |
| 144 | digest = release.tree_digest(fixture_release) |
| 145 | destination = release.RELEASES / digest[:16] |
| 146 | assert not destination.exists() |
| 147 | fixture_release.rename(destination) |
| 148 | fixture_release = destination |
| 149 | (fixture_release / ".studio-release.json").write_text(json.dumps({"id": digest[:16], "digest": digest, "version": 2})) |
| 150 | value = json.loads(stage_file.read_text()) |
| 151 | stage_file.write_text(json.dumps({**value, "release": digest[:16]})) |
| 152 | active = call("deploy.start", action="destroy", target=stage_id) |
| 153 | owned_runs.append(active["id"]) |
| 154 | call("deploy.start", action="deploy", target=current, status=409) |
| 155 | shell("systemctl", "restart", "studio-host-boundary-test.service") |
| 156 | assert call("deploy.last")["id"] == active["id"] |
| 157 | result = done(active["id"]) |
| 158 | assert result["code"] == 0 and "approved fixture worker" in result["lines"], result |
| 159 | (fixture_release / "tools/studio.py").write_text("print('unapproved changed worker')") |
| 160 | call("deploy.start", action="destroy", target=stage_id, status=502) |
| 161 | assert call("deploy.last")["id"] == active["id"] |
| 162 | stage_file.write_text(json.dumps(value)) |
| 163 | mount.mkdir() |
| 164 | (mount / "fixture").write_text("disposable") |
| 165 | |
| 166 | shell("podman", "run", "-d", "--name=" + container, "--pull=never", "--user=65534:65534", |
| 167 | "--read-only", "--cap-drop=all", "--security-opt=no-new-privileges", "--pids-limit=128", |
| 168 | "--memory=512m", "--cpus=2", "--tmpfs=/tmp:rw,noexec,nosuid,nodev,size=64m", |
| 169 | "--publish=127.0.0.1:7076:7072", "--volume=" + str(args.socket.parent) + ":/run/studio-host:ro", |
| 170 | "--volume=" + str(proxy_token) + ":/run/secrets/dashboard-proxy.token:ro", |
| 171 | "--volume=" + str(data) + ":/data:rw,nosuid,nodev", "--env=STUDIO_DATA_DIR=/data", args.image) |
| 172 | overview = ready() |
| 173 | assert overview["current"] == current and overview["recorded"] |
| 174 | assert len(overview["history"]) == len(json.loads(history_bytes)) |
| 175 | assert any(stage["id"] == stage_id for stage in overview["stages"]) |
| 176 | assert overview["run"]["id"] == active["id"] and overview["run"]["code"] == 0 |
| 177 | with ThreadPoolExecutor(max_workers=40) as workers: |
| 178 | latencies = [] |
| 179 | def overview_check(_): |
| 180 | start = time.monotonic() |
| 181 | assert http("")["current"] == current |
| 182 | return (time.monotonic() - start) * 1000 |
| 183 | latencies = list(workers.map(overview_check, range(1000))) |
| 184 | latencies.sort() |
| 185 | http("/" + service_id + "/promote", {}, status=400, section="services") |
| 186 | http("/" + service_id + "/stop", {}, status=404, section="services") |
| 187 | managed_file.write_text(json.dumps([*json.loads(managed_bytes), service_id])) |
| 188 | service_created = True |
| 189 | nomad("jobs", {"Job": {"ID": service_id, "Name": service_id, "Namespace": "default", "Datacenters": ["*"], "Type": "service", |
| 190 | "Meta": {"studio_secrets": json.dumps([{"name": "fixture", "generated": True, "bytes": 16}, {"name": "outside", "generated": False}])}, |
| 191 | "TaskGroups": [{"Name": "fixture", "Count": 1, "Tasks": [{"Name": "idle", "Driver": "podman", |
| 192 | "Config": {"image": "docker.io/library/alpine:3.22", "command": "/bin/sleep", "args": ["600"]}, |
| 193 | "Resources": {"CPU": 25, "MemoryMB": 32}}]}]}}) |
| 194 | deadline = time.monotonic() + 40 |
| 195 | while time.monotonic() < deadline: |
| 196 | allocations = nomad("job/" + service_id + "/allocations") |
| 197 | deployments = nomad("job/" + service_id + "/deployments") |
| 198 | if (any(allocation["ClientStatus"] == "running" for allocation in allocations) |
| 199 | and any(deployment["Status"] == "successful" for deployment in deployments)): |
| 200 | break |
| 201 | time.sleep(.1) |
| 202 | else: |
| 203 | raise AssertionError("service fixture did not start") |
| 204 | http("/" + service_id + "/stop", {}, status=403, section="services", groups="") |
| 205 | nomad("var/nomad/jobs/" + service_id, {"Namespace": "default", "Path": "nomad/jobs/" + service_id, |
| 206 | "Items": {"fixture": "old", "outside": "keep", "undeclared": "private"}}, "PUT") |
| 207 | prefix = "/" + service_id + "/secrets/" |
| 208 | assert http(prefix + "fixture", section="services")["value"] == "old" |
| 209 | http(prefix + "fixture", status=403, section="services", groups="") |
| 210 | with ThreadPoolExecutor(max_workers=20) as workers: |
| 211 | assert all(workers.map(lambda _: http(prefix + "fixture", section="services")["value"] == "old", range(100))) |
| 212 | http(prefix + "undeclared", status=404, section="services") |
| 213 | http(prefix + "fixture", {"value": "line\nbreak"}, method="PUT", status=400, section="services") |
| 214 | secret_value = "disposable-" + suffix + '-"$();☃' |
| 215 | http(prefix + "fixture", {"value": secret_value}, method="PUT", status=204, section="services") |
| 216 | saved = call("deploy.last") |
| 217 | owned_runs.append(saved["id"]) |
| 218 | assert saved["action"] == "secret-set" and saved["key"] == "fixture" and saved["code"] == 0 |
| 219 | assert secret_value not in "\n".join(call("deploy.run", id=saved["id"])["lines"]) |
| 220 | assert http(prefix + "fixture", section="services")["value"] == secret_value |
| 221 | assert nomad("var/nomad/jobs/" + service_id)["Items"]["outside"] == "keep" |
| 222 | assert not list(Path("/var/lib/studio/host-boundary-test/runs").glob("*.input")) |
| 223 | http(prefix + "fixture/rotate", {}, status=204, section="services") |
| 224 | saved = call("deploy.last") |
| 225 | owned_runs.append(saved["id"]) |
| 226 | rotated = http(prefix + "fixture", section="services")["value"] |
| 227 | assert len(rotated) == 32 and all(c in "0123456789abcdef" for c in rotated) |
| 228 | http(prefix + "outside/rotate", {}, status=502, section="services") |
| 229 | owned_runs.append(call("deploy.last")["id"]) |
| 230 | assert nomad("var/nomad/jobs/" + service_id)["Items"]["outside"] == "keep" |
| 231 | for action, status in [("restart", 204), ("stop", 204), ("start", 502)]: |
| 232 | http("/" + service_id + "/" + action, {}, status=status, section="services") |
| 233 | saved = call("deploy.last") |
| 234 | owned_runs.append(saved["id"]) |
| 235 | assert saved["action"] == action and saved["target"] == service_id |
| 236 | assert http("")["run"]["id"] == saved["id"] |
| 237 | if action == "restart": |
| 238 | allocations = nomad("job/" + service_id + "/allocations") |
| 239 | assert any(allocation["TaskStates"]["idle"]["Restarts"] > 0 for allocation in allocations), allocations |
| 240 | elif action == "stop": |
| 241 | assert nomad("job/" + service_id)["Stop"] is True |
| 242 | else: |
| 243 | assert saved["code"] != 0 |
| 244 | nomad("job/" + service_id + "?purge=true", method="DELETE") |
| 245 | nomad("var/nomad/jobs/" + service_id, method="DELETE") |
| 246 | managed_file.write_bytes(managed_bytes) |
| 247 | service_created = False |
| 248 | http("/openspeedtest/start", {}, status=204, section="services") |
| 249 | saved = call("deploy.last") |
| 250 | owned_runs.append(saved["id"]) |
| 251 | assert saved["action"] == "start" and saved["target"] == "openspeedtest" and saved["code"] == 0 |
| 252 | assert nomad("job/openspeedtest")["Stop"] is False |
| 253 | launched = http("/stages/" + stage_id + "/destroy", {}) |
| 254 | owned_runs.append(launched["id"]) |
| 255 | result = done(launched["id"]) |
| 256 | assert result["code"] == 0, result |
| 257 | assert not stage_file.exists() and not mount.exists() |
| 258 | events = http("/runs/" + launched["id"]).decode() |
| 259 | assert "event: exit\ndata: 0" in events, events |
| 260 | shell("podman", "stop", "--time=8", container) |
| 261 | shell("podman", "start", container) |
| 262 | assert ready()["run"]["id"] == launched["id"] |
| 263 | assert release.current_release() == current and release.HISTORY.read_bytes() == history_bytes |
| 264 | assert {path.name for path in (release.STATE / "stages").glob("*.json")} == stages_before |
| 265 | result = {"fixed_deployment_request_boundary": "passed", "manifest_tampering_refused": "passed", |
| 266 | "main_guards_and_noop_rollback": "passed", "host_restart_preserves_worker": "passed", |
| 267 | "concurrent_deployment_refused": "passed", "container_deployment_metadata": "passed", |
| 268 | "container_stage_destroy": "passed", "container_run_stream_and_restart": "passed", |
| 269 | "container_managed_service_restart_and_stop": "passed", "service_start_failure_preserves_host_run": "passed", |
| 270 | "container_start_approved_existing_stateless_service": "passed", |
| 271 | "container_declared_secret_reads_updates_and_rotation": "passed", "secret_siblings_and_private_input_cleanup": "passed", |
| 272 | "concurrent_named_secret_reads": 100, "secret_read_clients": 20, |
| 273 | "undeclared_secret_read_and_external_rotation_refused": "passed", |
| 274 | "unmanaged_service_and_non_admin_control_refused": "passed", "managed_service_registry_preserved": "passed", |
| 275 | "current_release_history_and_existing_stages_preserved": "passed", |
| 276 | "concurrent_overview_requests": 1000, "clients": 40, |
| 277 | "overview_p95_ms": round(latencies[949], 2), "overview_max_ms": round(latencies[-1], 2)} |
| 278 | if args.output: |
| 279 | args.output.write_text(json.dumps(result, indent=2) + "\n") |
| 280 | print(json.dumps(result)) |
| 281 | except BaseException: |
| 282 | print(subprocess.run(["podman", "logs", "--tail=10", container], capture_output=True, text=True).stderr) |
| 283 | raise |
| 284 | finally: |
| 285 | subprocess.run(["podman", "rm", "--force", container], capture_output=True) |
| 286 | for identity in owned_runs: |
| 287 | subprocess.run(["systemctl", "stop", "studio-run-" + identity], capture_output=True, timeout=15) |
| 288 | if service_created: |
| 289 | nomad("job/" + service_id + "?purge=true", method="DELETE") |
| 290 | nomad("var/nomad/jobs/" + service_id, method="DELETE") |
| 291 | jobs = json.loads(managed_file.read_bytes()) |
| 292 | if jobs == [*json.loads(managed_bytes), service_id]: |
| 293 | managed_file.write_bytes(managed_bytes) |
| 294 | else: |
| 295 | managed_file.write_text(json.dumps([job for job in jobs if job != service_id])) |
| 296 | stage_file.unlink(missing_ok=True) |
| 297 | shutil.rmtree(mount, ignore_errors=True) |
| 298 | if fixture_release: |
| 299 | shutil.rmtree(fixture_release) |
| 300 | |
| 301 | |
| 302 | if __name__ == "__main__": |
| 303 | main() |