crowncode-backend / tests /test_analysis_jobs.py
Rthur2003's picture
feat: core analysis pipeline ve routes için audio inspection ve XAI verdicts eklendi
e6e7656
Raw History Blame Contribute Delete
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": []}
@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