1#!/usr/bin/env python3
2import fcntl
3import json
4import os
5from pathlib import Path
6import re
7import selectors
8import secrets
9import signal
10import stat
11import subprocess
12import sys
13import time
14import uuid
15import urllib.error
16import urllib.parse
17import urllib.request
18
19import release
20
21
22FIELDS = {
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}
30ACTIONS = {"deploy", "destroy", "rollback", "start", "stop", "restart", "secret-set", "secret-rotate"}
31MAX_LOG = 1024 * 1024
32MAX_METADATA = 65536
33HOST_STATE = Path(os.environ.get("STUDIO_HOST_STATE_ROOT", "/var/lib/studio/host"))
34IAM_ROLES = {"infra-admin", "media", "media-manage"}
35IAM_ACTIONS = {"UPDATE_PASSWORD", "VERIFY_EMAIL", "UPDATE_PROFILE", "CONFIGURE_TOTP", "webauthn-register", "webauthn-register-passwordless"}
36
37
38class Error(Exception):
39 def __init__(self, status, message):
40 super().__init__(message)
41 self.status = status
42
43
44def 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
58def 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
65def 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
72def 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
81def 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
102def 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
135def 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
144def 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
153def 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
187def last():
188 text = read(HOST_STATE / "last-run.json")
189 if text is None:
190 return None
191 return json.loads(text)
192
193
194def 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
220def 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
256def 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
376def 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
435def 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
478def 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
546if __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:]))