| """ |
| 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 |
|
|
|
|
| |
| _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), |
| } |
|
|