Files
matrix-screen-controller/核桃派软件源代码/app/ota/manager.py
T

277 lines
12 KiB
Python

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 app.system.kernel_recovery import KernelRecovery
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,
"kernel_recovery": self.kernel_recovery_status(),
}
def _production_recovery(self) -> KernelRecovery | None:
if os.name == "nt" or self.data_root != Path("/var/lib/matrix-screen-controller"):
return None
return KernelRecovery(self.data_root)
def kernel_recovery_status(self) -> dict[str, str]:
recovery = self._production_recovery()
return recovery.status() if recovery else {"state": "not_needed", "reason": ""}
def schedule_kernel_recovery_after_ota(self) -> None:
recovery = self._production_recovery()
if recovery is None:
return
result = read_last_result(self.data_root)
if not result or result.get("status") != "success" or result.get("target_version") != str(self.software_version):
return
result_id = f"{result.get('target_version')}:{result.get('installed_at')}"
if not recovery.prepare_attempt(result_id):
return
command = [
"systemd-run", "--quiet", "--collect", "--unit=matrix-kernel-recovery",
"--on-active=5s", "/usr/bin/python3",
"/opt/matrix-screen-controller/app/system/kernel_recovery.py", "arm-and-reboot",
]
try:
subprocess.run(command, check=True, capture_output=True, timeout=10)
except (OSError, subprocess.CalledProcessError, subprocess.TimeoutExpired) as exc:
recovery.write_record("failed", f"could not schedule kernel recovery: {exc}", attempts=0)
def finalize_kernel_recovery_after_boot(self) -> None:
recovery = self._production_recovery()
if recovery is None:
return
thread = threading.Thread(target=recovery.finalize_after_boot, name="kernel-recovery-finalize", daemon=True)
thread.start()
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())