o2o-site-AEO/solution/backend/worker/handlers.py

109 lines
4.7 KiB
Python

"""잡 핸들러 레지스트리 — JobType 별로 '무엇을 하는가'.
핸들러 규약: `async def handler(job: dict) -> dict`
job : {"job_id", "job_type", "payload", "attempts", "max_attempts"}
반환값 : jobs.result 에 JSONB 로 저장된다(관측·디버깅용)
예외 : 워커가 잡아 백오프 재큐(소진 시 DEAD). 재시도해도 소용없는 실패는
예외 메시지에 이유를 남긴다 — last_error 로 남아 운영자가 본다
핸들러는 **재시도 안전(멱등)** 해야 한다. lease 만료·워커 재시작으로 같은 잡이 다시 돌 수 있다.
수집은 이미 확보한 fact 를 다시 덮어쓰지 않고(특히 CORRECTED), 사진은 origin_url 로 중복을 거른다.
현재 등록된 핸들러:
JobType.COLLECT ✓ services/collect_service.run_collect — Phase 1 은 MockAdapter 만 등록돼 있다
JobType.VISION ✓ services/vision_service.run_vision — Gemini Vision 사진 분류 + alt
JobType.COPY ✓ services/copy_service.run_copy — 소개문·FAQ (확보된 fact 만 근거)
JobType.BUILD ✓ services/build_service.run_build — 정적 빌드 + 발행 검수 게이트
JobType.LOCAL_SYNC ✓ services/story_service.run_local_sync — 지역 이야기 생성(지역당 1회)
JobType.SONG ✓ services/song_service.run_song — 이 숙소의 노래 한 곡(가사 Gemini → 작곡 Suno)
JobType.ROLLBACK ✓ services/rollback_service.run_rollback — 예전 버전으로 공개 주소를 되돌림
JobType.AI_CHECK → reports 모듈이 붙을 때
"""
from common.enums import JobType
from common.logger import LOG
class UnknownJobType(RuntimeError):
"""등록되지 않은 JobType — 재시도해도 소용없다(코드 배포 누락 신호)."""
# JobType -> async def handler(job) -> dict
HANDLERS: dict[int, object] = {}
def register(job_type: JobType):
"""핸들러 등록 데코레이터. 모듈이 붙을 때 자기 핸들러를 여기에 건다.
@register(JobType.COLLECT)
async def handle_collect(job: dict) -> dict:
...
"""
def _deco(fn):
if job_type.value in HANDLERS:
raise RuntimeError(f"JobType.{job_type.name} 핸들러가 이미 등록돼 있다")
HANDLERS[job_type.value] = fn
return fn
return _deco
def build_handler():
"""등록된 핸들러로 디스패처를 만든다. 워커에 주입한다."""
async def dispatch(job: dict) -> dict:
fn = HANDLERS.get(job["job_type"])
if fn is None:
name = JobType(job["job_type"]).name if job["job_type"] in {t.value for t in JobType} else job["job_type"]
raise UnknownJobType(f"JobType.{name} 핸들러 미등록 — 이 워커 이미지에 해당 모듈이 없다")
return await fn(job)
return dispatch
def registered_types() -> list[str]:
"""기동 로그용 — 이 워커가 처리할 수 있는 잡 종류."""
return [JobType(v).name for v in sorted(HANDLERS)]
def log_registry():
names = registered_types()
if names:
LOG.i(f"[worker] 처리 가능 잡: {', '.join(names)}")
else:
LOG.w("[worker] 등록된 핸들러가 없습니다 — 적재되는 잡은 모두 UnknownJobType 으로 DEAD 됩니다")
# ---- 등록 ------------------------------------------------------------------
# import 부작용으로 등록한다(모듈을 읽는 것만으로 워커가 처리 능력을 갖는다).
# 순환 import 를 피하려고 파일 맨 아래에서 붙인다.
def _register_builtin():
from services.collect_service import run_collect
from services.build_service import run_build
from services.copy_service import run_copy
from services.story_service import run_local_sync
from services.song_service import run_song
from services.social_service import run_draft, run_post
HANDLERS[JobType.SOCIAL_DRAFT.value] = run_draft
HANDLERS[JobType.SOCIAL_POST.value] = run_post
from services.rollback_service import run_rollback
from services.vision_service import run_vision
if JobType.COLLECT.value not in HANDLERS:
HANDLERS[JobType.COLLECT.value] = run_collect
if JobType.VISION.value not in HANDLERS:
HANDLERS[JobType.VISION.value] = run_vision
if JobType.COPY.value not in HANDLERS:
HANDLERS[JobType.COPY.value] = run_copy
if JobType.BUILD.value not in HANDLERS:
HANDLERS[JobType.BUILD.value] = run_build
if JobType.LOCAL_SYNC.value not in HANDLERS:
HANDLERS[JobType.LOCAL_SYNC.value] = run_local_sync
if JobType.SONG.value not in HANDLERS:
HANDLERS[JobType.SONG.value] = run_song
if JobType.ROLLBACK.value not in HANDLERS:
HANDLERS[JobType.ROLLBACK.value] = run_rollback
_register_builtin()