File size: 4,311 Bytes
23d337e
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
9bd3ee0
 
 
 
 
 
23d337e
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
"""
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),
            }