Spaces:
Running
Running
feat: background analysis job queue servisi ve corresponding tests eklendi
Browse files
app/services/analysis_jobs.py
CHANGED
|
@@ -234,7 +234,10 @@ class Job:
|
|
| 234 |
data: Dict[str, Any] = {
|
| 235 |
"jobId": self.id,
|
| 236 |
"sourceType": self.source_type,
|
| 237 |
-
"
|
|
|
|
|
|
|
|
|
|
| 238 |
"steps": self.progress.snapshot(),
|
| 239 |
"elapsedSec": round((self.updated if self.finished else time.time()) - self.created, 2),
|
| 240 |
"response": self.response,
|
|
|
|
| 234 |
data: Dict[str, Any] = {
|
| 235 |
"jobId": self.id,
|
| 236 |
"sourceType": self.source_type,
|
| 237 |
+
# "running" until settled, whether waiting or working, so clients
|
| 238 |
+
# that predate the queue keep polling; `phase` tells them apart.
|
| 239 |
+
"status": self.status if self.finished else "running",
|
| 240 |
+
"phase": self.status if self.finished else self.progress.phase,
|
| 241 |
"steps": self.progress.snapshot(),
|
| 242 |
"elapsedSec": round((self.updated if self.finished else time.time()) - self.created, 2),
|
| 243 |
"response": self.response,
|
tests/test_analysis_jobs.py
CHANGED
|
@@ -39,8 +39,9 @@ def test_jobs_take_the_cpu_one_at_a_time_in_order() -> None:
|
|
| 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]["
|
| 43 |
-
assert snaps[1]["status"] == "
|
|
|
|
| 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)
|
|
|
|
| 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]["phase"] == "running" and "queue" not in snaps[0]
|
| 43 |
+
assert snaps[1]["status"] == "running" and snaps[1]["phase"] == "queued"
|
| 44 |
+
assert snaps[1]["queue"]["position"] == 1 and snaps[1]["queue"]["ahead"] == 1
|
| 45 |
assert snaps[2]["queue"]["position"] == 2
|
| 46 |
await asyncio.sleep(0.25)
|
| 47 |
assert all(j.to_dict()["status"] == "done" for j in jobs)
|
tests/test_analyze_queue.py
CHANGED
|
@@ -70,7 +70,7 @@ def test_ten_visitors_at_once_queue_and_finish(stub_pipeline) -> None:
|
|
| 70 |
assert len(set(ids)) == 10
|
| 71 |
|
| 72 |
snaps = [(await client.get(f"/api/analyze/jobs/{i}")).json() for i in ids]
|
| 73 |
-
positions = sorted(s["queue"]["position"] for s in snaps if s["
|
| 74 |
assert positions == list(range(1, len(positions) + 1))
|
| 75 |
assert len(positions) >= 8
|
| 76 |
|
|
|
|
| 70 |
assert len(set(ids)) == 10
|
| 71 |
|
| 72 |
snaps = [(await client.get(f"/api/analyze/jobs/{i}")).json() for i in ids]
|
| 73 |
+
positions = sorted(s["queue"]["position"] for s in snaps if s["phase"] == "queued")
|
| 74 |
assert positions == list(range(1, len(positions) + 1))
|
| 75 |
assert len(positions) >= 8
|
| 76 |
|