"""Ten visitors at once, through the real endpoints and queue. The analysis itself is replaced by a stand-in that holds the CPU slot for a moment, so this checks the wiring (queue positions, polling, cancel, shared uploads) without running any model. """ from __future__ import annotations import asyncio from collections import defaultdict import io import math import struct import httpx import pytest from app.main import app from app.routes import analyze as route from app.routes.analyze_schemas import AnalyzeResponse from app.services import analysis_jobs as aj def _wav(seed: int) -> bytes: """A short, distinct WAV per visitor (distinct content, distinct job).""" sr, n = 8000, 8000 frames = b"".join(struct.pack(" httpx.Response: return await client.post( "/api/analyze/jobs", data={"sourceType": "file"}, files={"file": (f"t{seed}.wav", io.BytesIO(_wav(seed)), "audio/wav")}, headers={"x-forwarded-for": ip}, ) def test_ten_visitors_at_once_queue_and_finish(stub_pipeline) -> None: async def scenario() -> None: transport = httpx.ASGITransport(app=app) async with httpx.AsyncClient(transport=transport, base_url="http://test") as client: responses = await asyncio.gather(*(_post(client, i, f"10.0.0.{i}") for i in range(10))) assert all(r.status_code == 202 for r in responses) ids = [r.json()["jobId"] for r in responses] assert len(set(ids)) == 10 snaps = [(await client.get(f"/api/analyze/jobs/{i}")).json() for i in ids] positions = sorted(s["queue"]["position"] for s in snaps if s["phase"] == "queued") assert positions == list(range(1, len(positions) + 1)) assert len(positions) >= 8 # One visitor gives up; the others keep their order. cancelled = ids[-1] assert (await client.delete(f"/api/analyze/jobs/{cancelled}")).status_code == 200 for _ in range(80): snaps = [(await client.get(f"/api/analyze/jobs/{i}")).json() for i in ids] if all(s["status"] in ("done", "error") for s in snaps): break await asyncio.sleep(0.05) by_id = dict(zip(ids, snaps)) assert by_id[cancelled]["response"]["errors"] == ["cancelled"] assert all(by_id[i]["status"] == "done" for i in ids[:-1]) health = (await client.get("/api/health")).json() assert health["models"]["jobs"]["pending"] == 0 asyncio.run(scenario()) def test_same_upload_twice_shares_one_job(stub_pipeline) -> None: async def scenario() -> None: transport = httpx.ASGITransport(app=app) async with httpx.AsyncClient(transport=transport, base_url="http://test") as client: a, b = await asyncio.gather(_post(client, 1, "10.0.1.1"), _post(client, 1, "10.0.1.2")) assert a.json()["jobId"] == b.json()["jobId"] asyncio.run(scenario())