"""Job queue behaviour.""" from __future__ import annotations import asyncio from pathlib import Path import pytest from app.config import Settings from app.errors import LimitExceededError, QueueFullError from app.models import ConvertRequest, Source from app.services.jobs import JobManager from app.services.pipeline import ConversionResult from app.services.storage import Storage pytestmark = pytest.mark.asyncio class StubPipeline: def __init__(self, behaviour="ok", delay=0.0) -> None: self.behaviour = behaviour self.delay = delay async def run(self, request, workdir: Path, progress=None) -> ConversionResult: if progress is not None: progress(1, 1, 1, 1) if self.delay: await asyncio.sleep(self.delay) if self.behaviour == "conversion_error": raise LimitExceededError("Prekrocen limit stranek.", {"limit": 10}) if self.behaviour == "crash": raise RuntimeError("necekana chyba enginu") workdir.mkdir(parents=True, exist_ok=True) produced = workdir / "out.pdf" produced.write_bytes(b"%PDF-1.7\n%fake\n") return ConversionResult(path=produced, page_count=1, engine_used="stub") def make_manager(tmp_path, behaviour="ok", delay=0.0, **overrides) -> JobManager: settings = Settings() for key, value in overrides.items(): object.__setattr__(settings, key, value) object.__setattr__(settings, "storage_dir", str(tmp_path)) storage = Storage(str(tmp_path)) return JobManager(StubPipeline(behaviour, delay), storage, settings) def simple_request() -> ConvertRequest: return ConvertRequest(source=Source(html="

ahoj

")) async def test_successful_job_publishes_a_result(tmp_path) -> None: manager = make_manager(tmp_path) await manager.start() try: job = manager.submit(simple_request()) await asyncio.wait_for(job.done.wait(), timeout=5) assert job.state.status == "done" assert job.state.page_count == 1 assert (Path(tmp_path) / job.id / "result.pdf").exists() finally: await manager.stop() async def test_conversion_error_keeps_its_code(tmp_path) -> None: manager = make_manager(tmp_path, behaviour="conversion_error") await manager.start() try: job = manager.submit(simple_request()) await asyncio.wait_for(job.done.wait(), timeout=5) assert job.state.status == "failed" assert job.state.error is not None assert job.state.error.error_code == "limit_exceeded" finally: await manager.stop() async def test_unexpected_error_is_reported_not_swallowed(tmp_path) -> None: manager = make_manager(tmp_path, behaviour="crash") await manager.start() try: job = manager.submit(simple_request()) await asyncio.wait_for(job.done.wait(), timeout=5) assert job.state.status == "failed" assert job.state.error.error_code == "internal_error" finally: await manager.stop() async def test_running_job_can_be_cancelled(tmp_path) -> None: manager = make_manager(tmp_path, delay=5.0) await manager.start() try: job = manager.submit(simple_request()) await asyncio.sleep(0.2) manager.cancel(job.id) await asyncio.wait_for(job.done.wait(), timeout=5) assert job.state.status == "cancelled" assert not (Path(tmp_path) / job.id).exists() finally: await manager.stop() async def test_full_queue_is_rejected(tmp_path) -> None: manager = make_manager(tmp_path, delay=5.0, queue_max_size=1, workers=1) await manager.start() try: manager.submit(simple_request()) await asyncio.sleep(0.1) manager.submit(simple_request()) with pytest.raises(QueueFullError): manager.submit(simple_request()) finally: await manager.stop() async def test_shutdown_marks_pending_jobs_as_failed(tmp_path) -> None: manager = make_manager(tmp_path, delay=5.0, workers=1) await manager.start() job = manager.submit(simple_request()) await asyncio.sleep(0.1) await manager.stop() assert job.state.status == "failed" assert job.state.error.error_code == "service_restarted"