| 1 | #!/usr/bin/env python3 |
| 2 | import importlib |
| 3 | import json |
| 4 | import os |
| 5 | from pathlib import Path |
| 6 | import subprocess |
| 7 | import tempfile |
| 8 | import unittest |
| 9 | import uuid |
| 10 | from unittest.mock import patch |
| 11 | |
| 12 | import release |
| 13 | |
| 14 | runs = importlib.import_module("dashboard-run") |
| 15 | host = importlib.import_module("dashboard-host") |
| 16 | |
| 17 | |
| 18 | class DeploymentBoundary(unittest.TestCase): |
| 19 | def setUp(self): |
| 20 | self.temporary = tempfile.TemporaryDirectory(prefix="studio-deploy-boundary-") |
| 21 | self.addCleanup(self.temporary.cleanup) |
| 22 | self.root = Path(self.temporary.name).resolve() |
| 23 | self.state = self.root / "state" |
| 24 | self.state.mkdir() |
| 25 | self.releases = self.root / "releases" |
| 26 | self.releases.mkdir() |
| 27 | self.history = self.state / "deployments.json" |
| 28 | self.history.write_text("[]") |
| 29 | for name, value in [("ROOT", self.root), ("STATE", self.state), ("RELEASES", self.releases), ("HISTORY", self.history)]: |
| 30 | context = patch.object(release, name, value) |
| 31 | context.start() |
| 32 | self.addCleanup(context.stop) |
| 33 | self.host = host.Host("studio-demo") |
| 34 | state_context = patch.object(runs, "HOST_STATE", self.state / "host") |
| 35 | state_context.start() |
| 36 | self.addCleanup(state_context.stop) |
| 37 | self.version = "a" * 16 |
| 38 | self.snapshot = self.releases / self.version |
| 39 | self.snapshot.mkdir() |
| 40 | for name in release.SOURCES: |
| 41 | path = self.snapshot / name |
| 42 | if "." in name: |
| 43 | path.write_text("fixture") |
| 44 | else: |
| 45 | path.mkdir() |
| 46 | (self.snapshot / "tools/studio.py").write_text("print('fixture')") |
| 47 | (self.snapshot / ".studio-release.json").write_text(json.dumps({"id": self.version, "digest": release.tree_digest(self.snapshot), "version": 2, "main": {"commit": "b" * 40, "description": "Fixture main"}})) |
| 48 | (self.state / "stages").mkdir() |
| 49 | self.stage = "fixture-preview" |
| 50 | self.stage_file = self.state / "stages" / (self.stage + ".json") |
| 51 | self.stage_file.write_text(json.dumps({"sourceId": "fixture", "release": self.version, "ready": True, |
| 52 | "clone": None, "mount": "/srv/staging/fixture", "overrides": []})) |
| 53 | self.history.write_text(json.dumps([{"release": self.version, "source": "fixture", "time": 123, "legacy": False}])) |
| 54 | (self.state / "managed-jobs.json").write_text('["fixture"]') |
| 55 | (self.root / "current").symlink_to(self.snapshot) |
| 56 | (self.root / "main").symlink_to(self.snapshot) |
| 57 | |
| 58 | def call(self, operation, **kwargs): |
| 59 | return self.host.handle({"operation": operation, **kwargs}) |
| 60 | |
| 61 | def test_iam_requests_are_bounded_before_any_worker_starts(self): |
| 62 | user = "/users/" + str(uuid.uuid4()) |
| 63 | bad = [ |
| 64 | ("/clients", "GET", None), ("/roles", "POST", {}), |
| 65 | ("/users?username=fixture&exact=true&first=0", "GET", None), |
| 66 | ("/users?username=fixture&exact=true&exact=true", "GET", None), |
| 67 | (user + "/../../clients", "GET", None), (user + "/users", "POST", {}), |
| 68 | (user, "PUT", {"realmRoles": ["admin"]}), |
| 69 | (user, "PUT", {"attributes": {"admin": ["true"]}}), |
| 70 | (user + "/role-mappings/realm", "POST", [{"id": "fixture", "name": "admin"}]), |
| 71 | (user + "/role-mappings/realm", "POST", [{"id": "fixture", "name": []}]), |
| 72 | (user + "/reset-password", "PUT", {"type": "password", "value": "password", "temporary": 1}), |
| 73 | (user + "/execute-actions-email", "PUT", ["arbitrary"]), |
| 74 | (user + "/logout", "POST", {"redirect": "https://elsewhere.invalid"}), |
| 75 | ] |
| 76 | with patch.object(runs.subprocess, "run") as process: |
| 77 | for path, method, body in bad: |
| 78 | with self.subTest(path=path), self.assertRaises(runs.Error): |
| 79 | self.call("iam.request", path=path, method=method, body=body) |
| 80 | for operation in ["get", "set", "rotate"]: |
| 81 | fields = {"service": "keycloak", "key": "password"} |
| 82 | if operation == "set": |
| 83 | fields["value"] = "known-password" |
| 84 | with self.assertRaises(runs.Error) as denied: |
| 85 | self.call("deploy.secret." + operation, **fields) |
| 86 | self.assertEqual(denied.exception.status, 403) |
| 87 | process.assert_not_called() |
| 88 | |
| 89 | def test_requests_cannot_select_commands_paths_or_environment(self): |
| 90 | requests = [{"operation": "deploy.start", "action": "destroy", "target": self.stage, "argv": ["sh"]}, |
| 91 | {"operation": "deploy.start", "action": "deploy", "target": self.version, "env": {"PATH": "/tmp"}}, |
| 92 | {"operation": "deploy.release", "release": "../private"}, |
| 93 | {"operation": "deploy.run", "id": "../../nomad.token"}, |
| 94 | {"operation": "deploy.history", "path": "/etc/shadow"}] |
| 95 | for action, target in [("sh", self.stage), ([], self.stage), ("destroy", "../escape"), |
| 96 | ("destroy", "--root"), ("rollback", self.version), ("rollback", "0"), ("rollback", [])]: |
| 97 | requests.append({"operation": "deploy.start", "action": action, "target": target}) |
| 98 | with patch.object(runs.subprocess, "run") as executed: |
| 99 | for request in requests: |
| 100 | with self.subTest(request=request), self.assertRaises((host.Rejected, runs.Error)): |
| 101 | self.host.handle(request) |
| 102 | executed.assert_not_called() |
| 103 | |
| 104 | def test_commands_come_from_verified_releases_and_history(self): |
| 105 | self.assertEqual(runs.command("destroy", self.stage)[1:], [str(self.snapshot / "tools/studio.py"), "destroy", self.stage]) |
| 106 | self.assertEqual(runs.command("deploy", self.version)[1:], [str(Path(runs.__file__).with_name("release.py")), "deploy", self.version]) |
| 107 | self.assertEqual(runs.command("rollback", "1")[1:], [str(Path(runs.__file__).with_name("release.py")), "rollback", self.version]) |
| 108 | self.assertEqual(runs.command("start", "fixture")[1:], [str(self.snapshot / "tools/studio.py"), "deploy", "fixture"]) |
| 109 | self.assertEqual(runs.command("stop", "fixture"), ["nomad", "job", "stop", "-yes", "fixture"]) |
| 110 | self.assertEqual(runs.command("restart", "fixture"), ["nomad", "job", "restart", "-yes", "fixture"]) |
| 111 | (self.snapshot / "tools/studio.py").write_text("changed") |
| 112 | for action, target in [("destroy", self.stage), ("deploy", self.version), ("rollback", "1"), ("start", "fixture")]: |
| 113 | with self.subTest(action=action), self.assertRaisesRegex(ValueError, "contents changed"): |
| 114 | runs.command(action, target) |
| 115 | |
| 116 | def test_service_control_is_limited_to_managed_names(self): |
| 117 | self.assertEqual(self.call("deploy.managed"), ["fixture"]) |
| 118 | for action in ["start", "stop", "restart"]: |
| 119 | for target, status in [("unmanaged", 404), ("../escape", 400), ("--purge", 400), ([], 400)]: |
| 120 | with self.subTest(action=action, target=target), self.assertRaises(runs.Error) as error: |
| 121 | runs.command(action, target) |
| 122 | self.assertEqual(error.exception.status, status) |
| 123 | (self.state / "managed-jobs.json").write_text('["--purge"]') |
| 124 | with self.assertRaisesRegex(ValueError, "invalid managed"): |
| 125 | self.call("deploy.managed") |
| 126 | for data in [b'{"argv":["sh"]}', b'[1]', b'x' * (runs.MAX_METADATA + 1)]: |
| 127 | (self.state / "managed-jobs.json").write_bytes(data) |
| 128 | with self.subTest(data=data[:30]), self.assertRaises(ValueError): |
| 129 | self.call("deploy.managed") |
| 130 | (self.state / "managed-jobs.json").unlink() |
| 131 | self.assertEqual(self.call("deploy.managed"), []) |
| 132 | |
| 133 | def test_main_deployment_ignores_stage_state_and_rejects_stale_candidates(self): |
| 134 | self.stage_file.write_text(json.dumps({"sourceId": "fixture", "ready": False, "overrides": ["PASSWORD=private"]})) |
| 135 | self.assertEqual(runs.command("deploy", self.version)[-2:], ["deploy", self.version]) |
| 136 | for target in [self.stage, "c" * 16, "../escape", []]: |
| 137 | with self.subTest(target=target), self.assertRaises(runs.Error) as error: |
| 138 | runs.command("deploy", target) |
| 139 | self.assertEqual(error.exception.status, 409) |
| 140 | with self.assertRaises(runs.Error): |
| 141 | runs.command("promote", self.stage) |
| 142 | (self.root / "main").unlink() |
| 143 | self.assertIsNone(self.call("deploy.main")) |
| 144 | with self.assertRaises(runs.Error): |
| 145 | runs.command("deploy", self.version) |
| 146 | |
| 147 | def test_main_metadata_and_deployment_are_checked_again_under_lock(self): |
| 148 | import sys |
| 149 | with patch.object(sys, "argv", ["release.py", "deploy", "c" * 16]), patch.object(release, "activate") as activate: |
| 150 | with self.assertRaisesRegex(ValueError, "main changed"): |
| 151 | release.main() |
| 152 | activate.assert_not_called() |
| 153 | manifest = self.snapshot / ".studio-release.json" |
| 154 | original = json.loads(manifest.read_text()) |
| 155 | for revision in [None, {"commit": "b" * 40, "description": " "}, {"commit": "escape", "description": "main"}]: |
| 156 | manifest.write_text(json.dumps({**original, "main": revision})) |
| 157 | with self.subTest(revision=revision), self.assertRaises(ValueError): |
| 158 | self.call("deploy.main") |
| 159 | |
| 160 | def test_metadata_and_legacy_logs_are_read_without_running_commands(self): |
| 161 | directory = self.snapshot / "service/fixture" |
| 162 | directory.mkdir() |
| 163 | (directory / "service.pkl").write_text("fixture definition") |
| 164 | identity = str(uuid.uuid4()) |
| 165 | directory = self.state / "runs" |
| 166 | directory.mkdir() |
| 167 | (directory / (identity + ".log")).write_bytes(b"legacy\n\xff\n") |
| 168 | (directory / (identity + ".exit")).write_text("0") |
| 169 | with patch.object(runs.subprocess, "run") as executed: |
| 170 | self.assertEqual(self.call("deploy.current"), self.version) |
| 171 | self.assertEqual(len(self.call("deploy.history")), 1) |
| 172 | self.assertEqual(self.call("deploy.release", release=self.version), {"fixture": "fixture definition"}) |
| 173 | self.assertEqual(self.call("deploy.run", id=identity), {"lines": ["legacy", "\ufffd"], "code": 0}) |
| 174 | executed.assert_not_called() |
| 175 | |
| 176 | def test_first_activation_exposes_current_before_starting_host_units(self): |
| 177 | (self.root / "current").unlink() |
| 178 | |
| 179 | def execute(argv, **kwargs): |
| 180 | if argv[:2] == ["nixos-rebuild", "switch"]: |
| 181 | self.assertEqual((self.root / "current").resolve(), self.snapshot) |
| 182 | output = 'job "fixture" {\n' if argv[-1] == "render" else '"leader"' |
| 183 | return subprocess.CompletedProcess(argv, 0, stdout=output, stderr="") |
| 184 | |
| 185 | with patch.object(release.subprocess, "run", side_effect=execute): |
| 186 | self.assertIsNone(release.activate(self.version, initial=True)) |
| 187 | self.assertEqual(release.current_release(), self.version) |
| 188 | |
| 189 | def test_state_reads_reject_symlinks_and_large_files(self): |
| 190 | source = self.root / "fixture" |
| 191 | source.write_bytes(b"x" * 200) |
| 192 | link = self.root / "link" |
| 193 | link.symlink_to(source) |
| 194 | with self.assertRaises(OSError): |
| 195 | runs.read(link) |
| 196 | with self.assertRaisesRegex(ValueError, "too large"): |
| 197 | runs.read(source, 100) |
| 198 | |
| 199 | def test_output_selection_ignores_unrelated_and_old_large_logs(self): |
| 200 | directory = self.state / "runs" |
| 201 | directory.mkdir() |
| 202 | (directory / "unrelated.log").write_bytes(b"x" * (runs.MAX_LOG + 1)) |
| 203 | old = directory / (str(uuid.uuid4()) + ".log") |
| 204 | old.write_bytes(b"x" * (runs.MAX_LOG + 1)) |
| 205 | selected = directory / "new-prod-fixture.log" |
| 206 | selected.write_text("approved deployment\n") |
| 207 | os.utime(selected, (123, 123)) |
| 208 | self.assertEqual(self.call("deploy.output", kind="history", target="1"), ["approved deployment"]) |
| 209 | self.assertIsNone(self.call("deploy.output", kind="stages", target=self.stage)) |
| 210 | |
| 211 | def test_run_state_survives_restart_and_never_forwards_request_environment(self): |
| 212 | with patch.object(runs.subprocess, "run", return_value=subprocess.CompletedProcess([], 0, stdout="", stderr="")) as executed: |
| 213 | started = self.call("deploy.start", action="destroy", target=self.stage) |
| 214 | argv = executed.call_args.args[0] |
| 215 | self.assertIn("--property=RuntimeMaxSec=3600", argv) |
| 216 | self.assertEqual(argv[-3:], [started["id"], "destroy", self.stage]) |
| 217 | self.assertNotIn("--setenv=NOMAD_TOKEN", " ".join(argv)) |
| 218 | root = runs.HOST_STATE |
| 219 | (root / "runs").mkdir(exist_ok=True) |
| 220 | (root / "runs" / (started["id"] + ".log")).write_text("fixture\n") |
| 221 | restored = host.Host("studio-demo") |
| 222 | active = subprocess.CompletedProcess([], 0, stdout="ActiveState=active\nExecMainStatus=0\nLoadState=loaded\n") |
| 223 | with patch.object(runs.subprocess, "run", return_value=active): |
| 224 | self.assertEqual(restored.handle({"operation": "deploy.last"}), {**started, "code": None}) |
| 225 | with self.assertRaises(runs.Error) as error: |
| 226 | restored.handle({"operation": "deploy.start", "action": "deploy", "target": self.version}) |
| 227 | self.assertEqual(error.exception.status, 409) |
| 228 | (root / "runs" / (started["id"] + ".exit")).write_text("0") |
| 229 | self.assertEqual(restored.handle({"operation": "deploy.last"}), {**started, "code": 0}) |
| 230 | (root / "runs" / (started["id"] + ".exit")).unlink() |
| 231 | missing = subprocess.CompletedProcess([], 1, stdout="LoadState=not-found\nActiveState=inactive\nExecMainStatus=0\n") |
| 232 | with patch.object(runs.subprocess, "run", return_value=missing): |
| 233 | self.assertEqual(restored.handle({"operation": "deploy.run", "id": started["id"]})["code"], 1) |
| 234 | |
| 235 | def test_worker_output_and_child_cleanup_are_bounded(self): |
| 236 | identity = str(uuid.uuid4()) |
| 237 | root = runs.HOST_STATE |
| 238 | root.mkdir() |
| 239 | (root / "last-run.json").write_text(json.dumps({"id": identity, "action": "destroy", "target": self.stage})) |
| 240 | import sys |
| 241 | argv = [sys.executable, "-c", "import os,time; os.write(1,b'x'*200000); time.sleep(10)"] |
| 242 | with patch.object(runs, "command", return_value=argv), patch.object(runs, "MAX_LOG", 4096), patch.object(runs.signal, "signal"): |
| 243 | self.assertEqual(runs.worker(identity, "destroy", self.stage), 1) |
| 244 | self.assertEqual((root / "runs" / (identity + ".exit")).read_text(), "1") |
| 245 | log = (root / "runs" / (identity + ".log")).read_text() |
| 246 | self.assertLess(len(log), 8192) |
| 247 | self.assertIn("output exceeded its size limit", log) |
| 248 | |
| 249 | def test_management_token_is_loaded_only_in_the_host_worker(self): |
| 250 | identity = str(uuid.uuid4()) |
| 251 | root = runs.HOST_STATE |
| 252 | root.mkdir() |
| 253 | (root / "last-run.json").write_text(json.dumps({"id": identity, "action": "stop", "target": "fixture"})) |
| 254 | (self.state / "nomad.token").write_text("private-fixture-token\n") |
| 255 | import sys |
| 256 | argv = [sys.executable, "-c", "import os; assert os.environ['NOMAD_TOKEN']==open(" + repr(str(self.state / "nomad.token")) + ").read().strip(); print('stopped')"] |
| 257 | with patch.object(runs, "command", return_value=argv), patch.object(runs.signal, "signal"): |
| 258 | self.assertEqual(runs.worker(identity, "stop", "fixture"), 0) |
| 259 | log = (root / "runs" / (identity + ".log")).read_text() |
| 260 | self.assertIn("stopped", log) |
| 261 | self.assertNotIn("private-fixture-token", log) |
| 262 | |
| 263 | def test_secret_scope_and_compare_and_set_preserve_siblings(self): |
| 264 | (self.state / "nomad.token").write_text("private-fixture-token") |
| 265 | job = {"Meta": {"studio_secrets": json.dumps([{"name": "token", "generated": True, "bytes": 16}])}} |
| 266 | existing = {"ModifyIndex": 7, "Items": {"token": "old", "sibling": "keep", "undeclared": "private"}} |
| 267 | with patch.object(runs, "nomad", side_effect=[job, existing]) as requested: |
| 268 | self.assertEqual(runs.secret("get", "fixture", "token"), "old") |
| 269 | with patch.object(runs, "nomad", return_value=job) as requested, self.assertRaises(runs.Error): |
| 270 | runs.secret("get", "fixture", "undeclared") |
| 271 | requested.assert_called_once_with("job/fixture") |
| 272 | with patch.object(runs, "nomad", side_effect=[job, existing, {}]) as requested, patch.object(runs.subprocess, "run", return_value=subprocess.CompletedProcess([], 0)): |
| 273 | runs.secret("set", "fixture", "token", "new-secret") |
| 274 | self.assertEqual(requested.call_args.args, ("var/nomad/jobs/fixture?cas=7", "PUT", {"Namespace": "default", "Path": "nomad/jobs/fixture", "Items": {**existing["Items"], "token": "new-secret"}})) |
| 275 | with patch.object(runs, "nomad", side_effect=[job, existing, runs.Error(409, "conflict")]), patch.object(runs.subprocess, "run") as restarted, self.assertRaises(runs.Error): |
| 276 | runs.secret("set", "fixture", "token", "new-secret") |
| 277 | restarted.assert_not_called() |
| 278 | |
| 279 | def test_rotations_use_the_recorded_size_and_refuse_external_secrets(self): |
| 280 | (self.state / "nomad.token").write_text("private-fixture-token") |
| 281 | for generated, count in [(False, 16), (True, True), (True, 4097), (True, 0)]: |
| 282 | job = {"Meta": {"studio_secrets": json.dumps([{"name": "token", "generated": generated, "bytes": count}])}} |
| 283 | with patch.object(runs, "nomad", side_effect=[job, None]) as requested, self.assertRaises(runs.Error): |
| 284 | runs.secret("rotate", "fixture", "token") |
| 285 | self.assertEqual(requested.call_count, 2) |
| 286 | job = {"Meta": {"studio_secrets": '[{"name":"token","generated":true,"bytes":16}]'}} |
| 287 | with patch.object(runs, "nomad", side_effect=[job, None, {}]) as requested, patch.object(runs.subprocess, "run", return_value=subprocess.CompletedProcess([], 0)): |
| 288 | runs.secret("rotate", "fixture", "token") |
| 289 | value = requested.call_args.args[2]["Items"]["token"] |
| 290 | self.assertEqual(len(value), 32) |
| 291 | self.assertRegex(value, "^[a-f0-9]+$") |
| 292 | |
| 293 | def test_secret_input_is_private_never_forwarded_and_removed_after_consumption(self): |
| 294 | import stat |
| 295 | import sys |
| 296 | with patch.object(runs.subprocess, "run", return_value=subprocess.CompletedProcess([], 0)) as dispatched: |
| 297 | saved = self.call("deploy.secret.set", service="fixture", key="token", value="private-fixture-value") |
| 298 | self.assertNotIn("private-fixture-value", " ".join(dispatched.call_args.args[0])) |
| 299 | self.assertNotIn("value", saved) |
| 300 | source = runs.HOST_STATE / "runs" / (saved["id"] + ".input") |
| 301 | self.assertEqual(stat.S_IMODE(source.stat().st_mode), 0o600) |
| 302 | argv = [sys.executable, "-c", "import sys; assert len(sys.stdin.read())==21; print('consumed')"] |
| 303 | with patch.object(runs, "command", return_value=argv), patch.object(runs.signal, "signal"): |
| 304 | self.assertEqual(runs.worker(saved["id"], "secret-set", "fixture", "token"), 0) |
| 305 | self.assertFalse(source.exists()) |
| 306 | self.assertNotIn("private-fixture-value", (source.with_suffix(".log")).read_text()) |
| 307 | with patch.object(runs.subprocess, "run", side_effect=OSError("fixture dispatch failed")), self.assertRaises(OSError): |
| 308 | self.call("deploy.secret.set", service="fixture", key="token", value="private-fixture-value") |
| 309 | self.assertFalse(list(source.parent.glob("*.input"))) |
| 310 | |
| 311 | |
| 312 | if __name__ == "__main__": |
| 313 | unittest.main() |