| """ |
| Recognition service — runs recognition jobs against the reference |
| gallery. |
| """ |
|
|
| from __future__ import annotations |
|
|
| import time |
| from typing import Optional |
|
|
| from confidence.engine import ConfidenceEngine |
| from confidence.conflicts import ConflictDetector |
| from metrics.collector import MetricsCollector |
| from models.jobs import JobRequest |
| from models.providers import ProviderCapability |
| from normalization.merger import ReportMerger |
| from orchestrator.runner import Orchestrator |
| from pipeline import ( |
| InputValidator, |
| ImagePreprocessor, |
| ImageHasher, |
| FeatureExtractor, |
| ) |
| from storage.reference_store import ReferenceStore |
| from utils.logging import execution_context, new_execution_id |
|
|
|
|
| class RecognitionService: |
| """Handles recognition-only jobs.""" |
|
|
| def __init__( |
| self, |
| orchestrator: Orchestrator, |
| metrics: MetricsCollector, |
| validator: InputValidator, |
| preprocessor: ImagePreprocessor, |
| hasher: ImageHasher, |
| feature_extractor: FeatureExtractor, |
| reference_store: ReferenceStore, |
| 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._gallery = reference_store |
| self._merger = ReportMerger(confidence_engine, conflict_detector) |
|
|
| async def recognize(self, request: JobRequest) -> dict: |
| eid = new_execution_id() |
| with execution_context(execution_id=eid, provider_id="recognition_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, |
| ) |
|
|
| |
| pipeline_output.gallery = self._gallery.get_all() |
|
|
| results = await self._orchestrator.run( |
| pipeline_output=pipeline_output, |
| capabilities=[ProviderCapability.RECOGNITION], |
| 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="recognition", |
| ) |
| self._metrics.timings.record("job.recognition", elapsed) |
| self._metrics.counters.inc("jobs.recognition.completed") |
| return { |
| "success": True, |
| "report": report.model_dump(), |
| "elapsed_ms": round(elapsed, 3), |
| } |
|
|
| def list_known_persons(self) -> list[dict]: |
| return self._gallery.list_persons() |
|
|
| def add_known_person(self, name: str, embedding_bytes: bytes) -> dict: |
| import numpy as np |
| emb = np.frombuffer(embedding_bytes, dtype=np.float32) |
| self._gallery.add(name, emb) |
| return {"name": name, "added": True} |
|
|
| def remove_known_person(self, name: str) -> dict: |
| removed = self._gallery.remove(name) |
| return {"name": name, "removed": removed} |
|
|