o2o-negosium-original/lps/tests/test_job_queue.py
민헌 da9bd7750c feat(lps): 작업 큐 엔진 — PostgreSQL 원자적 claim + lease + dead-letter
레퍼런스의 DB-큐 반면교사를 전부 뒤집는 큐 엔진을 구축한다.
 - 원자적 claim: FOR UPDATE SKIP LOCKED 서브쿼리 + 같은 UPDATE + RETURNING
   → 워커 다수여도 이중 할당 원천 불가(fetch/claim 미분리)
 - 모든 전이는 조건부 CAS(status/worker_id 가드) + RETURNING
 - 복구는 timeout 추측이 아닌 lease 만료 소유권 기반(reaper 회수)
 - 재시도 지수백오프 + max_attempts + dead-letter 를 큐에 내장(스크립트 난립 제거)
 - dedupe_key 부분 유니크로 활성 중복 차단, counts()로 관측(수동 psql 대체)

- enums: JobStatus(PENDING/RUNNING/DONE/DEAD)·JobType(SEARCH/OUTBOX)
- models: job 테이블(무FK·SMALLINT코드·TIMESTAMPTZ, claim/lease/dedupe 인덱스)
- crud/job_crud: JobQueue(enqueue/claim/complete/fail/renew_lease/reap/counts)
- tests: 동시 8워커 이중할당0·lease회수·재시도→dead·소유권가드 등 8건(실 lps_db)

부수: 포트 9400→9600 (negodata 9400 충돌 회피)

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-08 16:37:32 +09:00

94 lines
3.6 KiB
Python

"""작업 큐 엔진 테스트 — 원자적 claim(이중할당 불가)·lease 회수·재시도/dead-letter·소유권 가드.
실제 lps_db 에 붙어 검증한다(db_engine 이 스키마 보장)."""
import asyncio
import pytest_asyncio
from sqlalchemy import text
from common.enums import JobStatus, JobType
from crud.job_crud import JobQueue, compute_backoff
@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_enqueue_claim_complete(q):
jid = await q.enqueue(JobType.SEARCH.value, {"query": "커피"}, dedupe_key="search-커피")
assert jid
job = await q.claim("w1")
assert job and job["job_id"] == jid
assert job["payload"]["query"] == "커피" and job["attempts"] == 1
assert await q.complete(jid, "w1", {"count": 3}) is True
counts = await q.counts()
assert counts["DONE"] == 1 and counts["PENDING"] == 0
async def test_dedupe_blocks_active_duplicate(q):
a = await q.enqueue(JobType.SEARCH.value, {"q": 1}, dedupe_key="k")
b = await q.enqueue(JobType.SEARCH.value, {"q": 1}, dedupe_key="k")
assert a and b is None # 활성 중복 차단
# 완료로 빠지면 같은 키 재적재 가능
job = await q.claim("w1")
await q.complete(job["job_id"], "w1")
c = await q.enqueue(JobType.SEARCH.value, {"q": 1}, dedupe_key="k")
assert c
async def test_atomic_claim_no_double_assignment(q):
N = 12
for i in range(N):
await q.enqueue(JobType.SEARCH.value, {"i": i})
# 8개 워커가 동시에 claim → 서로 다른 잡만, 이중 할당 0
results = await asyncio.gather(*[q.claim(f"w{i}") for i in range(8)])
claimed = [r["job_id"] for r in results if r]
assert len(claimed) == 8
assert len(set(claimed)) == 8
async def test_priority_and_order(q):
await q.enqueue(JobType.SEARCH.value, {"n": "low"}, priority=100)
await q.enqueue(JobType.SEARCH.value, {"n": "high"}, priority=1)
job = await q.claim("w1")
assert job["payload"]["n"] == "high" # priority 낮은 값 우선
async def test_lease_reclaim_by_reaper(q):
jid = await q.enqueue(JobType.SEARCH.value, {"q": "x"})
job = await q.claim("w1", lease_sec=1)
assert job["job_id"] == jid and job["attempts"] == 1
assert await q.reap() == [] # 아직 lease 유효 → 회수 없음
await asyncio.sleep(1.3) # lease 만료
assert jid in await q.reap() # 회수됨(워커 사망 시나리오)
job2 = await q.claim("w2") # 다시 claim 가능, attempts 누적
assert job2["job_id"] == jid and job2["attempts"] == 2
async def test_retry_then_dead_letter(q):
jid = await q.enqueue(JobType.SEARCH.value, {"q": "x"}, max_attempts=2)
await q.claim("w1")
assert await q.fail(jid, "w1", "boom", backoff_sec=0) == JobStatus.PENDING.value # 1/2 → 재큐
job2 = await q.claim("w1")
assert job2["attempts"] == 2
assert await q.fail(jid, "w1", "boom2", backoff_sec=0) == JobStatus.DEAD.value # 2/2 → dead-letter
assert (await q.counts())["DEAD"] == 1
async def test_transitions_require_ownership(q):
jid = await q.enqueue(JobType.SEARCH.value, {"q": "x"})
await q.claim("w1")
assert await q.complete(jid, "intruder") is False # 소유 아님 → 거부(CAS 가드)
assert await q.fail(jid, "intruder", "no") is None
assert await q.complete(jid, "w1") is True
def test_backoff_is_exponential_capped():
assert compute_backoff(1, base=5) == 5
assert compute_backoff(2, base=5) == 10
assert compute_backoff(3, base=5) == 20
assert compute_backoff(100, base=5, cap=600) == 600