diff --git a/lps/crud/job_crud.py b/lps/crud/job_crud.py index 0c7c31c..c221d24 100644 --- a/lps/crud/job_crud.py +++ b/lps/crud/job_crud.py @@ -203,7 +203,9 @@ class JobQueue: async def ops(self) -> dict: """운영 스냅샷(모니터링·알림용): 상태별 카운트 + 큐 지연(가장 오래된 PENDING 나이) + - 최근 1시간 DEAD + lease 만료 stuck(reaper 가 회수 못한 좀비 신호).""" + 최근 1시간 DEAD + stuck(좀비 신호). stuck 은 두 축 — lease 만료(워커 사망인데 reaper + 미회수) OR 실행 10분 초과(핸들러 행 — heartbeat 가 lease 를 계속 갱신해 lease 축엔 안 + 잡히므로 run_started_at 로 따로 본다. 잡 데드라인 300s 가 정상 작동하면 여기 안 온다).""" sql = text(""" SELECT count(*) FILTER (WHERE status = 1) AS pending, @@ -211,7 +213,10 @@ class JobQueue: count(*) FILTER (WHERE status = 3) AS done, count(*) FILTER (WHERE status = 4) AS dead, count(*) FILTER (WHERE status = 4 AND updated_at > now() - interval '1 hour') AS dead_1h, - count(*) FILTER (WHERE status = 2 AND lease_until IS NOT NULL AND lease_until < now()) AS stuck_running, + count(*) FILTER (WHERE status = 2 AND ( + (lease_until IS NOT NULL AND lease_until < now()) + OR run_started_at < now() - interval '10 minutes' + )) AS stuck_running, COALESCE(EXTRACT(EPOCH FROM (now() - min(created_at) FILTER (WHERE status = 1)))::int, 0) AS oldest_pending_sec FROM job """) diff --git a/lps/tests/test_worker.py b/lps/tests/test_worker.py index c8396ae..e4b7ac8 100644 --- a/lps/tests/test_worker.py +++ b/lps/tests/test_worker.py @@ -59,6 +59,76 @@ async def test_reaper_reclaims_then_worker_reprocesses(q): assert (await q.counts())["DONE"] == 1 +async def test_job_deadline_cancels_hung_handler(q): + """핸들러 행 → 데드라인 초과 시 취소·fail 처리(재큐/DEAD)돼야 한다. 없으면 heartbeat 가 + lease 를 계속 갱신해 워커 슬롯이 영구 점유된다(2026-07-10 부하테스트 실측).""" + jid = await q.enqueue(JobType.SEARCH.value, {"product_name": "hang"}, max_attempts=1) + + async def hang(job): + await asyncio.sleep(3600) + + w = Worker("w1", q, hang, backoff_fn=lambda a: 0, job_deadline_sec=0.2) + assert await w.process_one() is True # 행이어도 데드라인에 끊겨 반환된다 + assert (await q.counts())["DEAD"] == 1 # max_attempts=1 → 즉시 DEAD + job = await q.get(jid) + assert "JobDeadlineExceeded" in job["last_error"] + + +async def test_job_deadline_retries_before_dead(q): + """데드라인 초과도 일반 실패처럼 백오프 재큐를 탄다(시도 소진 전까지).""" + await q.enqueue(JobType.SEARCH.value, {"product_name": "hang"}, max_attempts=2) + + async def hang(job): + await asyncio.sleep(3600) + + w = Worker("w1", q, hang, backoff_fn=lambda a: 0, job_deadline_sec=0.2) + assert await w.process_one() is True + assert (await q.counts())["PENDING"] == 1 # 1/2 → 재큐 + assert await w.process_one() is True + assert (await q.counts())["DEAD"] == 1 # 2/2 → DEAD + + +async def test_ops_counts_long_running_as_stuck(q, db_engine): + """lease 가 계속 갱신돼도(행 상태의 heartbeat) 실행 10분 초과면 stuck_running 에 잡혀야 한다.""" + await q.enqueue(JobType.SEARCH.value, {"product_name": "x"}) + await q.claim("w1", lease_sec=3600) # lease 는 멀쩡(만료 안 됨) + assert (await q.ops())["stuck_running"] == 0 + async with db_engine.begin() as conn: + await conn.execute(text("UPDATE job SET run_started_at = now() - interval '11 minutes' WHERE status = 2")) + assert (await q.ops())["stuck_running"] == 1 + + +async def test_browser_reaper_survives_hung_adapter(): + """한 어댑터의 close 행이 정리 루프 전체를 멈추면 안 된다 — 타임아웃 후 다음 어댑터로.""" + from worker_main import run_browser_reaper + + class HungAdapter: + source = "hung" + async def close_if_idle(self, idle_sec): + await asyncio.sleep(3600) + + class OkAdapter: + source = "ok" + closed = False + async def close_if_idle(self, idle_sec): + self.closed = True + + ok = OkAdapter() + stop = asyncio.Event() + task = asyncio.create_task(run_browser_reaper( + [HungAdapter(), ok], stop, idle_sec=0, interval=0.01, close_timeout=0.05)) + try: + await asyncio.wait_for(_until(lambda: ok.closed), timeout=3.0) # 행 어댑터를 지나 ok 까지 도달 + finally: + stop.set() + await task + + +async def _until(cond, poll: float = 0.02): + while not cond(): + await asyncio.sleep(poll) + + async def test_enqueue_notifies_listener(q): listener = JobListener() await listener.start() diff --git a/lps/worker/runner.py b/lps/worker/runner.py index 3ccd8db..77f6b30 100644 --- a/lps/worker/runner.py +++ b/lps/worker/runner.py @@ -2,6 +2,9 @@ 워커는 큐에서 잡을 원자적으로 claim → 핸들러 실행 → complete/fail 한다. - 처리 중 heartbeat 로 lease 를 갱신(긴 잡이 reaper 에 회수되지 않게). +- 핸들러엔 데드라인(job_deadline_sec)을 건다 — heartbeat 가 lease 를 계속 갱신하므로 + 핸들러가 행하면 reaper 로는 영원히 회수 불가(2026-07-10 부하테스트에서 크롤 15분 행 실측). + 초과 시 취소 후 fail 처리 → 백오프 재큐(소진 시 DEAD), 워커 슬롯은 즉시 다음 잡으로. - claim 이 비면 LISTEN 알림 또는 짧은 타임아웃까지 대기(폴링 최소화). - 핸들러는 주입식(async def(job)->dict) — 프로덕션은 검색 파이프라인, 테스트는 fake. """ @@ -14,12 +17,14 @@ 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): + def __init__(self, worker_id: str, queue: JobQueue, handler, lease_sec: int = 120, backoff_fn=compute_backoff, + job_deadline_sec: float = 300.0): self.worker_id = worker_id self.queue = queue self.handler = handler self.lease_sec = lease_sec self.backoff_fn = backoff_fn + self.job_deadline_sec = job_deadline_sec # 잡 1건 처리 시간 상한(0 이면 무제한 — 테스트용) async def process_one(self) -> bool: """대기 잡 1건을 claim·처리. 처리했으면 True, 없으면 False.""" @@ -51,9 +56,19 @@ class Worker: jid = job["job_id"] hb = asyncio.create_task(self._heartbeat(jid)) try: - result = await self.handler(job) + 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 TimeoutError: + # 데드라인 초과 — wait_for 가 핸들러 태스크를 취소한 뒤 여기로 온다. in-flight 크롤이 + # 취소되며 브라우저가 어중간한 상태로 남을 수 있지만, 어댑터가 다음 검색에서 재기동으로 + # 회복한다. 행이 워커 슬롯을 영구 점유하는 것보다 낫다. + backoff = self.backoff_fn(job["attempts"]) + st = await self.queue.fail(jid, self.worker_id, f"JobDeadlineExceeded: {self.job_deadline_sec:.0f}s", backoff) + LOG.w(f"[{self.worker_id}] deadline {jid} → {JobStatus(st).name if st else '?'} ({self.job_deadline_sec:.0f}s 초과, 핸들러 취소)") 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) diff --git a/lps/worker_main.py b/lps/worker_main.py index ba3d805..c0195fc 100644 --- a/lps/worker_main.py +++ b/lps/worker_main.py @@ -88,16 +88,18 @@ def _build_worker(i: int, concurrency: int, has_openai: bool, neg_cache, history return handler, list(adapters.values()) + list(fallback_adapters.values()) -async def _warmup_worker(worker_adapters, tries: int = 3): +async def _warmup_worker(worker_adapters, tries: int = 3, attempt_timeout: float = 60.0): """워커의 챌린지 소스(Turnstile/Akamai)를 미리 풀어 쿠키(cf_clearance 등)를 확보한다. 콜드 비용을 시작 시 몰아, 이후 실 작업은 웜(빠름). 백그라운드로 돌려 잡 처리를 막지 않는다. - 나쁜 IP 는 인터랙티브 Turnstile 로 에스컬레이션되므로, 실패 시 **다른 IP 로 회전 재시도**한다.""" + 나쁜 IP 는 인터랙티브 Turnstile 로 에스컬레이션되므로, 실패 시 **다른 IP 로 회전 재시도**한다. + 시도당 타임아웃 필수 — 웜업은 search 중 어댑터 락을 쥐므로, 여기서 행하면 그 워커의 + 모든 실 검색이 락 대기로 함께 멈춘다(2026-07-10 부하테스트에서 15분 행 실측).""" for ad in worker_adapters: if ad.source not in ("gmarket", "auction", "coupang"): continue for attempt in range(tries): try: - await ad.search("생수", limit=1) + await asyncio.wait_for(ad.search("생수", limit=1), timeout=attempt_timeout) LOG.i(f"[warmup:{ad.source}] 챌린지 통과·쿠키 확보 (시도 {attempt + 1})") break except Exception as ex: @@ -151,9 +153,12 @@ async def run_ops_monitor(queue, bot_log, stop, interval: float = 30.0): pass -async def run_browser_reaper(adapters, stop, idle_sec: float = 120.0, interval: float = 30.0): +async def run_browser_reaper(adapters, stop, idle_sec: float = 120.0, interval: float = 30.0, + close_timeout: float = 60.0): """유휴 브라우저 정리 루프 — 일정 시간 검색 없는 어댑터의 Chrome 을 닫아 메모리를 회수한다. - 쿠키는 user_data_dir 에 남아, 다음 검색 때 재기동해도 (같은 IP면) 웜 유지.""" + 쿠키는 user_data_dir 에 남아, 다음 검색 때 재기동해도 (같은 IP면) 웜 유지. + 순차 순회라 close 1건에도 타임아웃을 건다 — 한 어댑터의 close 행이 루프 전체를 멈춰 + 다른 워커의 브라우저까지 못 닫게 되는 것을 실측(2026-07-10 부하테스트)했다.""" while not stop.is_set(): try: await asyncio.wait_for(stop.wait(), timeout=interval) @@ -164,7 +169,9 @@ async def run_browser_reaper(adapters, stop, idle_sec: float = 120.0, interval: if close_if_idle is None: # 네이버(httpx) 등 브라우저 없는 어댑터는 정리 대상 아님 continue try: - await close_if_idle(idle_sec) + await asyncio.wait_for(close_if_idle(idle_sec), timeout=close_timeout) + except asyncio.TimeoutError: + LOG.w(f"[browser-reaper] {getattr(ad, 'source', '?')} 정리 {close_timeout:.0f}s 초과 — 취소·스킵(다음 틱 재시도)") except Exception as ex: LOG.e_no_callstack(f"[browser-reaper] 정리 실패(무시): {ex}") @@ -210,6 +217,8 @@ async def main(concurrency: int = 1): for sig in (signal.SIGINT, signal.SIGTERM): loop.add_signal_handler(sig, _request_stop, sig.name) + # 잡 1건 데드라인 — 정상 검색은 폴백 포함 수분 내 끝난다(실측 15~22s). 크롤 행 실측(15분) 대비 상한. + job_deadline = float(os.environ.get("LPS_JOB_DEADLINE_SEC", "300")) for i in range(concurrency): handler, worker_adapters = _build_worker(i, concurrency, has_openai, neg_cache, history) all_adapters += worker_adapters @@ -217,7 +226,7 @@ async def main(concurrency: int = 1): listener = JobListener() await listener.start() listeners.append(listener) - worker = Worker(f"worker-{i}", queue, handler) + worker = Worker(f"worker-{i}", queue, handler, job_deadline_sec=job_deadline) tasks.append(asyncio.create_task(worker.run(listener, stop))) tasks.append(asyncio.create_task(run_reaper(queue, stop)))