Sluzba prijme adresu HTML dokumentu nebo HTML v tele requestu a vrati PDF. Navrzena pro dokumenty o stovkach az tisicich stranek. Rendering: - WeasyPrint jako vychozi engine, spravne CSS Paged Media, nizka pametova narocnost, bez JavaScriptu - Chromium pres Playwright pro dokumenty dokreslovane skripty - rezim auto s detekci skriptu a fallbackem pri selhani WeasyPrintu Velke dokumenty: - deleni na casti na strukturalnich hranicich, rez nikdy uvnitr tabulky nebo odstavce - dvoupruchodovy render obsahu se skutecnymi cisly stranek, pozice nadpisu se ctou z kotev hlasenych u kazde stranky - cislovani stranek bud pres CSS countery, nebo pres cislovaci vrstvu nastampovanou na hotove PDF, rozmer stranky se cte z vysledneho souboru - Chromium se restartuje po N jobech, nikdy vsak behem beziciho renderu API: - POST /convert synchronne, POST /jobs asynchronne se sledovanim stavu, stahovanim vysledku, rusenim a volitelnym callbackem - GET /health s overenim dostupnosti obou enginu a stavem fronty - OpenAPI respektuje prefix reverse proxy pres root_path Bezpecnost a provoz: - SSRF kontrola po DNS resolvu, na kazdem presmerovani a u vsech pozadavku prohlizece - nedostupne assety render nezastavi, ale hlasi se v odpovedi i v logu - fronta s omezenym poctem workeru, rozpracovane joby se pri ukonceni oznaci jako failed, nezmizi potichu - strukturovane JSON logovani s job_id - vsechny limity vypnute ve vychozim stavu Dockerfile je dvoufazovy, obsahuje zavislosti WeasyPrintu, Chromium a fonty s ceskou diakritikou. Autentizace zamerne neni implementovana, zpusob predavani neni domluveny. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
239 lines
9.0 KiB
Python
239 lines
9.0 KiB
Python
"""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})
|