큐를 실제로 돌린다: 적재 → 워커 claim → 핸들러 실행 → 결과 저장 → DONE. API(적재)와 워커(소비)를 분리 프로세스로(코드베이스 공유, 독립 스케일). - worker/notify: 전용 asyncpg LISTEN 리스너. enqueue 에서 pg_notify → 유휴 워커 즉시 기상(폴링 제거) - worker/runner: Worker(claim→처리, 처리중 heartbeat 로 lease 갱신, complete/fail) + run_reaper - worker/handlers: job_type 별 핸들러(주입식). SEARCH=소스 어댑터 검색→정규화 결과 - worker_main: API 분리 워커 진입점(브라우저 무거워 기본 동시성 1) - job_crud.enqueue: 삽입 시 pg_notify (중복 스킵 시엔 미발생) - tests: drain→DONE·실패→재시도→dead·reaper 회수 후 재처리·NOTIFY 기상 4건 (전체 20/20) - 라이브 E2E 확인: 적재→워커가 실제 쿠팡 검색(8건)→DONE Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
70 lines
2.4 KiB
Python
70 lines
2.4 KiB
Python
"""워커 루프 테스트 — drain→DONE, 실패→재시도→dead, reaper 회수 후 재처리, NOTIFY 깨움.
|
|
핸들러는 fake(브라우저 없이) — 워커 로직만 결정론적으로 검증한다."""
|
|
|
|
import asyncio
|
|
|
|
import pytest_asyncio
|
|
from sqlalchemy import text
|
|
|
|
from common.enums import JobStatus, JobType
|
|
from crud.job_crud import JobQueue
|
|
from worker.notify import JobListener
|
|
from worker.runner import Worker
|
|
|
|
|
|
@pytest_asyncio.fixture
|
|
async def q(db_engine):
|
|
async with db_engine.begin() as conn:
|
|
await conn.execute(text("TRUNCATE job"))
|
|
return JobQueue()
|
|
|
|
|
|
async def test_worker_drains_all_to_done(q):
|
|
for i in range(3):
|
|
await q.enqueue(JobType.SEARCH.value, {"product_name": f"item{i}"})
|
|
|
|
async def handler(job):
|
|
return {"ok": True, "q": job["payload"]["product_name"]}
|
|
|
|
processed = await Worker("w1", q, handler).drain()
|
|
assert processed == 3
|
|
assert (await q.counts())["DONE"] == 3
|
|
|
|
|
|
async def test_worker_failure_retries_then_dead(q):
|
|
await q.enqueue(JobType.SEARCH.value, {"product_name": "x"}, max_attempts=2)
|
|
|
|
async def boom(job):
|
|
raise RuntimeError("nope")
|
|
|
|
w = Worker("w1", q, boom, backoff_fn=lambda a: 0) # 백오프 0 → 즉시 재시도 가능
|
|
assert await w.process_one() is True # 1/2 실패 → PENDING
|
|
assert (await q.counts())["PENDING"] == 1
|
|
assert await w.process_one() is True # 2/2 실패 → DEAD
|
|
counts = await q.counts()
|
|
assert counts["DEAD"] == 1 and counts["PENDING"] == 0
|
|
assert await w.process_one() is False # DEAD 는 claim 대상 아님
|
|
|
|
|
|
async def test_reaper_reclaims_then_worker_reprocesses(q):
|
|
jid = await q.enqueue(JobType.SEARCH.value, {"product_name": "x"})
|
|
await q.claim("dead-worker", lease_sec=1) # 점유 후 사망 흉내
|
|
await asyncio.sleep(1.3)
|
|
assert jid in await q.reap() # 회수 → PENDING
|
|
|
|
async def handler(job):
|
|
return {"ok": True}
|
|
|
|
assert await Worker("w2", q, handler).process_one() is True
|
|
assert (await q.counts())["DONE"] == 1
|
|
|
|
|
|
async def test_enqueue_notifies_listener(q):
|
|
listener = JobListener()
|
|
await listener.start()
|
|
try:
|
|
await q.enqueue(JobType.SEARCH.value, {"product_name": "x"}) # pg_notify 발생
|
|
assert await listener.wait(3.0) is True # 즉시 깨어남
|
|
finally:
|
|
await listener.close()
|