# o2o-web4ai 워커 프로세스 진입점 (API 와 분리 실행 — 코드베이스 공유, 독립 스케일). # python worker_main.py # WORKER_CONCURRENCY=3 python worker_main.py # # 수집·비전분석·빌드는 몇 분씩 걸려 동기 요청으로 처리할 수 없다. API 는 잡만 적재하고 즉시 응답하며, # 실제 처리는 이 프로세스가 한다. 큐는 PostgreSQL(job.jobs) — 원자적 claim + lease 소유권이라 # 워커를 몇 개 띄우든(docker compose --scale) 같은 잡이 두 번 돌지 않는다. import asyncio import os import signal from common.database.db_session_manager import DB_SESSION_MNG from common.logger import LOG from config.server_configs import web_server_config from crud.job_crud import JobQueue from worker.handlers import build_handler, log_registry from worker.notify import JobListener from worker.runner import Worker, run_reaper LOG.SetPrefix(f"{web_server_config.server_name}-worker") # 잡 1건 처리 시간 상한. Perplexity(10~30s) + 크롤링 + Vision(사진 20~50장 배치)을 감안한 값. JOB_DEADLINE_SEC = float(os.environ.get("JOB_DEADLINE_SEC", "900")) # lease 임대 시간. heartbeat 가 lease_sec/3 마다 갱신하므로 짧아도 되지만, # 워커가 죽었을 때 이만큼 지나야 reaper 가 회수한다. LEASE_SEC = int(os.environ.get("JOB_LEASE_SEC", "120")) async def main(concurrency: int = 1): queue = JobQueue() handler = build_handler() log_registry() stop = asyncio.Event() listeners: list[JobListener] = [] tasks: list[asyncio.Task] = [] # ── graceful shutdown: SIGINT(Ctrl+C)/SIGTERM(docker stop) → stop 이벤트 ── # 하던 잡은 마무리하고 새 잡은 받지 않는다. 강제 종료돼도 lease 만료 후 reaper 가 재큐한다. def _request_stop(sig_name: str): if not stop.is_set(): LOG.i(f"{sig_name} 수신 — graceful 종료: 새 잡 중단, 하던 잡 마무리 (한 번 더 = 강제 종료)") stop.set() else: LOG.w(f"{sig_name} 재수신 — 강제 종료(실행 중 잡은 lease 만료 후 reaper 가 재큐)") for t in tasks: t.cancel() loop = asyncio.get_running_loop() for sig in (signal.SIGINT, signal.SIGTERM): loop.add_signal_handler(sig, _request_stop, sig.name) for i in range(concurrency): listener = JobListener() await listener.start() listeners.append(listener) worker = Worker(f"worker-{i}", queue, handler, lease_sec=LEASE_SEC, job_deadline_sec=JOB_DEADLINE_SEC) tasks.append(asyncio.create_task(worker.run(listener, stop))) tasks.append(asyncio.create_task(run_reaper(queue, stop))) LOG.i(f"워커 {concurrency}개 + reaper 기동 (lease {LEASE_SEC}s · 잡 데드라인 {JOB_DEADLINE_SEC:.0f}s)") gathered = asyncio.gather(*tasks) try: await gathered except asyncio.CancelledError: LOG.w("강제 종료 — 남은 리소스 정리 후 종료") finally: stop.set() for t in tasks: t.cancel() await asyncio.gather(gathered, return_exceptions=True) for listener in listeners: try: await listener.close() except Exception as ex: LOG.e_no_callstack(f"[shutdown] 리스너 정리 실패(무시): {ex}") await DB_SESSION_MNG.dispose_all() LOG.i("워커 종료 완료") if __name__ == "__main__": asyncio.run(main(int(os.environ.get("WORKER_CONCURRENCY", "1"))))