o2o-site-AEO/solution/backend/worker/runner.py
Mina Choi b4085a0e0f [fix] solution: 온보딩 생성·크롤링 진단·장애 알림 묶음
운영 번들 자동 로그인 자격증명 유출, 온보딩 COPY 잡이 Gemini 429 로 죽던 것,
크롤링 실패가 로그에만 남던 것을 한 번에 정리한다. 실측(2026-09-15 밤, 킹서버):
사진분석 배치가 Gemini 분당 쿼터를 다 써서 같은 키를 쓰는 온보딩 COPY 잡도 같이
429 를 맞고 DEAD 로 갔다 — 확인된 fact 만으로도 편집·발행이 되는데 잡을 죽일
이유가 없었다.

- solution/frontend: `VITE_AUTO_LOGIN_ID`·`PW` 를 운영 진입점에 안 넘긴다(자동 로그인은
  dev 서버 전용) + `Step5Generating` 겉모습을 이전 카드 스타일로, 데이터는 실제 잡
  진행(useGenerationJob) 그대로
- solution/backend: copy_service — Gemini 호출 실패해도 잡을 안 죽이고 fact 만으로 계속.
  db_session_manager — 유니크 제약 충돌(정상 경로) 로그를 ERROR → WARN.
  worker/runner + alert_service + teams_webhook — 잡 dead-letter·발행 실패·큐 정체를
  Teams 로 알림(영구 저장 + 재시도 + dedupe). `/readyz` 추가.
  collect_diagnostics(신규) — 크롤링 채널별 실패를 jobs.result 에 구조화해서 싣는다.
- postgres-init: 0015(users token_version) · 0016(alert_outbox) 마이그레이션

검증: 백엔드 pytest 759 passed. tsc(solution/frontend) 통과. Teams 알림 실채널 수신 확인.
2026-09-16 16:25:02 +09:00

172 lines
8.2 KiB
Python

