Spaces:
Running
Running
Download tests/test_analysis_jobs.py from Rthur2003/crowncode-backend: direct link, hf CLI and curl.
- Browser
- Download file 6.11 kB
-
https://huggingface.co/spaces/Rthur2003/crowncode-backend/resolve/main/tests/test_analysis_jobs.py
- Command line
-
hf download hf://spaces/Rthur2003/crowncode-backend/tests/test_analysis_jobs.py
-
curl -L -o test_analysis_jobs.py https://huggingface.co/spaces/Rthur2003/crowncode-backend/resolve/main/tests/test_analysis_jobs.py
6.11 kB
| """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": []} | |
| 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 | |