from fastapi import Depends from common.enums import ErrorType, JobType from common.logger import LOG from crud.job_crud import JobQueue from router.v1.job.protocol import JobData, Res_Job, Res_JobOps class JobService: """작업 큐 조회·관리. 적재는 각 도메인 서비스가 JobQueue 로 직접 한다. 수집·비전분석·빌드는 몇 분 걸린다 — API 는 잡만 넣고 즉시 응답하고, 클라이언트는 GET /v1/job/{id} 를 폴링한다.""" def __init__(self, queue: JobQueue = Depends(JobQueue)): self.queue = queue async def get_job(self, job_id: str) -> Res_Job: res = Res_Job() row = await self.queue.get(job_id) if row is None: res.result.SetResult(ErrorType.JOB_NOT_FOUND) return res res.job = JobData(**{k: v for k, v in row.items() if k in JobData.model_fields}) return res async def ops(self) -> Res_JobOps: snap = await self.queue.ops() return Res_JobOps(**{k: v for k, v in snap.items() if k in Res_JobOps.model_fields}) async def requeue(self, job_id: str) -> Res_Job: """DEAD 잡 재큐(운영자 액션).""" try: requeued = await self.queue.requeue(job_id) except Exception as ex: # 같은 dedupe_key 의 활성 잡이 이미 있으면 부분 유니크 위반. LOG.w(f"[job] requeue 실패 {job_id}: {type(ex).__name__}") res = Res_Job() res.result.SetResult(ErrorType.JOB_ALREADY_QUEUED) return res if requeued is None: res = Res_Job() res.result.SetResult(ErrorType.JOB_NOT_DEAD) return res return await self.get_job(job_id) async def enqueue_job( queue: JobQueue, job_type: JobType, payload: dict, dedupe_key: str | None = None, priority: int = 100, ) -> tuple[str | None, bool]: """도메인 서비스가 잡을 넣을 때 쓰는 공용 진입점. 반환: (job_id, 새로 만들었는가). 활성 중복이면 기존 잡의 id 와 False 를 돌려준다 — 같은 사업장 수집을 두 번 눌러도 잡이 두 번 돌지 않는다.""" job_id = await queue.enqueue(job_type.value, payload, priority=priority, dedupe_key=dedupe_key) if job_id is not None: return job_id, True if dedupe_key: existing = await queue.find_active(dedupe_key) if existing: return existing["job_id"], False return None, False