o2o-site-AEO/backend/worker_main.py
Mina Choi 6784e59ca5 최초 커밋 — 기존 코드 전체 + 문서 체계 신설
git 저장소가 없어 히스토리·협업 기반이 아예 없던 상태를 연다.
함께 문서를 재편했다. 그동안 문서가 있어도 "이 제품이 뭘 푸는가"와
"어떻게 도는가"를 담은 문서가 없어서, 목표 문장이 backend/frontend
README 두 곳에 복붙돼 있었다 — 상위 문서가 없어 아래로 샌 것이다.

신설
  README.md               레포 진입점 + 문서 지도 + 문서 규칙 4가지
  AGENTS.md               에이전트·신규 합류자용 함정 목록과 규약
                          (CLAUDE.md 는 여기로 걸린 심볼릭 링크)
  docs/PRODUCT.md         제품 정의 — 문제·사용자·원칙·**non-goals**·성공 기준
  docs/ARCHITECTURE.md    payload 경계·발행 파이프라인·서빙 결정·앱 분리 설계

이동
  backend/docs/DECISIONS.md → docs/DECISIONS.md
    백엔드만의 결정이 아니다. 게다가 코드 주석 ~25곳이 이미
    `docs/DECISIONS.md` 로 적고 있어 레포 루트 기준으로는 그게 맞다.

갱신
  docs/DEPLOY.md          서빙 결정 반영 — nginx 정적 서빙이 지금 경로(3절),
                          Azure 는 나중에 켤 때(4절)로 분리
  docs/ARCHITECTURE.md    사이트 = 한 장(2026-08-31) 구조 반영
  docs/COLLECTION_SEO_AEO_FLOW.md
                          robots.txt·sitemap.xml 은 오리진 루트에만 굽는다는 점 명시
  frontend/site/scripts/prerender.ts
                          헤더 주석의 렌더 보고서 경로가 실제(422줄)와 달라 수정

.gitignore
  ★ CLAUDE.md 를 더 이상 무시하지 않는다. 에이전트 지침은 팀과 모든
    에이전트가 공유하는 규약이라 커밋해야 한다 — 무시하면 클론한 사람이
    "배포 후 republish_all.py 필수" 같은 함정을 전달받지 못한다.
    개인용 오버라이드는 ~/.claude/CLAUDE.md 에 둔다.
2026-08-31 13:57:59 +09:00

85 lines
3.5 KiB
Python

# 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"))))