init
This commit is contained in:
@@ -0,0 +1,463 @@
|
||||
"""Channel build storage, locking, and builder orchestration."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import ctypes
|
||||
import errno
|
||||
import fcntl
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
import re
|
||||
import shutil
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
import uuid
|
||||
from contextlib import contextmanager
|
||||
from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
from typing import Any, Iterator
|
||||
|
||||
CHANNEL_RE = re.compile(r"^[0-9a-z]{32}$")
|
||||
DOMAIN_RE = re.compile(
|
||||
r"^(?=.{1,253}\.?$)(?:[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?\.)*"
|
||||
r"[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?\.?$"
|
||||
)
|
||||
|
||||
|
||||
class ApiError(Exception):
|
||||
def __init__(self, status: int, code: str, message: str):
|
||||
super().__init__(message)
|
||||
self.status = status
|
||||
self.code = code
|
||||
self.message = message
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class Settings:
|
||||
artifact_root: Path
|
||||
bearer_token: str
|
||||
project_script: Path
|
||||
python: str = sys.executable
|
||||
build_timeout: int = 900
|
||||
|
||||
@classmethod
|
||||
def from_env(cls) -> "Settings":
|
||||
repo_root = Path(__file__).resolve().parent.parent
|
||||
token = os.environ.get("BUILD_API_TOKEN", "")
|
||||
if not token:
|
||||
raise RuntimeError("BUILD_API_TOKEN is required")
|
||||
timeout = int(os.environ.get("BUILD_API_TIMEOUT", "900"))
|
||||
if timeout < 1:
|
||||
raise RuntimeError("BUILD_API_TIMEOUT must be positive")
|
||||
return cls(
|
||||
artifact_root=Path(
|
||||
os.environ.get("BUILD_API_ARTIFACT_ROOT", repo_root / "artifacts")
|
||||
),
|
||||
bearer_token=token,
|
||||
project_script=Path(
|
||||
os.environ.get(
|
||||
"BUILD_API_PROJECT_SCRIPT",
|
||||
repo_root / "frontend" / "tools" / "new_project.py",
|
||||
)
|
||||
),
|
||||
python=os.environ.get("BUILD_API_PYTHON", sys.executable),
|
||||
build_timeout=timeout,
|
||||
)
|
||||
|
||||
|
||||
class ChannelLock:
|
||||
def __init__(self, descriptor: int):
|
||||
self.descriptor = descriptor
|
||||
|
||||
def close(self) -> None:
|
||||
try:
|
||||
fcntl.flock(self.descriptor, fcntl.LOCK_UN)
|
||||
finally:
|
||||
os.close(self.descriptor)
|
||||
|
||||
|
||||
class BuildService:
|
||||
def __init__(self, settings: Settings):
|
||||
self.settings = settings
|
||||
configured_root = settings.artifact_root.expanduser().absolute()
|
||||
configured_root.mkdir(parents=True, exist_ok=True)
|
||||
if configured_root.is_symlink() or not configured_root.is_dir():
|
||||
raise RuntimeError("artifact root must be a real directory")
|
||||
self.root = configured_root.resolve(strict=True)
|
||||
self.channels = self.root / "channel"
|
||||
self.staging = self.root / "staging"
|
||||
self.locks = self.root / "locks"
|
||||
self._prepare_storage()
|
||||
|
||||
def _prepare_storage(self) -> None:
|
||||
for path in (self.channels, self.staging, self.locks):
|
||||
path.mkdir(mode=0o750, exist_ok=True)
|
||||
if path.is_symlink() or not path.is_dir():
|
||||
raise RuntimeError(f"managed path must be a real directory: {path}")
|
||||
root_device = self.root.stat().st_dev
|
||||
if any(path.stat().st_dev != root_device for path in (self.channels, self.staging)):
|
||||
raise RuntimeError("staging and channel directories must share a filesystem")
|
||||
|
||||
def build(
|
||||
self, channel_id: str, request_input: dict[str, Any], *, force: bool
|
||||
) -> tuple[int, dict[str, Any], list[tuple[str, str]]]:
|
||||
self.validate_channel(channel_id)
|
||||
digest = self._digest(request_input)
|
||||
request_id = uuid.uuid4().hex
|
||||
|
||||
with self._channel_lock(channel_id):
|
||||
destination = self._channel_path(channel_id)
|
||||
if self._path_exists(destination):
|
||||
self._assert_safe_tree(destination)
|
||||
current = self._read_manifest(destination)
|
||||
if current is None:
|
||||
raise ApiError(
|
||||
409,
|
||||
"existing_build_invalid",
|
||||
"channel release exists without a valid manifest",
|
||||
)
|
||||
if current.get("input_digest") == digest and not force:
|
||||
return 200, self._status_payload(channel_id, current), []
|
||||
if not force:
|
||||
raise ApiError(
|
||||
409,
|
||||
"build_conflict",
|
||||
"channel already exists with different build input",
|
||||
)
|
||||
|
||||
stage = self._staging_path(request_id)
|
||||
stage.mkdir(mode=0o750)
|
||||
try:
|
||||
completed = self._run_builder(stage, channel_id, request_input)
|
||||
self._validate_output(stage, channel_id)
|
||||
release = self._read_manifest(stage) or {}
|
||||
manifest = {
|
||||
**release,
|
||||
"schema_version": 1,
|
||||
"channel_id": channel_id,
|
||||
"request_id": request_id,
|
||||
"input_digest": digest,
|
||||
"input": request_input,
|
||||
"built_at": int(time.time()),
|
||||
"support_path": release.get(
|
||||
"support_path", f"/channel/{channel_id}/web/support.html"
|
||||
),
|
||||
"daily_path": release.get(
|
||||
"daily_path", f"/channel/{channel_id}/sync/daily.html"
|
||||
),
|
||||
"builder": {
|
||||
"exit_code": completed.returncode,
|
||||
"stdout_sha256": hashlib.sha256(
|
||||
completed.stdout.encode("utf-8", "replace")
|
||||
).hexdigest(),
|
||||
},
|
||||
}
|
||||
self._replace_json(stage / "manifest.json", manifest)
|
||||
self._publish(stage, destination, request_id)
|
||||
except subprocess.TimeoutExpired as exc:
|
||||
raise ApiError(500, "build_timeout", "builder timed out") from exc
|
||||
except subprocess.CalledProcessError as exc:
|
||||
message = self._builder_error(exc)
|
||||
raise ApiError(500, "build_failed", message) from exc
|
||||
finally:
|
||||
if stage.exists():
|
||||
self._remove_staging(stage)
|
||||
|
||||
return 201, self._status_payload(channel_id, manifest), [
|
||||
("Location", f"/v1/channels/{channel_id}")
|
||||
]
|
||||
|
||||
def status(self, channel_id: str) -> dict[str, Any]:
|
||||
self.validate_channel(channel_id)
|
||||
destination = self._channel_path(channel_id)
|
||||
if not self._path_exists(destination):
|
||||
raise ApiError(404, "not_found", "channel release not found")
|
||||
self._assert_safe_tree(destination)
|
||||
manifest = self._read_manifest(destination)
|
||||
if manifest is None:
|
||||
raise ApiError(500, "invalid_release", "channel manifest is invalid")
|
||||
return self._status_payload(channel_id, manifest)
|
||||
|
||||
def delete(self, channel_id: str) -> None:
|
||||
self.validate_channel(channel_id)
|
||||
with self._channel_lock(channel_id):
|
||||
destination = self._channel_path(channel_id)
|
||||
if self._path_exists(destination):
|
||||
self._safe_rmtree(destination)
|
||||
|
||||
def _run_builder(
|
||||
self, stage: Path, channel_id: str, request_input: dict[str, Any]
|
||||
) -> subprocess.CompletedProcess[str]:
|
||||
script = self.settings.project_script
|
||||
if script.is_symlink() or not script.is_file():
|
||||
raise ApiError(500, "builder_unavailable", "builder script is unavailable")
|
||||
# Staging is created empty by the API; new_project refuses an existing
|
||||
# --root unless --force is set.
|
||||
command = [
|
||||
self.settings.python,
|
||||
str(script),
|
||||
"--root",
|
||||
str(stage),
|
||||
"--channel-id",
|
||||
channel_id,
|
||||
"--force",
|
||||
]
|
||||
for domain in request_input["deployment_domains"]:
|
||||
command.extend(["--deployment-domains", domain])
|
||||
for domain in request_input["reporting_domains"]:
|
||||
command.extend(["--reporting-domains", domain])
|
||||
return subprocess.run(
|
||||
command,
|
||||
cwd=str(Path(__file__).resolve().parent.parent),
|
||||
check=True,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=self.settings.build_timeout,
|
||||
env={**os.environ, "PYTHONUNBUFFERED": "1"},
|
||||
)
|
||||
|
||||
def _publish(self, stage: Path, destination: Path, request_id: str) -> None:
|
||||
self._assert_safe_tree(stage)
|
||||
backup = self.channels / f".replaced-{destination.name}-{request_id}"
|
||||
moved_old = False
|
||||
try:
|
||||
if self._path_exists(destination):
|
||||
self._assert_safe_tree(destination)
|
||||
if self._exchange_directories(stage, destination):
|
||||
self._safe_rmtree(stage)
|
||||
self._fsync_directory(self.channels)
|
||||
return
|
||||
os.rename(destination, backup)
|
||||
moved_old = True
|
||||
os.rename(stage, destination)
|
||||
except Exception:
|
||||
if (
|
||||
moved_old
|
||||
and not self._path_exists(destination)
|
||||
and self._path_exists(backup)
|
||||
):
|
||||
os.rename(backup, destination)
|
||||
raise
|
||||
self._fsync_directory(self.channels)
|
||||
if self._path_exists(backup):
|
||||
self._safe_rmtree(backup)
|
||||
|
||||
@staticmethod
|
||||
def _exchange_directories(first: Path, second: Path) -> bool:
|
||||
"""Atomically exchange directories on Linux; return false if unsupported."""
|
||||
if not sys.platform.startswith("linux"):
|
||||
return False
|
||||
libc = ctypes.CDLL(None, use_errno=True)
|
||||
renameat2 = getattr(libc, "renameat2", None)
|
||||
if renameat2 is None:
|
||||
return False
|
||||
renameat2.argtypes = [
|
||||
ctypes.c_int,
|
||||
ctypes.c_char_p,
|
||||
ctypes.c_int,
|
||||
ctypes.c_char_p,
|
||||
ctypes.c_uint,
|
||||
]
|
||||
renameat2.restype = ctypes.c_int
|
||||
result = renameat2(
|
||||
-100,
|
||||
os.fsencode(first),
|
||||
-100,
|
||||
os.fsencode(second),
|
||||
2,
|
||||
)
|
||||
if result == 0:
|
||||
return True
|
||||
error = ctypes.get_errno()
|
||||
if error in (errno.ENOSYS, errno.EINVAL, errno.ENOTSUP):
|
||||
return False
|
||||
raise OSError(error, os.strerror(error))
|
||||
|
||||
def _validate_output(self, stage: Path, channel_id: str) -> None:
|
||||
required_dirs = (stage / "web", stage / "sync", stage / "out")
|
||||
if any(not path.is_dir() or path.is_symlink() for path in required_dirs):
|
||||
raise ApiError(
|
||||
500,
|
||||
"invalid_build_output",
|
||||
"builder must produce web, sync, and out directories",
|
||||
)
|
||||
required_files = (
|
||||
stage / "web" / "support.html",
|
||||
stage / "sync" / "daily.html",
|
||||
stage / "manifest.json",
|
||||
)
|
||||
if any(not path.is_file() or path.is_symlink() for path in required_files):
|
||||
raise ApiError(
|
||||
500,
|
||||
"invalid_build_output",
|
||||
"builder must produce support.html, daily.html, and manifest.json",
|
||||
)
|
||||
release = self._read_manifest(stage)
|
||||
if not release or release.get("channel_id") != channel_id:
|
||||
raise ApiError(
|
||||
500,
|
||||
"invalid_build_output",
|
||||
"builder manifest channel_id does not match request",
|
||||
)
|
||||
self._assert_safe_tree(stage)
|
||||
|
||||
@staticmethod
|
||||
def normalize_domains(value: list[str], field: str) -> list[str]:
|
||||
if not value or len(value) > 8:
|
||||
raise ApiError(
|
||||
422, "validation_error", f"{field} must contain 1 to 8 domains"
|
||||
)
|
||||
result: list[str] = []
|
||||
for item in value:
|
||||
domain = item.strip().lower()
|
||||
if not DOMAIN_RE.fullmatch(domain):
|
||||
raise ApiError(
|
||||
422, "validation_error", f"{field} contains an invalid domain"
|
||||
)
|
||||
result.append(domain.rstrip("."))
|
||||
if len(set(result)) != len(result):
|
||||
raise ApiError(
|
||||
422, "validation_error", f"{field} must not contain duplicates"
|
||||
)
|
||||
return result
|
||||
|
||||
@staticmethod
|
||||
def validate_channel(channel_id: str) -> None:
|
||||
if not CHANNEL_RE.fullmatch(channel_id):
|
||||
raise ApiError(
|
||||
422,
|
||||
"validation_error",
|
||||
"channel id must be 32 lowercase alphanumeric characters",
|
||||
)
|
||||
|
||||
@contextmanager
|
||||
def _channel_lock(self, channel_id: str) -> Iterator[ChannelLock]:
|
||||
path = self.locks / f"{channel_id}.lock"
|
||||
flags = os.O_CREAT | os.O_RDWR
|
||||
if hasattr(os, "O_NOFOLLOW"):
|
||||
flags |= os.O_NOFOLLOW
|
||||
try:
|
||||
descriptor = os.open(path, flags, 0o640)
|
||||
except OSError as exc:
|
||||
raise ApiError(409, "channel_locked", "channel lock is unavailable") from exc
|
||||
lock = ChannelLock(descriptor)
|
||||
try:
|
||||
try:
|
||||
fcntl.flock(descriptor, fcntl.LOCK_EX | fcntl.LOCK_NB)
|
||||
except OSError as exc:
|
||||
if exc.errno in (errno.EACCES, errno.EAGAIN):
|
||||
raise ApiError(
|
||||
409, "channel_locked", "another channel operation is in progress"
|
||||
) from exc
|
||||
raise
|
||||
yield lock
|
||||
finally:
|
||||
lock.close()
|
||||
|
||||
def _channel_path(self, channel_id: str) -> Path:
|
||||
return self._contained(self.channels, self.channels / channel_id)
|
||||
|
||||
def _staging_path(self, request_id: str) -> Path:
|
||||
return self._contained(self.staging, self.staging / request_id)
|
||||
|
||||
@staticmethod
|
||||
def _contained(parent: Path, child: Path) -> Path:
|
||||
if child.parent != parent or not child.is_relative_to(parent):
|
||||
raise ApiError(500, "unsafe_path", "managed path escaped its root")
|
||||
return child
|
||||
|
||||
@staticmethod
|
||||
def _assert_safe_tree(path: Path) -> None:
|
||||
if path.is_symlink() or not path.is_dir():
|
||||
raise ApiError(409, "unsafe_release", "managed release is not a real directory")
|
||||
for root, directories, files in os.walk(path, followlinks=False):
|
||||
for name in (*directories, *files):
|
||||
entry = Path(root) / name
|
||||
if entry.is_symlink():
|
||||
raise ApiError(
|
||||
409, "unsafe_release", "managed release contains a symbolic link"
|
||||
)
|
||||
|
||||
def _safe_rmtree(self, path: Path) -> None:
|
||||
self._assert_safe_tree(path)
|
||||
shutil.rmtree(path)
|
||||
|
||||
def _remove_staging(self, path: Path) -> None:
|
||||
"""Remove an unpublished builder tree without following child symlinks."""
|
||||
self._contained(self.staging, path)
|
||||
if path.is_symlink() or not path.is_dir():
|
||||
raise ApiError(500, "unsafe_path", "staging root is not a real directory")
|
||||
shutil.rmtree(path)
|
||||
|
||||
@staticmethod
|
||||
def _path_exists(path: Path) -> bool:
|
||||
return os.path.lexists(path)
|
||||
|
||||
@staticmethod
|
||||
def _fsync_directory(path: Path) -> None:
|
||||
descriptor = os.open(path, os.O_RDONLY)
|
||||
try:
|
||||
os.fsync(descriptor)
|
||||
finally:
|
||||
os.close(descriptor)
|
||||
|
||||
@staticmethod
|
||||
def _replace_json(path: Path, value: dict[str, Any]) -> None:
|
||||
"""Atomically overwrite JSON that the builder may already have created."""
|
||||
data = json.dumps(value, sort_keys=True, indent=2).encode() + b"\n"
|
||||
temporary = path.with_name(f".{path.name}.{uuid.uuid4().hex}.tmp")
|
||||
descriptor = os.open(temporary, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o640)
|
||||
try:
|
||||
with os.fdopen(descriptor, "wb") as stream:
|
||||
stream.write(data)
|
||||
stream.flush()
|
||||
os.fsync(stream.fileno())
|
||||
os.replace(temporary, path)
|
||||
except Exception:
|
||||
if temporary.exists():
|
||||
temporary.unlink()
|
||||
raise
|
||||
|
||||
@staticmethod
|
||||
def _read_manifest(destination: Path) -> dict[str, Any] | None:
|
||||
if destination.is_symlink() or not destination.is_dir():
|
||||
return None
|
||||
manifest_path = destination / "manifest.json"
|
||||
if manifest_path.is_symlink() or not manifest_path.is_file():
|
||||
return None
|
||||
try:
|
||||
value = json.loads(manifest_path.read_text(encoding="utf-8"))
|
||||
except (OSError, UnicodeDecodeError, json.JSONDecodeError):
|
||||
return None
|
||||
return value if isinstance(value, dict) else None
|
||||
|
||||
@staticmethod
|
||||
def _digest(value: dict[str, Any]) -> str:
|
||||
canonical = json.dumps(
|
||||
value, sort_keys=True, separators=(",", ":"), ensure_ascii=True
|
||||
).encode()
|
||||
return hashlib.sha256(canonical).hexdigest()
|
||||
|
||||
@staticmethod
|
||||
def _status_payload(channel_id: str, manifest: dict[str, Any]) -> dict[str, Any]:
|
||||
return {
|
||||
"status": "built",
|
||||
"channel_id": channel_id,
|
||||
"request_id": manifest.get("request_id"),
|
||||
"release_path": f"channel/{channel_id}",
|
||||
"support_path": manifest.get(
|
||||
"support_path", f"/channel/{channel_id}/web/support.html"
|
||||
),
|
||||
"daily_path": manifest.get(
|
||||
"daily_path", f"/channel/{channel_id}/sync/daily.html"
|
||||
),
|
||||
"manifest": manifest,
|
||||
}
|
||||
|
||||
@staticmethod
|
||||
def _builder_error(exc: subprocess.CalledProcessError) -> str:
|
||||
stderr = (exc.stderr or "").strip().splitlines()
|
||||
detail = stderr[-1][:500] if stderr else "builder exited unsuccessfully"
|
||||
return f"builder failed: {detail}"
|
||||
Reference in New Issue
Block a user