| 1 | #!/usr/bin/env python3 |
| 2 | import fcntl |
| 3 | import json |
| 4 | import os |
| 5 | from pathlib import Path |
| 6 | import re |
| 7 | import selectors |
| 8 | import secrets |
| 9 | import signal |
| 10 | import stat |
| 11 | import subprocess |
| 12 | import sys |
| 13 | import time |
| 14 | import uuid |
| 15 | import urllib.error |
| 16 | import urllib.request |
| 17 | |
| 18 | import release |
| 19 | |
| 20 | |
| 21 | FIELDS = { |
| 22 | "deploy.current": set(), "deploy.main": set(), "deploy.history": set(), "deploy.stages": set(), "deploy.managed": set(), |
| 23 | "deploy.release": {"release"}, "deploy.output": {"kind", "target"}, |
| 24 | "deploy.last": set(), "deploy.run": {"id"}, "deploy.start": {"action", "target"}, |
| 25 | "deploy.secret.get": {"service", "key"}, "deploy.secret.set": {"service", "key", "value"}, |
| 26 | "deploy.secret.rotate": {"service", "key"}, |
| 27 | } |
| 28 | ACTIONS = {"deploy", "destroy", "rollback", "start", "stop", "restart", "secret-set", "secret-rotate"} |
| 29 | MAX_LOG = 1024 * 1024 |
| 30 | MAX_METADATA = 65536 |
| 31 | HOST_STATE = Path(os.environ.get("STUDIO_HOST_STATE_ROOT", "/var/lib/studio/host")) |
| 32 | |
| 33 | |
| 34 | class Error(Exception): |
| 35 | def __init__(self, status, message): |
| 36 | super().__init__(message) |
| 37 | self.status = status |
| 38 | |
| 39 | |
| 40 | def read(path, limit=MAX_METADATA, missing=None): |
| 41 | try: |
| 42 | fd = os.open(path, os.O_RDONLY | os.O_NOFOLLOW) |
| 43 | except FileNotFoundError: |
| 44 | return missing |
| 45 | with os.fdopen(fd, "rb") as source: |
| 46 | if not stat.S_ISREG(os.fstat(source.fileno()).st_mode): |
| 47 | raise ValueError("deployment state must be a regular file") |
| 48 | data = source.read(limit + 1) |
| 49 | if len(data) > limit: |
| 50 | raise ValueError("deployment state is too large") |
| 51 | return data.decode(errors="replace") |
| 52 | |
| 53 | |
| 54 | def history(): |
| 55 | entries = json.loads(read(release.HISTORY, missing="[]")) |
| 56 | if not isinstance(entries, list) or any(not isinstance(entry, dict) for entry in entries): |
| 57 | raise ValueError("invalid deployment history") |
| 58 | return entries |
| 59 | |
| 60 | |
| 61 | def managed(): |
| 62 | jobs = json.loads(read(release.STATE / "managed-jobs.json", missing="[]")) |
| 63 | if not isinstance(jobs, list) or any(not isinstance(job, str) or not release.STAGE_ID.fullmatch(job) for job in jobs): |
| 64 | raise ValueError("invalid managed job list") |
| 65 | return jobs |
| 66 | |
| 67 | |
| 68 | def service(target, key=None): |
| 69 | if not isinstance(target, str) or not release.STAGE_ID.fullmatch(target): |
| 70 | raise Error(400, "Choose a service from the list.") |
| 71 | if target not in managed(): |
| 72 | raise Error(404, "That service is no longer managed. Reload services.") |
| 73 | if key is not None and (not isinstance(key, str) or not re.fullmatch(r"[A-Za-z_][A-Za-z0-9_]{0,127}", key)): |
| 74 | raise Error(400, "Choose a secret from this service's list.") |
| 75 | |
| 76 | |
| 77 | def nomad(path, method="GET", body=None): |
| 78 | token = read(release.STATE / "nomad.token", missing="").strip() |
| 79 | if not token: |
| 80 | raise Error(503, "Nomad credentials are unavailable. Check the host service configuration.") |
| 81 | request = urllib.request.Request("http://127.0.0.1:4646/v1/" + path, method=method, |
| 82 | data=json.dumps(body).encode() if body is not None else None, |
| 83 | headers={"X-Nomad-Token": token, "Content-Type": "application/json"}) |
| 84 | try: |
| 85 | with urllib.request.urlopen(request, timeout=10) as response: |
| 86 | data = response.read(MAX_LOG + 1) |
| 87 | except urllib.error.HTTPError as error: |
| 88 | if error.code == 404 and method == "GET": |
| 89 | return None |
| 90 | if error.code == 409: |
| 91 | raise Error(409, "The secret changed during this update. Reload it before retrying.") from None |
| 92 | raise Error(502, "Nomad refused the secret request. Check its service logs.") from None |
| 93 | if len(data) > MAX_LOG: |
| 94 | raise Error(502, "The secret response is too large. Check the service configuration.") |
| 95 | return json.loads(data) |
| 96 | |
| 97 | |
| 98 | def secret(action, target, key, value=None): |
| 99 | service(target, key) |
| 100 | job = nomad("job/" + target) |
| 101 | specs = json.loads((job or {}).get("Meta", {}).get("studio_secrets", "[]")) |
| 102 | if not isinstance(specs, list): |
| 103 | raise Error(502, "The service's secret names couldn't be read. Deploy its current release again.") |
| 104 | matches = [spec for spec in specs if isinstance(spec, dict) and spec.get("name") == key] |
| 105 | if len(matches) != 1: |
| 106 | raise Error(404, "That secret is unavailable. Reload the service's secret list.") |
| 107 | path = "nomad/jobs/" + target |
| 108 | existing = nomad("var/" + path) |
| 109 | items = dict(existing["Items"]) if existing else {} |
| 110 | if action == "get": |
| 111 | if key not in items: |
| 112 | raise Error(404, "That secret has no value. Set it from the service page.") |
| 113 | return items[key] |
| 114 | if action == "rotate": |
| 115 | count = matches[0].get("bytes") |
| 116 | if matches[0].get("generated") is not True or type(count) is not int or not 1 <= count <= 4096: |
| 117 | raise Error(409, "This secret comes from outside the server. Set it instead of generating it.") |
| 118 | value = secrets.token_hex(count) |
| 119 | if not isinstance(value, str) or not value or len(value.encode()) > 8192 or any(c in value for c in "\r\n\0"): |
| 120 | raise Error(400, "Enter a single-line secret under 8 KiB.") |
| 121 | items[key] = value |
| 122 | nomad("var/" + path + "?cas=" + str(existing["ModifyIndex"] if existing else 0), "PUT", |
| 123 | {"Namespace": "default", "Path": path, "Items": items}) |
| 124 | restarted = subprocess.run(["nomad", "job", "restart", "-yes", target], capture_output=True, text=True, timeout=55, |
| 125 | env={**os.environ, "NOMAD_TOKEN": read(release.STATE / "nomad.token").strip()}) |
| 126 | if restarted.returncode: |
| 127 | raise Error(502, "The secret was saved, but the service couldn't restart. Check its run before retrying.") |
| 128 | print("Updated " + target + "/" + key, flush=True) |
| 129 | |
| 130 | |
| 131 | def stage(target): |
| 132 | if not isinstance(target, str) or not release.STAGE_ID.fullmatch(target): |
| 133 | raise Error(400, "Choose a stage from the list.") |
| 134 | text = read(release.STATE / "stages" / (target + ".json")) |
| 135 | if text is None: |
| 136 | raise Error(404, "That stage is no longer available. Reload deploys.") |
| 137 | return json.loads(text) |
| 138 | |
| 139 | |
| 140 | def entry(target): |
| 141 | if not isinstance(target, str) or not re.fullmatch(r"[1-9][0-9]{0,9}", target): |
| 142 | raise Error(400, "Choose a deployment from history.") |
| 143 | entries = history() |
| 144 | if int(target) > len(entries): |
| 145 | raise Error(404, "That deployment is outside history. Reload deploys.") |
| 146 | return entries[int(target) - 1] |
| 147 | |
| 148 | |
| 149 | def command(action, target, key=None): |
| 150 | if not isinstance(action, str) or action not in ACTIONS: |
| 151 | raise Error(400, "Choose a supported deployment action.") |
| 152 | if action in {"start", "stop", "restart", "secret-set", "secret-rotate"}: |
| 153 | service(target, key) |
| 154 | if action.startswith("secret-"): |
| 155 | if key is None: |
| 156 | raise Error(400, "Choose a secret from this service's list.") |
| 157 | return [sys.executable, str(Path(__file__)), "--secret", action.removeprefix("secret-"), target, key] |
| 158 | if action in {"stop", "restart"}: |
| 159 | return ["nomad", "job", action, "-yes", target] |
| 160 | version = release.current_release() |
| 161 | if version is None: |
| 162 | raise Error(409, "No release is running. Deploy a release before starting this service.") |
| 163 | return [sys.executable, str(release.check_release(version) / "tools/studio.py"), "deploy", target] |
| 164 | if action == "rollback": |
| 165 | saved = entry(target) |
| 166 | version = saved["release"] |
| 167 | release.check_release(version, legacy=saved.get("legacy") is True) |
| 168 | return [sys.executable, str(Path(__file__).with_name("release.py")), "rollback", version] |
| 169 | if action == "deploy": |
| 170 | candidate = release.main_release() |
| 171 | if not candidate or target != candidate["release"]: |
| 172 | raise Error(409, "Main changed or isn't uploaded. Reload deploys before deploying.") |
| 173 | release.check_release(target) |
| 174 | return [sys.executable, str(Path(__file__).with_name("release.py")), "deploy", target] |
| 175 | metadata = stage(target) |
| 176 | version = metadata.get("release") |
| 177 | if not isinstance(version, str): |
| 178 | raise Error(409, "This stage has no release. Stage it again.") |
| 179 | root = release.check_release(version) |
| 180 | return [sys.executable, str(root / "tools/studio.py"), "destroy", target] |
| 181 | |
| 182 | |
| 183 | def last(): |
| 184 | text = read(HOST_STATE / "last-run.json") |
| 185 | if text is None: |
| 186 | return None |
| 187 | return json.loads(text) |
| 188 | |
| 189 | |
| 190 | def run(identity): |
| 191 | if not isinstance(identity, str) or not re.fullmatch(r"[0-9a-f]{8}(?:-[0-9a-f]{4}){3}-[0-9a-f]{12}", identity): |
| 192 | raise Error(400, "Choose a deployment run from the list.") |
| 193 | root = HOST_STATE / "runs" |
| 194 | text = read(root / (identity + ".log"), MAX_LOG + 4096) |
| 195 | if text is None: |
| 196 | legacy = read(release.STATE / "runs" / (identity + ".log"), MAX_LOG) |
| 197 | if legacy is not None: |
| 198 | root, text = release.STATE / "runs", legacy |
| 199 | latest = last() |
| 200 | if text is None and (not latest or latest.get("id") != identity): |
| 201 | raise Error(404, "That run is no longer available.") |
| 202 | code = read(root / (identity + ".exit"), 64) |
| 203 | if code is None: |
| 204 | shown = subprocess.run(["systemctl", "show", "--property=ActiveState,ExecMainStatus,LoadState", "studio-run-" + identity], |
| 205 | capture_output=True, text=True, timeout=5) |
| 206 | state = dict(line.split("=", 1) for line in shown.stdout.splitlines() if "=" in line) |
| 207 | if shown.returncode and state.get("LoadState") != "not-found": |
| 208 | shown.check_returncode() |
| 209 | if state.get("ActiveState") not in {"active", "activating", "deactivating"}: |
| 210 | code = read(root / (identity + ".exit"), 64) |
| 211 | if code is None: |
| 212 | code = int(state.get("ExecMainStatus", 0)) or 1 |
| 213 | return {"lines": [line for line in (text or "").splitlines() if line], "code": int(code) if code is not None else None} |
| 214 | |
| 215 | |
| 216 | def start(action, target, key=None, value=None): |
| 217 | if action == "secret-set" and (not isinstance(value, str) or not value or len(value.encode()) > 8192 or any(c in value for c in "\r\n\0")): |
| 218 | raise Error(400, "Enter a single-line secret under 8 KiB.") |
| 219 | HOST_STATE.mkdir(mode=0o700, parents=True, exist_ok=True) |
| 220 | with (HOST_STATE / "run.lock").open("a") as lock: |
| 221 | fcntl.flock(lock, fcntl.LOCK_EX) |
| 222 | previous = last() |
| 223 | if previous and run(previous["id"])["code"] is None: |
| 224 | raise Error(409, "A deployment is running. Wait for it to finish, then retry.") |
| 225 | command(action, target, key) |
| 226 | identity = str(uuid.uuid4()) |
| 227 | metadata = {"id": identity, "action": action, "target": target} |
| 228 | if key is not None: |
| 229 | metadata["key"] = key |
| 230 | root = HOST_STATE / "runs" |
| 231 | root.mkdir(mode=0o700, parents=True, exist_ok=True) |
| 232 | input_file = root / (identity + ".input") |
| 233 | try: |
| 234 | if action == "secret-set": |
| 235 | with os.fdopen(os.open(input_file, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600), "w") as source: |
| 236 | source.write(value) |
| 237 | pending = HOST_STATE / "last-run.pending" |
| 238 | pending.write_text(json.dumps(metadata) + "\n") |
| 239 | pending.replace(HOST_STATE / "last-run.json") |
| 240 | args = ["systemd-run", "--unit=studio-run-" + identity, "--collect", "--quiet", |
| 241 | "--property=RuntimeMaxSec=3600", "--property=TimeoutStopSec=10", |
| 242 | "--property=MemoryMax=4G", "--property=TasksMax=1024", "--property=CPUWeight=10"] |
| 243 | args.extend("--setenv=" + key + "=" + value for key, value in os.environ.items() if key == "PATH" or key.startswith("STUDIO_")) |
| 244 | subprocess.run([*args, "--", sys.executable, str(Path(__file__)), identity, action, target, *([key] if key is not None else [])], |
| 245 | check=True, capture_output=True, timeout=10) |
| 246 | except BaseException: |
| 247 | input_file.unlink(missing_ok=True) |
| 248 | raise |
| 249 | return metadata |
| 250 | |
| 251 | |
| 252 | def handle(request): |
| 253 | operation = request["operation"] |
| 254 | if operation == "deploy.secret.get": |
| 255 | if not isinstance(request["key"], str): |
| 256 | raise Error(400, "Choose a secret from this service's list.") |
| 257 | service(request["service"], request["key"]) |
| 258 | process = subprocess.run(["systemd-run", "--pipe", "--wait", "--collect", "--quiet", |
| 259 | "--unit=studio-secret-read-" + str(uuid.uuid4()), "--property=RuntimeMaxSec=25", |
| 260 | "--property=MemoryMax=128M", "--property=TasksMax=8", "--property=ProtectSystem=strict", |
| 261 | "--property=ProtectHome=yes", "--property=NoNewPrivileges=yes", "--property=CapabilityBoundingSet=", |
| 262 | "--property=RestrictAddressFamilies=AF_INET AF_UNIX", "--property=IPAddressDeny=any", |
| 263 | "--property=IPAddressAllow=localhost", "--", sys.executable, |
| 264 | str(Path(__file__)), "--secret", "get", request["service"], request["key"]], |
| 265 | capture_output=True, text=True, timeout=30, check=True) |
| 266 | response = json.loads(process.stdout) |
| 267 | if "error" in response: |
| 268 | raise Error(response["status"], response["error"]) |
| 269 | return response["value"] |
| 270 | if operation in {"deploy.secret.set", "deploy.secret.rotate"}: |
| 271 | return start("secret-" + operation.rpartition(".")[2], request["service"], request["key"], request.get("value")) |
| 272 | if operation == "deploy.current": |
| 273 | return release.current_release() |
| 274 | if operation == "deploy.main": |
| 275 | return release.main_release() |
| 276 | if operation == "deploy.history": |
| 277 | return history() |
| 278 | if operation == "deploy.managed": |
| 279 | return managed() |
| 280 | if operation == "deploy.stages": |
| 281 | values = [] |
| 282 | for path in (release.STATE / "stages").glob("*.json"): |
| 283 | metadata = stage(path.stem) |
| 284 | values.append({"id": path.stem, "service": metadata["sourceId"], "release": metadata.get("release"), |
| 285 | "ready": metadata.get("ready", True), "created": path.stat().st_mtime, |
| 286 | "overrides": [value.partition("=")[0] for value in metadata.get("overrides", [])], |
| 287 | "clone": metadata.get("clone"), "mount": metadata.get("mount")}) |
| 288 | return values |
| 289 | if operation == "deploy.release": |
| 290 | version = request["release"] |
| 291 | if not isinstance(version, str) or not release.RELEASE_ID.fullmatch(version): |
| 292 | raise Error(400, "Choose a release from deployment history.") |
| 293 | root = release.RELEASES / version / "service" |
| 294 | if root.parent.is_symlink(): |
| 295 | raise ValueError("release directory is a symlink") |
| 296 | if not root.is_dir(): |
| 297 | return None |
| 298 | values = {} |
| 299 | for path in root.iterdir(): |
| 300 | if path.is_symlink(): |
| 301 | raise ValueError("release service directory is a symlink") |
| 302 | if path.is_dir(): |
| 303 | text = read(path / "service.pkl") |
| 304 | if text is not None: |
| 305 | values[path.name] = text |
| 306 | return values |
| 307 | if operation == "deploy.output": |
| 308 | kind, target = request["kind"], request["target"] |
| 309 | if kind == "stages": |
| 310 | value = stage(target) |
| 311 | elif kind == "history": |
| 312 | value = entry(target) |
| 313 | else: |
| 314 | raise Error(400, "Choose a stage or deployment from history.") |
| 315 | found = [] |
| 316 | for path in (release.STATE / "runs").glob("*.log"): |
| 317 | name = path.name |
| 318 | if kind == "stages": |
| 319 | if not name.endswith("-stage-" + value["sourceId"] + ".log"): |
| 320 | continue |
| 321 | text = read(path, MAX_LOG) |
| 322 | if "stage=" + target in text.splitlines(): |
| 323 | found.append((name, 0, text)) |
| 324 | else: |
| 325 | source, version = value["source"], value["release"] |
| 326 | selected = name.endswith(("-prod-" + source + ".log", "-prod-" + version + ".log", "-rollback-" + version + ".log")) |
| 327 | if not selected and not re.fullmatch(r"[0-9a-f-]{36}\.log", name): |
| 328 | continue |
| 329 | delta = abs(path.stat().st_mtime - value["time"]) |
| 330 | if delta >= 600: |
| 331 | continue |
| 332 | text = read(path, MAX_LOG) |
| 333 | first = next(iter(text.splitlines()), "") |
| 334 | if selected or first.endswith(" deploy " + version) or first.endswith(" promote " + source) or first.endswith(" rollback " + version): |
| 335 | found.append((name, delta, text)) |
| 336 | if kind == "history": |
| 337 | for path in (HOST_STATE / "runs").glob("*.log"): |
| 338 | delta = abs(path.stat().st_mtime - value["time"]) |
| 339 | if delta >= 600: |
| 340 | continue |
| 341 | text = read(path, MAX_LOG + 4096) |
| 342 | first = next(iter(text.splitlines()), "") |
| 343 | if first.endswith(" deploy " + value["release"]) or first.endswith(" promote " + value["source"]) or first.endswith(" rollback " + value["release"]): |
| 344 | found.append((path.name, delta, text)) |
| 345 | found.sort(key=lambda item: item[0] if kind == "stages" else item[1], reverse=kind == "stages") |
| 346 | return [line for line in found[0][2].splitlines() if line] if found else None |
| 347 | if operation == "deploy.last": |
| 348 | value = last() |
| 349 | return {**value, "code": run(value["id"])["code"]} if value else None |
| 350 | if operation == "deploy.run": |
| 351 | return run(request["id"]) |
| 352 | return start(request["action"], request["target"]) |
| 353 | |
| 354 | |
| 355 | def worker(identity, action, target, key=None): |
| 356 | if not re.fullmatch(r"[0-9a-f]{8}(?:-[0-9a-f]{4}){3}-[0-9a-f]{12}", identity): |
| 357 | raise ValueError("incorrect deployment run ID") |
| 358 | saved = last() |
| 359 | expected = {"id": identity, "action": action, "target": target} |
| 360 | if key is not None: |
| 361 | expected["key"] = key |
| 362 | if saved != expected: |
| 363 | raise ValueError("deployment run is outside host state") |
| 364 | root = HOST_STATE / "runs" |
| 365 | root.mkdir(mode=0o700, parents=True, exist_ok=True) |
| 366 | os.umask(0o077) |
| 367 | code = 1 |
| 368 | def interrupted(signum, frame): |
| 369 | raise RuntimeError("Deployment stopped before completion.") |
| 370 | signal.signal(signal.SIGTERM, interrupted) |
| 371 | with (root / (identity + ".log")).open("xb", buffering=0) as output: |
| 372 | try: |
| 373 | argv = command(action, target, key) |
| 374 | output.write(("$ " + " ".join(argv) + "\n").encode()) |
| 375 | environment = None |
| 376 | if action in {"stop", "restart"}: |
| 377 | environment = {**os.environ, "NOMAD_TOKEN": read(release.STATE / "nomad.token", missing="").strip()} |
| 378 | if not environment["NOMAD_TOKEN"]: |
| 379 | raise RuntimeError("Nomad credentials are unavailable. Check the host service configuration.") |
| 380 | input_file = root / (identity + ".input") |
| 381 | with subprocess.Popen(argv, stdin=subprocess.PIPE if action == "secret-set" else subprocess.DEVNULL, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, |
| 382 | start_new_session=True, env=environment) as process: |
| 383 | try: |
| 384 | if action == "secret-set": |
| 385 | process.stdin.write(read(input_file, 8192).encode()) |
| 386 | process.stdin.close() |
| 387 | input_file.unlink() |
| 388 | total = 0 |
| 389 | deadline = time.monotonic() + 3500 |
| 390 | with selectors.DefaultSelector() as selector: |
| 391 | selector.register(process.stdout, selectors.EVENT_READ) |
| 392 | while selector.get_map(): |
| 393 | remaining = deadline - time.monotonic() |
| 394 | if remaining <= 0: |
| 395 | raise TimeoutError("Deployment exceeded its time limit.") |
| 396 | for key, _ in selector.select(remaining): |
| 397 | chunk = os.read(key.fd, 65536) |
| 398 | if not chunk: |
| 399 | selector.unregister(key.fileobj) |
| 400 | continue |
| 401 | total += len(chunk) |
| 402 | if total > MAX_LOG: |
| 403 | raise RuntimeError("Deployment output exceeded its size limit.") |
| 404 | output.write(chunk) |
| 405 | code = process.wait(timeout=max(.001, deadline - time.monotonic())) |
| 406 | except BaseException: |
| 407 | try: |
| 408 | os.killpg(process.pid, signal.SIGKILL) |
| 409 | except ProcessLookupError: |
| 410 | pass |
| 411 | raise |
| 412 | except Exception as error: |
| 413 | output.write((str(error) + "\n").encode()) |
| 414 | finally: |
| 415 | (root / (identity + ".input")).unlink(missing_ok=True) |
| 416 | output.write(("exit " + str(code) + "\n").encode()) |
| 417 | pending = root / (identity + ".exit.tmp") |
| 418 | pending.write_text(str(code)) |
| 419 | pending.replace(root / (identity + ".exit")) |
| 420 | return code |
| 421 | |
| 422 | |
| 423 | if __name__ == "__main__": |
| 424 | if len(sys.argv) == 5 and sys.argv[1] == "--secret" and sys.argv[2] in {"get", "set", "rotate"}: |
| 425 | action, target, key = sys.argv[2:] |
| 426 | try: |
| 427 | result = secret(action, target, key, sys.stdin.read(8193) if action == "set" else None) |
| 428 | except Error as error: |
| 429 | if action != "get": |
| 430 | raise SystemExit(str(error)) |
| 431 | result = {"error": str(error), "status": error.status} |
| 432 | else: |
| 433 | result = {"value": result} |
| 434 | if action == "get": |
| 435 | print(json.dumps(result)) |
| 436 | raise SystemExit(0) |
| 437 | if len(sys.argv) not in {4, 5}: |
| 438 | raise SystemExit("Expected a run ID, action, and target") |
| 439 | raise SystemExit(worker(*sys.argv[1:])) |