"""워커 루프 + reaper. 워커는 큐에서 잡을 원자적으로 claim → 핸들러 실행 → complete/fail 한다. - 처리 중 heartbeat 로 lease 를 갱신(긴 잡이 reaper 에 회수되지 않게). - claim 이 비면 LISTEN 알림 또는 짧은 타임아웃까지 대기(폴링 최소화). - 핸들러는 주입식(async def(job)->dict) — 프로덕션은 검색 파이프라인, 테스트는 fake. """ import asyncio from common.enums import JobStatus from common.logger import LOG from crud.job_crud import JobQueue, compute_backoff class Worker: def __init__(self, worker_id: str, queue: JobQueue, handler, lease_sec: int = 120, backoff_fn=compute_backoff): self.worker_id = worker_id self.queue = queue self.handler = handler self.lease_sec = lease_sec self.backoff_fn = backoff_fn 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(): worked = await self.process_one() 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: result = await self.handler(job) await self.queue.complete(jid, self.worker_id, result) LOG.d(f"[{self.worker_id}] done {jid}") except Exception as ex: backoff = self.backoff_fn(job["attempts"]) st = await self.queue.fail(jid, self.worker_id, f"{type(ex).__name__}: {ex}", backoff) LOG.w(f"[{self.worker_id}] fail {jid} → {JobStatus(st).name if st else '?'} ({type(ex).__name__}: {ex})") 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) await self.queue.renew_lease(jid, self.worker_id, self.lease_sec) async def run_reaper(queue: JobQueue, stop: asyncio.Event, interval: float = 30.0): """만료 lease(워커 사망) 잡을 주기적으로 회수. 재시도 남으면 재큐, 소진되면 DEAD.""" while not stop.is_set(): reclaimed = await queue.reap() if reclaimed: LOG.w(f"[reaper] reclaimed {len(reclaimed)} stale job(s)") try: await asyncio.wait_for(stop.wait(), interval) except asyncio.TimeoutError: pass