| """ |
| 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: |
| |
| 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: |
| |
| 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() |
| self._db.save_job(job_dict) |
| return True |
|
|