Files
hashbro 9153a4f557 admin
2026-08-09 02:42:45 +08:00

466 lines
18 KiB
Python

"""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])
support_template = str(request_input.get("support_template") or "test")
command.extend(["--support-template", support_template])
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}"