""" Job service — full-pipeline orchestration + job persistence. Wraps DetectionService + RecognitionService + SearchService + AnalysisService to run any job kind, and persists the result to the Database. """ from __future__ import annotations import asyncio import time from typing import Optional from loguru import logger from metrics.collector import MetricsCollector from models.jobs import Job, JobKind, JobRequest, JobResult, JobStatus from services.analysis_service import AnalysisService from services.detection_service import DetectionService from services.recognition_service import RecognitionService from services.search_service import SearchService from storage.database import Database from utils.logging import execution_context, new_execution_id class JobService: """Full-pipeline job orchestration + persistence.""" def __init__( self, detection: DetectionService, recognition: RecognitionService, search: SearchService, analysis: AnalysisService, database: Database, metrics: MetricsCollector, job_timeout_seconds: float = 300.0, ) -> None: self._detection = detection self._recognition = recognition self._search = search self._analysis = analysis self._db = database self._metrics = metrics self._timeout = job_timeout_seconds async def create_and_run(self, request: JobRequest) -> dict: eid = new_execution_id() job = Job(id=eid, kind=request.kind, request=request) job.mark_running() self._db.save_job(job.model_dump(mode="json")) with execution_context(execution_id=eid, provider_id="job_service"): t0 = time.perf_counter() try: # Run with overall job timeout result = await asyncio.wait_for( self._dispatch(request), timeout=self._timeout, ) elapsed = (time.perf_counter() - t0) * 1000.0 job.mark_completed() self._db.save_job(job.model_dump(mode="json")) self._db.save_result( job_id=eid, status=job.status.value, report=result.get("report") if isinstance(result, dict) else result, error=None, elapsed_ms=elapsed, ) self._metrics.counters.inc(f"jobs.{request.kind.value}.completed") return {"job_id": eid, "status": job.status.value, "result": result, "elapsed_ms": round(elapsed, 3)} except asyncio.TimeoutError: elapsed = (time.perf_counter() - t0) * 1000.0 job.mark_timeout() self._db.save_job(job.model_dump(mode="json")) self._db.save_result( job_id=eid, status=job.status.value, report=None, error=f"Job timeout after {self._timeout}s", elapsed_ms=elapsed, ) self._metrics.counters.inc(f"jobs.{request.kind.value}.timeout") logger.error(f"Job {eid} timed out after {self._timeout}s") return {"job_id": eid, "status": job.status.value, "error": f"Job timeout after {self._timeout}s"} except Exception as e: logger.exception(f"Job {eid} failed") job.mark_failed(str(e)) self._db.save_job(job.model_dump(mode="json")) self._db.save_result( job_id=eid, status=job.status.value, report=None, error=str(e), elapsed_ms=0.0, ) self._metrics.counters.inc(f"jobs.{request.kind.value}.failed") return {"job_id": eid, "status": job.status.value, "error": str(e), "error_type": type(e).__name__} async def _dispatch(self, request: JobRequest) -> dict: """Route the request to the correct service.""" kind = request.kind if kind == JobKind.DETECTION: return await self._detection.detect(request) elif kind == JobKind.RECOGNITION: return await self._recognition.recognize(request) elif kind == JobKind.SEARCH: return await self._search.search(request) elif kind in (JobKind.IMAGE_ANALYSIS, JobKind.METADATA, JobKind.FORENSICS, JobKind.OCR, JobKind.OBJECT_DETECTION, JobKind.SCENE_RECOGNITION, JobKind.NSFW_DETECTION, JobKind.AI_IMAGE_DETECTION, JobKind.EMBEDDING): return await self._analysis.analyze(request) elif kind == JobKind.FULL_PIPELINE: # Run all capabilities concurrently async def _safe(coro, name): try: return await coro except Exception as e: return {"success": False, "error": str(e), "error_type": type(e).__name__} det, rec, ana, met, forr, srch = await asyncio.gather( _safe(self._detection.detect(request), "detection"), _safe(self._recognition.recognize(request), "recognition"), _safe(self._analysis.analyze(JobRequest( kind=JobKind.IMAGE_ANALYSIS, image_url=request.image_url, image_base64=request.image_base64, providers=request.providers)), "image_analysis"), _safe(self._analysis.analyze(JobRequest( kind=JobKind.METADATA, image_url=request.image_url, image_base64=request.image_base64, providers=request.providers)), "metadata"), _safe(self._analysis.analyze(JobRequest( kind=JobKind.FORENSICS, image_url=request.image_url, image_base64=request.image_base64, providers=request.providers)), "forensics"), _safe(self._search.search(request), "search"), return_exceptions=True, ) return { "success": True, "detection": det, "recognition": rec, "image_analysis": ana, "metadata": met, "forensics": forr, "search": srch, } else: return {"success": False, "error": f"Unknown job kind: {kind}", "error_type": "ValidationError"} def get_job(self, job_id: str) -> Optional[dict]: return self._db.get_job(job_id) def get_result(self, job_id: str) -> Optional[dict]: return self._db.get_result(job_id) def list_jobs(self, limit: int = 50, status: Optional[str] = None) -> list[dict]: return self._db.list_jobs(limit=limit, status=status) def cancel_job(self, job_id: str) -> bool: """Mark a job as cancelled (best-effort; doesn't interrupt running work).""" job_dict = self._db.get_job(job_id) if not job_dict: return False if job_dict["status"] in (JobStatus.COMPLETED.value, JobStatus.FAILED.value, JobStatus.CANCELLED.value): return False job_dict["status"] = JobStatus.CANCELLED.value job_dict["completed_at"] = time.perf_counter() # approx self._db.save_job(job_dict) return True