from __future__ import annotations import hashlib import os import traceback from typing import Any from app.config import settings from app.core.checksum import sha256_file from app.core.command_runner import CommandRunner from app.core.downloader import Downloader from app.core.docker_installer import DockerInstaller, image_reference from app.core.installer import AptInstaller, DebInstaller from app.core.manifest_client import ManifestClient from app.core.manifest_validator import ManifestValidator from app.core.service_manager import ServiceManager from app.models.schemas import InstallRequest, RemoveRequest, UpdateRequest from app.storage.repository import Repository, utc_now APT_SERVICE_POLICIES: dict[str, dict[str, Any]] = { "postgresql": { "managedServices": ["postgresql.service"], "readiness": "postgresql", }, } class InstalledComponentVerificationError(RuntimeError): """The package changed on disk, but its required runtime verification failed.""" class TaskRunner: def __init__(self, repository: Repository) -> None: self.repository = repository self.manifest_client = ManifestClient() self.manifest_validator = ManifestValidator() def run_install(self, task_id: str, request: InstallRequest | UpdateRequest, task_type: str = "install") -> None: manifest: dict[str, Any] | None = None try: self._mark_started(task_id, f"starting {task_type}") self._require_root_if_available() manifest = self._resolve_manifest(request) self.repository.add_log(task_id, "info", f"Installing {manifest['appId']} {manifest['version']}") self._install_manifest(task_id, manifest) manifest_hash = hashlib.sha256( self.repository.export_manifest_hash(manifest).encode("utf-8") ).hexdigest() self.repository.upsert_installed_app( manifest["appId"], manifest["appName"], manifest["version"], manifest_hash, manifest.get("openUrl"), ) self.repository.update_task( task_id, status="success", progress=100, current_step="completed", finished_at=utc_now(), ) self.repository.add_log(task_id, "info", f"Task {task_id} completed") except InstalledComponentVerificationError as error: self._track_attention_install(task_id, manifest) self._fail_task(task_id, error) except Exception as error: self._fail_task(task_id, error) def run_remove(self, task_id: str, request: RemoveRequest) -> None: try: if not settings.allow_remove: raise ValueError("Remove is disabled on this Agent") self._mark_started(task_id, "starting remove") self._require_root_if_available() effective_purge = request.purge and settings.allow_purge if request.purge and not effective_purge: self.repository.add_log( task_id, "warning", "Purge cleanup was requested but ALLOW_PURGE is disabled; falling back to remove cleanup", ) components = self.repository.list_installed_components(request.app_id) if not components and request.package_name: components = [ { "component_id": request.package_name, "type": "deb", "install_order": 10, "package_name": request.package_name, "service_name": request.service_name, } ] if not components: raise ValueError("No installed components found for this app") command_runner = CommandRunner(self.repository, task_id) installer = DebInstaller(command_runner) services = ServiceManager(command_runner) ordered = sorted(components, key=lambda item: item["install_order"], reverse=True) total = len(ordered) removed_apt_package = False for index, component in enumerate(ordered, start=1): progress = int((index - 1) / total * 80) + 10 component_id = component["component_id"] self.repository.update_task( task_id, progress=progress, current_step=f"removing {component_id}", current_component_id=component_id, ) service_name = component.get("service_name") if service_name: self.repository.add_log(task_id, "info", f"Stopping service {service_name}") self._best_effort(task_id, f"stop service {service_name}", lambda: services.stop_service(service_name)) self._best_effort(task_id, f"disable service {service_name}", lambda: services.disable_service(service_name)) package_name = component.get("package_name") if component["type"] in {"deb", "apt"} and package_name: self.repository.add_log(task_id, "info", f"Removing package {package_name}") installer.remove_package(package_name, purge=effective_purge) if component["type"] == "deb": self._clean_cached_package_files(task_id, package_name, component_id) if service_name: self._best_effort(task_id, f"reset failed state for {service_name}", lambda: services.reset_failed(service_name)) removed_apt_package = True elif component["type"] == "docker": container_name = component.get("container_name") or component_id self.repository.add_log(task_id, "info", f"Removing Docker container {container_name}") docker_installer = DockerInstaller(command_runner) docker_installer.ensure_runtime(auto_install=settings.auto_install_docker) docker_installer.remove_labeled_containers( request.app_id, component_id, remove_volumes=effective_purge, ) docker_installer.remove_container(container_name, remove_volumes=effective_purge) image = component.get("docker_image") or component.get("image") if effective_purge and image: self._best_effort(task_id, f"remove Docker image {image}", lambda: docker_installer.remove_image(image)) else: raise ValueError(f"Unsupported installed component type: {component['type']}") if removed_apt_package: self.repository.update_task(task_id, progress=92, current_step="cleaning package leftovers") self._best_effort( task_id, "clean unused packages and apt cache", lambda: installer.cleanup_after_remove(purge=effective_purge), ) self.repository.delete_installed_app(request.app_id) self.repository.update_task( task_id, status="success", progress=100, current_step="completed", finished_at=utc_now(), ) except Exception as error: self._fail_task(task_id, error) def _resolve_manifest(self, request: InstallRequest | UpdateRequest) -> dict[str, Any]: if request.download_url: digest = request.sha256 or request.checksum return self.manifest_validator.validate( { "schemaVersion": "1.0", "appId": request.app_id, "appName": request.app_name or request.app_id, "version": request.version, "architecture": "amd64", "components": [ { "componentId": request.package_name, "type": "deb", "installOrder": 10, "required": True, "packageName": request.package_name, "version": request.version, "downloadUrl": request.download_url, "sha256": digest, "serviceName": request.service_name, } ], } ) payload = self.manifest_client.fetch_manifest(request.app_id, request.version) return self.manifest_validator.validate(payload) def _install_manifest(self, task_id: str, manifest: dict[str, Any]) -> None: components = manifest["components"] for component in components: self.repository.create_task_component( task_id, manifest["appId"], component["componentId"], component["type"], component.get("installOrder", 10), ) total = len(components) if total == 0: raise ValueError("Manifest has no installable components") for index, component in enumerate(components, start=1): base_progress = int((index - 1) / total * 80) + 10 component_id = component["componentId"] self.repository.update_task( task_id, progress=base_progress, current_step=f"installing {component_id}", current_component_id=component_id, ) self.repository.update_task_component( task_id, component_id, status="running", progress=5, current_step="preparing", started_at=utc_now(), ) if component["type"] == "deb": self._install_deb_component(task_id, manifest["appId"], component) elif component["type"] == "apt": self._install_apt_component(task_id, manifest["appId"], component) elif component["type"] == "docker": self._install_docker_component(task_id, manifest["appId"], component) else: raise ValueError(f"Unsupported component type: {component['type']}") self.repository.update_task_component( task_id, component_id, status="success", progress=100, current_step="completed", finished_at=utc_now(), ) def _install_deb_component(self, task_id: str, app_id: str, component: dict[str, Any]) -> None: component_id = component["componentId"] downloader = Downloader(self.repository, task_id) command_runner = CommandRunner(self.repository, task_id) installer = DebInstaller(command_runner) services = ServiceManager(command_runner) self.repository.update_task_component(task_id, component_id, progress=10, current_step="downloading") package_path = downloader.download(component["downloadUrl"]) self.repository.update_task_component(task_id, component_id, progress=35, current_step="verifying checksum") actual_sha256 = sha256_file(package_path) expected_sha256 = component["sha256"].lower() if actual_sha256.lower() != expected_sha256: raise ValueError( f"Checksum mismatch for {component_id}: expected {expected_sha256}, got {actual_sha256}" ) self.repository.add_log(task_id, "info", f"Checksum verified for {component_id}") self.repository.update_task_component(task_id, component_id, progress=50, current_step="validating package metadata") deb_metadata = installer.get_deb_metadata(package_path) expected_package_name = component["packageName"] actual_package_name = deb_metadata["package"] if actual_package_name != expected_package_name: raise ValueError( f"Package metadata mismatch for {component_id}: manifest packageName is " f"{expected_package_name}, but .deb Package is {actual_package_name}. " f"Create or update the package in the web server with Package code {actual_package_name}." ) expected_version = component.get("version") or "" actual_version = deb_metadata["version"] if expected_version and actual_version != expected_version: raise ValueError( f"Package metadata mismatch for {component_id}: manifest version is " f"{expected_version}, but .deb Version is {actual_version}." ) self.repository.add_log( task_id, "info", f"Package metadata verified for {actual_package_name} {actual_version}", ) self.repository.update_task_component(task_id, component_id, progress=60, current_step="installing package") installer.install_deb(package_path) self.repository.update_task_component(task_id, component_id, progress=75, current_step="verifying package") installed_version = installer.get_package_version(component["packageName"]) self.repository.add_log( task_id, "info", f"Package {component['packageName']} installed with version {installed_version}", ) service_name = component.get("serviceName") if service_name: self.repository.update_task_component(task_id, component_id, progress=90, current_step="starting service") services.enable_service(service_name) services.start_service(service_name) self.repository.upsert_installed_component(app_id, component) def _install_apt_component(self, task_id: str, app_id: str, component: dict[str, Any]) -> None: component_id = component["componentId"] package_name = component["packageName"] command_runner = CommandRunner(self.repository, task_id) installer = AptInstaller(command_runner) services = ServiceManager(command_runner) self.repository.update_task_component( task_id, component_id, progress=10, current_step="refreshing APT package index", ) self.repository.add_log(task_id, "info", "Refreshing APT package index") installer.update_package_index() self.repository.update_task_component( task_id, component_id, progress=35, current_step=f"installing APT package {package_name}", ) self.repository.add_log(task_id, "info", f"Installing trusted APT package {package_name}") installer.install_package(package_name) self.repository.update_task_component( task_id, component_id, progress=70, current_step="verifying installed package", ) installed_version = installer.get_package_version(package_name) if not installed_version: raise RuntimeError(f"APT package was not installed: {package_name}") self.repository.add_log( task_id, "info", f"APT package {package_name} installed with version {installed_version}", ) installed_component = dict(component) installed_component["version"] = installed_version policy = APT_SERVICE_POLICIES.get(package_name, {}) managed_services = list(policy.get("managedServices", [])) discovered_services = installer.discover_service_units(package_name) service_names = list(dict.fromkeys([*managed_services, *discovered_services])) self.repository.update_task_component( task_id, component_id, progress=80, current_step="checking package services", service_check_status="checking", service_checks=[], ) verification_errors: list[str] = [] if managed_services: self.repository.update_task_component( task_id, component_id, progress=85, current_step="starting managed package services", ) for service_name in managed_services: try: services.enable_service(service_name) services.start_service(service_name) except Exception as error: verification_errors.append(f"Could not start {service_name}: {error}") readiness_status: str | None = None if policy.get("readiness") == "postgresql": self.repository.update_task_component( task_id, component_id, progress=90, current_step="checking PostgreSQL readiness", ) try: installer.wait_for_postgresql() readiness_status = "ready" except Exception as error: readiness_status = "failed" verification_errors.append(f"PostgreSQL is not accepting connections: {error}") service_checks: list[dict[str, Any]] = [] self.repository.update_task_component( task_id, component_id, progress=95, current_step="reading service status", ) for service_name in service_names: check = services.get_service_status(service_name) if service_name in managed_services and readiness_status is not None: check["readinessType"] = str(policy["readiness"]) check["readinessStatus"] = readiness_status check["healthy"] = bool(check.get("healthy")) and readiness_status == "ready" service_checks.append(check) if check.get("healthy"): self.repository.add_log( task_id, "info", f"Service {service_name} is active ({check.get('subState', 'unknown')})", ) else: self.repository.add_log( task_id, "warning", f"Service {service_name} is not healthy " f"(load={check.get('loadState', 'unknown')}, " f"active={check.get('activeState', 'unknown')}, " f"sub={check.get('subState', 'unknown')})", ) if not service_names: service_check_status = "not-applicable" self.repository.add_log( task_id, "info", f"APT package {package_name} does not expose a concrete systemd service unit", ) elif all(bool(check.get("healthy")) for check in service_checks): service_check_status = "healthy" else: service_check_status = "unhealthy" self.repository.update_task_component( task_id, component_id, progress=98, current_step="service checks completed", service_check_status=service_check_status, service_checks=service_checks, ) installed_component["serviceCheckStatus"] = service_check_status installed_component["serviceChecks"] = service_checks if managed_services: installed_component["serviceName"] = managed_services[0] self.repository.upsert_installed_component(app_id, installed_component) unhealthy_managed_services = [ str(check.get("serviceName")) for check in service_checks if check.get("serviceName") in managed_services and not check.get("healthy") ] if unhealthy_managed_services: verification_errors.append( f"Required service is not healthy: {', '.join(unhealthy_managed_services)}" ) if verification_errors: raise InstalledComponentVerificationError( "; ".join(dict.fromkeys(verification_errors)) ) def _install_docker_component(self, task_id: str, app_id: str, component: dict[str, Any]) -> None: component_id = component["componentId"] container_name = component["containerName"] reference = image_reference(component) command_runner = CommandRunner(self.repository, task_id) installer = DockerInstaller(command_runner) self.repository.update_task_component(task_id, component_id, progress=15, current_step="checking Docker runtime") installer.ensure_runtime(auto_install=settings.auto_install_docker) self.repository.update_task_component(task_id, component_id, progress=35, current_step="pulling image") self.repository.add_log(task_id, "info", f"Pulling Docker image {reference}") installer.pull_image(reference) self.repository.update_task_component(task_id, component_id, progress=70, current_step="recreating container") self.repository.add_log(task_id, "info", f"Recreating Docker container {container_name}") installer.recreate_container(app_id, component) self.repository.update_task_component(task_id, component_id, progress=90, current_step="verifying container") installer.assert_container_running(container_name) self.repository.add_log(task_id, "info", f"Docker container {container_name} is running") installed_component = dict(component) installed_component["image"] = reference self.repository.upsert_installed_component(app_id, installed_component) def _mark_started(self, task_id: str, step: str) -> None: self.repository.update_task( task_id, status="running", progress=5, current_step=step, started_at=utc_now(), ) self.repository.add_log(task_id, "info", step) def _fail_task(self, task_id: str, error: Exception) -> None: task = self.repository.get_task(task_id) component_id = task.get("current_component_id") if task else None finished_at = utc_now() if component_id: self.repository.update_task_component( task_id, component_id, status="failed", current_step="failed", error_message=str(error), finished_at=finished_at, ) self.repository.update_task( task_id, status="failed", current_step="failed", error_message=str(error), finished_at=finished_at, ) self.repository.add_log(task_id, "error", str(error)) self.repository.add_log(task_id, "debug", traceback.format_exc()) def _require_root_if_available(self) -> None: geteuid = getattr(os, "geteuid", None) if callable(geteuid) and geteuid() != 0: raise PermissionError("Agent must run as root to call apt and systemctl") def _best_effort(self, task_id: str, action: str, callback: Any) -> None: try: callback() except Exception as error: self.repository.add_log(task_id, "warning", f"Could not {action}: {error}") def _track_attention_install(self, task_id: str, manifest: dict[str, Any] | None) -> None: if not manifest: return try: manifest_hash = hashlib.sha256( self.repository.export_manifest_hash(manifest).encode("utf-8") ).hexdigest() self.repository.upsert_installed_app( manifest["appId"], manifest["appName"], manifest["version"], manifest_hash, manifest.get("openUrl"), status="attention", ) self.repository.add_log( task_id, "warning", "The package was installed but its service needs attention; it remains available for status and removal", ) except Exception as tracking_error: self.repository.add_log( task_id, "warning", f"Could not persist installed package attention state: {tracking_error}", ) def _clean_cached_package_files(self, task_id: str, *identifiers: str | None) -> None: cache_dir = settings.cache_dir if not cache_dir.exists() or not cache_dir.is_dir(): return patterns: list[str] = [] seen_patterns: set[str] = set() for identifier in identifiers: name = (identifier or "").strip() if not name: continue for pattern in (f"{name}.deb", f"{name}_*.deb"): if pattern not in seen_patterns: patterns.append(pattern) seen_patterns.add(pattern) removed_files: set[str] = set() for pattern in patterns: for file_path in cache_dir.glob(pattern): if not file_path.is_file(): continue try: file_path.unlink() removed_files.add(str(file_path)) except Exception as error: self.repository.add_log(task_id, "warning", f"Could not remove cached package {file_path}: {error}") for file_path in sorted(removed_files): self.repository.add_log(task_id, "info", f"Removed cached package {file_path}")