| 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.parse |
| 17 | import urllib.request |
| 18 | |
| 19 | import release |
| 20 | |
| 21 | |
| 22 | FIELDS = { |
| 23 | "iam.request": {"path", "method", "body"}, |
| 24 | "deploy.current": set(), "deploy.main": set(), "deploy.history": set(), "deploy.stages": set(), "deploy.managed": set(), |
| 25 | "deploy.release": {"release"}, "deploy.output": {"kind", "target"}, |
| 26 | "deploy.last": set(), "deploy.run": {"id"}, "deploy.start": {"action", "target"}, |
| 27 | "deploy.secret.get": {"service", "key"}, "deploy.secret.set": {"service", "key", "value"}, |
| 28 | "deploy.secret.rotate": {"service", "key"}, |
| 29 | } |
| 30 | ACTIONS = {"deploy", "destroy", "rollback", "start", "stop", "restart", "secret-set", "secret-rotate"} |
| 31 | MAX_LOG = 1024 * 1024 |
| 32 | MAX_METADATA = 65536 |
| 33 | HOST_STATE = Path(os.environ.get("STUDIO_HOST_STATE_ROOT", "/var/lib/studio/host")) |
| 34 | IAM_ROLES = {"infra-admin", "media", "media-manage"} |
| 35 | IAM_ACTIONS = {"UPDATE_PASSWORD", "VERIFY_EMAIL", "UPDATE_PROFILE", "CONFIGURE_TOTP", "webauthn-register", "webauthn-register-passwordless"} |
| 36 | |
| 37 | |
| 38 | class Error(Exception): |
| 39 | def __init__(self, status, message): |
| 40 | super().__init__(message) |
| 41 | self.status = status |
| 42 | |
| 43 | |
| 44 | def read(path, limit=MAX_METADATA, missing=None): |
| 45 | try: |
| 46 | fd = os.open(path, os.O_RDONLY | os.O_NOFOLLOW) |
| 47 | except FileNotFoundError: |
| 48 | return missing |
| 49 | with os.fdopen(fd, "rb") as source: |
| 50 | if not stat.S_ISREG(os.fstat(source.fileno()).st_mode): |
| 51 | raise ValueError("deployment state must be a regular file") |
| 52 | data = source.read(limit + 1) |
| 53 | if len(data) > limit: |
| 54 | raise ValueError("deployment state is too large") |
| 55 | return data.decode(errors="replace") |
| 56 | |
| 57 | |
| 58 | def history(): |
| 59 | entries = json.loads(read(release.HISTORY, missing="[]")) |
| 60 | if not isinstance(entries, list) or any(not isinstance(entry, dict) for entry in entries): |
| 61 | raise ValueError("invalid deployment history") |
| 62 | return entries |
| 63 | |
| 64 | |
| 65 | def managed(): |
| 66 | jobs = json.loads(read(release.STATE / "managed-jobs.json", missing="[]")) |
| 67 | if not isinstance(jobs, list) or any(not isinstance(job, str) or not release.STAGE_ID.fullmatch(job) for job in jobs): |
| 68 | raise ValueError("invalid managed job list") |
| 69 | return jobs |
| 70 | |
| 71 | |
| 72 | def service(target, key=None): |
| 73 | if not isinstance(target, str) or not release.STAGE_ID.fullmatch(target): |
| 74 | raise Error(400, "Choose a service from the list.") |
| 75 | if target not in managed(): |
| 76 | raise Error(404, "That service is no longer managed. Reload services.") |
| 77 | 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)): |
| 78 | raise Error(400, "Choose a secret from this service's list.") |
| 79 | |
| 80 | |
| 81 | def nomad(path, method="GET", body=None): |
| 82 | token = read(release.STATE / "nomad.token", missing="").strip() |
| 83 | if not token: |
| 84 | raise Error(503, "Nomad credentials are unavailable. Check the host service configuration.") |
| 85 | request = urllib.request.Request("http://127.0.0.1:4646/v1/" + path, method=method, |
| 86 | data=json.dumps(body).encode() if body is not None else None, |
| 87 | headers={"X-Nomad-Token": token, "Content-Type": "application/json"}) |
| 88 | try: |
| 89 | with urllib.request.urlopen(request, timeout=10) as response: |
| 90 | data = response.read(MAX_LOG + 1) |
| 91 | except urllib.error.HTTPError as error: |
| 92 | if error.code == 404 and method == "GET": |
| 93 | return None |
| 94 | if error.code == 409: |
| 95 | raise Error(409, "The secret changed during this update. Reload it before retrying.") from None |
| 96 | raise Error(502, "Nomad refused the secret request. Check its service logs.") from None |
| 97 | if len(data) > MAX_LOG: |
| 98 | raise Error(502, "The secret response is too large. Check the service configuration.") |
| 99 | return json.loads(data) |
| 100 | |
| 101 | |
| 102 | def secret(action, target, key, value=None): |
| 103 | service(target, key) |
| 104 | job = nomad("job/" + target) |
| 105 | specs = json.loads((job or {}).get("Meta", {}).get("studio_secrets", "[]")) |
| 106 | if not isinstance(specs, list): |
| 107 | raise Error(502, "The service's secret names couldn't be read. Deploy its current release again.") |
| 108 | matches = [spec for spec in specs if isinstance(spec, dict) and spec.get("name") == key] |
| 109 | if len(matches) != 1: |
| 110 | raise Error(404, "That secret is unavailable. Reload the service's secret list.") |
| 111 | path = "nomad/jobs/" + target |
| 112 | existing = nomad("var/" + path) |
| 113 | items = dict(existing["Items"]) if existing else {} |
| 114 | if action == "get": |
| 115 | if key not in items: |
| 116 | raise Error(404, "That secret has no value. Set it from the service page.") |
| 117 | return items[key] |
| 118 | if action == "rotate": |
| 119 | count = matches[0].get("bytes") |
| 120 | if matches[0].get("generated") is not True or type(count) is not int or not 1 <= count <= 4096: |
| 121 | raise Error(409, "This secret comes from outside the server. Set it instead of generating it.") |
| 122 | value = secrets.token_hex(count) |
| 123 | if not isinstance(value, str) or not value or len(value.encode()) > 8192 or any(c in value for c in "\r\n\0"): |
| 124 | raise Error(400, "Enter a single-line secret under 8 KiB.") |
| 125 | items[key] = value |
| 126 | nomad("var/" + path + "?cas=" + str(existing["ModifyIndex"] if existing else 0), "PUT", |
| 127 | {"Namespace": "default", "Path": path, "Items": items}) |
| 128 | restarted = subprocess.run(["nomad", "job", "restart", "-yes", target], capture_output=True, text=True, timeout=55, |
| 129 | env={**os.environ, "NOMAD_TOKEN": read(release.STATE / "nomad.token").strip()}) |
| 130 | if restarted.returncode: |
| 131 | raise Error(502, "The secret was saved, but the service couldn't restart. Check its run before retrying.") |
| 132 | print("Updated " + target + "/" + key, flush=True) |
| 133 | |
| 134 | |
| 135 | def stage(target): |
| 136 | if not isinstance(target, str) or not release.STAGE_ID.fullmatch(target): |
| 137 | raise Error(400, "Choose a stage from the list.") |
| 138 | text = read(release.STATE / "stages" / (target + ".json")) |
| 139 | if text is None: |
| 140 | raise Error(404, "That stage is no longer available. Reload deploys.") |
| 141 | return json.loads(text) |
| 142 | |
| 143 | |
| 144 | def entry(target): |
| 145 | if not isinstance(target, str) or not re.fullmatch(r"[1-9][0-9]{0,9}", target): |
| 146 | raise Error(400, "Choose a deployment from history.") |
| 147 | entries = history() |
| 148 | if int(target) > len(entries): |
| 149 | raise Error(404, "That deployment is outside history. Reload deploys.") |
| 150 | return entries[int(target) - 1] |
| 151 | |
| 152 | |
| 153 | def command(action, target, key=None): |
| 154 | if not isinstance(action, str) or action not in ACTIONS: |
| 155 | raise Error(400, "Choose a supported deployment action.") |
| 156 | if action in {"start", "stop", "restart", "secret-set", "secret-rotate"}: |
| 157 | service(target, key) |
| 158 | if action.startswith("secret-"): |
| 159 | if key is None: |
| 160 | raise Error(400, "Choose a secret from this service's list.") |
| 161 | return [sys.executable, str(Path(__file__)), "--secret", action.removeprefix("secret-"), target, key] |
| 162 | if action in {"stop", "restart"}: |
| 163 | return ["nomad", "job", action, "-yes", target] |
| 164 | version = release.current_release() |
| 165 | if version is None: |
| 166 | raise Error(409, "No release is running. Deploy a release before starting this service.") |
| 167 | return [sys.executable, str(release.check_release(version) / "tools/studio.py"), "deploy", target] |
| 168 | if action == "rollback": |
| 169 | saved = entry(target) |
| 170 | version = saved["release"] |
| 171 | release.check_release(version, legacy=saved.get("legacy") is True) |
| 172 | return [sys.executable, str(Path(__file__).with_name("release.py")), "rollback", version] |
| 173 | if action == "deploy": |
| 174 | candidate = release.main_release() |
| 175 | if not candidate or target != candidate["release"]: |
| 176 | raise Error(409, "Main changed or isn't uploaded. Reload deploys before deploying.") |
| 177 | release.check_release(target) |
| 178 | return [sys.executable, str(Path(__file__).with_name("release.py")), "deploy", target] |
| 179 | metadata = stage(target) |
| 180 | version = metadata.get("release") |
| 181 | if not isinstance(version, str): |
| 182 | raise Error(409, "This stage has no release. Stage it again.") |
| 183 | root = release.check_release(version) |
| 184 | return [sys.executable, str(root / "tools/studio.py"), "destroy", target] |
| 185 | |
| 186 | |
| 187 | def last(): |
| 188 | text = read(HOST_STATE / "last-run.json") |
| 189 | if text is None: |
| 190 | return None |
| 191 | return json.loads(text) |
| 192 | |
| 193 | |
| 194 | def run(identity): |
| 195 | if not isinstance(identity, str) or not re.fullmatch(r"[0-9a-f]{8}(?:-[0-9a-f]{4}){3}-[0-9a-f]{12}", identity): |
| 196 | raise Error(400, "Choose a deployment run from the list.") |
| 197 | root = HOST_STATE / "runs" |
| 198 | text = read(root / (identity + ".log"), MAX_LOG + 4096) |
| 199 | if text is None: |
| 200 | legacy = read(release.STATE / "runs" / (identity + ".log"), MAX_LOG) |
| 201 | if legacy is not None: |
| 202 | root, text = release.STATE / "runs", legacy |
| 203 | latest = last() |
| 204 | if text is None and (not latest or latest.get("id") != identity): |
| 205 | raise Error(404, "That run is no longer available.") |
| 206 | code = read(root / (identity + ".exit"), 64) |
| 207 | if code is None: |
| 208 | shown = subprocess.run(["systemctl", "show", "--property=ActiveState,ExecMainStatus,LoadState", "studio-run-" + identity], |
| 209 | capture_output=True, text=True, timeout=5) |
| 210 | state = dict(line.split("=", 1) for line in shown.stdout.splitlines() if "=" in line) |
| 211 | if shown.returncode and state.get("LoadState") != "not-found": |
| 212 | shown.check_returncode() |
| 213 | if state.get("ActiveState") not in {"active", "activating", "deactivating"}: |
| 214 | code = read(root / (identity + ".exit"), 64) |
| 215 | if code is None: |
| 216 | code = int(state.get("ExecMainStatus", 0)) or 1 |
| 217 | return {"lines": [line for line in (text or "").splitlines() if line], "code": int(code) if code is not None else None} |
| 218 | |
| 219 | |
| 220 | def start(action, target, key=None, value=None): |
| 221 | 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")): |
| 222 | raise Error(400, "Enter a single-line secret under 8 KiB.") |
| 223 | HOST_STATE.mkdir(mode=0o700, parents=True, exist_ok=True) |
| 224 | with (HOST_STATE / "run.lock").open("a") as lock: |
| 225 | fcntl.flock(lock, fcntl.LOCK_EX) |
| 226 | previous = last() |
| 227 | if previous and run(previous["id"])["code"] is None: |
| 228 | raise Error(409, "A deployment is running. Wait for it to finish, then retry.") |
| 229 | command(action, target, key) |
| 230 | identity = str(uuid.uuid4()) |
| 231 | metadata = {"id": identity, "action": action, "target": target} |
| 232 | if key is not None: |
| 233 | metadata["key"] = key |
| 234 | root = HOST_STATE / "runs" |
| 235 | root.mkdir(mode=0o700, parents=True, exist_ok=True) |
| 236 | input_file = root / (identity + ".input") |
| 237 | try: |
| 238 | if action == "secret-set": |
| 239 | with os.fdopen(os.open(input_file, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600), "w") as source: |
| 240 | source.write(value) |
| 241 | pending = HOST_STATE / "last-run.pending" |
| 242 | pending.write_text(json.dumps(metadata) + "\n") |
| 243 | pending.replace(HOST_STATE / "last-run.json") |
| 244 | args = ["systemd-run", "--unit=studio-run-" + identity, "--collect", "--quiet", |
| 245 | "--property=RuntimeMaxSec=3600", "--property=TimeoutStopSec=10", |
| 246 | "--property=MemoryMax=4G", "--property=TasksMax=1024", "--property=CPUWeight=10"] |
| 247 | args.extend("--setenv=" + key + "=" + value for key, value in os.environ.items() if key == "PATH" or key.startswith("STUDIO_")) |
| 248 | subprocess.run([*args, "--", sys.executable, str(Path(__file__)), identity, action, target, *([key] if key is not None else [])], |
| 249 | check=True, capture_output=True, timeout=10) |
| 250 | except BaseException: |
| 251 | input_file.unlink(missing_ok=True) |
| 252 | raise |
| 253 | return metadata |
| 254 | |
| 255 | |
| 256 | def handle(request): |
| 257 | operation = request["operation"] |
| 258 | if operation == "iam.request": |
| 259 | iam_validate(request) |
| 260 | process = subprocess.run([ |
| 261 | "systemd-run", "--pipe", "--wait", "--collect", "--quiet", |
| 262 | "--unit=studio-iam-" + str(uuid.uuid4()), "--property=RuntimeMaxSec=25", |
| 263 | "--property=MemoryMax=128M", "--property=TasksMax=8", "--property=ProtectSystem=strict", |
| 264 | "--property=ProtectHome=yes", "--property=NoNewPrivileges=yes", "--property=CapabilityBoundingSet=", |
| 265 | "--property=RestrictAddressFamilies=AF_INET AF_INET6 AF_UNIX", "--property=IPAddressDeny=any", |
| 266 | "--property=IPAddressAllow=localhost", "--setenv=STUDIO_DOMAIN=" + os.environ["STUDIO_DOMAIN"], |
| 267 | "--setenv=STUDIO_API_TIMEOUT=5", "--", sys.executable, str(Path(__file__)), "--iam"], |
| 268 | input=json.dumps(request), capture_output=True, text=True, timeout=30, check=True) |
| 269 | response = json.loads(process.stdout) |
| 270 | if "error" in response: |
| 271 | raise Error(response["status"], response["error"]) |
| 272 | return response["value"] |
| 273 | if operation.startswith("deploy.secret.") and request["service"] == "keycloak": |
| 274 | raise Error(403, "Keycloak credentials are managed by the host.") |
| 275 | if operation == "deploy.secret.get": |
| 276 | if not isinstance(request["key"], str): |
| 277 | raise Error(400, "Choose a secret from this service's list.") |
| 278 | service(request["service"], request["key"]) |
| 279 | process = subprocess.run(["systemd-run", "--pipe", "--wait", "--collect", "--quiet", |
| 280 | "--unit=studio-secret-read-" + str(uuid.uuid4()), "--property=RuntimeMaxSec=25", |
| 281 | "--property=MemoryMax=128M", "--property=TasksMax=8", "--property=ProtectSystem=strict", |
| 282 | "--property=ProtectHome=yes", "--property=NoNewPrivileges=yes", "--property=CapabilityBoundingSet=", |
| 283 | "--property=RestrictAddressFamilies=AF_INET AF_UNIX", "--property=IPAddressDeny=any", |
| 284 | "--property=IPAddressAllow=localhost", "--", sys.executable, |
| 285 | str(Path(__file__)), "--secret", "get", request["service"], request["key"]], |
| 286 | capture_output=True, text=True, timeout=30, check=True) |
| 287 | response = json.loads(process.stdout) |
| 288 | if "error" in response: |
| 289 | raise Error(response["status"], response["error"]) |
| 290 | return response["value"] |
| 291 | if operation in {"deploy.secret.set", "deploy.secret.rotate"}: |
| 292 | return start("secret-" + operation.rpartition(".")[2], request["service"], request["key"], request.get("value")) |
| 293 | if operation == "deploy.current": |
| 294 | return release.current_release() |
| 295 | if operation == "deploy.main": |
| 296 | return release.main_release() |
| 297 | if operation == "deploy.history": |
| 298 | return history() |
| 299 | if operation == "deploy.managed": |
| 300 | return managed() |
| 301 | if operation == "deploy.stages": |
| 302 | values = [] |
| 303 | for path in (release.STATE / "stages").glob("*.json"): |
| 304 | metadata = stage(path.stem) |
| 305 | values.append({"id": path.stem, "service": metadata["sourceId"], "release": metadata.get("release"), |
| 306 | "ready": metadata.get("ready", True), "created": path.stat().st_mtime, |
| 307 | "overrides": [value.partition("=")[0] for value in metadata.get("overrides", [])], |
| 308 | "clone": metadata.get("clone"), "mount": metadata.get("mount")}) |
| 309 | return values |
| 310 | if operation == "deploy.release": |
| 311 | version = request["release"] |
| 312 | if not isinstance(version, str) or not release.RELEASE_ID.fullmatch(version): |
| 313 | raise Error(400, "Choose a release from deployment history.") |
| 314 | root = release.RELEASES / version / "service" |
| 315 | if root.parent.is_symlink(): |
| 316 | raise ValueError("release directory is a symlink") |
| 317 | if not root.is_dir(): |
| 318 | return None |
| 319 | values = {} |
| 320 | for path in root.iterdir(): |
| 321 | if path.is_symlink(): |
| 322 | raise ValueError("release service directory is a symlink") |
| 323 | if path.is_dir(): |
| 324 | text = read(path / "service.pkl") |
| 325 | if text is not None: |
| 326 | values[path.name] = text |
| 327 | return values |
| 328 | if operation == "deploy.output": |
| 329 | kind, target = request["kind"], request["target"] |
| 330 | if kind == "stages": |
| 331 | value = stage(target) |
| 332 | elif kind == "history": |
| 333 | value = entry(target) |
| 334 | else: |
| 335 | raise Error(400, "Choose a stage or deployment from history.") |
| 336 | found = [] |
| 337 | for path in (release.STATE / "runs").glob("*.log"): |
| 338 | name = path.name |
| 339 | if kind == "stages": |
| 340 | if not name.endswith("-stage-" + value["sourceId"] + ".log"): |
| 341 | continue |
| 342 | text = read(path, MAX_LOG) |
| 343 | if "stage=" + target in text.splitlines(): |
| 344 | found.append((name, 0, text)) |
| 345 | else: |
| 346 | source, version = value["source"], value["release"] |
| 347 | selected = name.endswith(("-prod-" + source + ".log", "-prod-" + version + ".log", "-rollback-" + version + ".log")) |
| 348 | if not selected and not re.fullmatch(r"[0-9a-f-]{36}\.log", name): |
| 349 | continue |
| 350 | delta = abs(path.stat().st_mtime - value["time"]) |
| 351 | if delta >= 600: |
| 352 | continue |
| 353 | text = read(path, MAX_LOG) |
| 354 | first = next(iter(text.splitlines()), "") |
| 355 | if selected or first.endswith(" deploy " + version) or first.endswith(" promote " + source) or first.endswith(" rollback " + version): |
| 356 | found.append((name, delta, text)) |
| 357 | if kind == "history": |
| 358 | for path in (HOST_STATE / "runs").glob("*.log"): |
| 359 | delta = abs(path.stat().st_mtime - value["time"]) |
| 360 | if delta >= 600: |
| 361 | continue |
| 362 | text = read(path, MAX_LOG + 4096) |
| 363 | first = next(iter(text.splitlines()), "") |
| 364 | if first.endswith(" deploy " + value["release"]) or first.endswith(" promote " + value["source"]) or first.endswith(" rollback " + value["release"]): |
| 365 | found.append((path.name, delta, text)) |
| 366 | found.sort(key=lambda item: item[0] if kind == "stages" else item[1], reverse=kind == "stages") |
| 367 | return [line for line in found[0][2].splitlines() if line] if found else None |
| 368 | if operation == "deploy.last": |
| 369 | value = last() |
| 370 | return {**value, "code": run(value["id"])["code"]} if value else None |
| 371 | if operation == "deploy.run": |
| 372 | return run(request["id"]) |
| 373 | return start(request["action"], request["target"]) |
| 374 | |
| 375 | |
| 376 | def iam_validate(request): |
| 377 | path, method, body = request["path"], request["method"], request["body"] |
| 378 | if not isinstance(path, str) or not isinstance(method, str): |
| 379 | raise Error(400, "Choose a supported user operation.") |
| 380 | allowed = { |
| 381 | "/roles": {"GET"}, "/users?max=1000": {"GET"}, "/users": {"POST"}, |
| 382 | "": {"GET", "PUT", "DELETE"}, "/sessions": {"GET"}, "/credentials": {"GET"}, |
| 383 | "/role-mappings/realm": {"GET", "POST", "DELETE"}, "/logout": {"POST"}, |
| 384 | "/shale-username": {"GET"}, |
| 385 | "/execute-actions-email": {"PUT"}, "/reset-password": {"PUT"}, |
| 386 | } |
| 387 | user = re.fullmatch(r"/users/([0-9a-fA-F]{8}(?:-[0-9a-fA-F]{4}){3}-[0-9a-fA-F]{12})(.*)", path) |
| 388 | suffix = user[2] if user else path |
| 389 | if user and suffix not in {"", "/sessions", "/credentials", "/role-mappings/realm", "/logout", "/execute-actions-email", "/reset-password", "/shale-username"}: |
| 390 | raise Error(400, "Choose a supported user operation.") |
| 391 | if path.startswith("/users?username="): |
| 392 | try: |
| 393 | query = urllib.parse.parse_qs(path.partition("?")[2], keep_blank_values=True, strict_parsing=True) |
| 394 | except ValueError: |
| 395 | raise Error(400, "Choose a username.") from None |
| 396 | if (set(query) != {"username", "exact"} or query["exact"] != ["true"] |
| 397 | or len(query["username"]) != 1 or not re.fullmatch(r"[a-z0-9][a-z0-9._@-]{0,254}", query["username"][0])): |
| 398 | raise Error(400, "Choose a username.") |
| 399 | suffix = "/users?max=1000" |
| 400 | if (method not in allowed.get(suffix, set()) or not user and suffix == "" |
| 401 | or method in {"GET", "DELETE", "POST"} and suffix not in {"/users", "/role-mappings/realm"} and body is not None |
| 402 | or method == "GET" and body is not None or suffix == "/shale-username" and not user): |
| 403 | raise Error(400, "Choose a supported user operation.") |
| 404 | if method == "PUT" and suffix == "/reset-password": |
| 405 | if (not isinstance(body, dict) or set(body) != {"type", "value", "temporary"} |
| 406 | or body["type"] != "password" or not isinstance(body["value"], str) |
| 407 | or not 8 <= len(body["value"]) <= 8192 or type(body["temporary"]) is not bool): |
| 408 | raise Error(400, "Enter a password and choose whether it is temporary.") |
| 409 | elif method == "PUT" and suffix == "/execute-actions-email": |
| 410 | if not isinstance(body, list) or not body or not all(isinstance(action, str) and action in IAM_ACTIONS for action in body): |
| 411 | raise Error(400, "Choose a sign-in action.") |
| 412 | elif suffix == "/role-mappings/realm" and method != "GET": |
| 413 | if (not isinstance(body, list) or len(body) != 1 or not isinstance(body[0], dict) |
| 414 | or not isinstance(body[0].get("name"), str) or body[0]["name"] not in IAM_ROLES or not isinstance(body[0].get("id"), str)): |
| 415 | raise Error(403, "Choose a dashboard group.") |
| 416 | elif method == "POST" and suffix == "/users" or method == "PUT" and suffix == "": |
| 417 | if not isinstance(body, dict) or not body or not set(body) <= {"username", "email", "firstName", "lastName", "enabled", "emailVerified", "requiredActions", "attributes"}: |
| 418 | raise Error(400, "Enter a user profile.") |
| 419 | for key, value in body.items(): |
| 420 | if key in {"enabled", "emailVerified"}: |
| 421 | valid = type(value) is bool |
| 422 | elif key == "requiredActions": |
| 423 | valid = isinstance(value, list) and all(isinstance(action, str) and action in IAM_ACTIONS for action in value) |
| 424 | elif key == "attributes": |
| 425 | valid = isinstance(value, dict) and set(value) == {"picture"} and (value["picture"] is None or isinstance(value["picture"], list) and len(value["picture"]) == 1 and isinstance(value["picture"][0], str) and len(value["picture"][0]) <= 8192) |
| 426 | else: |
| 427 | valid = value is None and key != "username" or isinstance(value, str) and len(value) <= 8192 |
| 428 | if key == "username": |
| 429 | valid = isinstance(value, str) and bool(re.fullmatch(r"[a-z0-9][a-z0-9._@-]{0,254}", value)) and value != "admin" |
| 430 | if not valid: |
| 431 | raise Error(400, "Enter a valid user profile.") |
| 432 | return user |
| 433 | |
| 434 | |
| 435 | def iam(request): |
| 436 | user = iam_validate(request) |
| 437 | sys.path.insert(0, str(Path(__file__).resolve().parent.parent / "service/keycloak")) |
| 438 | from api import Keycloak |
| 439 | client = Keycloak("keycloak." + os.environ["STUDIO_DOMAIN"], secret("get", "keycloak", "password"), |
| 440 | attempts=1, cafile="/var/lib/studio/ca-bundle.crt") |
| 441 | path, method, body = request["path"], request["method"], request["body"] |
| 442 | if user: |
| 443 | profile = client.request("/admin/realms/master/users/" + user[1]) |
| 444 | if profile["username"] == "admin": |
| 445 | raise Error(403, "The Keycloak administrator is managed outside the dashboard.") |
| 446 | if user[2] == "/shale-username": |
| 447 | clients = client.request("/admin/realms/master/clients?clientId=shale") |
| 448 | clients = [item for item in clients if item["clientId"] == "shale"] |
| 449 | if len(clients) != 1: |
| 450 | raise Error(502, "Shale's sign-in client is unavailable.") |
| 451 | aliases = client.request(f"/admin/realms/master/users/{user[1]}/role-mappings/clients/{clients[0]['id']}/composite") |
| 452 | if len(aliases) > 1: |
| 453 | raise Error(502, "This account has multiple Shale usernames.") |
| 454 | return {"body": {"username": aliases[0]["name"] if aliases else profile["username"], "enabled": profile["enabled"]}} |
| 455 | if isinstance(body, dict) and "attributes" in body: |
| 456 | attributes = dict(profile.get("attributes", {})) |
| 457 | picture = body["attributes"]["picture"] |
| 458 | if picture is None: |
| 459 | attributes.pop("picture", None) |
| 460 | else: |
| 461 | attributes["picture"] = picture |
| 462 | body = {**{key: profile[key] for key in ["username", "email", "firstName", "lastName"] if key in profile}, |
| 463 | **body, "attributes": attributes} |
| 464 | if path.endswith("/role-mappings/realm") and method != "GET": |
| 465 | roles = client.request("/admin/realms/master/roles") |
| 466 | role = next((role for role in roles if role["name"] in IAM_ROLES and role["id"] == body[0]["id"] and role["name"] == body[0]["name"]), None) |
| 467 | if role is None: |
| 468 | raise Error(403, "Choose a dashboard group.") |
| 469 | body = [{"id": role["id"], "name": role["name"]}] |
| 470 | result = client.request("/admin/realms/master" + path, method, body, full=True) |
| 471 | if method == "GET" and path.startswith("/users?"): |
| 472 | result["body"] = [user for user in result["body"] if user["username"] != "admin"] |
| 473 | elif method == "GET" and (path == "/roles" or path.endswith("/role-mappings/realm")): |
| 474 | result["body"] = [role for role in result["body"] if role["name"] in IAM_ROLES] |
| 475 | return result |
| 476 | |
| 477 | |
| 478 | def worker(identity, action, target, key=None): |
| 479 | if not re.fullmatch(r"[0-9a-f]{8}(?:-[0-9a-f]{4}){3}-[0-9a-f]{12}", identity): |
| 480 | raise ValueError("incorrect deployment run ID") |
| 481 | saved = last() |
| 482 | expected = {"id": identity, "action": action, "target": target} |
| 483 | if key is not None: |
| 484 | expected["key"] = key |
| 485 | if saved != expected: |
| 486 | raise ValueError("deployment run is outside host state") |
| 487 | root = HOST_STATE / "runs" |
| 488 | root.mkdir(mode=0o700, parents=True, exist_ok=True) |
| 489 | os.umask(0o077) |
| 490 | code = 1 |
| 491 | def interrupted(signum, frame): |
| 492 | raise RuntimeError("Deployment stopped before completion.") |
| 493 | signal.signal(signal.SIGTERM, interrupted) |
| 494 | with (root / (identity + ".log")).open("xb", buffering=0) as output: |
| 495 | try: |
| 496 | argv = command(action, target, key) |
| 497 | output.write(("$ " + " ".join(argv) + "\n").encode()) |
| 498 | environment = None |
| 499 | if action in {"stop", "restart"}: |
| 500 | environment = {**os.environ, "NOMAD_TOKEN": read(release.STATE / "nomad.token", missing="").strip()} |
| 501 | if not environment["NOMAD_TOKEN"]: |
| 502 | raise RuntimeError("Nomad credentials are unavailable. Check the host service configuration.") |
| 503 | input_file = root / (identity + ".input") |
| 504 | with subprocess.Popen(argv, stdin=subprocess.PIPE if action == "secret-set" else subprocess.DEVNULL, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, |
| 505 | start_new_session=True, env=environment) as process: |
| 506 | try: |
| 507 | if action == "secret-set": |
| 508 | process.stdin.write(read(input_file, 8192).encode()) |
| 509 | process.stdin.close() |
| 510 | input_file.unlink() |
| 511 | total = 0 |
| 512 | deadline = time.monotonic() + 3500 |
| 513 | with selectors.DefaultSelector() as selector: |
| 514 | selector.register(process.stdout, selectors.EVENT_READ) |
| 515 | while selector.get_map(): |
| 516 | remaining = deadline - time.monotonic() |
| 517 | if remaining <= 0: |
| 518 | raise TimeoutError("Deployment exceeded its time limit.") |
| 519 | for key, _ in selector.select(remaining): |
| 520 | chunk = os.read(key.fd, 65536) |
| 521 | if not chunk: |
| 522 | selector.unregister(key.fileobj) |
| 523 | continue |
| 524 | total += len(chunk) |
| 525 | if total > MAX_LOG: |
| 526 | raise RuntimeError("Deployment output exceeded its size limit.") |
| 527 | output.write(chunk) |
| 528 | code = process.wait(timeout=max(.001, deadline - time.monotonic())) |
| 529 | except BaseException: |
| 530 | try: |
| 531 | os.killpg(process.pid, signal.SIGKILL) |
| 532 | except ProcessLookupError: |
| 533 | pass |
| 534 | raise |
| 535 | except Exception as error: |
| 536 | output.write((str(error) + "\n").encode()) |
| 537 | finally: |
| 538 | (root / (identity + ".input")).unlink(missing_ok=True) |
| 539 | output.write(("exit " + str(code) + "\n").encode()) |
| 540 | pending = root / (identity + ".exit.tmp") |
| 541 | pending.write_text(str(code)) |
| 542 | pending.replace(root / (identity + ".exit")) |
| 543 | return code |
| 544 | |
| 545 | |
| 546 | if __name__ == "__main__": |
| 547 | if sys.argv[1:] == ["--iam"]: |
| 548 | try: |
| 549 | payload = sys.stdin.read(MAX_METADATA + 1) |
| 550 | if len(payload.encode()) > MAX_METADATA: |
| 551 | raise Error(400, "The user request is too large. Narrow the selection.") |
| 552 | result = {"value": iam(json.loads(payload))} |
| 553 | except Error as error: |
| 554 | result = {"error": str(error), "status": error.status} |
| 555 | except urllib.error.HTTPError as error: |
| 556 | result = {"error": "Keycloak refused this change. Reload the page and retry.", "status": error.code if error.code in {404, 409} else 502} |
| 557 | except Exception: |
| 558 | result = {"error": "Keycloak is unavailable. Check its service logs.", "status": 502} |
| 559 | print(json.dumps(result)) |
| 560 | raise SystemExit(0) |
| 561 | if len(sys.argv) == 5 and sys.argv[1] == "--secret" and sys.argv[2] in {"get", "set", "rotate"}: |
| 562 | action, target, key = sys.argv[2:] |
| 563 | try: |
| 564 | result = secret(action, target, key, sys.stdin.read(8193) if action == "set" else None) |
| 565 | except Error as error: |
| 566 | if action != "get": |
| 567 | raise SystemExit(str(error)) |
| 568 | result = {"error": str(error), "status": error.status} |
| 569 | else: |
| 570 | result = {"value": result} |
| 571 | if action == "get": |
| 572 | print(json.dumps(result)) |
| 573 | raise SystemExit(0) |
| 574 | if len(sys.argv) not in {4, 5}: |
| 575 | raise SystemExit("Expected a run ID, action, and target") |
| 576 | raise SystemExit(worker(*sys.argv[1:])) |