o2o-site-AEO/solution/backend/tests/test_job_queue.py
Mina Choi 9d25ed613e 구조: 사장님(solution)과 내부 운영(admin)을 두 앱으로 가른다
최상단을 프로젝트 단위로 평평하게 둔다 — o2o-negosium 과 같은 규약이고, 이 레포만
다르게 갈 이유가 없다. negodata/{backend,front} 가 프로젝트 안에서 f/b 를 가르는 선례,
lps-admin/ 이 백엔드 없이 프론트만 가진 최상단 폴더의 선례다.

  backend/ frontend/{admin,site,shared}  →  solution/{backend,front,site,shared} + admin/

## 왜

내부 라우트(/local-content, /places/:id/seo)의 이름과 화면 코드가 사장님 번들에
그대로 실려 나가고 있었다. UserRole.DEVELOPER 주석의 "고객사에 존재를 노출하지 않는다"를
번들이 깨고 있었다 — 라우트 가드는 화면을 가리지 번들은 못 가린다.
번들을 갈라 확인했다: 사장님 dist 에서 local-content · /places · SeoAudit 이 전부 0건이다.

그 과정에서 두 곳이 더 새고 있었다.
- AppShell 의 NAV 배열이 내부 메뉴를 하드코딩하고 있었다. 앱을 가른 뒤에도 dist 에
  local-content 가 남아서 찾았다. 메뉴는 이제 앱이 prop 으로 들고 온다.
- EditorHeader·BuilderPage·LoginPage 가 /places 로 링크하고 있었다. 그 화면이 admin 으로
  나갔으니 사장님 앱에서는 404 다. 링크를 걷어내고 LoginPage 기본 도착지는 '/' 로 바꿨다
  (앱마다 홈이 다르고 각 라우터의 '/' 가 이미 그걸 안다).

## admin 에 백엔드를 두지 않았다

내부 화면이 부르는 훅이 전부 router/v1/{place,fact,local,validator} 에 이미 있다.
자체 백엔드를 두면 place·fact·link 를 같은 DB 에 대고 두 번 구현하게 된다.
대가는 solution/backend 가 죽으면 admin 도 멈추는 것 — 내부 도구라 감수한다.

## admin 의 `@` 는 solution/front/src 를 가리킨다

내부 화면이 쓰는 API 클라이언트·UI·수집 배선이 solution 에 한 벌만 있고 그 파일들끼리도
`@/...` 로 서로를 부른다. admin 에서 `@` 를 자기 src 로 잡으면 그 참조가 전부 깨진다
(실측 TS2307 14건). 복제하는 길도 있지만 RecollectPanel 주석이 금지한다 —
"수집 경로를 두 벌 만들면 확정 게이트"가 갈라진다.
admin 자기 파일만 `@admin` 이고, 의존 방향은 admin → solution 한 쪽뿐이다.

admin 이 여는 빌더는 다른 오리진이라 절대 URL + 새 탭이다(admin/src/lib/solutionUrl.ts).
react-router Link 로 두면 admin 안에서 라우트를 찾다 404 다.

## 그 밖

- npm 워크스페이스 루트를 레포 루트로 올렸다(admin 이 solution 밖이라).
- docker-compose 를 255→174줄로 줄이고 admin(:3002) 서비스를 넣었다. ADMIN_BIND 기본값은
  127.0.0.1 — 0.0.0.0 으로 열면 앱을 가른 의미가 없다.
- 발행 호스트를 프론트 .env 에 따로 적지 않는다. compose 가 루트의 SITE_PUBLIC_HOST 를
  VITE_PUBLISH_HOST 로 흘려보낸다 — 두 곳에 적으면 canonical 과 화면 주소가 조용히 갈라진다.
- nginx/site.conf 를 git 에서 빼고 .example 만 남겼다(.env·*.toml 과 같은 규약).
  compose 가 bind mount 하므로 클론 직후 복사해야 한다 — 없으면 Docker 가 그 자리에
  디렉토리를 만들어 nginx 가 설정 없이 뜬다.
- config.test.toml.example 을 추가했다. 없으면 클론한 사람이 pytest 를 아예 못 돌린다
  (conftest import 단계에서 죽는다). 외부 API 키는 전부 빈값이다 —
  APP_ENV=test 가 .env 를 안 읽는 이유를 여기서 우회하면 안 된다.
- 경로가 한 칸 깊어져 test_schema_ddl(parents[2]→[3]) 과 test_site_theme 을 고쳤다.

검증: front·admin·site 전부 lint 0 / build 0. 백엔드 514 passed.
남은 4건(test_build_publish 3 · test_snapshot 1)은 이 변경 전부터 실패하던 것으로,
손대지 않은 메인 체크아웃에서 같은 4건이 같게 실패하는 것을 확인했다.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_019uYhHQdssRubirPirrdJJC
2026-08-31 15:12:09 +09:00

229 lines
8.9 KiB
Python

"""작업 큐 — 도커에서 워커를 여러 개 띄웠을 때 지켜져야 하는 것들.
수집·비전분석·빌드는 몇 분 걸려 동기 요청으로 못 한다. 그 큐가 실제로:
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