1#!/usr/bin/env python3
2import argparse
3from concurrent.futures import ThreadPoolExecutor
4import json
5import os
6from pathlib import Path
7import shutil
8import subprocess
9import sys
10import tempfile
11import time
12import urllib.error
13import urllib.request
14import uuid
15
16import release
17
18
19def 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
302if __name__ == "__main__":
303 main()