"""Job runner: executes update jobs in background tasks, streams logs to WS + DB.""" import asyncio from datetime import UTC, datetime from app.core.database import async_session_factory from app.core.logging import get_logger from app.models.credential import Credential, CredentialType from app.models.server import Server, ServerType from app.models.update_job import JobStatus, UpdateJob, UpdateLog from app.services.cau import CAUService from app.services.ssh import SSHCredentials, SSHService from app.services.winrm import WinRMCredentials, WinRMService from app.websocket import handlers as ws logger = get_logger(__name__) class JobRunner: """Manages asyncio tasks for running update jobs.""" def __init__(self) -> None: self._tasks: dict[int, asyncio.Task] = {} # type: ignore[type-arg] async def start(self, job_id: int) -> None: task = asyncio.create_task(self._run(job_id), name=f"job-{job_id}") self._tasks[job_id] = task task.add_done_callback(lambda _t: self._tasks.pop(job_id, None)) async def cancel(self, job_id: int) -> bool: task = self._tasks.get(job_id) if not task: return False task.cancel() return True # ------------------------------------------------------------------ async def _run(self, job_id: int) -> None: async with async_session_factory() as db: job = await db.get(UpdateJob, job_id) if not job: logger.error("job.not_found", job_id=job_id) return server = await db.get(Server, job.server_id) if not server: await self._finish(db, job, JobStatus.FAILED, "Server nicht gefunden") return job.status = JobStatus.RUNNING job.started_at = datetime.now(UTC) await db.commit() await ws.emit_job_start(job.id, server.id, job.type.value) started = datetime.now(UTC) try: await self._dispatch(db, job, server) await self._finish(db, job, JobStatus.SUCCESS, None, started) except asyncio.CancelledError: await self._finish(db, job, JobStatus.CANCELLED, "Vom Benutzer abgebrochen", started) raise except Exception as exc: # noqa: BLE001 logger.error("job.failed", job_id=job_id, error=str(exc)) await self._finish(db, job, JobStatus.FAILED, str(exc), started) async def _dispatch(self, db, job: UpdateJob, server: Server) -> None: # type: ignore[no-untyped-def] credential = await db.get(Credential, server.credential_id) if server.credential_id else None if server.type == ServerType.LINUX: creds = self._ssh_creds(credential) service = SSHService(server.hostname, port=server.port, credentials=creds) async for line in service.stream_updates(): await self._log(db, job, line) elif server.type == ServerType.WINDOWS: creds = self._winrm_creds(credential) service = WinRMService(server.hostname, port=server.port, credentials=creds) async for line in service.install_updates(): await self._log(db, job, line) elif server.type == ServerType.CAU_CLUSTER: creds = self._winrm_creds(credential) service = CAUService(server.hostname, access_node=server.hostname, port=server.port, credentials=creds) async for line in service.invoke_cau_run(): await self._log(db, job, line) async def _log(self, db, job: UpdateJob, line: str, level: str = "info") -> None: # type: ignore[no-untyped-def] db.add(UpdateLog(job_id=job.id, level=level, line=line)) await db.commit() await ws.emit_job_log(job.id, line, level) async def _finish( # type: ignore[no-untyped-def] self, db, job: UpdateJob, status: JobStatus, error: str | None, started: datetime | None = None ) -> None: job.status = status job.error = error job.finished_at = datetime.now(UTC) if status == JobStatus.SUCCESS: job.progress_percent = 100 await db.commit() duration = ( (job.finished_at - started).total_seconds() if started else None ) await ws.emit_job_complete(job.id, status.value, duration) # ------------------------------------------------------------------ @staticmethod def _winrm_creds(credential: Credential | None) -> WinRMCredentials | None: if not credential or credential.type != CredentialType.WINRM_USERPASS: return None from app.core.security import decrypt return WinRMCredentials( username=credential.username, password=decrypt(credential.password_encrypted or ""), ) @staticmethod def _ssh_creds(credential: Credential | None) -> SSHCredentials | None: if not credential: return None from app.core.security import decrypt return SSHCredentials( username=credential.username, password=decrypt(credential.password_encrypted) if credential.password_encrypted else None, private_key=decrypt(credential.private_key_encrypted) if credential.private_key_encrypted else None, passphrase=decrypt(credential.key_passphrase_encrypted) if credential.key_passphrase_encrypted else None, ) job_runner = JobRunner()