"""워커 루프 + reaper. (LPS `worker/runner.py` 이식)
워커는 큐에서 잡을 원자적으로 claim → 핸들러 실행 → complete/fail 한다.
- 처리 중 heartbeat 로 lease 를 갱신(긴 잡이 reaper 에 회수되지 않게).
수집·비전분석은 몇 분이 정상이라 heartbeat 없이는 lease 가 먼저 만료된다.
- 핸들러엔 데드라인(job_deadline_sec)을 건다 — heartbeat 가 lease 를 계속 갱신하므로
핸들러가 행하면 reaper 로는 영원히 회수 불가.
초과 시 취소 후 fail 처리 → 백오프 재큐(소진 시 DEAD), 워커 슬롯은 즉시 다음 잡으로.
- claim 이 비면 LISTEN 알림 또는 짧은 타임아웃까지 대기(폴링 최소화).
- 핸들러는 주입식(async def(job)->dict) — 프로덕션은 수집 파이프라인, 테스트는 fake.
"""
import asyncio
from common.enums import JobStatus, JobType
from common.job_errors import PermanentJobError
from common.logger import LOG
from crud.job_crud import JobQueue, compute_backoff
from services import alert_service
def _job_type_name(job_type: int) -> str:
try:
return JobType(job_type).name
except ValueError:
return str(job_type)
async def _alert_job_dead(job: dict, error: str) -> None:
"""잡이 dead-letter 로 떨어졌다 — 재시도를 소진했다는 뜻이다(수동 개입 대상).
★ dedupe_key 는 "같은 대상이 반복해서 죽는가" 를 잡는다. job_id 는 잡마다 새로 생기므로
쓰지 않는다 — payload 의 place_id(대부분의 잡이 갖는 자연키)가 있으면 그걸 쓰고,
없으면 job_type 만으로 묶는다(어느 쪽이든 완벽하진 않지만, 없는 것보다는 낫다)."""
job_type = job.get("job_type")
type_name = _job_type_name(job_type)
payload = job.get("payload") or {}
target = payload.get("place_id") or payload.get("region_code") or job.get("job_id", "")
await alert_service.send_alert(
kind="job_dead",
title=f"[{type_name}] 잡이 재시도를 소진했다(DEAD)",
detail=f"job_id={job.get('job_id')} attempts={job.get('attempts')}\n{error}",
dedupe_key=f"job_dead:{type_name}:{target}",
)
class Worker:
"""큐에서 잡을 원자적으로 claim → 핸들러 실행 → complete/fail 하는 워커 루프."""
def __init__(
self,
worker_id: str,
queue: JobQueue,
handler,
lease_sec: int = 120,
backoff_fn=compute_backoff,
job_deadline_sec: float = 900.0,
):
self.worker_id = worker_id
self.queue = queue
self.handler = handler
self.lease_sec = lease_sec
self.backoff_fn = backoff_fn
# 잡 1건 처리 시간 상한(0 이면 무제한 — 테스트용).
# 수집 파이프라인은 Perplexity(10~30s) + 크롤링 + Vision(사진 20~50장) 이라 15분을 기본으로 둔다.
self.job_deadline_sec = job_deadline_sec
async def process_one(self) -> bool:
"""대기 잡 1건을 claim·처리. 처리했으면 True, 없으면 False."""
job = await self.queue.claim(self.worker_id, self.lease_sec)
if not job:
return False
await self._process(job)
return True
async def drain(self) -> int:
"""큐가 빌 때까지 처리(테스트/일회성 배치용). 처리한 잡 수 반환."""
n = 0
while await self.process_one():
n += 1
return n
async def run(self, listener=None, stop: asyncio.Event | None = None, idle_timeout: float = 5.0):
"""상시 루프. stop 이 설정될 때까지 처리하고, 유휴 시 알림/타임아웃까지 대기."""
stop = stop or asyncio.Event()
while not stop.is_set():
try:
worked = await self.process_one()
except Exception as ex:
# claim 자체가 실패(DB 순단 등) — 루프를 죽이지 않는다. 잠깐 쉬고 다시 시도.
LOG.e_no_callstack(f"[{self.worker_id}] claim 실패(계속): {type(ex).__name__}: {ex}")
worked = False
if not worked:
if listener is not None:
await listener.wait(idle_timeout)
else:
await asyncio.sleep(idle_timeout)
async def _process(self, job: dict):
jid = job["job_id"]
hb = asyncio.create_task(self._heartbeat(jid))
try:
if self.job_deadline_sec > 0:
result = await asyncio.wait_for(self.handler(job), timeout=self.job_deadline_sec)
else:
result = await self.handler(job)
await self.queue.complete(jid, self.worker_id, result)
LOG.d(f"[{self.worker_id}] done {jid}")
except (asyncio.TimeoutError, TimeoutError):
# 데드라인 초과 — wait_for 가 핸들러 태스크를 취소한 뒤 여기로 온다.
# 행이 워커 슬롯을 영구 점유하는 것보다 낫다.
backoff = self.backoff_fn(job["attempts"])
reason = f"JobDeadlineExceeded: {self.job_deadline_sec:.0f}s"
st = await self.queue.fail(jid, self.worker_id, reason, backoff)
LOG.w(f"[{self.worker_id}] deadline {jid}{JobStatus(st).name if st else '?'} "
f"({self.job_deadline_sec:.0f}s 초과, 핸들러 취소)")
if st == JobStatus.DEAD.value:
await _alert_job_dead(job, reason)
except PermanentJobError as ex:
# ★ 재시도하지 않는다 — 다시 해도 같은 결과다(common/job_errors 주석).
# 실측(2026-09-15): 사업장이 지워진 뒤 남은 소개문 잡이 "사업장을 찾을 수 없다" 로
# 세 번 돌고 DEAD 로 갔다. 결과는 같고 큐 지연과 알림만 늘었다.
# ★ fail_permanent 는 재시도 없이 곧장 DEAD 다 — 위 일반 실패 경로처럼 상태를 다시
# 조회해 분기할 필요 없이 바로 알린다(job_crud.fail_permanent 주석 참고).
reason = f"{type(ex).__name__}: {ex}"
await self.queue.fail_permanent(jid, self.worker_id, reason)
LOG.w(f"[{self.worker_id}] fail {jid} → DEAD (재시도 안 함: {reason})")
await _alert_job_dead(job, reason)
except Exception as ex:
backoff = self.backoff_fn(job["attempts"])
reason = f"{type(ex).__name__}: {ex}"
st = await self.queue.fail(jid, self.worker_id, reason, backoff)
LOG.w(f"[{self.worker_id}] fail {jid}{JobStatus(st).name if st else '?'} ({type(ex).__name__}: {ex})")
if st == JobStatus.DEAD.value:
await _alert_job_dead(job, reason)
finally:
hb.cancel()
try:
await hb
except asyncio.CancelledError:
pass
async def _heartbeat(self, jid: str):
interval = max(1, self.lease_sec // 3)
while True:
await asyncio.sleep(interval)
try:
await self.queue.renew_lease(jid, self.worker_id, self.lease_sec)
except Exception as ex:
# 갱신 실패는 치명적이지 않다(다음 틱 재시도). 계속 실패하면 lease 만료 → reaper 회수.
LOG.w(f"[{self.worker_id}] lease 갱신 실패(계속): {type(ex).__name__}")
async def run_reaper(queue: JobQueue, stop: asyncio.Event, interval: float = 30.0):
"""만료 lease(워커 사망) 잡을 주기적으로 회수. 재시도 남으면 재큐, 소진되면 DEAD.
도커에서 워커 컨테이너가 재시작되면 진행 중이던 잡이 여기서 되살아난다."""
while not stop.is_set():
try:
reclaimed = await queue.reap()
if reclaimed:
LOG.w(f"[reaper] reclaimed {len(reclaimed)} stale job(s)")
for job in reclaimed:
if job["status"] == JobStatus.DEAD.value:
await _alert_job_dead(job, job.get("last_error") or "lease 만료 후 재시도 소진")
except Exception as ex:
LOG.e_no_callstack(f"[reaper] {type(ex).__name__}: {ex}")
try:
await asyncio.wait_for(stop.wait(), interval)
except asyncio.TimeoutError:
pass