"""In memory job queue. The queue is deliberately in process. There is no broker and no database, which means a restart loses queued work. Such jobs are marked as failed with an explicit reason instead of silently disappearing. """ from __future__ import annotations import asyncio import logging import uuid from datetime import datetime, timedelta, timezone import httpx from ..config import Settings from ..errors import ConversionError, JobNotFoundError, QueueFullError from ..logging_setup import current_job_id from ..models import ConvertRequest, ErrorInfo, JobProgress, JobState from .pipeline import ConversionPipeline from .storage import Storage logger = logging.getLogger(__name__) def _now() -> datetime: return datetime.now(timezone.utc) class Job: def __init__(self, job_id: str, request: ConvertRequest) -> None: self.id = job_id self.request = request self.state = JobState(job_id=job_id, status="queued", created_at=_now()) self.task = None self.done = asyncio.Event() class JobManager: def __init__(self, pipeline: ConversionPipeline, storage: Storage, settings: Settings) -> None: self._pipeline = pipeline self._storage = storage self._settings = settings self._jobs = {} self._queue = asyncio.Queue(maxsize=settings.queue_max_size) self._workers = [] self._cleaner = None async def start(self) -> None: self._storage.sweep_orphans(set(self._jobs)) for index in range(max(1, self._settings.workers)): self._workers.append(asyncio.create_task(self._worker(index), name=f"htp-worker-{index}")) self._cleaner = asyncio.create_task(self._cleanup_loop(), name="htp-cleanup") logger.info("Job manager started", extra={"workers": len(self._workers)}) async def stop(self) -> None: for task in self._workers: task.cancel() if self._cleaner is not None: self._cleaner.cancel() for job in self._jobs.values(): if job.state.status in ("queued", "running"): self._fail( job, ErrorInfo( error_code="service_restarted", message="Sluzba byla ukoncena drive, nez job dobehl. Odeslete pozadavek znovu.", ), ) logger.info("Job manager stopped") def submit(self, request: ConvertRequest) -> Job: job = Job(str(uuid.uuid4()), request) self._jobs[job.id] = job try: self._queue.put_nowait(job.id) except asyncio.QueueFull as exc: del self._jobs[job.id] raise QueueFullError( "Fronta je plna, zkuste to prosim za chvili.", {"queue_max_size": self._settings.queue_max_size}, ) from exc logger.info("Job queued", extra={"job_id": job.id, "queue_size": self._queue.qsize()}) return job def get(self, job_id: str) -> Job: job = self._jobs.get(job_id) if job is None: raise JobNotFoundError("Job s timto identifikatorem neexistuje nebo uz expiroval.") return job def cancel(self, job_id: str) -> Job: job = self.get(job_id) if job.task is not None and not job.task.done(): job.task.cancel() else: job.state.status = "cancelled" job.state.finished_at = _now() self._storage.discard(job.id) job.done.set() job.state.expires_at = _now() + timedelta(seconds=self._settings.job_result_ttl_seconds) logger.info("Job cancelled", extra={"job_id": job_id}) return job def stats(self) -> dict: counts = {"queued": 0, "running": 0, "done": 0, "failed": 0, "cancelled": 0, "expired": 0} for job in self._jobs.values(): counts[job.state.status] = counts.get(job.state.status, 0) + 1 counts["workers"] = len(self._workers) return counts async def _worker(self, index: int) -> None: while True: job_id = await self._queue.get() job = self._jobs.get(job_id) if job is None or job.state.status != "queued": self._queue.task_done() continue job.task = asyncio.current_task() token = current_job_id.set(job.id) try: await self._execute(job) except asyncio.CancelledError: job.state.status = "cancelled" job.state.finished_at = _now() job.state.expires_at = _now() + timedelta(seconds=self._settings.job_result_ttl_seconds) self._storage.discard(job.id) logger.info("Job execution cancelled", extra={"job_id": job.id}) finally: current_job_id.reset(token) job.task = None job.done.set() self._queue.task_done() await self._notify_callback(job) async def _execute(self, job: Job) -> None: job.state.status = "running" job.state.started_at = _now() logger.info("Job started") def progress(pages, chunks_done, chunks_total, pass_number): job.state.progress = JobProgress( pages_rendered=pages, chunks_done=chunks_done, chunks_total=chunks_total, pass_number=pass_number, ) workdir = self._storage.job_dir(job.id) try: result = await self._pipeline.run(job.request, workdir, progress) except ConversionError as exc: logger.error("Job failed", extra={"error_code": exc.error_code}, exc_info=exc) self._fail(job, ErrorInfo(error_code=exc.error_code, message=exc.message, detail=exc.detail or None)) return except asyncio.CancelledError: raise except Exception as exc: logger.exception("Job failed with an unexpected error") self._fail( job, ErrorInfo( error_code="internal_error", message="Pri generovani PDF doslo k neocekavane chybe.", detail={"reason": str(exc)}, ), ) return self._storage.publish(job.id, result.path) job.state.status = "done" job.state.finished_at = _now() job.state.expires_at = _now() + timedelta(seconds=self._settings.job_result_ttl_seconds) job.state.page_count = result.page_count job.state.engine_used = result.engine_used job.state.missing_assets = result.missing_assets job.state.warnings = result.warnings logger.info( "Job finished", extra={ "pages": result.page_count, "engine": result.engine_used, "missing_assets": len(result.missing_assets), }, ) def _fail(self, job: Job, error: ErrorInfo) -> None: job.state.status = "failed" job.state.finished_at = _now() job.state.expires_at = _now() + timedelta(seconds=self._settings.job_result_ttl_seconds) job.state.error = error self._storage.discard(job.id) async def _notify_callback(self, job: Job) -> None: url = job.request.callback_url if not url or job.state.status not in ("done", "failed"): return payload = job.state.model_dump(mode="json") for attempt in range(1, max(1, self._settings.callback_retries) + 1): try: async with httpx.AsyncClient(timeout=self._settings.callback_timeout_seconds) as client: response = await client.post(url, json=payload) if response.status_code < 400: logger.info("Callback delivered", extra={"attempt": attempt}) return logger.warning( "Callback returned an error status", extra={"attempt": attempt, "status_code": response.status_code}, ) except Exception as exc: logger.warning("Callback delivery failed", extra={"attempt": attempt}, exc_info=exc) await asyncio.sleep(min(2 ** attempt, 10)) logger.error("Callback could not be delivered, job result stays available over the API") async def _cleanup_loop(self) -> None: while True: await asyncio.sleep(60) try: self._expire_old_jobs() except Exception: logger.exception("Cleanup loop failed") def _expire_old_jobs(self) -> None: now = _now() for job_id, job in list(self._jobs.items()): expires_at = job.state.expires_at if expires_at is None or expires_at > now: continue job.state.status = "expired" self._storage.discard(job_id) del self._jobs[job_id] logger.info("Job result expired and was removed", extra={"job_id": job_id})