| """ |
| OSINT service — orchestrates reverse-image-search across multiple |
| providers and merges results into a unified evidence set. |
| |
| This is NOT a provider itself — it's a service that: |
| 1. Runs the existing reverse-search providers (serpapi, social_lookup) |
| via the orchestrator. |
| 2. Collects raw results from each provider. |
| 3. Merges + deduplicates via cores.osint.merger. |
| 4. Classifies each source (social/news/forum/marketplace/etc.) via |
| cores.osint.sources. |
| 5. Returns an OSINTResult with confidence-scored, classified matches. |
| |
| Pure stdlib for the merging logic; relies on existing providers for |
| network calls. |
| """ |
|
|
| from __future__ import annotations |
|
|
| import time |
| from typing import Optional |
|
|
| from cores.osint import merge_reverse_results, classify_source, SourceType |
| from metrics.collector import MetricsCollector |
| from models.jobs import JobRequest |
| from models.reports import OSINTResult, OSINTMatch |
| from models.providers import ProviderCapability |
| from orchestrator.runner import Orchestrator |
| from pipeline import InputValidator, ImagePreprocessor, ImageHasher, FeatureExtractor |
| from utils.logging import execution_context, new_execution_id |
|
|
|
|
| class OSINTService: |
| """Orchestrates reverse-image-search across providers + merges results.""" |
|
|
| def __init__( |
| self, |
| orchestrator: Orchestrator, |
| metrics: MetricsCollector, |
| validator: InputValidator, |
| preprocessor: ImagePreprocessor, |
| hasher: ImageHasher, |
| feature_extractor: FeatureExtractor, |
| ) -> None: |
| self._orchestrator = orchestrator |
| self._metrics = metrics |
| self._validator = validator |
| self._preprocessor = preprocessor |
| self._hasher = hasher |
| self._feature_extractor = feature_extractor |
|
|
| async def reverse_search(self, request: JobRequest) -> dict: |
| """Run reverse-image-search across all enabled providers + merge.""" |
| eid = new_execution_id() |
| with execution_context(execution_id=eid, provider_id="osint_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, |
| ) |
|
|
| |
| if request.image_url: |
| pipeline_output.scrape_url = request.image_url |
|
|
| |
| results = await self._orchestrator.run( |
| pipeline_output=pipeline_output, |
| capabilities=[ProviderCapability.REVERSE_SEARCH], |
| provider_whitelist=request.providers or None, |
| execution_id=eid, |
| ) |
|
|
| |
| provider_results: dict[str, list[dict]] = {} |
| providers_invoked: list[str] = [] |
| providers_succeeded: list[str] = [] |
| providers_failed: list[str] = [] |
|
|
| for name, result in results.items(): |
| providers_invoked.append(name) |
| if result.success: |
| providers_succeeded.append(name) |
| raw_matches = result.normalized.get("results", []) |
| if raw_matches: |
| provider_results[name] = raw_matches |
| else: |
| providers_failed.append(name) |
|
|
| |
| merged = merge_reverse_results(provider_results) |
|
|
| |
| breakdown: dict[str, int] = {} |
| for m in merged: |
| st = m.get("source_type", "unknown") |
| breakdown[st] = breakdown.get(st, 0) + 1 |
|
|
| elapsed = (time.perf_counter() - t0) * 1000.0 |
| osint_result = OSINTResult( |
| providers_invoked=providers_invoked, |
| providers_succeeded=providers_succeeded, |
| providers_failed=providers_failed, |
| total_matches=len(merged), |
| matches=[OSINTMatch(**m) for m in merged], |
| source_type_breakdown=breakdown, |
| elapsed_ms=round(elapsed, 3), |
| ) |
|
|
| self._metrics.timings.record("job.osint", elapsed) |
| self._metrics.counters.inc("jobs.osint.completed") |
|
|
| return { |
| "success": True, |
| "osint": osint_result.model_dump(), |
| "elapsed_ms": round(elapsed, 3), |
| } |
|
|