from __future__ import annotations import asyncio import os from pathlib import Path import re import shutil import subprocess import threading from typing import Any, AsyncIterable, Callable from uuid import uuid4 from app.display.service import DisplayService from .diagnostics import failure_log_metadata, read_failure_log from .package import MAX_OTA_UPLOAD_BYTES, OtaPackageError, inspect_package from .state import read_json, read_last_result, utc_now, write_json from .versioning import SoftwareVersion _OTA_FILENAME = re.compile(r"^matrix-screen-controller-[0-9]+\.[0-9]+\.[0-9]+\.ota$") class OtaUploadError(RuntimeError): def __init__(self, message: str, *, status_code: int = 422) -> None: super().__init__(message) self.status_code = status_code class OtaBusyError(OtaUploadError): def __init__(self, message: str = "an OTA upload or update is already active") -> None: super().__init__(message, status_code=409) class OtaManager: def __init__( self, *, code_root: Path, data_root: Path, runtime_root: Path, software_version: SoftwareVersion, display: DisplayService, worker_starter: Callable[[], None] | None = None, staging_root: Path | None = None, ) -> None: self.code_root = Path(code_root) self.data_root = Path(data_root) self.runtime_root = Path(runtime_root) self.software_version = software_version self.display = display configured_staging = os.environ.get("MATRIX_OTA_STAGING_DIR") if staging_root is not None: self.staging_root = Path(staging_root) elif configured_staging: self.staging_root = Path(configured_staging) elif os.name != "nt" and str(self.code_root).startswith("/opt/"): self.staging_root = Path("/opt/matrix-screen-controller-ota") else: self.staging_root = self.runtime_root / "ota-staging" self.status_path = self.runtime_root / "ota-status.json" self.request_path = self.runtime_root / "ota-request.json" self._worker_starter = worker_starter or self._start_systemd_worker self._upload_lock = asyncio.Lock() self._monitor_stop = threading.Event() self._monitor_thread: threading.Thread | None = None self._completion_callback: Callable[[], None] = lambda: None def set_completion_callback(self, callback: Callable[[], None]) -> None: self._completion_callback = callback def status(self) -> dict[str, Any]: active_document = read_json(self.status_path) job = active_document.get("job") if isinstance(active_document, dict) else None component_pending = (self.data_root / "ota/component-transaction").exists() active = component_pending or bool(active_document and active_document.get("active") is True and isinstance(job, dict)) last_result = read_last_result(self.data_root) failure_log_available, failure_log_bytes = failure_log_metadata(self.data_root, last_result) return { "current_version": str(self.software_version), "max_upload_bytes": MAX_OTA_UPLOAD_BYTES, "active": active, "component_recovery_pending": component_pending, "job": dict(job) if active or isinstance(job, dict) else None, "last_result": last_result, "failure_log_available": failure_log_available, "failure_log_bytes": failure_log_bytes, } def failure_log(self) -> bytes | None: last_result = read_last_result(self.data_root) return read_failure_log(self.data_root, last_result) def is_active(self) -> bool: return self.status()["active"] or self._upload_lock.locked() async def accept_upload( self, *, filename: str, chunks: AsyncIterable[bytes], content_length: int | None, orientation: int, ) -> dict[str, Any]: if Path(filename).name != filename or not _OTA_FILENAME.fullmatch(filename): raise OtaUploadError( "OTA filename must be matrix-screen-controller-MAJOR.MINOR.PATCH.ota", status_code=400, ) if content_length is not None and (content_length <= 0 or content_length > MAX_OTA_UPLOAD_BYTES): raise OtaUploadError("OTA package exceeds the 256 MiB upload limit", status_code=413) if self._upload_lock.locked() or self.status()["active"]: raise OtaBusyError() async with self._upload_lock: if self.status()["active"]: raise OtaBusyError() upload_root = self.staging_root / "uploads" upload_root.mkdir(parents=True, exist_ok=True) temporary = upload_root / f".{uuid4().hex}.part" final: Path | None = None total = 0 try: with temporary.open("xb") as handle: os.chmod(temporary, 0o600) async for chunk in chunks: if not chunk: continue total += len(chunk) if total > MAX_OTA_UPLOAD_BYTES: raise OtaUploadError("OTA package exceeds the 256 MiB upload limit", status_code=413) handle.write(chunk) handle.flush() os.fsync(handle.fileno()) if total == 0: raise OtaUploadError("OTA package is empty", status_code=400) if content_length is not None and total != content_length: raise OtaUploadError("OTA upload length does not match Content-Length", status_code=400) try: package = inspect_package(temporary, current_version=self.software_version) except OtaPackageError as exc: message = str(exc) status_code = 409 if "must be newer" in message else 422 raise OtaUploadError(message, status_code=status_code) from exc expected_name = f"matrix-screen-controller-{package.target_version}.ota" if filename != expected_name: raise OtaUploadError("OTA filename version does not match its manifest", status_code=422) job_id = uuid4().hex final = upload_root / f"{job_id}.ota" os.replace(temporary, final) request = { "schema_version": 1, "job_id": job_id, "package_path": str(final), "status_path": str(self.status_path), "request_path": str(self.request_path), "current_version": str(self.software_version), "target_version": str(package.target_version), "packaged_at": package.created_at, "release_notes": package.release_notes, "orientation": int(orientation), "accepted_at": utc_now(), } job = { "id": job_id, "target_version": str(package.target_version), "packaged_at": package.created_at, "phase": "queued", "percent": 5, "message": "软件安装包已校验,准备检查并安装系统组件" if package.target_version.patch == 0 else "更新包已校验,准备安装", "started_at": request["accepted_at"], "finished_at": None, "error": None, } write_json(self.request_path, request) write_json(self.status_path, {"schema_version": 1, "active": True, "job": job}) self.display.start_ota_indicator(self.status_path) try: self._worker_starter() except Exception as exc: failed = dict(job) failed.update( phase="failed", percent=0, message="无法启动独立 OTA 工作服务", finished_at=utc_now(), error=str(exc), ) write_json(self.status_path, {"schema_version": 1, "active": False, "job": failed}) self.display.stop_ota_indicator() raise OtaUploadError("independent OTA worker could not be started", status_code=500) from exc self.resume_or_start_monitor(self._completion_callback) return {"accepted": True, "job": job} except BaseException: temporary.unlink(missing_ok=True) if final is not None and not self.status()["active"]: final.unlink(missing_ok=True) raise def resume_or_start_monitor(self, on_finished: Callable[[], None]) -> bool: if not self.status()["active"]: return False self.display.start_ota_indicator(self.status_path) if self._monitor_thread is not None and self._monitor_thread.is_alive(): return True self._monitor_stop.clear() def monitor() -> None: while not self._monitor_stop.wait(0.25): if self.status()["active"]: continue try: self.display.stop_ota_indicator() on_finished() except Exception: pass return self._monitor_thread = threading.Thread(target=monitor, name="ota-status-monitor", daemon=True) self._monitor_thread.start() return True def close(self) -> None: self._monitor_stop.set() thread = self._monitor_thread if thread is not None and thread is not threading.current_thread(): thread.join(timeout=2.0) @staticmethod def _start_systemd_worker() -> None: result = subprocess.run( ["systemctl", "start", "--no-block", "matrix-screen-controller-ota.service"], check=False, capture_output=True, text=True, timeout=10, ) if result.returncode != 0: raise RuntimeError((result.stderr or result.stdout or "systemctl start failed").strip())