face-intel / services /job_service.py
Marwan
Restructure + add reverse face search (PimEyes-style)
f5eeb1c
Raw
History Blame Contribute Delete
7.42 kB
"""
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