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() if item["provider"] != "snowglobe"}
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 endpoint.get("loopback") or (not task["hostNetwork"] and kind == "http" and endpoint.get("hostPort") is None):
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 native = request["provider"] == "snowglobe"
684 provider = {"id": "snowglobe", "provide": True, "containers": {}} if native else definitions[request["provider"]]
685 if not provider.get("provide"):
686 raise ValueError(f"{provider['id']} does not provide inputs")
687 path = f"nomad/jobs/{data['id']}/inputs/{alias}"
688 variable = get_variable(path)
689 own = (get_variable("nomad/jobs/" + provider["id"]) or {}).get("Items", {})
690 http = next((task["http"] for task in provider["containers"].values()
691 if task.get("http") and task["http"].get("hostname")), None)
692 if http:
693 check_http(http["hostname"], http["checkPath"], http["tlsInternal"])
694 payload = {
695 "request": request,
696 "existing": variable["Items"] if variable else None,
697 "providerSecrets": own,
698 "host": http["hostname"] if http else None,
699 }
700 if stage_id:
701 payload["stageId"] = stage_id
702 if postgres_source and provider["id"] == "postgres":
703 payload["sourceContainer"] = postgres_source
704 script = REPO / "tools/oidc-provider.py" if native else config_path(provider["id"], provider["provide"])
705 result = command("python3", str(script), input=json.dumps(payload), capture=True)
706 values = json.loads(result.stdout)
707 if not isinstance(values, dict) or not values or any(
708 not re.fullmatch(r"[A-Za-z_][A-Za-z0-9_]*", key) or not isinstance(value, str) or not value
709 for key, value in values.items()
710 ):
711 raise ValueError(f"invalid outputs from {provider['id']}")
712 if not variable or variable["Items"] != values:
713 put_variable(path, values, variable["ModifyIndex"] if variable else None)
714 print(f"Allocated {data['id']}.{alias} from {provider['id']}")
715
716
717def bootstrap(all_data):
718 if os.geteuid() != 0:
719 raise PermissionError("bootstrap requires root")
720 STATE.mkdir(parents=True, exist_ok=True)
721 if not (STATE / "nomad.token").exists():
722 result = command("nomad", "acl", "bootstrap", "-json", capture=True)
723 write_private(STATE / "nomad.token", json.loads(result.stdout)["SecretID"] + "\n")
724 token()
725 command(
726 "nomad",
727 "acl",
728 "policy",
729 "apply",
730 "studio-router",
731 str(REPO / "config/policies/router.hcl"),
732 )
733 command("nomad", "acl", "policy", "apply", "studio-dashboard",
734 str(REPO / "config/policies/dashboard.hcl"))
735 router_token = STATE / "router.token"
736 if not router_token.exists():
737 result = command(
738 "nomad",
739 "acl",
740 "token",
741 "create",
742 "-policy=studio-router",
743 "-name=studio-router",
744 "-json",
745 capture=True,
746 )
747 write_private(router_token, json.loads(result.stdout)["SecretID"] + "\n")
748 dashboard_token = STATE / "dashboard.token"
749 if not dashboard_token.exists():
750 result = command("nomad", "acl", "token", "create", "-policy=studio-dashboard",
751 "-name=studio-dashboard", "-json", capture=True)
752 write_private(dashboard_token, json.loads(result.stdout)["SecretID"] + "\n")
753 for data in all_data:
754 paths = {
755 f"nomad/jobs/{data['id']}/inputs/{alias}" for alias in data["inputs"]
756 }
757 if paths and data["containers"]:
758 policy = (
759 'namespace "default" {\n variables {\n'
760 + "".join(
761 f' path {q(path)} {{ capabilities = ["read"] }}\n'
762 for path in sorted(paths)
763 )
764 + " }\n}\n"
765 )
766 with tempfile.NamedTemporaryFile("w", delete=False) as file:
767 file.write(policy)
768 path = file.name
769 try:
770 for task_name in data["containers"]:
771 command(
772 "nomad", "acl", "policy", "apply", "-namespace", "default",
773 "-job", data["id"], "-group", "app", "-task", task_name,
774 data["id"] + "-" + task_name + "-imports", path,
775 )
776 finally:
777 Path(path).unlink()
778 if (Path("/opt/studio/current")).is_symlink():
779 command("systemctl", "start", "studio-router.service")
780
781
782def destroy_stage(stage, metadata, definitions):
783 subprocess.run(["nomad", "job", "stop", "-purge", "-yes", stage], check=False)
784 for alias, request in metadata["inputs"].items():
785 native = request["provider"] == "snowglobe"
786 provider = {"id": "snowglobe", "provide": True, "containers": {}} if native else definitions[request["provider"]]
787 own = (get_variable("nomad/jobs/" + provider["id"]) or {}).get("Items", {})
788 variable = get_variable(f"nomad/jobs/{stage}/inputs/{alias}")
789 existing = variable["Items"] if variable else None
790 hosts = [task["http"]["hostname"] for task in provider["containers"].values()
791 if task.get("http") and task["http"].get("hostname")]
792 script = REPO / "tools/oidc-provider.py" if native else config_path(provider["id"], provider["provide"])
793 command("python3", str(script), input=json.dumps({
794 "operation": "delete", "request": request, "stageId": stage,
795 "providerSecrets": own, "host": hosts[0] if hosts else None,
796 "existing": existing,
797 }))
798 datasets = list(metadata.get("external", {}).values())
799 if metadata.get("clone"):
800 datasets.append({"clone": metadata["clone"], "snapshot": metadata.get("snapshot")})
801 for dataset in datasets:
802 for attempt in range(30):
803 if subprocess.run(["zfs", "destroy", dataset["clone"]], capture_output=True).returncode == 0:
804 break
805 time.sleep(1)
806 else:
807 raise RuntimeError("preview dataset is still in use")
808 if dataset.get("snapshot"):
809 command("zfs", "destroy", dataset["snapshot"])
810 if not metadata.get("clone"):
811 shutil.rmtree(metadata["mount"], ignore_errors=True)
812 for alias in metadata["inputs"]:
813 subprocess.run(["nomad", "var", "purge", f"nomad/jobs/{stage}/inputs/{alias}"], check=False)
814 subprocess.run(["nomad", "var", "purge", "nomad/jobs/" + stage], check=False)
815 for task in metadata["tasks"]:
816 subprocess.run(["nomad", "acl", "policy", "delete", stage + "-" + task + "-imports"], check=False)
817 for path in (
818 STATE / "stages" / f"{stage}.json",
819 STATE / "routes" / f"{stage}.json",
820 STATE / "secrets" / f"{stage}.nv.hcl",
821 ):
822 path.unlink(missing_ok=True)
823 shutil.rmtree(ASSETS / stage, ignore_errors=True)
824
825
826def check_pool(data):
827 if os.geteuid() != 0:
828 raise PermissionError("pool setup requires root")
829 clover = subprocess.run(
830 ["findmnt", "-n", "-o", "SOURCE,FSTYPE", "--mountpoint", data[0]["cloverRoot"]],
831 capture_output=True, text=True,
832 )
833 clover_info = clover.stdout.split()
834 if clover.returncode or len(clover_info) != 2 or clover_info[1] != "zfs":
835 raise ValueError(f"clover must be a mounted ZFS dataset: {data[0]['cloverRoot']}")
836 media = data[0]["mediaRoot"]
837 mounted = subprocess.run(
838 ["findmnt", "-n", "-o", "SOURCE,FSTYPE", "--mountpoint", media],
839 capture_output=True, text=True,
840 )
841 media_info = mounted.stdout.split()
842 expected = {"fuse"} if os.environ.get("STUDIO_MEDIA_READ_ONLY") == "true" else {"zfs"}
843 if mounted.returncode or len(media_info) != 2 or media_info[1] not in expected or media_info[0] == clover_info[0]:
844 raise ValueError(f"media must be a separate mounted dataset: {media}")
845 for path, info in ((data[0]["cloverRoot"], clover_info), (media, media_info)):
846 if info[1] == "zfs":
847 acltype = command("zfs", "get", "-H", "-o", "value", "acltype", info[0], capture=True).stdout.strip()
848 if acltype == "nfsv4":
849 raise ValueError(f"{path} uses NFSv4 ACLs, which NixOS cannot enforce. Rehearse a POSIX permission conversion on a ZFS clone before migration.")
850 pool = data[0]["pool"]
851 if subprocess.run(["zpool", "list", pool], capture_output=True).returncode:
852 raise ValueError(f"ZFS pool is unavailable: {pool}")
853 prod = pool + "/prod"
854 acl_source = prod if subprocess.run(["zfs", "list", prod], capture_output=True).returncode == 0 else pool
855 if command("zfs", "get", "-H", "-o", "value", "acltype", acl_source, capture=True).stdout.strip() == "nfsv4":
856 raise ValueError(f"{acl_source} uses NFSv4 ACLs; service datasets need POSIX ACLs on NixOS.")
857 encrypted = any(
858 command("zfs", "get", "-H", "-o", "value", "encryption", info[0], capture=True).stdout.strip() != "off"
859 for info in (clover_info, media_info) if info[1] == "zfs"
860 )
861 if encrypted:
862 root = str(Path(data[0]["hostRoot"]).parent)
863 mounted = subprocess.run(["findmnt", "-n", "-o", "SOURCE", "--mountpoint", root], capture_output=True, text=True)
864 encryption = subprocess.run(["zfs", "get", "-H", "-o", "value", "encryption", prod], capture_output=True, text=True)
865 if encryption.returncode or encryption.stdout.strip() == "off" or mounted.stdout.strip() != prod:
866 raise ValueError(f"Create and mount encrypted {prod} at {root} before deployment; the existing storage is encrypted.")
867
868
869def check_staging_pool(data):
870 check_pool([data])
871 prod = data["pool"] + "/prod"
872 staging = data["pool"] + "/staging"
873 acl_source = staging if subprocess.run(["zfs", "list", staging], capture_output=True).returncode == 0 else data["pool"]
874 if command("zfs", "get", "-H", "-o", "value", "acltype", acl_source, capture=True).stdout.strip() == "nfsv4":
875 raise ValueError(f"{acl_source} uses NFSv4 ACLs; staged datasets need POSIX ACLs on NixOS.")
876 encryption = subprocess.run(["zfs", "get", "-H", "-o", "value", "encryption", prod], capture_output=True, text=True)
877 if encryption.returncode or encryption.stdout.strip() == "off":
878 return
879 mounted = subprocess.run(["findmnt", "-n", "-o", "SOURCE", "--mountpoint", data["stagingRoot"]], capture_output=True, text=True)
880 staging_encryption = subprocess.run(["zfs", "get", "-H", "-o", "value", "encryption", staging], capture_output=True, text=True)
881 if staging_encryption.returncode or staging_encryption.stdout.strip() == "off" or mounted.stdout.strip() != staging:
882 raise ValueError(f"Create and mount encrypted {staging} at {data['stagingRoot']} before staging.")
883
884
885def ensure_pool(data):
886 check_pool(data)
887 STATE.mkdir(parents=True, exist_ok=True)
888 pool = data[0]["pool"]
889 for item in data:
890 if item["hasManagedVolumes"]:
891 dataset = pool + "/prod/" + item["id"]
892 if subprocess.run(["zfs", "list", dataset], capture_output=True).returncode:
893 root = Path(item["hostRoot"])
894 if root.exists() and any(root.iterdir()):
895 raise ValueError(f"refusing to mount over {root}")
896 if subprocess.run(["zfs", "list", pool + "/prod"], capture_output=True).returncode:
897 parent = str(root.parent)
898 if Path(parent).exists() and any(Path(parent).iterdir()):
899 raise ValueError(f"refusing to mount over {parent}")
900 command("zfs", "create", "-o", "mountpoint=" + parent, pool + "/prod")
901 command("zfs", "create", "-o", "mountpoint=" + str(root), dataset)
902
903
904def main():
905 parser = argparse.ArgumentParser()
906 parser.add_argument(
907 "mode",
908 choices=[
909 "render",
910 "secrets",
911 "validate",
912 "check",
913 "bootstrap",
914 "deploy",
915 "preflight",
916 "pool",
917 "allocate",
918 "stage",
919 "destroy",
920 ],
921 )
922 parser.add_argument("name", nargs="?")
923 parser.add_argument("--base-domain")
924 parser.add_argument("--root")
925 parser.add_argument("--pool")
926 parser.add_argument("--media-root")
927 parser.add_argument("--env", action="append", default=[])
928 parser.add_argument("--key", action="append", default=[])
929 args = parser.parse_args()
930 if args.mode == "allocate" and not args.name:
931 parser.error("allocate requires a service name")
932 if args.key and args.mode != "secrets":
933 parser.error("--key is available only for secrets")
934 if args.name and args.mode != "destroy" and args.name not in service_files():
935 raise ValueError(f"unknown service: {args.name}")
936 properties = {
937 key: value
938 for key, value in (("domain", args.base_domain), ("root", args.root), ("pool", args.pool), ("mediaRoot", args.media_root))
939 if value
940 }
941 names = [args.name] if args.name else services()
942 if args.mode == "secrets":
943 if not args.name or not NAME.fullmatch(args.name):
944 parser.error("secrets requires a service name")
945 data = load(args.name, properties)
946 declared = [*data["requiredSecrets"], *data["secrets"]]
947 fields = args.key or declared
948 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):
949 raise ValueError(f"invalid required secrets for {args.name}")
950 if len(set(fields)) != len(fields) or any(field not in declared for field in fields):
951 raise ValueError(f"unknown or duplicate secret key for {args.name}")
952 values = sys.stdin.read().splitlines()
953 if len(values) != len(fields) or any(not value for value in values):
954 raise ValueError(f"expected {len(fields)} nonempty secret lines for {args.name}")
955 token()
956 path = "nomad/jobs/" + args.name
957 existing = get_variable(path)
958 items = dict(existing["Items"]) if existing else {}
959 items.update(zip(fields, values))
960 put_variable(path, items, existing["ModifyIndex"] if existing else 0)
961 print(f"Imported {len(fields)} secrets for {args.name}")
962 return
963 if args.mode == "destroy":
964 if not args.name or not NAME.fullmatch(args.name):
965 parser.error("destroy requires a stage ID")
966 token()
967 metadata = json.loads((STATE / "stages" / f"{args.name}.json").read_text())
968 names = {metadata["sourceId"]} | {request["provider"] for request in metadata["inputs"].values() if request["provider"] != "snowglobe"}
969 definitions = {name: load(name, properties) for name in names}
970 destroy_stage(args.name, metadata, definitions)
971 return
972 if args.mode == "stage":
973 if not args.name:
974 parser.error("stage requires a service name")
975 token()
976 data = load(args.name, properties)
977 check_staging_pool(data)
978 definitions = {args.name: data}
979 definitions.update({provider: load(provider, properties) for provider in dependencies(data)})
980 if not data["containers"]:
981 raise ValueError(f"service has no container to stage: {args.name}")
982 stage_dir = STATE / "stages"
983 previous = []
984 for path in stage_dir.glob("*.json"):
985 saved = json.loads(path.read_text())
986 if saved.get("sourceId") == args.name:
987 previous.append((path, saved))
988 if len(previous) > 1:
989 raise ValueError(f"multiple stages exist for {args.name}; destroy one first")
990 created = not previous
991 stage = f"{args.name}-preview-{secrets.token_hex(4)}" if created else previous[0][0].stem
992 http = [task["http"] for task in data["containers"].values() if task.get("http") and task["http"].get("hostname")]
993 if len(http) > 1:
994 raise ValueError("preview requires at most one HTTP route")
995 original_host = http[0].get("hostname") if http else None
996 hostname = stage + "." + original_host.split(".", 1)[1] if original_host else None
997 data = load(args.name, properties, stage)
998 data["rollout"] = "simple"
999 data["hostRoot"] = str(Path(data["stagingRoot"]) / stage)
1000 if data["stageIsolation"] == "fresh" and any(
1001 volume.get("src") and not volume["readOnly"]
1002 for task in data["containers"].values()
1003 for volume in task["volumes"].values()
1004 ):
1005 raise ValueError("fresh stages cannot use writable external mounts")
1006 for task in data["containers"].values():
1007 if task.get("tcp") and not task["hostNetwork"]:
1008 task["tcp"]["hostPort"] = None
1009 if hostname:
1010 for task in data["containers"].values():
1011 if task.get("http") and task["http"].get("hostname") == original_host:
1012 task["http"]["hostname"] = hostname
1013 task["http"]["plainHostnames"] = [stage + "-" + host for host in task["http"]["plainHostnames"]]
1014 task["env"] = {key: value.replace(original_host, hostname) for key, value in task["env"].items()}
1015 for request in data["inputs"].values():
1016 if request["provider"] == "snowglobe":
1017 request["redirectUris"] = [uri.replace(original_host, hostname) for uri in request["redirectUris"]]
1018 for assignment in args.env:
1019 key, separator, value = assignment.partition("=")
1020 if not separator or len(data["containers"]) != 1:
1021 parser.error(f"unknown environment override: {key}")
1022 env = next(iter(data["containers"].values()))["env"]
1023 if key not in env:
1024 parser.error(f"unknown environment override: {key}")
1025 env[key] = value
1026 metadata = (
1027 {"mount": data["hostRoot"], "inputs": data["inputs"],
1028 "tasks": list(data["containers"]), "sourceId": args.name,
1029 "stageIsolation": data["stageIsolation"]}
1030 if created else previous[0][1]
1031 )
1032 if not created:
1033 if metadata.get("stageIsolation", "clone") != data["stageIsolation"]:
1034 raise ValueError("stage isolation changed; destroy the stage before recreating it")
1035 if bool(metadata.get("clone")) != data["hasManagedVolumes"]:
1036 raise ValueError("storage layout changed; destroy the stage before recreating it")
1037 for alias, request in data["inputs"].items():
1038 if alias in metadata["inputs"] and metadata["inputs"][alias] != request:
1039 raise ValueError(f"requirement {alias} changed; destroy the stage before recreating it")
1040 if alias not in metadata["inputs"] and request["provider"] == "postgres":
1041 raise ValueError("new database requirement needs a new stage")
1042 metadata["inputs"].update(data["inputs"])
1043 metadata["tasks"] = sorted(set(metadata["tasks"]) | set(data["containers"]))
1044 stage_dir.mkdir(parents=True, exist_ok=True)
1045 postgres_snapshot = None
1046 pending_snapshots = set()
1047 try:
1048 postgres_inputs = created and data["stageIsolation"] == "clone" and any(
1049 request["provider"] == "postgres" and request["kind"] == "database"
1050 for request in data["inputs"].values()
1051 )
1052 snapshots = []
1053 source = None
1054 if created and data["hasManagedVolumes"]:
1055 pool = data["pool"]
1056 source = subprocess.run(
1057 ["findmnt", "-n", "-o", "SOURCE", "--mountpoint", definitions[args.name]["hostRoot"]],
1058 capture_output=True, text=True,
1059 ).stdout.strip()
1060 dataset = pool + "/staging/" + stage
1061 if subprocess.run(["zfs", "list", pool + "/staging"], capture_output=True).returncode:
1062 command("zfs", "create", "-o", "mountpoint=none", pool + "/staging")
1063 if source and data["stageIsolation"] == "clone":
1064 if source.split("/", 1)[0] != pool:
1065 raise ValueError(f"production dataset belongs to another pool: {source}")
1066 snapshot = source + "@" + stage
1067 snapshots.append(snapshot)
1068 if postgres_inputs:
1069 postgres_dataset = dataset_for("postgres")
1070 if not postgres_dataset:
1071 raise ValueError("Postgres needs a ZFS dataset for consistent stages")
1072 postgres_snapshot = postgres_dataset + "@" + stage
1073 snapshots.append(postgres_snapshot)
1074 external = {}
1075 for task in data["containers"].values():
1076 for volume in task["volumes"].values():
1077 if volume.get("src") and not volume["readOnly"]:
1078 external.setdefault(volume["src"], []).append(volume)
1079 external_mounts = {}
1080 for external_source in external:
1081 path = Path(external_source).resolve(strict=True)
1082 mount = json.loads(command(
1083 "findmnt", "--json", "-o", "SOURCE,TARGET,FSTYPE,OPTIONS", "--target", str(path), capture=True,
1084 ).stdout)["filesystems"][0]
1085 if "ro" not in mount["options"].split(",") and mount["fstype"] != "zfs":
1086 raise ValueError(f"writable external mount is not ZFS: {external_source}")
1087 external_mounts[external_source] = (path, mount)
1088 if created and mount["fstype"] == "zfs" and "ro" not in mount["options"].split(","):
1089 target = mount["source"] + "@" + stage
1090 if target not in snapshots:
1091 snapshots.append(target)
1092 if snapshots:
1093 if len({target.split("/", 1)[0] for target in snapshots}) != 1:
1094 raise ValueError("staged datasets must belong to one ZFS pool")
1095 command("zfs", "snapshot", *snapshots)
1096 pending_snapshots.update(snapshots)
1097 if created and data["hasManagedVolumes"]:
1098 if source and data["stageIsolation"] == "clone":
1099 try:
1100 clone_dataset(source, snapshot, dataset, data["hostRoot"])
1101 except Exception:
1102 command("zfs", "destroy", snapshot)
1103 pending_snapshots.remove(snapshot)
1104 raise
1105 metadata["snapshot"] = snapshot
1106 pending_snapshots.remove(snapshot)
1107 else:
1108 command("zfs", "create", "-o", "mountpoint=" + data["hostRoot"], dataset)
1109 metadata["clone"] = dataset
1110 bindings = {}
1111 for source, volumes in sorted(external.items()):
1112 path, mount = external_mounts[source]
1113 if "ro" in mount["options"].split(","):
1114 binding = {"src": source, "readOnly": True}
1115 elif mount["fstype"] == "zfs":
1116 dataset = mount["source"]
1117 if metadata.get("snapshot") == dataset + "@" + stage:
1118 clone_mount = data["hostRoot"]
1119 else:
1120 staged = metadata.setdefault("external", {})
1121 if dataset not in staged:
1122 if not created:
1123 raise ValueError("external storage changed; destroy the stage before recreating it")
1124 suffix = hashlib.sha256(dataset.encode()).hexdigest()[:8]
1125 pool = dataset.split("/", 1)[0]
1126 clone = pool + "/staging/" + stage + "-external-" + suffix
1127 clone_mount = str(Path(data["stagingRoot"]) / (stage + "-external-" + suffix))
1128 snapshot = dataset + "@" + stage
1129 if subprocess.run(["zfs", "list", pool + "/staging"], capture_output=True).returncode:
1130 command("zfs", "create", "-o", "mountpoint=none", pool + "/staging")
1131 try:
1132 clone_dataset(dataset, snapshot, clone, clone_mount)
1133 except Exception:
1134 command("zfs", "destroy", snapshot)
1135 pending_snapshots.remove(snapshot)
1136 raise
1137 staged[dataset] = {"clone": clone, "snapshot": snapshot, "mount": clone_mount}
1138 pending_snapshots.remove(snapshot)
1139 clone_mount = staged[dataset]["mount"]
1140 relative = path.relative_to(Path(mount["target"]).resolve())
1141 binding = {"src": str(Path(clone_mount) / relative), "readOnly": False}
1142 else:
1143 raise ValueError(f"writable external mount is not ZFS: {source}")
1144 bindings[source] = binding
1145 for volume in volumes:
1146 if path.is_relative_to(Path(data["cloverRoot"]).resolve()):
1147 volume["clover"] = True
1148 volume.update(binding)
1149 if not created and bindings != metadata.get("externalSources", {}):
1150 raise ValueError("external storage changed; destroy the stage before recreating it")
1151 metadata["externalSources"] = bindings
1152 if REPO.parent == Path("/opt/studio/releases"):
1153 metadata["release"] = REPO.name
1154 metadata["ready"] = False
1155 write_private(stage_dir / f"{stage}.json", json.dumps(metadata))
1156 bootstrap([data])
1157 if created and data["stageIsolation"] == "clone" and (data["secrets"] or data["requiredSecrets"]):
1158 source = get_variable("nomad/jobs/" + args.name)
1159 if source:
1160 put_variable("nomad/jobs/" + stage, source["Items"], 0)
1161 if postgres_snapshot:
1162 postgres_clone = postgres_dataset.rsplit("/prod/", 1)[0] + "/staging/" + stage + "-postgres"
1163 postgres_mount = Path(data["stagingRoot"]) / (stage + "-postgres")
1164 source_context = cloned_postgres(postgres_snapshot, postgres_clone, postgres_mount)
1165 else:
1166 source_context = nullcontext(None)
1167 with source_context as postgres_source:
1168 provision_inputs(data, definitions, stage, postgres_source)
1169 provision(data)
1170 prepare(data)
1171 build_images(data)
1172 submit(data, "run", definitions)
1173 except Exception:
1174 if created:
1175 destroy_stage(stage, metadata, definitions)
1176 raise
1177 finally:
1178 for snapshot in sorted(pending_snapshots):
1179 command("zfs", "destroy", snapshot)
1180 if hostname:
1181 check_http(hostname, http[0]["checkPath"], http[0]["tlsInternal"])
1182 configure(data)
1183 if REPO.parent == Path("/opt/studio/releases"):
1184 metadata["ready"] = True
1185 metadata["overrides"] = [assignment.partition("=")[0] for assignment in args.env]
1186 write_private(stage_dir / f"{stage}.json", json.dumps(metadata))
1187 print(f"stage={stage}\nrelease={metadata.get('release', '')}\nhost={hostname}\ndataset={metadata.get('clone', '')}")
1188 return
1189 definitions = {name: load(name, properties) for name in services()}
1190 if args.name and args.name not in definitions:
1191 raise ValueError(f"unknown service: {args.name}")
1192 if not args.name and args.mode == "deploy":
1193 names = [name for name in names if definitions[name]["containers"]]
1194 selected = set(names)
1195 while True:
1196 required = selected | set().union(*(dependencies(definitions[name]) for name in selected))
1197 missing = required - definitions.keys()
1198 if missing:
1199 raise ValueError(f"unknown dependency: {sorted(missing)}")
1200 if required == selected:
1201 break
1202 selected = required
1203 data = ordered([definitions[name] for name in sorted(selected)])
1204 if args.mode == "render":
1205 for item in data:
1206 if item["containers"]:
1207 print(render(item, definitions))
1208 elif args.mode == "validate":
1209 token()
1210 for item in data:
1211 submit(item, "validate", definitions)
1212 elif args.mode == "check":
1213 for item in data:
1214 for task in item["containers"].values():
1215 if task.get("http") and task["http"].get("hostname"):
1216 check_http(task["http"]["hostname"], task["http"]["checkPath"], task["http"]["tlsInternal"])
1217 elif args.mode == "preflight":
1218 check_pool(data)
1219 token()
1220 missing_secrets = {}
1221 for item in data:
1222 if item["containers"] and item["requiredSecrets"]:
1223 variable = get_variable("nomad/jobs/" + item["id"])
1224 values = variable["Items"] if variable else {}
1225 missing = [name for name in item["requiredSecrets"] if not values.get(name)]
1226 if missing:
1227 missing_secrets[item["id"]] = missing
1228 if missing_secrets:
1229 raise ValueError("missing required secrets: " + "; ".join(
1230 f"{name}: {', '.join(keys)}" for name, keys in missing_secrets.items()))
1231 elif args.mode == "pool":
1232 ensure_pool(data)
1233 elif args.mode == "allocate":
1234 ensure_pool(data)
1235 bootstrap(data)
1236 for item in data:
1237 provision_inputs(item, definitions)
1238 provision(item)
1239 prepare(item)
1240 build_images(item)
1241 elif args.mode == "bootstrap":
1242 bootstrap(data)
1243 elif args.mode == "deploy":
1244 bootstrap(data)
1245 for item in data:
1246 provision_inputs(item, definitions)
1247 provision(item)
1248 prepare(item)
1249 build_images(item)
1250 submit(item, "run", definitions)
1251 configure(item)
1252
1253
1254if __name__ == "__main__":
1255 try:
1256 main()
1257 except (OSError, RuntimeError, ValueError, subprocess.CalledProcessError) as error:
1258 print(error, file=sys.stderr)
1259 sys.exit(1)