""" 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, ) # Attach gallery to pipeline output so recognition providers can read it 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}