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.request
17
18import release
19
20
21FIELDS = {
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}
28ACTIONS = {"deploy", "destroy", "rollback", "start", "stop", "restart", "secret-set", "secret-rotate"}
29MAX_LOG = 1024 * 1024
30MAX_METADATA = 65536
31HOST_STATE = Path(os.environ.get("STUDIO_HOST_STATE_ROOT", "/var/lib/studio/host"))
32
33
34class Error(Exception):
35 def __init__(self, status, message):
36 super().__init__(message)
37 self.status = status
38
39
40def 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
54def 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
61def 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
68def 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
77def 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
98def 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
131def 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
140def 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
149def 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
183def last():
184 text = read(HOST_STATE / "last-run.json")
185 if text is None:
186 return None
187 return json.loads(text)
188
189
190def 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
216def 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
252def 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
355def 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
423if __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:]))