"""워커 루프 + reaper.""" 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: """dedupe_key 는 "같은 대상이 반복해서 죽는가" 를 잡는다.""" 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 이면 무제한 — 테스트용). self.job_deadline_sec = job_deadline_sec async def process_one(self) -> bool: """대기 잡 1건을 claim·처리.""" 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 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 주석). 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: # 갱신 실패는 치명적이지 않다(다음 틱 재시도). LOG.w(f"[{self.worker_id}] lease 갱신 실패(계속): {type(ex).__name__}") async def run_reaper(queue: JobQueue, stop: asyncio.Event, interval: float = 30.0): """만료 lease(워커 사망) 잡을 주기적으로 회수.""" 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