""" Analysis service — runs image-analysis, metadata, and forensics jobs. These capabilities share the same pipeline shape (no gallery/scrape context needed), so one service handles all three. """ from __future__ import annotations import time from confidence.engine import ConfidenceEngine from confidence.conflicts import ConflictDetector from metrics.collector import MetricsCollector from models.jobs import JobKind, JobRequest from models.providers import ProviderCapability from normalization.merger import ReportMerger from orchestrator.runner import Orchestrator from pipeline import ( InputValidator, ImagePreprocessor, ImageHasher, FeatureExtractor, ) from utils.logging import execution_context, new_execution_id # Mapping of JobKind -> capability to invoke _KIND_TO_CAPABILITY = { JobKind.IMAGE_ANALYSIS: ProviderCapability.IMAGE_ANALYSIS, JobKind.METADATA: ProviderCapability.METADATA, JobKind.FORENSICS: ProviderCapability.FORENSICS, JobKind.OCR: ProviderCapability.OCR, JobKind.OBJECT_DETECTION: ProviderCapability.OBJECT_DETECTION, JobKind.SCENE_RECOGNITION: ProviderCapability.SCENE_RECOGNITION, JobKind.NSFW_DETECTION: ProviderCapability.NSFW_DETECTION, JobKind.AI_IMAGE_DETECTION: ProviderCapability.AI_IMAGE_DETECTION, JobKind.EMBEDDING: ProviderCapability.EMBEDDING, } class AnalysisService: """Handles image-analysis / metadata / forensics jobs.""" def __init__( self, orchestrator: Orchestrator, metrics: MetricsCollector, validator: InputValidator, preprocessor: ImagePreprocessor, hasher: ImageHasher, feature_extractor: FeatureExtractor, confidence_engine: ConfidenceEngine, conflict_detector: ConflictDetector, ) -> None: self._orchestrator = orchestrator self._metrics = metrics self._validator = validator self._preprocessor = preprocessor self._hasher = hasher self._feature_extractor = feature_extractor self._merger = ReportMerger(confidence_engine, conflict_detector) async def analyze(self, request: JobRequest) -> dict: """Run an analysis job. `request.kind` determines the capability.""" if request.kind not in _KIND_TO_CAPABILITY: return {"success": False, "error": f"Unsupported job kind: {request.kind}", "error_type": "ValidationError"} capability = _KIND_TO_CAPABILITY[request.kind] eid = new_execution_id() with execution_context(execution_id=eid, provider_id=f"{request.kind.value}_service"): t0 = time.perf_counter() vr = self._validator.validate( image_url=request.image_url, image_base64=request.image_base64, ) if not vr.valid: return {"success": False, "error": vr.error, "error_type": "ValidationError"} if vr.source == "url": pre = self._preprocessor.from_url(request.image_url) else: pre = self._preprocessor.from_bytes(vr.image_bytes, vr.source) img_hash = self._hasher.hash(pre.image) pipeline_output = self._feature_extractor.extract( pre.image, img_hash, pre.width, pre.height, pre.source, original_bytes=pre.original_bytes, original_format=pre.original_format, ) results = await self._orchestrator.run( pipeline_output=pipeline_output, capabilities=[capability], provider_whitelist=request.providers or None, execution_id=eid, ) elapsed = (time.perf_counter() - t0) * 1000.0 report = self._merger.merge( results=results, image_hash=img_hash, job_id=eid, total_elapsed_ms=elapsed, kind=request.kind.value, ) self._metrics.timings.record(f"job.{request.kind.value}", elapsed) self._metrics.counters.inc(f"jobs.{request.kind.value}.completed") return { "success": True, "report": report.model_dump(), "elapsed_ms": round(elapsed, 3), }