Spaces:
Running
Running
test: unit tests için analysis job queue cancellation ve önbellek eklendi
Browse files- tests/test_analysis_jobs.py +120 -21
tests/test_analysis_jobs.py
CHANGED
|
@@ -1,33 +1,127 @@
|
|
| 1 |
-
"""Tests for the
|
| 2 |
|
| 3 |
from __future__ import annotations
|
| 4 |
|
| 5 |
import asyncio
|
| 6 |
|
| 7 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 8 |
from app.services.inference_xai import XAIInferenceService
|
| 9 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 10 |
|
| 11 |
-
def
|
| 12 |
async def scenario() -> None:
|
| 13 |
store = JobStore()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 14 |
|
| 15 |
-
|
| 16 |
-
|
| 17 |
-
|
| 18 |
-
|
| 19 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 20 |
|
| 21 |
job = store.start("file", FILE_STEPS, work)
|
| 22 |
-
|
| 23 |
-
|
|
|
|
|
|
|
| 24 |
|
| 25 |
-
|
| 26 |
-
|
| 27 |
-
|
| 28 |
-
|
| 29 |
-
|
| 30 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 31 |
|
| 32 |
asyncio.run(scenario())
|
| 33 |
|
|
@@ -36,7 +130,7 @@ def test_failed_job_settles_as_error() -> None:
|
|
| 36 |
async def scenario() -> None:
|
| 37 |
store = JobStore()
|
| 38 |
|
| 39 |
-
async def boom(
|
| 40 |
raise RuntimeError("boom")
|
| 41 |
|
| 42 |
job = store.start("file", FILE_STEPS, boom)
|
|
@@ -47,14 +141,19 @@ def test_failed_job_settles_as_error() -> None:
|
|
| 47 |
asyncio.run(scenario())
|
| 48 |
|
| 49 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 50 |
def test_confidence_band_is_measured_from_the_threshold() -> None:
|
| 51 |
svc = XAIInferenceService.__new__(XAIInferenceService)
|
| 52 |
svc.threshold = 0.431577
|
| 53 |
-
# Just above the threshold: an AI verdict, but barely.
|
| 54 |
assert svc._confidence_band(0.45).tier == "uncertain"
|
| 55 |
-
# Below 0.5 yet clearly under the threshold's human side? No — close to it.
|
| 56 |
assert svc._confidence_band(0.40).tier == "uncertain"
|
| 57 |
assert svc._confidence_band(0.97).tier == "very_strong"
|
| 58 |
assert svc._confidence_band(0.02).tier == "very_strong"
|
| 59 |
-
|
| 60 |
-
assert band.margin == 0.0
|
|
|
|
| 1 |
+
"""Tests for the job queue, cancellation, cache and confidence tiers (no model inference)."""
|
| 2 |
|
| 3 |
from __future__ import annotations
|
| 4 |
|
| 5 |
import asyncio
|
| 6 |
|
| 7 |
+
import pytest
|
| 8 |
+
|
| 9 |
+
from app.services import analysis_jobs as aj
|
| 10 |
+
from app.services.analysis_jobs import (
|
| 11 |
+
FILE_STEPS, STEP_DONE, CpuGate, JobCancelled, JobStore, Progress, ResultCache,
|
| 12 |
+
)
|
| 13 |
from app.services.inference_xai import XAIInferenceService
|
| 14 |
|
| 15 |
+
OK = {"result": {"ok": True}, "warnings": [], "errors": []}
|
| 16 |
+
|
| 17 |
+
|
| 18 |
+
@pytest.fixture(autouse=True)
|
| 19 |
+
def fresh_state(monkeypatch):
|
| 20 |
+
"""Each test gets its own gate and cache."""
|
| 21 |
+
monkeypatch.setattr(aj, "cpu_gate", CpuGate(1))
|
| 22 |
+
monkeypatch.setattr(aj, "result_cache", ResultCache(8, 3600))
|
| 23 |
+
|
| 24 |
+
|
| 25 |
+
def gated(hold: float = 0.05):
|
| 26 |
+
"""A job that takes the CPU slot like the real pipeline does."""
|
| 27 |
+
async def work(job):
|
| 28 |
+
async with aj.cpu_gate.slot(job.id, job.progress):
|
| 29 |
+
job.progress.running("features")
|
| 30 |
+
await asyncio.sleep(hold)
|
| 31 |
+
job.progress.set("features", STEP_DONE)
|
| 32 |
+
return OK
|
| 33 |
+
return work
|
| 34 |
+
|
| 35 |
+
|
| 36 |
+
def test_jobs_take_the_cpu_one_at_a_time_in_order() -> None:
|
| 37 |
+
async def scenario() -> None:
|
| 38 |
+
store = JobStore()
|
| 39 |
+
jobs = [store.start("file", FILE_STEPS, gated(0.05)) for _ in range(3)]
|
| 40 |
+
await asyncio.sleep(0.01)
|
| 41 |
+
snaps = [j.to_dict() for j in jobs]
|
| 42 |
+
assert snaps[0]["status"] == "running"
|
| 43 |
+
assert snaps[1]["status"] == "queued" and snaps[1]["queue"]["position"] == 1
|
| 44 |
+
assert snaps[2]["queue"]["position"] == 2
|
| 45 |
+
await asyncio.sleep(0.25)
|
| 46 |
+
assert all(j.to_dict()["status"] == "done" for j in jobs)
|
| 47 |
+
# The gate now knows how long a job holds the CPU... only past 1 s, so still None here.
|
| 48 |
+
assert aj.cpu_gate.estimated_wait(1) is None
|
| 49 |
+
|
| 50 |
+
asyncio.run(scenario())
|
| 51 |
+
|
| 52 |
|
| 53 |
+
def test_cancelling_a_queued_job_frees_its_place() -> None:
|
| 54 |
async def scenario() -> None:
|
| 55 |
store = JobStore()
|
| 56 |
+
first = store.start("file", FILE_STEPS, gated(0.1))
|
| 57 |
+
second = store.start("file", FILE_STEPS, gated(0.1))
|
| 58 |
+
third = store.start("file", FILE_STEPS, gated(0.1))
|
| 59 |
+
await asyncio.sleep(0.01)
|
| 60 |
+
store.cancel(second.id)
|
| 61 |
+
await asyncio.sleep(0.01)
|
| 62 |
+
assert second.status == "error" and second.response["errors"] == ["cancelled"]
|
| 63 |
+
assert third.to_dict()["queue"]["position"] == 1
|
| 64 |
+
await asyncio.sleep(0.3)
|
| 65 |
+
assert first.status == "done" and third.status == "done"
|
| 66 |
+
|
| 67 |
+
asyncio.run(scenario())
|
| 68 |
|
| 69 |
+
|
| 70 |
+
def test_cancelling_a_running_job_stops_it_at_the_next_check() -> None:
|
| 71 |
+
async def scenario() -> None:
|
| 72 |
+
store = JobStore()
|
| 73 |
+
|
| 74 |
+
async def work(job):
|
| 75 |
+
async with aj.cpu_gate.slot(job.id, job.progress):
|
| 76 |
+
await asyncio.sleep(0.05)
|
| 77 |
+
job.progress.check()
|
| 78 |
+
return OK
|
| 79 |
|
| 80 |
job = store.start("file", FILE_STEPS, work)
|
| 81 |
+
await asyncio.sleep(0.01)
|
| 82 |
+
store.cancel(job.id)
|
| 83 |
+
await asyncio.sleep(0.1)
|
| 84 |
+
assert job.status == "error" and job.response["errors"] == ["cancelled"]
|
| 85 |
|
| 86 |
+
asyncio.run(scenario())
|
| 87 |
+
|
| 88 |
+
|
| 89 |
+
def test_a_job_nobody_polls_is_dropped_when_its_turn_comes(monkeypatch) -> None:
|
| 90 |
+
monkeypatch.setattr(aj, "ABANDON_SEC", 0.02)
|
| 91 |
+
|
| 92 |
+
async def scenario() -> None:
|
| 93 |
+
store = JobStore()
|
| 94 |
+
store.start("file", FILE_STEPS, gated(0.08))
|
| 95 |
+
waiting = store.start("file", FILE_STEPS, gated(0.01))
|
| 96 |
+
await asyncio.sleep(0.2)
|
| 97 |
+
assert waiting.status == "error" and waiting.response["errors"] == ["job_abandoned"]
|
| 98 |
+
|
| 99 |
+
asyncio.run(scenario())
|
| 100 |
+
|
| 101 |
+
|
| 102 |
+
def test_same_file_shares_the_running_job_then_hits_the_cache() -> None:
|
| 103 |
+
async def scenario() -> None:
|
| 104 |
+
store = JobStore()
|
| 105 |
+
a = store.start("file", FILE_STEPS, gated(0.05), key="file:abc")
|
| 106 |
+
b = store.start("file", FILE_STEPS, gated(0.05), key="file:abc")
|
| 107 |
+
assert a is b
|
| 108 |
+
await asyncio.sleep(0.15)
|
| 109 |
+
c = store.start("file", FILE_STEPS, gated(0.05), key="file:abc")
|
| 110 |
+
assert c is not a and c.cached and c.status == "done" and c.response == OK
|
| 111 |
+
|
| 112 |
+
asyncio.run(scenario())
|
| 113 |
+
|
| 114 |
+
|
| 115 |
+
def test_full_queue_is_refused(monkeypatch) -> None:
|
| 116 |
+
monkeypatch.setattr(aj, "MAX_PENDING", 2)
|
| 117 |
+
|
| 118 |
+
async def scenario() -> None:
|
| 119 |
+
store = JobStore()
|
| 120 |
+
store.start("file", FILE_STEPS, gated(0.05))
|
| 121 |
+
store.start("file", FILE_STEPS, gated(0.05))
|
| 122 |
+
with pytest.raises(aj.QueueFull):
|
| 123 |
+
store.start("file", FILE_STEPS, gated(0.05))
|
| 124 |
+
await asyncio.sleep(0.2)
|
| 125 |
|
| 126 |
asyncio.run(scenario())
|
| 127 |
|
|
|
|
| 130 |
async def scenario() -> None:
|
| 131 |
store = JobStore()
|
| 132 |
|
| 133 |
+
async def boom(job):
|
| 134 |
raise RuntimeError("boom")
|
| 135 |
|
| 136 |
job = store.start("file", FILE_STEPS, boom)
|
|
|
|
| 141 |
asyncio.run(scenario())
|
| 142 |
|
| 143 |
|
| 144 |
+
def test_progress_check_raises_when_cancelled() -> None:
|
| 145 |
+
progress = Progress(FILE_STEPS)
|
| 146 |
+
progress.check()
|
| 147 |
+
progress.cancelled = True
|
| 148 |
+
with pytest.raises(JobCancelled):
|
| 149 |
+
progress.check()
|
| 150 |
+
|
| 151 |
+
|
| 152 |
def test_confidence_band_is_measured_from_the_threshold() -> None:
|
| 153 |
svc = XAIInferenceService.__new__(XAIInferenceService)
|
| 154 |
svc.threshold = 0.431577
|
|
|
|
| 155 |
assert svc._confidence_band(0.45).tier == "uncertain"
|
|
|
|
| 156 |
assert svc._confidence_band(0.40).tier == "uncertain"
|
| 157 |
assert svc._confidence_band(0.97).tier == "very_strong"
|
| 158 |
assert svc._confidence_band(0.02).tier == "very_strong"
|
| 159 |
+
assert svc._confidence_band(0.431577).margin == 0.0
|
|
|