"""Collector baseline and coverage helpers for BGP observations.""" from __future__ import annotations from collections import defaultdict from datetime import UTC, datetime, timedelta from typing import Any from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from app.core.time import to_iso8601_utc from app.models.bgp_observation import BGPObservation from app.services.collectors.bgp_common import RIPE_RIS_COLLECTOR_COORDS async def build_bgp_collector_coverage( db: AsyncSession, *, source_filter: tuple[str, ...] | None = None, ) -> list[dict[str, Any]]: now = datetime.now(UTC) recent_15m_threshold = now - timedelta(minutes=15) recent_24h_threshold = now - timedelta(hours=24) recent_7d_threshold = now - timedelta(days=7) stmt = select(BGPObservation).order_by(BGPObservation.observed_at.desc(), BGPObservation.id.desc()) if source_filter: stmt = stmt.where(BGPObservation.source.in_(source_filter)) result = await db.execute(stmt) records = list(result.scalars().all()) by_collector: dict[str, dict[str, Any]] = {} for record in records: collector = str(record.collector or "").strip() if not collector: continue coverage = by_collector.get(collector) if coverage is None: location = record.collector_geo or RIPE_RIS_COLLECTOR_COORDS.get(collector, {}) coverage = { "collector": collector, "city": location.get("city"), "country": location.get("country"), "latitude": location.get("latitude"), "longitude": location.get("longitude"), "observation_count": 0, "prefixes": set(), "origin_asns": set(), "peer_asns": set(), "event_types": defaultdict(int), "countries": set(), "cities": set(), "recent_15m_observation_count": 0, "recent_24h_observation_count": 0, "recent_7d_observation_count": 0, "recent_15m_prefixes": set(), "recent_24h_prefixes": set(), "recent_7d_prefixes": set(), "latest_observed_at": None, "latest_event_type": None, } by_collector[collector] = coverage coverage["observation_count"] += 1 if record.prefix: coverage["prefixes"].add(record.prefix) if record.origin_asn is not None: coverage["origin_asns"].add(record.origin_asn) if record.peer_asn is not None: coverage["peer_asns"].add(record.peer_asn) if record.event_type: coverage["event_types"][record.event_type] += 1 observed_at = record.observed_at if observed_at is not None: aware_observed_at = ( observed_at.astimezone(UTC) if observed_at.tzinfo else observed_at.replace(tzinfo=UTC) ) if aware_observed_at >= recent_15m_threshold: coverage["recent_15m_observation_count"] += 1 if record.prefix: coverage["recent_15m_prefixes"].add(record.prefix) if aware_observed_at >= recent_24h_threshold: coverage["recent_24h_observation_count"] += 1 if record.prefix: coverage["recent_24h_prefixes"].add(record.prefix) if aware_observed_at >= recent_7d_threshold: coverage["recent_7d_observation_count"] += 1 if record.prefix: coverage["recent_7d_prefixes"].add(record.prefix) geo = record.collector_geo or {} if geo.get("country"): coverage["countries"].add(geo["country"]) if geo.get("city"): coverage["cities"].add(geo["city"]) current_latest = coverage["latest_observed_at"] if current_latest is None or ( record.observed_at is not None and record.observed_at > current_latest ): coverage["latest_observed_at"] = record.observed_at coverage["latest_event_type"] = record.event_type for collector, location in RIPE_RIS_COLLECTOR_COORDS.items(): if collector in by_collector: continue by_collector[collector] = { "collector": collector, "city": location.get("city"), "country": location.get("country"), "latitude": location.get("latitude"), "longitude": location.get("longitude"), "observation_count": 0, "prefixes": set(), "origin_asns": set(), "peer_asns": set(), "event_types": defaultdict(int), "countries": {location.get("country")} if location.get("country") else set(), "cities": {location.get("city")} if location.get("city") else set(), "recent_15m_observation_count": 0, "recent_24h_observation_count": 0, "recent_7d_observation_count": 0, "recent_15m_prefixes": set(), "recent_24h_prefixes": set(), "recent_7d_prefixes": set(), "latest_observed_at": None, "latest_event_type": None, } results: list[dict[str, Any]] = [] for collector in sorted(by_collector.keys()): item = by_collector[collector] top_event_types = sorted( item["event_types"].items(), key=lambda pair: (-pair[1], pair[0]), ) results.append( { "collector": item["collector"], "city": item["city"], "country": item["country"], "latitude": item["latitude"], "longitude": item["longitude"], "observation_count": item["observation_count"], "prefix_count": len(item["prefixes"]), "origin_asn_count": len(item["origin_asns"]), "peer_asn_count": len(item["peer_asns"]), "recent_15m_observation_count": item["recent_15m_observation_count"], "recent_24h_observation_count": item["recent_24h_observation_count"], "recent_7d_observation_count": item["recent_7d_observation_count"], "recent_15m_prefix_count": len(item["recent_15m_prefixes"]), "recent_24h_prefix_count": len(item["recent_24h_prefixes"]), "recent_7d_prefix_count": len(item["recent_7d_prefixes"]), "top_event_types": [ {"event_type": event_type, "count": count} for event_type, count in top_event_types[:3] ], "latest_observed_at": to_iso8601_utc(item["latest_observed_at"]), "latest_event_type": item["latest_event_type"], "baseline_scope": { "countries": sorted(country for country in item["countries"] if country), "cities": sorted(city for city in item["cities"] if city), }, } ) return results