o2o-site-AEO/solution/backend/worker/runner.py
Mina Choi 11d30bb3d1 [chore] solution,admin,ontology: 코드 주석을 한 줄로 — 히스토리 주석 삭제
여러 줄 주석이 설명보다 경위(예전·실측·지적)를 적고 있어 읽는 사람이 결론을 찾기 어려웠다.

- ts·tsx·js·mjs·css·py 478개: 여러 줄 주석은 첫 문장 한 줄로, 과거형·날짜 문장은 삭제
- 주석 위치는 TypeScript 파서·파이썬 tokenize/ast 로 찾는다 — 문자열 안의 # · /* 는 건드리지 않는다
- eslint·ts·noqa·type: ignore 같은 지시 주석은 그대로 둔다

파이썬 275개 정리 전후 AST 동일, TS 298개 주석 뺀 토큰 동일(빈 JSX 주석 10곳만 차이).
site·frontend·admin tsc, site vitest 105 passed

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
2026-09-28 16:05:19 +09:00

150 lines
6.0 KiB
Python

"""워커 루프 + 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