o2o-negosium-original/lps/tests/test_worker.py
민헌 abcb58ba05 feat(lps): 워커 루프 + LISTEN/NOTIFY — 큐 소비 파이프라인 가동
큐를 실제로 돌린다: 적재 → 워커 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>
2026-07-08 16:58:58 +09:00

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()