"""작업 큐 — 도커에서 워커를 여러 개 띄웠을 때 지켜져야 하는 것들. 수집·비전분석·빌드는 몇 분 걸려 동기 요청으로 못 한다. 그 큐가 실제로: 1. 같은 잡을 두 워커에 이중 할당하지 않는가 (FOR UPDATE SKIP LOCKED) 2. 워커가 죽어도 잡이 증발하지 않는가 (lease 만료 → reaper 회수) 3. 같은 요청을 두 번 눌러도 잡이 두 번 돌지 않는가 (dedupe_key) 4. 실패가 백오프 재시도 → dead-letter 로 흐르는가 """ import asyncio import uuid from sqlalchemy import text from common.enums import JobStatus, JobType from crud.job_crud import JobQueue, compute_backoff from worker.handlers import UnknownJobType, build_handler from worker.runner import Worker, run_reaper async def test_enqueue_and_claim(db_engine): """검증: 잡을 넣고 워커가 claim 한다. 기대결과: RUNNING 으로 전이되고 attempts 가 1 올라간다.""" q = JobQueue() jid = await q.enqueue(JobType.COLLECT.value, {"place_id": "abc"}) assert jid claimed = await q.claim("w1") assert claimed["job_id"] == jid assert claimed["payload"] == {"place_id": "abc"} assert claimed["attempts"] == 1 row = await q.get(jid) assert row["status"] == JobStatus.RUNNING.value async def test_claim_is_atomic_across_workers(db_engine): """검증: 잡 1건에 워커 5개가 동시에 달려든다. 기대결과: 정확히 1명만 가져간다 — 도커에서 워커를 몇 개로 스케일하든 이중 실행이 없다.""" q = JobQueue() jid = await q.enqueue(JobType.COLLECT.value, {"place_id": "race"}) results = await asyncio.gather(*[q.claim(f"w{i}") for i in range(5)]) won = [r for r in results if r is not None] assert len(won) == 1, f"이중 할당 발생: {len(won)}명이 같은 잡을 가져갔다" assert won[0]["job_id"] == jid async def test_claim_returns_none_when_empty(db_engine): """검증: 빈 큐에서 claim. 기대결과: None — 워커는 LISTEN 대기로 넘어간다.""" assert await JobQueue().claim("w1") is None async def test_dedupe_key_blocks_duplicate_active_job(db_engine): """검증: 같은 dedupe_key 로 두 번 적재한다. 기대결과: 두 번째는 None — 같은 사업장 수집을 두 번 눌러도 잡은 하나다.""" q = JobQueue() key = f"collect:{uuid.uuid4()}" first = await q.enqueue(JobType.COLLECT.value, {}, dedupe_key=key) second = await q.enqueue(JobType.COLLECT.value, {}, dedupe_key=key) assert first is not None assert second is None active = await q.find_active(key) assert active["job_id"] == first async def test_dedupe_key_frees_after_completion(db_engine): """검증: 완료된 뒤 같은 dedupe_key 로 다시 적재한다. 기대결과: 새 잡이 생긴다 — 부분 유니크가 PENDING/RUNNING 에만 걸리기 때문.""" q = JobQueue() key = f"collect:{uuid.uuid4()}" first = await q.enqueue(JobType.COLLECT.value, {}, dedupe_key=key) claimed = await q.claim("w1") assert await q.complete(claimed["job_id"], "w1", {"ok": True}) is True second = await q.enqueue(JobType.COLLECT.value, {}, dedupe_key=key) assert second is not None and second != first async def test_complete_requires_ownership(db_engine): """검증: 자기가 점유하지 않은 잡을 완료 처리한다. 기대결과: False — 소유권 가드(worker_id) 때문에 남의 잡을 끝낼 수 없다.""" q = JobQueue() await q.enqueue(JobType.BUILD.value, {}) claimed = await q.claim("w1") assert await q.complete(claimed["job_id"], "other-worker", {}) is False assert await q.complete(claimed["job_id"], "w1", {}) is True async def test_fail_retries_then_dead_letters(db_engine): """검증: max_attempts 만큼 계속 실패시킨다. 기대결과: 시도가 남으면 PENDING(백오프), 소진되면 DEAD — 무한 재시도가 없다.""" q = JobQueue() jid = await q.enqueue(JobType.VISION.value, {}, max_attempts=2) claimed = await q.claim("w1") assert await q.fail(jid, "w1", "boom", backoff_sec=0) == JobStatus.PENDING.value claimed = await q.claim("w1") assert claimed is not None, "백오프 0 이면 즉시 다시 claim 가능해야 한다" assert await q.fail(jid, "w1", "boom again", backoff_sec=0) == JobStatus.DEAD.value row = await q.get(jid) assert row["status"] == JobStatus.DEAD.value assert "boom again" in row["last_error"] async def test_backoff_delays_next_claim(db_engine): """검증: 실패 후 백오프가 걸린 잡. 기대결과: run_after 가 지나기 전에는 claim 되지 않는다.""" q = JobQueue() jid = await q.enqueue(JobType.COPY.value, {}, max_attempts=5) await q.claim("w1") await q.fail(jid, "w1", "transient", backoff_sec=60) assert await q.claim("w1") is None, "백오프 중인 잡이 claim 됐다" async def test_reaper_reclaims_dead_worker_job(db_engine): """검증: 워커가 잡을 쥔 채 죽는다(도커 컨테이너 재시작). 기대결과: lease 가 만료되면 reaper 가 회수해 다시 대기로 돌린다 — 잡이 증발하지 않는다.""" q = JobQueue() jid = await q.enqueue(JobType.COLLECT.value, {}) await q.claim("dead-worker", lease_sec=120) # 워커 사망을 흉내 낸다: lease 를 과거로 돌린다. async with db_engine.begin() as conn: await conn.execute( text("UPDATE jobs SET lease_until = now() - interval '1 minute' WHERE job_id = CAST(:id AS uuid)"), {"id": jid}, ) reclaimed = await q.reap() assert jid in reclaimed row = await q.get(jid) assert row["status"] == JobStatus.PENDING.value assert "lease-expired reclaim" in row["last_error"] assert await q.claim("new-worker") is not None async def test_requeue_only_works_on_dead_jobs(db_engine): """검증: DEAD 잡과 PENDING 잡에 각각 재큐를 시도한다. 기대결과: DEAD 만 되살아나고 attempts 가 0으로 초기화된다.""" q = JobQueue() jid = await q.enqueue(JobType.COLLECT.value, {}, max_attempts=1) await q.claim("w1") assert await q.fail(jid, "w1", "fatal", backoff_sec=0) == JobStatus.DEAD.value assert await q.requeue(jid) == jid row = await q.get(jid) assert row["status"] == JobStatus.PENDING.value assert row["attempts"] == 0 assert await q.requeue(jid) is None, "PENDING 잡은 재큐 대상이 아니다" async def test_worker_runs_handler_and_completes(db_engine): """검증: 워커 루프가 잡을 집어 핸들러를 돌리고 완료 처리한다. 기대결과: 핸들러 반환값이 result 에 저장되고 status=DONE.""" q = JobQueue() jid = await q.enqueue(JobType.COLLECT.value, {"place_id": "p1"}) seen = [] async def handler(job): seen.append(job["payload"]) return {"collected": 3} worker = Worker("w-test", q, handler, job_deadline_sec=10) assert await worker.process_one() is True assert seen == [{"place_id": "p1"}] row = await q.get(jid) assert row["status"] == JobStatus.DONE.value assert row["result"] == {"collected": 3} async def test_worker_deadline_cancels_hung_handler(db_engine): """검증: 핸들러가 행(hang)에 빠진다. 기대결과: 데드라인에 취소되고 잡은 재큐된다 — 행이 워커 슬롯을 영구 점유하지 않는다.""" q = JobQueue() jid = await q.enqueue(JobType.COLLECT.value, {}, max_attempts=5) async def hung(job): await asyncio.sleep(30) worker = Worker("w-test", q, hung, job_deadline_sec=0.2, backoff_fn=lambda _a: 0) assert await worker.process_one() is True row = await q.get(jid) assert row["status"] == JobStatus.PENDING.value assert row["last_error"].startswith("JobDeadlineExceeded") async def test_unregistered_job_type_fails_loudly(db_engine): """검증: 핸들러가 등록되지 않은 JobType 을 처리한다. 기대결과: UnknownJobType 으로 실패하고 last_error 에 남는다(조용히 성공하지 않는다).""" q = JobQueue() jid = await q.enqueue(JobType.AI_CHECK.value, {}, max_attempts=1) worker = Worker("w-test", q, build_handler(), backoff_fn=lambda _a: 0) await worker.process_one() row = await q.get(jid) assert row["status"] == JobStatus.DEAD.value assert UnknownJobType.__name__ in row["last_error"] async def test_reaper_loop_stops_on_event(db_engine): """검증: reaper 루프에 stop 이벤트를 건다. 기대결과: 즉시 빠져나온다(graceful shutdown 이 매달리지 않는다).""" stop = asyncio.Event() task = asyncio.create_task(run_reaper(JobQueue(), stop, interval=0.05)) await asyncio.sleep(0.1) stop.set() await asyncio.wait_for(task, timeout=2) def test_backoff_is_exponential_and_capped(): """검증: 백오프 계산. 기대결과: 5초부터 2배씩 늘고 600초에서 멈춘다.""" assert compute_backoff(1) == 5 assert compute_backoff(3) == 20 assert compute_backoff(99) == 600