"""Tests for the job queue, cancellation, cache and confidence tiers (no model inference).""" from __future__ import annotations import asyncio import pytest from app.services import analysis_jobs as aj from app.services.analysis_jobs import ( FILE_STEPS, STEP_DONE, STEP_RUNNING, CpuGate, JobCancelled, JobStore, Progress, ResultCache, ) from app.services.inference_xai import XAIInferenceService OK = {"result": {"ok": True}, "warnings": [], "errors": []} @pytest.fixture(autouse=True) def fresh_state(monkeypatch): """Each test gets its own gate and cache.""" monkeypatch.setattr(aj, "cpu_gate", CpuGate(1)) monkeypatch.setattr(aj, "result_cache", ResultCache(8, 3600)) def gated(hold: float = 0.05): """A job that takes the CPU slot like the real pipeline does.""" async def work(job): async with aj.cpu_gate.slot(job.id, job.progress): job.progress.running("features") await asyncio.sleep(hold) job.progress.set("features", STEP_DONE) return OK return work def test_jobs_take_the_cpu_one_at_a_time_in_order() -> None: async def scenario() -> None: store = JobStore() jobs = [store.start("file", FILE_STEPS, gated(0.05)) for _ in range(3)] await asyncio.sleep(0.01) snaps = [j.to_dict() for j in jobs] assert snaps[0]["phase"] == "running" and "queue" not in snaps[0] assert snaps[1]["status"] == "running" and snaps[1]["phase"] == "queued" assert snaps[1]["queue"]["position"] == 1 and snaps[1]["queue"]["ahead"] == 1 assert snaps[2]["queue"]["position"] == 2 await asyncio.sleep(0.25) assert all(j.to_dict()["status"] == "done" for j in jobs) # The gate now knows how long a job holds the CPU... only past 1 s, so still None here. assert aj.cpu_gate.estimated_wait(1) is None asyncio.run(scenario()) def test_cancelling_a_queued_job_frees_its_place() -> None: async def scenario() -> None: store = JobStore() first = store.start("file", FILE_STEPS, gated(0.1)) second = store.start("file", FILE_STEPS, gated(0.1)) third = store.start("file", FILE_STEPS, gated(0.1)) await asyncio.sleep(0.01) store.cancel(second.id) await asyncio.sleep(0.01) assert second.status == "error" and second.response["errors"] == ["cancelled"] assert third.to_dict()["queue"]["position"] == 1 await asyncio.sleep(0.3) assert first.status == "done" and third.status == "done" asyncio.run(scenario()) def test_cancelling_a_running_job_stops_it_at_the_next_check() -> None: async def scenario() -> None: store = JobStore() async def work(job): async with aj.cpu_gate.slot(job.id, job.progress): await asyncio.sleep(0.05) job.progress.check() return OK job = store.start("file", FILE_STEPS, work) await asyncio.sleep(0.01) store.cancel(job.id) await asyncio.sleep(0.1) assert job.status == "error" and job.response["errors"] == ["cancelled"] asyncio.run(scenario()) def test_a_job_nobody_polls_is_dropped_when_its_turn_comes(monkeypatch) -> None: monkeypatch.setattr(aj, "ABANDON_SEC", 0.02) async def scenario() -> None: store = JobStore() store.start("file", FILE_STEPS, gated(0.08)) waiting = store.start("file", FILE_STEPS, gated(0.01)) await asyncio.sleep(0.2) assert waiting.status == "error" and waiting.response["errors"] == ["job_abandoned"] asyncio.run(scenario()) def test_same_file_shares_the_running_job_then_hits_the_cache() -> None: async def scenario() -> None: store = JobStore() a = store.start("file", FILE_STEPS, gated(0.05), key="file:abc") b = store.start("file", FILE_STEPS, gated(0.05), key="file:abc") assert a is b await asyncio.sleep(0.15) c = store.start("file", FILE_STEPS, gated(0.05), key="file:abc") assert c is not a and c.cached and c.status == "done" and c.response == OK asyncio.run(scenario()) def test_full_queue_is_refused(monkeypatch) -> None: monkeypatch.setattr(aj, "MAX_PENDING", 2) async def scenario() -> None: store = JobStore() store.start("file", FILE_STEPS, gated(0.05)) store.start("file", FILE_STEPS, gated(0.05)) with pytest.raises(aj.QueueFull): store.start("file", FILE_STEPS, gated(0.05)) await asyncio.sleep(0.2) asyncio.run(scenario()) def test_failed_job_settles_as_error() -> None: async def scenario() -> None: store = JobStore() async def boom(job): raise RuntimeError("boom") job = store.start("file", FILE_STEPS, boom) await asyncio.sleep(0.01) assert job.status == "error" assert job.response == {"result": None, "warnings": [], "errors": ["internal_error"]} asyncio.run(scenario()) def test_progress_check_raises_when_cancelled() -> None: progress = Progress(FILE_STEPS) progress.check() progress.cancelled = True with pytest.raises(JobCancelled): progress.check() def test_confidence_band_is_measured_from_the_threshold() -> None: svc = XAIInferenceService.__new__(XAIInferenceService) svc.threshold = 0.431577 assert svc._confidence_band(0.45).tier == "uncertain" assert svc._confidence_band(0.40).tier == "uncertain" assert svc._confidence_band(0.97).tier == "very_strong" assert svc._confidence_band(0.02).tier == "very_strong" assert svc._confidence_band(0.431577).margin == 0.0 def test_progress_reports_how_far_a_running_step_is() -> None: progress = Progress(FILE_STEPS) progress.set("timeline", STEP_RUNNING) progress.advance("timeline", 3, 12) item = next(s for s in progress.snapshot() if s["id"] == "timeline") assert item["progress"] == 0.25 progress.set("timeline", STEP_DONE) item = next(s for s in progress.snapshot() if s["id"] == "timeline") assert "progress" not in item