1#!/usr/bin/env python3
2import argparse
3from contextlib import nullcontext
4import fcntl
5import hashlib
6import json
7import os
8from pathlib import Path
9import re
10import secrets
11import shutil
12import subprocess
13import sys
14import tempfile
15import time
16
17from data import cloned_postgres, dataset_for
18from release import excluded_services
19
20
21REPO = Path(__file__).resolve().parent.parent
22SERVICES = REPO / "service"
23STATE = Path("/var/lib/studio")
24IDENTITIES = STATE / "identities.json"
25ASSETS = Path("/var/lib/caddy/studio")
26NAME = re.compile(r"[a-z][a-z0-9-]*\Z")
27
28
29def command(*args, capture=False, **kwargs):
30 return subprocess.run(args, check=True, text=True, capture_output=capture, **kwargs)
31
32
33def service_files():
34 excluded = excluded_services(REPO)
35 files = {}
36 for directory in SERVICES.iterdir():
37 if directory.name in excluded or not directory.is_dir():
38 continue
39 for path in directory.glob("*.pkl"):
40 name = directory.name if path.name == "service.pkl" else path.stem
41 if name in excluded:
42 continue
43 if name in files:
44 raise ValueError(f"duplicate service definition: {name}")
45 files[name] = path
46 return files
47
48
49def services():
50 return sorted(service_files())
51
52
53def service_dir(name):
54 path = service_files()[name]
55 return path.parent if path.name == "service.pkl" else path.parent / name
56
57
58def dependencies(data):
59 return set(data["dependsOn"]) | {item["provider"] for item in data["inputs"].values()}
60
61
62def ordered(data):
63 pending = {item["id"]: item for item in data}
64 result = []
65 while pending:
66 ready = [
67 item
68 for item in pending.values()
69 if not dependencies(item) & pending.keys()
70 ]
71 if not ready:
72 raise ValueError("cycle in service dependencies")
73 for item in ready:
74 result.append(item)
75 del pending[item["id"]]
76 return result
77
78
79def identity(name):
80 STATE.mkdir(parents=True, exist_ok=True)
81 if not IDENTITIES.exists():
82 write_private(IDENTITIES, (REPO / "config/identities.json").read_text())
83 with (IDENTITIES.parent / ".identities.lock").open("w") as lock:
84 fcntl.flock(lock, fcntl.LOCK_EX)
85 with IDENTITIES.open() as file:
86 ids = json.load(file)
87 if len(set(ids.values())) != len(ids):
88 raise ValueError("duplicate service UIDs")
89 if name not in ids:
90 used = set(ids.values())
91 ids[name] = next(uid for uid in range(3100, 60000) if uid not in used)
92 pending = IDENTITIES.with_suffix(".pending")
93 pending.write_text(json.dumps(ids, indent=2) + "\n")
94 pending.replace(IDENTITIES)
95 return ids[name]
96
97
98def load(name, properties, instance_id=None):
99 files = service_files()
100 if not NAME.fullmatch(name) or name not in files:
101 raise ValueError(f"unknown service: {name}")
102 instance_id = instance_id or name
103 if not NAME.fullmatch(instance_id):
104 raise ValueError(f"invalid instance ID: {instance_id}")
105 uid = identity(name)
106 result = command(
107 "pkl",
108 "eval",
109 *(
110 part
111 for key, value in {**properties, "serviceId": instance_id, "uid": uid, "preview": str(instance_id != name).lower()}.items()
112 for part in ("-p", f"{key}={value}")
113 ),
114 str(files[name]),
115 capture=True,
116 )
117 data = json.loads(result.stdout)
118 if data["id"] != instance_id:
119 raise ValueError(f"service ID must match deployment: {instance_id}")
120 inputs = {}
121 for requirement in data.pop("requirements"):
122 alias = requirement["alias"]
123 if not NAME.fullmatch(alias) or alias == "own" or alias in inputs:
124 raise ValueError(f"invalid or duplicate requirement alias: {alias}")
125 inputs[alias] = requirement
126 data["inputs"] = inputs
127 data["uid"] = uid
128 data["sourceId"] = name
129 for task_name, task in data["containers"].items():
130 if bool(task.get("image")) == bool(task.get("build")):
131 raise ValueError(f"{name}.{task_name} needs one image or build directory")
132 if task.get("build"):
133 if REPO.parent == Path("/opt/studio/releases"):
134 context = config_path(name, task["build"])
135 if not context.is_dir() or not (context / "Dockerfile").is_file():
136 raise ValueError(f"invalid build context: {name}.{task_name}")
137 digest = hashlib.sha256()
138 for path in sorted(context.rglob("*")):
139 if path.is_file():
140 digest.update(str(path.relative_to(context)).encode() + b"\0")
141 digest.update(path.read_bytes())
142 tag = digest.hexdigest()[:16]
143 elif not (service_dir(name) / task["build"] / "Dockerfile").is_file() and not (service_dir(name) / "build-source.json").is_file():
144 raise ValueError(f"missing build source: {name}.{task_name}")
145 else:
146 tag = REPO.name
147 task["image"] = f"localhost/studio/{name}-{task_name}:{tag}"
148 return data
149
150
151def q(value):
152 return json.dumps(value, ensure_ascii=False)
153
154
155def write_private(path, content):
156 pending = None
157 try:
158 with tempfile.NamedTemporaryFile("w", dir=path.parent, delete=False) as file:
159 pending = Path(file.name)
160 file.write(content)
161 pending.replace(path)
162 finally:
163 if pending:
164 pending.unlink(missing_ok=True)
165
166
167def host_path(root, relative):
168 if relative == ".":
169 return root
170 part = Path(relative)
171 if part.is_absolute() or ".." in part.parts or not part.parts:
172 raise ValueError(f"invalid relative path: {relative}")
173 return str(Path(root) / part)
174
175
176def config_path(service, relative):
177 root = service_dir(service).resolve()
178 source = (root / relative).resolve()
179 if not source.is_relative_to(root) or not source.exists():
180 raise ValueError(f"invalid config path: {relative}")
181 return source
182
183
184def config_assets(data):
185 paths = set()
186 for task in data["containers"].values():
187 paths.update(volume["config"] for volume in task["volumes"].values() if volume.get("config"))
188 if task.get("http"):
189 paths.update(task["http"]["overrideFiles"].values())
190 sources = {rel: config_path(data.get("sourceId", data["id"]), rel) for rel in paths}
191 digest = hashlib.sha256()
192 for rel, source in sorted(sources.items()):
193 digest.update(rel.encode())
194 files = sorted(source.rglob("*")) if source.is_dir() else [source]
195 for file in files:
196 if file.is_symlink():
197 raise ValueError(f"config asset cannot be a symlink: {file}")
198 if file.is_file():
199 digest.update(str(Path(rel) / file.relative_to(source)).encode() if source.is_dir() else rel.encode())
200 digest.update(file.read_bytes())
201 return ASSETS / data["id"] / digest.hexdigest()[:16], sources
202
203
204def render(data, definitions):
205 if not data["containers"]:
206 return ""
207 name = data["id"]
208 asset_dir, _ = config_assets(data)
209 prepare_digest = None
210 if data.get("prepare"):
211 source = service_dir(data.get("sourceId", name))
212 digest = hashlib.sha256()
213 for path in sorted(source.rglob("*")):
214 if not path.is_file() or "__pycache__" in path.parts or path.suffix == ".pyc":
215 continue
216 if path.name == "service.pkl" or path.name.startswith("icon"):
217 continue
218 digest.update(str(path.relative_to(source)).encode() + b"\0")
219 digest.update(path.read_bytes())
220 prepare_digest = digest.hexdigest()[:16]
221 routed = [task["http"] for task in data["containers"].values()
222 if task.get("http") and task["http"].get("hostname")]
223 primary = routed[0] if routed else {}
224 requirements = {provider: "startup" for provider in data["dependsOn"]}
225 requirements.update({item["provider"]: item["kind"] for item in data["inputs"].values()})
226 secrets_meta = [
227 {"name": name, "generated": True, "bytes": spec["bytes"]}
228 for name, spec in data["secrets"].items()
229 ] + [{"name": name, "generated": False} for name in data["requiredSecrets"]]
230 lines = [
231 f"job {q(name)} {{",
232 ' datacenters = ["clover"]',
233 ' type = "service"',
234 ' meta {',
235 f" studio_service = {q(data['sourceId'])}",
236 f" studio_name = {q(data['name'])}",
237 f" studio_tagline = {q(data['tagline'])}",
238 f" studio_hostname = {q(primary.get('hostname') or '')}",
239 f" studio_launcher = {q(str(data['launcher']).lower())}",
240 f" studio_auth_role = {q(data.get('access') or primary.get('authRole') or '')}",
241 f" studio_trace_service = {q(data.get('traceServiceName') or '')}",
242 f" studio_metrics_pushed = {q(str(data['metricsPushed']).lower())}",
243 f" studio_requires = {q(json.dumps(requirements))}",
244 f" studio_secrets = {q(json.dumps(secrets_meta))}",
245 ' }',
246 ' group "app" {',
247 ]
248 if dependencies(data):
249 lines += [
250 " restart {",
251 " attempts = 3",
252 ' delay = "1m"',
253 ' interval = "4m"',
254 ' mode = "delay"',
255 " }",
256 ]
257 if data["rollout"] == "overlapped":
258 for task in data["containers"].values():
259 if task["hostNetwork"] or any(
260 endpoint and endpoint.get("hostPort") is not None
261 for endpoint in (task.get("http"), task.get("tcp"))
262 ) or any(not volume["readOnly"] for volume in task["volumes"].values()):
263 raise ValueError(f"overlapped rollout requires dynamic ports and read-only mounts: {name}")
264 elif data["rollout"] != "simple":
265 raise ValueError(f"unknown rollout: {data['rollout']}")
266 if data["rollout"] == "overlapped" or data.get("healthyDeadline"):
267 lines.append(" update {")
268 if data["rollout"] == "overlapped":
269 lines += [" canary = 1", " auto_promote = true", " auto_revert = true"]
270 if data.get("healthyDeadline"):
271 lines += [f" healthy_deadline = {q(data['healthyDeadline'])}", ' progress_deadline = "0"']
272 lines.append(" }")
273 if data.get("healthyDeadline"):
274 lines += [" migrate {", f" healthy_deadline = {q(data['healthyDeadline'])}", " }"]
275 ports = {}
276 network = []
277 host_network = any(task["hostNetwork"] for task in data["containers"].values())
278 if host_network and not all(task["hostNetwork"] for task in data["containers"].values()):
279 raise ValueError("mixed host and bridge networking in one group")
280 for task_name, task in data["containers"].items():
281 for kind in ("http", "tcp"):
282 if task.get(kind):
283 endpoint = task[kind]
284 port_name = endpoint["name"] if kind == "tcp" else kind
285 label = port_name if len(data["containers"]) == 1 else f"{task_name}-{port_name}"
286 ports[(task_name, kind)] = label
287 network.append(f" port {q(label)} {{")
288 if endpoint.get("hostPort") is not None:
289 network.append(f" static = {endpoint['hostPort']}")
290 if not task["hostNetwork"]:
291 network.append(f" to = {endpoint['containerPort']}")
292 if (kind == "http" and endpoint.get("hostPort") is None) or (kind == "tcp" and endpoint["loopback"]):
293 network.append(' host_network = "loopback"')
294 network.append(" }")
295 if network:
296 lines.append(" network {")
297 if host_network:
298 lines.append(' mode = "host"')
299 lines += network + [" }"]
300 for provider_id in sorted(dependencies(data)):
301 provider = definitions[provider_id]
302 for provider_task, container in provider["containers"].items():
303 endpoint = container.get("http")
304 if not endpoint or not endpoint.get("hostname") or endpoint.get("authRole"):
305 continue
306 host = endpoint["hostname"]
307 task_name = f"wait-{provider_id}-{provider_task}"
308 lines += [
309 f" task {q(task_name)} {{",
310 ' driver = "podman"',
311 " lifecycle {",
312 ' hook = "prestart"',
313 " sidecar = false",
314 " }",
315 " config {",
316 ' image = "docker.io/curlimages/curl@sha256:43ebaa53d3806db6b1ce4353b6b26ae638ec1c167ee351524b05690f988bb20d"',
317 f" extra_hosts = {q([host + ':host-gateway'])}",
318 f" args = {q(['--insecure', '--fail', '--silent', '--show-error', '--retry', '999999', '--retry-all-errors', '--retry-delay', '5', '--connect-timeout', '3', '--max-time', '5', 'https://' + host + endpoint['checkPath']])}",
319 " }",
320 " resources {",
321 " cpu = 50",
322 " memory = 64",
323 " }",
324 " }",
325 ]
326 for task_name, task in data["containers"].items():
327 if task.get("http") or task.get("tcp"):
328 kind = "http" if task.get("http") else "tcp"
329 endpoint = task[kind]
330 lines += [" service {", f" name = {q(name if task_name == 'app' else name + '-' + task_name)}", ' provider = "nomad"', f" port = {q(ports[(task_name, kind)])}"]
331 if len(data["containers"]) == 1:
332 lines.append(f" task = {q(task_name)}")
333 tags = []
334 if kind == "http" and endpoint.get("metricsPath"):
335 tags.append("studio-metrics-path=" + endpoint["metricsPath"])
336 if kind == "http" and endpoint.get("hostname"):
337 tags += ["caddy-host=" + host for host in [endpoint["hostname"], *endpoint["plainHostnames"]]]
338 if endpoint.get("authRole"):
339 tags.append("caddy-auth-role=" + endpoint["authRole"])
340 if endpoint.get("userHeader"):
341 tags.append("caddy-user-header=" + endpoint["userHeader"])
342 if endpoint["forwardRealIp"]:
343 tags.append("caddy-real-ip=true")
344 if endpoint["tlsInternal"]:
345 tags.append("caddy-internal=true")
346 if tags:
347 lines.append(f" tags = {q(tags)}")
348 lines += [" check {", f" type = {q(kind)}"]
349 if kind == "http":
350 lines.append(f" path = {q(endpoint['checkPath'])}")
351 if endpoint["checkHeaders"]:
352 lines.append(" header {")
353 for key, value in endpoint["checkHeaders"].items():
354 if not re.fullmatch(r"[A-Za-z][A-Za-z0-9-]*", key):
355 raise ValueError(f"invalid health check header: {key}")
356 lines.append(f" {key} = [{q(value)}]")
357 lines.append(" }")
358 lines += [
359 ' interval = "10s"',
360 ' timeout = "15s"',
361 ' check_restart {',
362 ' limit = 12',
363 f" grace = {q(data.get('healthRestartGrace') or data.get('healthyDeadline') or '5m')}",
364 ' }',
365 ' }',
366 ' }',
367 ]
368 for task_name, task in data["containers"].items():
369 lines += [f" task {q(task_name)} {{", ' driver = "podman"']
370 if task["lifecycle"] != "main":
371 lines += [' lifecycle {', ' hook = "prestart"',
372 f" sidecar = {str(task['lifecycle'] == 'prestartSidecar').lower()}", ' }']
373 if not task["imageUser"]:
374 uses_clover = any(
375 volume.get("clover") or (
376 volume.get("src")
377 and Path(volume["src"]).resolve().is_relative_to(Path(data["cloverRoot"]).resolve())
378 )
379 for volume in task["volumes"].values()
380 )
381 gid = data["cloverGid"] if uses_clover else 0 if task["rootGroup"] else data["uid"]
382 lines.append(f" user = {q(str(data['uid']) + ':' + str(gid))}")
383 lines += [" config {", f" image = {q(task['image'])}"]
384 if task.get("entrypoint"):
385 lines.append(f" entrypoint = {q(task['entrypoint'])}")
386 if task["args"]:
387 lines.append(f" args = {q(task['args'])}")
388 if task["hostNetwork"]:
389 lines.append(' network_mode = "host"')
390 task_ports = [label for (owner, _), label in ports.items() if owner == task_name]
391 if task_ports and not task["hostNetwork"]:
392 lines.append(f" ports = {q(task_ports)}")
393 for field, values in (("cap_add", task["capAdd"]), ("devices", task["devices"]), ("extra_hosts", task["extraHosts"])):
394 if values:
395 lines.append(f" {field} = {q(values)}")
396 volumes = []
397 for target, volume in task["volumes"].items():
398 if not target.startswith("/") or ":" in target:
399 raise ValueError(f"invalid volume target: {target}")
400 if volume.get("config") is not None:
401 if volume.get("src") is not None or not volume["readOnly"]:
402 raise ValueError("config mounts must be read only and have no host source")
403 config_path(data.get("sourceId", name), volume["config"])
404 source = host_path(str(asset_dir), volume["config"])
405 elif volume.get("src") is not None:
406 source = volume["src"]
407 if not Path(source).is_absolute():
408 raise ValueError(f"volume source must be absolute: {source}")
409 else:
410 source = host_path(data["hostRoot"], target.lstrip("/"))
411 if ":" in source:
412 raise ValueError(f"invalid volume source: {source}")
413 volumes.append(f"{source}:{target}" + (":ro" if volume["readOnly"] else ""))
414 if volumes:
415 lines.append(f" volumes = {q(volumes)}")
416 if task["tmpfs"]:
417 lines.append(f" tmpfs = {q(task['tmpfs'])}")
418 lines.append(" }")
419 secret_env = []
420 plain_env = {}
421 for key, value in task["env"].items():
422 if not re.fullmatch(r"[A-Za-z_][A-Za-z0-9_]*", key):
423 raise ValueError(f"invalid env key: {key}")
424 parts = re.split(r"(\$\{secret\.[A-Za-z_][A-Za-z0-9_]*\.[A-Za-z_][A-Za-z0-9_]*\})", value)
425 if len(parts) == 1:
426 if "${secret." in value:
427 raise ValueError(f"invalid secret reference: {value}")
428 plain_env[key] = value
429 continue
430 args = []
431 paths = {}
432 for part in parts:
433 match = re.fullmatch(r"\$\{secret\.([A-Za-z_][A-Za-z0-9_]*)\.([A-Za-z_][A-Za-z0-9_]*)\}", part)
434 if match:
435 alias, field = match.groups()
436 if alias != "own" and alias not in data["inputs"]:
437 raise ValueError(f"unknown secret alias: {alias}")
438 path = f"nomad/jobs/{name}" + (f"/inputs/{alias}" if alias != "own" else "")
439 variable = paths.setdefault(path, f"$s{len(paths)}")
440 args.append(f"(index {variable} {q(field)}).Value")
441 elif part:
442 args.append(q(part))
443 start = "".join(f'{{{{ with {variable} := nomadVar {q(path)} }}}}' for path, variable in paths.items())
444 end = "{{ end }}" * len(paths)
445 secret_env.append(f'{start}{key}={{{{ print {" ".join(args)} | toJSON }}}}{end}')
446 if "STUDIO_PREPARE_DIGEST" in task["env"]:
447 raise ValueError("STUDIO_PREPARE_DIGEST is reserved")
448 lines.append(" env {")
449 for key, value in plain_env.items():
450 lines.append(f" {key} = {q(value)}")
451 if prepare_digest:
452 lines.append(f" STUDIO_PREPARE_DIGEST = {q(prepare_digest)}")
453 lines.append(" }")
454 if secret_env:
455 lines += [" template {", " data = <<EOF", *secret_env, "EOF", f' destination = {q("secrets/" + task_name + "-secret.env")}', " env = true", " error_on_missing_key = true", " }"]
456 if task.get("envTemplate"):
457 if "\nEOF\n" in task["envTemplate"]:
458 raise ValueError("invalid env template delimiter")
459 lines += [" template {", " data = <<EOF", task["envTemplate"].rstrip(), "EOF", f' destination = {q("secrets/" + task_name + ".env")}', " env = true", " }"]
460 lines += [" resources {", f" cpu = {task['cpu']}", f" memory = {task['memory']}", " }", " }"]
461 lines += [" }", "}"]
462 return "\n".join(lines) + "\n"
463
464
465def token():
466 if "NOMAD_TOKEN" not in os.environ:
467 os.environ["NOMAD_TOKEN"] = (STATE / "nomad.token").read_text().strip()
468
469
470def get_variable(path):
471 result = subprocess.run(["nomad", "var", "get", "-out=json", path], capture_output=True, text=True)
472 if result.returncode == 0:
473 return json.loads(result.stdout)
474 if "Variable not found" in result.stderr + result.stdout:
475 return None
476 raise RuntimeError(f"Nomad variable read failed for {path}: {result.stderr.strip()}")
477
478
479def provision(data):
480 if not data["containers"]:
481 return
482 if os.geteuid() != 0:
483 raise PermissionError("provisioning requires root")
484 root = Path(data["hostRoot"])
485 if data["hasManagedVolumes"]:
486 root.mkdir(parents=True, exist_ok=True)
487 if command(
488 "findmnt", "-n", "-o", "FSTYPE", "--mountpoint", str(root), capture=True
489 ).stdout.strip() != "zfs":
490 raise ValueError(f"expected ZFS dataset at {root}")
491 for task in data["containers"].values():
492 gid = 0 if task["rootGroup"] else data["uid"]
493 for target, volume in task["volumes"].items():
494 if volume.get("config") is not None or volume.get("src") is not None:
495 continue
496 path = Path(host_path(str(root), target.lstrip("/")))
497 path.mkdir(parents=True, exist_ok=True)
498 os.chown(path, data["uid"], gid)
499 path.chmod(0o750)
500 if data["sourceId"] == "ytdl" and target == "/data":
501 command("systemd-tmpfiles", "--create", "--prefix=" + str(path))
502 bundle = STATE / "ca-bundle.crt"
503 if any(
504 volume.get("src") == str(bundle)
505 for task in data["containers"].values()
506 for volume in task["volumes"].values()
507 ):
508 contents = Path("/etc/ssl/certs/ca-certificates.crt").read_bytes()
509 local_ca = Path("/var/lib/caddy/.local/share/caddy/pki/authorities/local/root.crt")
510 if local_ca.exists():
511 contents += b"\n" + local_ca.read_bytes()
512 if not bundle.exists() or bundle.read_bytes() != contents:
513 pending = bundle.with_suffix(".pending")
514 pending.write_bytes(contents)
515 pending.chmod(0o644)
516 pending.replace(bundle)
517 route_dir = STATE / "routes"
518 route_dir.mkdir(parents=True, exist_ok=True)
519 asset_dir, assets = config_assets(data)
520 static = {}
521 directories = {}
522 identity = {}
523 head_html = {}
524 metrics_paths = {}
525 for task in data["containers"].values():
526 if task.get("http"):
527 http = task["http"]
528 if http["identityHeaders"]:
529 identity[http["hostname"]] = http["identityHeaders"]
530 if http["headHtml"]:
531 head_html[http["hostname"]] = http["headHtml"]
532 if http.get("metricsPath") and http["hostname"]:
533 metrics_paths[http["hostname"]] = http["metricsPath"]
534 for request, rel in http["overrideFiles"].items():
535 origin = assets[rel]
536 target = directories if origin.is_dir() else static
537 target[request] = host_path(str(asset_dir), rel)
538 if assets:
539 asset_dir.parent.mkdir(parents=True, exist_ok=True)
540 asset_dir.parent.chmod(0o755)
541 if assets and not asset_dir.exists():
542 temporary = tempfile.mkdtemp(dir=asset_dir.parent)
543 try:
544 for rel, origin in assets.items():
545 destination = Path(temporary) / rel
546 destination.parent.mkdir(parents=True, exist_ok=True)
547 if origin.is_dir():
548 shutil.copytree(origin, destination)
549 for child in destination.rglob("*"):
550 child.chmod(0o755 if child.is_dir() else 0o644)
551 else:
552 shutil.copy2(origin, destination)
553 destination.chmod(0o644)
554 Path(temporary).chmod(0o755)
555 os.replace(temporary, asset_dir)
556 finally:
557 shutil.rmtree(temporary, ignore_errors=True)
558 secret_dir = STATE / "secrets"
559 secret_dir.mkdir(parents=True, exist_ok=True)
560 secret_file = secret_dir / f"{data['id']}.nv.hcl"
561 generated = data["secrets"]
562 required = data["requiredSecrets"]
563 secret_path = "nomad/jobs/" + data["id"]
564 existing = get_variable(secret_path) if generated or required else None
565 if generated and not existing and secret_file.exists():
566 command("nomad", "var", "put", "-in=hcl", "-out=none", secret_path, "@" + str(secret_file))
567 existing = get_variable(secret_path)
568 values = dict(existing["Items"]) if existing else {}
569 if generated:
570 for item, spec in generated.items():
571 if item not in values:
572 values[item] = secrets.token_hex(spec["bytes"])
573 if not existing or values != existing["Items"]:
574 put_variable(secret_path, values, existing["ModifyIndex"] if existing else 0)
575 backup = "items {\n" + "".join(f" {item} = {q(value)}\n" for item, value in values.items()) + "}\n"
576 if not secret_file.exists() or secret_file.read_text() != backup:
577 write_private(secret_file, backup)
578 if required:
579 missing = [name for name in required if not values.get(name)]
580 if missing:
581 raise ValueError(f"missing required secrets for {data['id']}: {', '.join(missing)}")
582 proof = {}
583 for task in data["containers"].values():
584 http = task.get("http")
585 if http and http.get("identityProofSecret"):
586 name = http["identityProofSecret"]
587 if not http["identityHeaders"] or name not in values:
588 raise ValueError(f"missing identity proof secret for {data['id']}")
589 proof[http["hostname"]] = "X-Studio-Idp-" + values[name]
590 write_private(route_dir / f"{data['id']}.json", json.dumps({
591 "files": static, "dirs": directories, "identity": identity,
592 "identityProof": proof, "headHtml": head_html, "metricsPaths": metrics_paths,
593 }))
594
595
596def submit(data, mode, definitions):
597 if not data["containers"]:
598 return
599 with tempfile.NamedTemporaryFile("w", suffix=".nomad.hcl", delete=False) as file:
600 file.write(render(data, definitions))
601 path = file.name
602 try:
603 command("nomad", "job", mode, path)
604 finally:
605 Path(path).unlink()
606
607
608def build_images(data):
609 for task in data["containers"].values():
610 if task.get("build"):
611 image = task["image"]
612 if subprocess.run(["podman", "image", "exists", image]).returncode:
613 command("podman", "build", "-t", image, str(config_path(data["sourceId"], task["build"])))
614
615
616def clone_dataset(source, snapshot, dataset, mount):
617 acltype = command("zfs", "get", "-H", "-o", "value", "acltype", source, capture=True).stdout.strip()
618 command("zfs", "clone", "-o", "mountpoint=" + mount, "-o", "acltype=" + acltype, snapshot, dataset)
619
620
621def check_http(host, path, tls_internal=False, attempts=120):
622 for attempt in range(attempts):
623 try:
624 if subprocess.run(
625 ["curl", "--noproxy", "*", "-sf", *(["--cacert", "/var/lib/caddy/.local/share/caddy/pki/authorities/local/root.crt"] if tls_internal else []), "--resolve", f"{host}:443:127.0.0.1", f"https://{host}{path}"],
626 capture_output=True, timeout=3,
627 ).returncode == 0:
628 print(f"Healthy {host}")
629 return
630 except subprocess.TimeoutExpired:
631 pass
632 time.sleep(2)
633 raise RuntimeError(f"HTTPS route is unhealthy: {host}")
634
635
636def put_variable(path, items, index=None):
637 with tempfile.NamedTemporaryFile("w", suffix=".nv.hcl", delete=False) as file:
638 file.write("items {\n" + "".join(f" {key} = {q(value)}\n" for key, value in items.items()) + "}\n")
639 filename = file.name
640 try:
641 command("nomad", "var", "put", "-in=hcl", "-out=none", *([f"-check-index={index}"] if index is not None else []), path, "@" + filename)
642 finally:
643 Path(filename).unlink()
644
645
646def configure(data):
647 if not data.get("setup"):
648 return
649 if data["setup"].startswith("tools/"):
650 script = (REPO / data["setup"]).resolve()
651 if not script.is_relative_to(REPO / "tools") or not script.is_file():
652 raise ValueError(f"invalid shared setup script: {data['setup']}")
653 else:
654 script = config_path(data.get("sourceId", data["id"]), data["setup"])
655 routes = [task["http"] for task in data["containers"].values()
656 if task.get("http") and task["http"].get("hostname")]
657 if len(routes) != 1:
658 raise ValueError("service setup requires one HTTP hostname")
659 route = routes[0]
660 check_http(route["hostname"], route["checkPath"], route["tlsInternal"], attempts=600)
661 variable = get_variable("nomad/jobs/" + data["id"])
662 own = variable["Items"] if variable else {}
663 command("python3", str(script), input=json.dumps({
664 "serviceId": data["id"], "host": route["hostname"], "hostRoot": data["hostRoot"],
665 "preview": data["id"] != data.get("sourceId", data["id"]),
666 "ownerEmail": data["ownerEmail"], **own,
667 }))
668 print(f"Configured {data['id']}")
669
670
671def prepare(data):
672 if data.get("prepare"):
673 script = config_path(data.get("sourceId", data["id"]), data["prepare"])
674 command("python3", str(script), input=json.dumps({
675 "hostRoot": data["hostRoot"], "uid": data["uid"],
676 }))
677
678
679def provision_inputs(data, definitions, stage_id=None, postgres_source=None):
680 for alias, request in data["inputs"].items():
681 if not NAME.fullmatch(alias):
682 raise ValueError(f"invalid input alias: {alias}")
683 provider = definitions[request["provider"]]
684 if not provider.get("provide"):
685 raise ValueError(f"{provider['id']} does not provide inputs")
686 path = f"nomad/jobs/{data['id']}/inputs/{alias}"
687 variable = get_variable(path)
688 own = (get_variable("nomad/jobs/" + provider["id"]) or {}).get("Items", {})
689 http = next((task["http"] for task in provider["containers"].values()
690 if task.get("http") and task["http"].get("hostname")), None)
691 if http:
692 check_http(http["hostname"], http["checkPath"], http["tlsInternal"])
693 payload = {
694 "request": request,
695 "existing": variable["Items"] if variable else None,
696 "providerSecrets": own,
697 "host": http["hostname"] if http else None,
698 }
699 if stage_id:
700 payload["stageId"] = stage_id
701 if postgres_source and provider["id"] == "postgres":
702 payload["sourceContainer"] = postgres_source
703 script = config_path(provider["id"], provider["provide"])
704 result = command("python3", str(script), input=json.dumps(payload), capture=True)
705 values = json.loads(result.stdout)
706 if not isinstance(values, dict) or not values or any(
707 not re.fullmatch(r"[A-Za-z_][A-Za-z0-9_]*", key) or not isinstance(value, str) or not value
708 for key, value in values.items()
709 ):
710 raise ValueError(f"invalid outputs from {provider['id']}")
711 if not variable or variable["Items"] != values:
712 put_variable(path, values, variable["ModifyIndex"] if variable else None)
713 print(f"Allocated {data['id']}.{alias} from {provider['id']}")
714
715
716def bootstrap(all_data):
717 if os.geteuid() != 0:
718 raise PermissionError("bootstrap requires root")
719 STATE.mkdir(parents=True, exist_ok=True)
720 if not (STATE / "nomad.token").exists():
721 result = command("nomad", "acl", "bootstrap", "-json", capture=True)
722 write_private(STATE / "nomad.token", json.loads(result.stdout)["SecretID"] + "\n")
723 token()
724 command(
725 "nomad",
726 "acl",
727 "policy",
728 "apply",
729 "studio-router",
730 str(REPO / "config/policies/router.hcl"),
731 )
732 command("nomad", "acl", "policy", "apply", "studio-dashboard",
733 str(REPO / "config/policies/dashboard.hcl"))
734 router_token = STATE / "router.token"
735 if not router_token.exists():
736 result = command(
737 "nomad",
738 "acl",
739 "token",
740 "create",
741 "-policy=studio-router",
742 "-name=studio-router",
743 "-json",
744 capture=True,
745 )
746 write_private(router_token, json.loads(result.stdout)["SecretID"] + "\n")
747 dashboard_token = STATE / "dashboard.token"
748 if not dashboard_token.exists():
749 result = command("nomad", "acl", "token", "create", "-policy=studio-dashboard",
750 "-name=studio-dashboard", "-json", capture=True)
751 write_private(dashboard_token, json.loads(result.stdout)["SecretID"] + "\n")
752 for data in all_data:
753 paths = {
754 f"nomad/jobs/{data['id']}/inputs/{alias}" for alias in data["inputs"]
755 }
756 if paths and data["containers"]:
757 policy = (
758 'namespace "default" {\n variables {\n'
759 + "".join(
760 f' path {q(path)} {{ capabilities = ["read"] }}\n'
761 for path in sorted(paths)
762 )
763 + " }\n}\n"
764 )
765 with tempfile.NamedTemporaryFile("w", delete=False) as file:
766 file.write(policy)
767 path = file.name
768 try:
769 for task_name in data["containers"]:
770 command(
771 "nomad", "acl", "policy", "apply", "-namespace", "default",
772 "-job", data["id"], "-group", "app", "-task", task_name,
773 data["id"] + "-" + task_name + "-imports", path,
774 )
775 finally:
776 Path(path).unlink()
777 if (Path("/opt/studio/current")).is_symlink():
778 command("systemctl", "start", "studio-router.service")
779
780
781def destroy_stage(stage, metadata, definitions):
782 subprocess.run(["nomad", "job", "stop", "-purge", "-yes", stage], check=False)
783 for alias, request in metadata["inputs"].items():
784 provider = definitions[request["provider"]]
785 own = (get_variable("nomad/jobs/" + provider["id"]) or {}).get("Items", {})
786 variable = get_variable(f"nomad/jobs/{stage}/inputs/{alias}")
787 existing = variable["Items"] if variable else None
788 hosts = [task["http"]["hostname"] for task in provider["containers"].values()
789 if task.get("http") and task["http"].get("hostname")]
790 script = config_path(provider["id"], provider["provide"])
791 command("python3", str(script), input=json.dumps({
792 "operation": "delete", "request": request, "stageId": stage,
793 "providerSecrets": own, "host": hosts[0] if hosts else None,
794 "existing": existing,
795 }))
796 datasets = list(metadata.get("external", {}).values())
797 if metadata.get("clone"):
798 datasets.append({"clone": metadata["clone"], "snapshot": metadata.get("snapshot")})
799 for dataset in datasets:
800 for attempt in range(30):
801 if subprocess.run(["zfs", "destroy", dataset["clone"]], capture_output=True).returncode == 0:
802 break
803 time.sleep(1)
804 else:
805 raise RuntimeError("preview dataset is still in use")
806 if dataset.get("snapshot"):
807 command("zfs", "destroy", dataset["snapshot"])
808 if not metadata.get("clone"):
809 shutil.rmtree(metadata["mount"], ignore_errors=True)
810 for alias in metadata["inputs"]:
811 subprocess.run(["nomad", "var", "purge", f"nomad/jobs/{stage}/inputs/{alias}"], check=False)
812 subprocess.run(["nomad", "var", "purge", "nomad/jobs/" + stage], check=False)
813 for task in metadata["tasks"]:
814 subprocess.run(["nomad", "acl", "policy", "delete", stage + "-" + task + "-imports"], check=False)
815 for path in (
816 STATE / "stages" / f"{stage}.json",
817 STATE / "routes" / f"{stage}.json",
818 STATE / "secrets" / f"{stage}.nv.hcl",
819 ):
820 path.unlink(missing_ok=True)
821 shutil.rmtree(ASSETS / stage, ignore_errors=True)
822
823
824def check_pool(data):
825 if os.geteuid() != 0:
826 raise PermissionError("pool setup requires root")
827 clover = subprocess.run(
828 ["findmnt", "-n", "-o", "SOURCE,FSTYPE", "--mountpoint", data[0]["cloverRoot"]],
829 capture_output=True, text=True,
830 )
831 clover_info = clover.stdout.split()
832 if clover.returncode or len(clover_info) != 2 or clover_info[1] != "zfs":
833 raise ValueError(f"clover must be a mounted ZFS dataset: {data[0]['cloverRoot']}")
834 media = data[0]["mediaRoot"]
835 mounted = subprocess.run(
836 ["findmnt", "-n", "-o", "SOURCE,FSTYPE", "--mountpoint", media],
837 capture_output=True, text=True,
838 )
839 media_info = mounted.stdout.split()
840 expected = {"fuse"} if os.environ.get("STUDIO_MEDIA_READ_ONLY") == "true" else {"zfs"}
841 if mounted.returncode or len(media_info) != 2 or media_info[1] not in expected or media_info[0] == clover_info[0]:
842 raise ValueError(f"media must be a separate mounted dataset: {media}")
843 for path, info in ((data[0]["cloverRoot"], clover_info), (media, media_info)):
844 if info[1] == "zfs":
845 acltype = command("zfs", "get", "-H", "-o", "value", "acltype", info[0], capture=True).stdout.strip()
846 if acltype == "nfsv4":
847 raise ValueError(f"{path} uses NFSv4 ACLs, which NixOS cannot enforce. Rehearse a POSIX permission conversion on a ZFS clone before migration.")
848 pool = data[0]["pool"]
849 if subprocess.run(["zpool", "list", pool], capture_output=True).returncode:
850 raise ValueError(f"ZFS pool is unavailable: {pool}")
851 prod = pool + "/prod"
852 acl_source = prod if subprocess.run(["zfs", "list", prod], capture_output=True).returncode == 0 else pool
853 if command("zfs", "get", "-H", "-o", "value", "acltype", acl_source, capture=True).stdout.strip() == "nfsv4":
854 raise ValueError(f"{acl_source} uses NFSv4 ACLs; service datasets need POSIX ACLs on NixOS.")
855 encrypted = any(
856 command("zfs", "get", "-H", "-o", "value", "encryption", info[0], capture=True).stdout.strip() != "off"
857 for info in (clover_info, media_info) if info[1] == "zfs"
858 )
859 if encrypted:
860 root = str(Path(data[0]["hostRoot"]).parent)
861 mounted = subprocess.run(["findmnt", "-n", "-o", "SOURCE", "--mountpoint", root], capture_output=True, text=True)
862 encryption = subprocess.run(["zfs", "get", "-H", "-o", "value", "encryption", prod], capture_output=True, text=True)
863 if encryption.returncode or encryption.stdout.strip() == "off" or mounted.stdout.strip() != prod:
864 raise ValueError(f"Create and mount encrypted {prod} at {root} before deployment; the existing storage is encrypted.")
865
866
867def check_staging_pool(data):
868 check_pool([data])
869 prod = data["pool"] + "/prod"
870 staging = data["pool"] + "/staging"
871 acl_source = staging if subprocess.run(["zfs", "list", staging], capture_output=True).returncode == 0 else data["pool"]
872 if command("zfs", "get", "-H", "-o", "value", "acltype", acl_source, capture=True).stdout.strip() == "nfsv4":
873 raise ValueError(f"{acl_source} uses NFSv4 ACLs; staged datasets need POSIX ACLs on NixOS.")
874 encryption = subprocess.run(["zfs", "get", "-H", "-o", "value", "encryption", prod], capture_output=True, text=True)
875 if encryption.returncode or encryption.stdout.strip() == "off":
876 return
877 mounted = subprocess.run(["findmnt", "-n", "-o", "SOURCE", "--mountpoint", data["stagingRoot"]], capture_output=True, text=True)
878 staging_encryption = subprocess.run(["zfs", "get", "-H", "-o", "value", "encryption", staging], capture_output=True, text=True)
879 if staging_encryption.returncode or staging_encryption.stdout.strip() == "off" or mounted.stdout.strip() != staging:
880 raise ValueError(f"Create and mount encrypted {staging} at {data['stagingRoot']} before staging.")
881
882
883def ensure_pool(data):
884 check_pool(data)
885 STATE.mkdir(parents=True, exist_ok=True)
886 pool = data[0]["pool"]
887 for item in data:
888 if item["hasManagedVolumes"]:
889 dataset = pool + "/prod/" + item["id"]
890 if subprocess.run(["zfs", "list", dataset], capture_output=True).returncode:
891 root = Path(item["hostRoot"])
892 if root.exists() and any(root.iterdir()):
893 raise ValueError(f"refusing to mount over {root}")
894 if subprocess.run(["zfs", "list", pool + "/prod"], capture_output=True).returncode:
895 parent = str(root.parent)
896 if Path(parent).exists() and any(Path(parent).iterdir()):
897 raise ValueError(f"refusing to mount over {parent}")
898 command("zfs", "create", "-o", "mountpoint=" + parent, pool + "/prod")
899 command("zfs", "create", "-o", "mountpoint=" + str(root), dataset)
900
901
902def main():
903 parser = argparse.ArgumentParser()
904 parser.add_argument(
905 "mode",
906 choices=[
907 "render",
908 "secrets",
909 "validate",
910 "check",
911 "bootstrap",
912 "deploy",
913 "preflight",
914 "pool",
915 "allocate",
916 "stage",
917 "destroy",
918 ],
919 )
920 parser.add_argument("name", nargs="?")
921 parser.add_argument("--base-domain")
922 parser.add_argument("--root")
923 parser.add_argument("--pool")
924 parser.add_argument("--media-root")
925 parser.add_argument("--env", action="append", default=[])
926 parser.add_argument("--key", action="append", default=[])
927 args = parser.parse_args()
928 if args.mode == "allocate" and not args.name:
929 parser.error("allocate requires a service name")
930 if args.key and args.mode != "secrets":
931 parser.error("--key is available only for secrets")
932 if args.name and args.mode != "destroy" and args.name not in service_files():
933 raise ValueError(f"unknown service: {args.name}")
934 properties = {
935 key: value
936 for key, value in (("domain", args.base_domain), ("root", args.root), ("pool", args.pool), ("mediaRoot", args.media_root))
937 if value
938 }
939 names = [args.name] if args.name else services()
940 if args.mode == "secrets":
941 if not args.name or not NAME.fullmatch(args.name):
942 parser.error("secrets requires a service name")
943 data = load(args.name, properties)
944 declared = [*data["requiredSecrets"], *data["secrets"]]
945 fields = args.key or declared
946 if not declared or len(set(declared)) != len(declared) or any(not re.fullmatch(r"[A-Za-z_][A-Za-z0-9_]*", field) for field in declared):
947 raise ValueError(f"invalid required secrets for {args.name}")
948 if len(set(fields)) != len(fields) or any(field not in declared for field in fields):
949 raise ValueError(f"unknown or duplicate secret key for {args.name}")
950 values = sys.stdin.read().splitlines()
951 if len(values) != len(fields) or any(not value for value in values):
952 raise ValueError(f"expected {len(fields)} nonempty secret lines for {args.name}")
953 token()
954 path = "nomad/jobs/" + args.name
955 existing = get_variable(path)
956 items = dict(existing["Items"]) if existing else {}
957 items.update(zip(fields, values))
958 put_variable(path, items, existing["ModifyIndex"] if existing else 0)
959 print(f"Imported {len(fields)} secrets for {args.name}")
960 return
961 if args.mode == "destroy":
962 if not args.name or not NAME.fullmatch(args.name):
963 parser.error("destroy requires a stage ID")
964 token()
965 metadata = json.loads((STATE / "stages" / f"{args.name}.json").read_text())
966 names = {metadata["sourceId"]} | {request["provider"] for request in metadata["inputs"].values()}
967 definitions = {name: load(name, properties) for name in names}
968 destroy_stage(args.name, metadata, definitions)
969 return
970 if args.mode == "stage":
971 if not args.name:
972 parser.error("stage requires a service name")
973 token()
974 data = load(args.name, properties)
975 check_staging_pool(data)
976 definitions = {args.name: data}
977 definitions.update({provider: load(provider, properties) for provider in dependencies(data)})
978 if not data["containers"]:
979 raise ValueError(f"service has no container to stage: {args.name}")
980 stage_dir = STATE / "stages"
981 previous = []
982 for path in stage_dir.glob("*.json"):
983 saved = json.loads(path.read_text())
984 if saved.get("sourceId") == args.name:
985 previous.append((path, saved))
986 if len(previous) > 1:
987 raise ValueError(f"multiple stages exist for {args.name}; destroy one first")
988 created = not previous
989 stage = f"{args.name}-preview-{secrets.token_hex(4)}" if created else previous[0][0].stem
990 http = [task["http"] for task in data["containers"].values() if task.get("http") and task["http"].get("hostname")]
991 if len(http) > 1:
992 raise ValueError("preview requires at most one HTTP route")
993 original_host = http[0].get("hostname") if http else None
994 hostname = stage + "." + original_host.split(".", 1)[1] if original_host else None
995 data = load(args.name, properties, stage)
996 data["rollout"] = "simple"
997 data["hostRoot"] = str(Path(data["stagingRoot"]) / stage)
998 if data["stageIsolation"] == "fresh" and any(
999 volume.get("src") and not volume["readOnly"]
1000 for task in data["containers"].values()
1001 for volume in task["volumes"].values()
1002 ):
1003 raise ValueError("fresh stages cannot use writable external mounts")
1004 for task in data["containers"].values():
1005 if task.get("tcp") and not task["hostNetwork"]:
1006 task["tcp"]["hostPort"] = None
1007 if hostname:
1008 for task in data["containers"].values():
1009 if task.get("http") and task["http"].get("hostname") == original_host:
1010 task["http"]["hostname"] = hostname
1011 task["http"]["plainHostnames"] = [stage + "-" + host for host in task["http"]["plainHostnames"]]
1012 task["env"] = {key: value.replace(original_host, hostname) for key, value in task["env"].items()}
1013 for assignment in args.env:
1014 key, separator, value = assignment.partition("=")
1015 if not separator or len(data["containers"]) != 1:
1016 parser.error(f"unknown environment override: {key}")
1017 env = next(iter(data["containers"].values()))["env"]
1018 if key not in env:
1019 parser.error(f"unknown environment override: {key}")
1020 env[key] = value
1021 metadata = (
1022 {"mount": data["hostRoot"], "inputs": data["inputs"],
1023 "tasks": list(data["containers"]), "sourceId": args.name,
1024 "stageIsolation": data["stageIsolation"]}
1025 if created else previous[0][1]
1026 )
1027 if not created:
1028 if metadata.get("stageIsolation", "clone") != data["stageIsolation"]:
1029 raise ValueError("stage isolation changed; destroy the stage before recreating it")
1030 if bool(metadata.get("clone")) != data["hasManagedVolumes"]:
1031 raise ValueError("storage layout changed; destroy the stage before recreating it")
1032 for alias, request in data["inputs"].items():
1033 if alias in metadata["inputs"] and metadata["inputs"][alias] != request:
1034 raise ValueError(f"requirement {alias} changed; destroy the stage before recreating it")
1035 if alias not in metadata["inputs"] and request["provider"] == "postgres":
1036 raise ValueError("new database requirement needs a new stage")
1037 metadata["inputs"].update(data["inputs"])
1038 metadata["tasks"] = sorted(set(metadata["tasks"]) | set(data["containers"]))
1039 stage_dir.mkdir(parents=True, exist_ok=True)
1040 postgres_snapshot = None
1041 pending_snapshots = set()
1042 try:
1043 postgres_inputs = created and data["stageIsolation"] == "clone" and any(
1044 request["provider"] == "postgres" and request["kind"] == "database"
1045 for request in data["inputs"].values()
1046 )
1047 snapshots = []
1048 source = None
1049 if created and data["hasManagedVolumes"]:
1050 pool = data["pool"]
1051 source = subprocess.run(
1052 ["findmnt", "-n", "-o", "SOURCE", "--mountpoint", definitions[args.name]["hostRoot"]],
1053 capture_output=True, text=True,
1054 ).stdout.strip()
1055 dataset = pool + "/staging/" + stage
1056 if subprocess.run(["zfs", "list", pool + "/staging"], capture_output=True).returncode:
1057 command("zfs", "create", "-o", "mountpoint=none", pool + "/staging")
1058 if source and data["stageIsolation"] == "clone":
1059 if source.split("/", 1)[0] != pool:
1060 raise ValueError(f"production dataset belongs to another pool: {source}")
1061 snapshot = source + "@" + stage
1062 snapshots.append(snapshot)
1063 if postgres_inputs:
1064 postgres_dataset = dataset_for("postgres")
1065 if not postgres_dataset:
1066 raise ValueError("Postgres needs a ZFS dataset for consistent stages")
1067 postgres_snapshot = postgres_dataset + "@" + stage
1068 snapshots.append(postgres_snapshot)
1069 external = {}
1070 for task in data["containers"].values():
1071 for volume in task["volumes"].values():
1072 if volume.get("src") and not volume["readOnly"]:
1073 external.setdefault(volume["src"], []).append(volume)
1074 external_mounts = {}
1075 for external_source in external:
1076 path = Path(external_source).resolve(strict=True)
1077 mount = json.loads(command(
1078 "findmnt", "--json", "-o", "SOURCE,TARGET,FSTYPE,OPTIONS", "--target", str(path), capture=True,
1079 ).stdout)["filesystems"][0]
1080 if "ro" not in mount["options"].split(",") and mount["fstype"] != "zfs":
1081 raise ValueError(f"writable external mount is not ZFS: {external_source}")
1082 external_mounts[external_source] = (path, mount)
1083 if created and mount["fstype"] == "zfs" and "ro" not in mount["options"].split(","):
1084 target = mount["source"] + "@" + stage
1085 if target not in snapshots:
1086 snapshots.append(target)
1087 if snapshots:
1088 if len({target.split("/", 1)[0] for target in snapshots}) != 1:
1089 raise ValueError("staged datasets must belong to one ZFS pool")
1090 command("zfs", "snapshot", *snapshots)
1091 pending_snapshots.update(snapshots)
1092 if created and data["hasManagedVolumes"]:
1093 if source and data["stageIsolation"] == "clone":
1094 try:
1095 clone_dataset(source, snapshot, dataset, data["hostRoot"])
1096 except Exception:
1097 command("zfs", "destroy", snapshot)
1098 pending_snapshots.remove(snapshot)
1099 raise
1100 metadata["snapshot"] = snapshot
1101 pending_snapshots.remove(snapshot)
1102 else:
1103 command("zfs", "create", "-o", "mountpoint=" + data["hostRoot"], dataset)
1104 metadata["clone"] = dataset
1105 bindings = {}
1106 for source, volumes in sorted(external.items()):
1107 path, mount = external_mounts[source]
1108 if "ro" in mount["options"].split(","):
1109 binding = {"src": source, "readOnly": True}
1110 elif mount["fstype"] == "zfs":
1111 dataset = mount["source"]
1112 if metadata.get("snapshot") == dataset + "@" + stage:
1113 clone_mount = data["hostRoot"]
1114 else:
1115 staged = metadata.setdefault("external", {})
1116 if dataset not in staged:
1117 if not created:
1118 raise ValueError("external storage changed; destroy the stage before recreating it")
1119 suffix = hashlib.sha256(dataset.encode()).hexdigest()[:8]
1120 pool = dataset.split("/", 1)[0]
1121 clone = pool + "/staging/" + stage + "-external-" + suffix
1122 clone_mount = str(Path(data["stagingRoot"]) / (stage + "-external-" + suffix))
1123 snapshot = dataset + "@" + stage
1124 if subprocess.run(["zfs", "list", pool + "/staging"], capture_output=True).returncode:
1125 command("zfs", "create", "-o", "mountpoint=none", pool + "/staging")
1126 try:
1127 clone_dataset(dataset, snapshot, clone, clone_mount)
1128 except Exception:
1129 command("zfs", "destroy", snapshot)
1130 pending_snapshots.remove(snapshot)
1131 raise
1132 staged[dataset] = {"clone": clone, "snapshot": snapshot, "mount": clone_mount}
1133 pending_snapshots.remove(snapshot)
1134 clone_mount = staged[dataset]["mount"]
1135 relative = path.relative_to(Path(mount["target"]).resolve())
1136 binding = {"src": str(Path(clone_mount) / relative), "readOnly": False}
1137 else:
1138 raise ValueError(f"writable external mount is not ZFS: {source}")
1139 bindings[source] = binding
1140 for volume in volumes:
1141 if path.is_relative_to(Path(data["cloverRoot"]).resolve()):
1142 volume["clover"] = True
1143 volume.update(binding)
1144 if not created and bindings != metadata.get("externalSources", {}):
1145 raise ValueError("external storage changed; destroy the stage before recreating it")
1146 metadata["externalSources"] = bindings
1147 if REPO.parent == Path("/opt/studio/releases"):
1148 metadata["release"] = REPO.name
1149 metadata["ready"] = False
1150 write_private(stage_dir / f"{stage}.json", json.dumps(metadata))
1151 bootstrap([data])
1152 if created and data["stageIsolation"] == "clone" and (data["secrets"] or data["requiredSecrets"]):
1153 source = get_variable("nomad/jobs/" + args.name)
1154 if source:
1155 put_variable("nomad/jobs/" + stage, source["Items"], 0)
1156 if postgres_snapshot:
1157 postgres_clone = postgres_dataset.rsplit("/prod/", 1)[0] + "/staging/" + stage + "-postgres"
1158 postgres_mount = Path(data["stagingRoot"]) / (stage + "-postgres")
1159 source_context = cloned_postgres(postgres_snapshot, postgres_clone, postgres_mount)
1160 else:
1161 source_context = nullcontext(None)
1162 with source_context as postgres_source:
1163 provision_inputs(data, definitions, stage, postgres_source)
1164 provision(data)
1165 prepare(data)
1166 build_images(data)
1167 submit(data, "run", definitions)
1168 except Exception:
1169 if created:
1170 destroy_stage(stage, metadata, definitions)
1171 raise
1172 finally:
1173 for snapshot in sorted(pending_snapshots):
1174 command("zfs", "destroy", snapshot)
1175 if hostname:
1176 check_http(hostname, http[0]["checkPath"], http[0]["tlsInternal"])
1177 configure(data)
1178 if REPO.parent == Path("/opt/studio/releases"):
1179 metadata["ready"] = True
1180 metadata["overrides"] = [assignment.partition("=")[0] for assignment in args.env]
1181 write_private(stage_dir / f"{stage}.json", json.dumps(metadata))
1182 print(f"stage={stage}\nrelease={metadata.get('release', '')}\nhost={hostname}\ndataset={metadata.get('clone', '')}")
1183 return
1184 definitions = {name: load(name, properties) for name in services()}
1185 if args.name and args.name not in definitions:
1186 raise ValueError(f"unknown service: {args.name}")
1187 if not args.name and args.mode == "deploy":
1188 names = [name for name in names if definitions[name]["containers"]]
1189 selected = set(names)
1190 while True:
1191 required = selected | set().union(*(dependencies(definitions[name]) for name in selected))
1192 missing = required - definitions.keys()
1193 if missing:
1194 raise ValueError(f"unknown dependency: {sorted(missing)}")
1195 if required == selected:
1196 break
1197 selected = required
1198 data = ordered([definitions[name] for name in sorted(selected)])
1199 if args.mode == "render":
1200 for item in data:
1201 if item["containers"]:
1202 print(render(item, definitions))
1203 elif args.mode == "validate":
1204 token()
1205 for item in data:
1206 submit(item, "validate", definitions)
1207 elif args.mode == "check":
1208 for item in data:
1209 for task in item["containers"].values():
1210 if task.get("http") and task["http"].get("hostname"):
1211 check_http(task["http"]["hostname"], task["http"]["checkPath"], task["http"]["tlsInternal"])
1212 elif args.mode == "preflight":
1213 check_pool(data)
1214 token()
1215 missing_secrets = {}
1216 for item in data:
1217 if item["containers"] and item["requiredSecrets"]:
1218 variable = get_variable("nomad/jobs/" + item["id"])
1219 values = variable["Items"] if variable else {}
1220 missing = [name for name in item["requiredSecrets"] if not values.get(name)]
1221 if missing:
1222 missing_secrets[item["id"]] = missing
1223 if missing_secrets:
1224 raise ValueError("missing required secrets: " + "; ".join(
1225 f"{name}: {', '.join(keys)}" for name, keys in missing_secrets.items()))
1226 elif args.mode == "pool":
1227 ensure_pool(data)
1228 elif args.mode == "allocate":
1229 ensure_pool(data)
1230 bootstrap(data)
1231 for item in data:
1232 provision_inputs(item, definitions)
1233 provision(item)
1234 prepare(item)
1235 build_images(item)
1236 elif args.mode == "bootstrap":
1237 bootstrap(data)
1238 elif args.mode == "deploy":
1239 bootstrap(data)
1240 for item in data:
1241 provision_inputs(item, definitions)
1242 provision(item)
1243 prepare(item)
1244 build_images(item)
1245 submit(item, "run", definitions)
1246 configure(item)
1247
1248
1249if __name__ == "__main__":
1250 try:
1251 main()
1252 except (OSError, RuntimeError, ValueError, subprocess.CalledProcessError) as error:
1253 print(error, file=sys.stderr)
1254 sys.exit(1)