o2o-site-AEO/solution/backend/tests/test_job_queue.py
Mina Choi 11d30bb3d1 [chore] solution,admin,ontology: 코드 주석을 한 줄로 — 히스토리 주석 삭제
여러 줄 주석이 설명보다 경위(예전·실측·지적)를 적고 있어 읽는 사람이 결론을 찾기 어려웠다.

- ts·tsx·js·mjs·css·py 478개: 여러 줄 주석은 첫 문장 한 줄로, 과거형·날짜 문장은 삭제
- 주석 위치는 TypeScript 파서·파이썬 tokenize/ast 로 찾는다 — 문자열 안의 # · /* 는 건드리지 않는다
- eslint·ts·noqa·type: ignore 같은 지시 주석은 그대로 둔다

파이썬 275개 정리 전후 AST 동일, TS 298개 주석 뺀 토큰 동일(빈 JSX 주석 10곳만 차이).
site·frontend·admin tsc, site vitest 105 passed

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
2026-09-28 16:05:19 +09:00

251 lines
8.8 KiB
Python

"""작업 큐 — 도커에서 워커를 여러 개 띄웠을 때 지켜져야 하는 것들."""
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 한다."""
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_progress_rejects_old_worker_and_resets_on_retry(db_engine):
from services.job_progress import JobProgress
q = JobQueue()
jid = await q.enqueue(JobType.COPY.value, {})
old = await q.claim("w1")
progress = JobProgress(old, ("prepare", "generate"))
async with progress.step("prepare"):
pass
try:
async with progress.step("generate"):
raise ValueError("upstream failed")
except ValueError:
pass
row = await q.get(jid)
assert row["progress"]["steps"][1]["status"] == "failed"
await q.fail(jid, "w1", "upstream failed", backoff_sec=0)
new = await q.claim("w2")
assert not await q.set_progress(old, {"steps": [], "attempt": 1})
retried = JobProgress(new, ("prepare", "generate"))
async with retried.step("prepare"):
row = await q.get(jid)
assert row["progress"]["attempt"] == 2
assert [s["status"] for s in row["progress"]["steps"]] == ["running", "pending"]
async def test_claim_is_atomic_across_workers(db_engine):
"""검증: 잡 1건에 워커 5개가 동시에 달려든다."""
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."""
assert await JobQueue().claim("w1") is None
async def test_dedupe_key_blocks_duplicate_active_job(db_engine):
"""검증: 같은 dedupe_key 로 두 번 적재한다."""
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 로 다시 적재한다."""
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):
"""검증: 자기가 점유하지 않은 잡을 완료 처리한다."""
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 만큼 계속 실패시킨다."""
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):
"""검증: 실패 후 백오프가 걸린 잡."""
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):
"""검증: 워커가 잡을 쥔 채 죽는다(도커 컨테이너 재시작)."""
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 [r["job_id"] for r 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 잡에 각각 재큐를 시도한다."""
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):
"""검증: 워커 루프가 잡을 집어 핸들러를 돌리고 완료 처리한다."""
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 을 처리한다."""
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_dead_letter_creates_an_alert(db_engine):
"""검증: 잡이 재시도를 소진해 DEAD 로 떨어진다."""
q = JobQueue()
jid = await q.enqueue(JobType.AI_CHECK.value, {"place_id": "p-alert-1"}, 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
async with db_engine.begin() as c:
result = await c.execute(text("SELECT kind, dedupe_key, title FROM alert_outbox WHERE kind = 'job_dead'"))
rows = [dict(r._mapping) for r in result]
assert len(rows) == 1
assert rows[0]["dedupe_key"] == "job_dead:AI_CHECK:p-alert-1"
async def test_reaper_loop_stops_on_event(db_engine):
"""검증: reaper 루프에 stop 이벤트를 건다."""
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():
"""검증: 백오프 계산."""
assert compute_backoff(1) == 5
assert compute_backoff(3) == 20
assert compute_backoff(99) == 600