여러 줄 주석이 설명보다 경위(예전·실측·지적)를 적고 있어 읽는 사람이 결론을 찾기 어려웠다. - 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>
77 lines
2.8 KiB
Python
77 lines
2.8 KiB
Python
# o2o-web4ai 워커 프로세스 진입점 (API 와 분리 실행 — 코드베이스 공유, 독립 스케일).
|
|
|
|
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건 처리 시간 상한.
|
|
JOB_DEADLINE_SEC = float(os.environ.get("JOB_DEADLINE_SEC", "900"))
|
|
# lease 임대 시간.
|
|
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 이벤트 ── 하던 잡은 마무리하고 새 잡은 받지 않는다.
|
|
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"))))
|