From 945786cee5f6e4f7d66602c83c16e4e139eafe7d Mon Sep 17 00:00:00 2001 From: linkong Date: Mon, 30 Mar 2026 16:13:36 +0800 Subject: [PATCH] fix: expand bgp pipeline and stabilize backend tests --- VERSION | 2 +- backend/app/api/v1/bgp.py | 67 ++- backend/app/api/v1/visualization.py | 79 ++- backend/app/db/session.py | 2 + backend/app/models/__init__.py | 4 + backend/app/models/bgp_incident.py | 64 +++ backend/app/models/bgp_observation.py | 62 +++ backend/app/services/bgp_detectors.py | 193 +++++++ backend/app/services/bgp_enrichment.py | 280 ++++++++++ backend/app/services/bgp_incidents.py | 151 ++++++ backend/app/services/collectors/bgp_common.py | 226 ++++---- backend/app/services/collectors/bgpstream.py | 14 +- backend/app/services/collectors/ris_live.py | 14 +- backend/tests/test_api.py | 56 +- backend/tests/test_bgp.py | 491 ++++++++++++++++++ backend/tests/test_collectors.py | 71 +-- backend/tests/test_models.py | 2 +- backend/tests/test_security.py | 20 +- docs/CHANGELOG.md | 33 ++ frontend/package.json | 2 +- frontend/public/earth/js/main.js | 180 +++++-- pyproject.toml | 2 +- 22 files changed, 1787 insertions(+), 228 deletions(-) create mode 100644 backend/app/models/bgp_incident.py create mode 100644 backend/app/models/bgp_observation.py create mode 100644 backend/app/services/bgp_detectors.py create mode 100644 backend/app/services/bgp_enrichment.py create mode 100644 backend/app/services/bgp_incidents.py diff --git a/VERSION b/VERSION index 4b9d1879..59f23022 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -0.21.8 +0.21.9 diff --git a/backend/app/api/v1/bgp.py b/backend/app/api/v1/bgp.py index 0a7818f5..11da3290 100644 --- a/backend/app/api/v1/bgp.py +++ b/backend/app/api/v1/bgp.py @@ -8,7 +8,8 @@ from sqlalchemy.ext.asyncio import AsyncSession from app.core.security import get_current_user from app.db.session import get_db from app.models.bgp_anomaly import BGPAnomaly -from app.models.collected_data import CollectedData +from app.models.bgp_incident import BGPIncident +from app.models.bgp_observation import BGPObservation from app.models.user import User router = APIRouter() @@ -48,12 +49,12 @@ async def list_bgp_events( db: AsyncSession = Depends(get_db), ): stmt = ( - select(CollectedData) - .where(CollectedData.source.in_(BGP_SOURCES)) - .order_by(CollectedData.reference_date.desc().nullslast(), CollectedData.id.desc()) + select(BGPObservation) + .where(BGPObservation.source.in_(BGP_SOURCES)) + .order_by(BGPObservation.observed_at.desc(), BGPObservation.id.desc()) ) if source: - stmt = stmt.where(CollectedData.source == source) + stmt = stmt.where(BGPObservation.source == source) result = await db.execute(stmt) records = result.scalars().all() @@ -62,18 +63,17 @@ async def list_bgp_events( filtered = [] for record in records: - metadata = record.extra_data or {} - if prefix and metadata.get("prefix") != prefix: + if prefix and record.prefix != prefix: continue - if origin_asn is not None and metadata.get("origin_asn") != origin_asn: + if origin_asn is not None and record.origin_asn != origin_asn: continue - if peer_asn is not None and metadata.get("peer_asn") != peer_asn: + if peer_asn is not None and record.peer_asn != peer_asn: continue - if collector and metadata.get("collector") != collector: + if collector and record.collector != collector: continue - if event_type and metadata.get("event_type") != event_type: + if event_type and record.event_type != event_type: continue - if (dt_from or dt_to) and not _matches_time(record.reference_date, dt_from, dt_to): + if (dt_from or dt_to) and not _matches_time(record.observed_at, dt_from, dt_to): continue filtered.append(record) @@ -92,7 +92,7 @@ async def get_bgp_event( current_user: User = Depends(get_current_user), db: AsyncSession = Depends(get_db), ): - record = await db.get(CollectedData, event_id) + record = await db.get(BGPObservation, event_id) if not record or record.source not in BGP_SOURCES: raise HTTPException(status_code=404, detail="BGP event not found") return record.to_dict() @@ -180,3 +180,44 @@ async def get_bgp_anomaly( if not record: raise HTTPException(status_code=404, detail="BGP anomaly not found") return record.to_dict() + + +@router.get("/incidents") +async def list_bgp_incidents( + severity: Optional[str] = Query(None), + incident_type: Optional[str] = Query(None), + status: Optional[str] = Query(None), + page: int = Query(1, ge=1), + page_size: int = Query(50, ge=1, le=200), + current_user: User = Depends(get_current_user), + db: AsyncSession = Depends(get_db), +): + stmt = select(BGPIncident).order_by(BGPIncident.created_at.desc(), BGPIncident.id.desc()) + if severity: + stmt = stmt.where(BGPIncident.severity == severity) + if incident_type: + stmt = stmt.where(BGPIncident.incident_type == incident_type) + if status: + stmt = stmt.where(BGPIncident.status == status) + + result = await db.execute(stmt) + records = result.scalars().all() + offset = (page - 1) * page_size + return { + "total": len(records), + "page": page, + "page_size": page_size, + "data": [record.to_dict() for record in records[offset : offset + page_size]], + } + + +@router.get("/incidents/{incident_id}") +async def get_bgp_incident( + incident_id: int, + current_user: User = Depends(get_current_user), + db: AsyncSession = Depends(get_db), +): + record = await db.get(BGPIncident, incident_id) + if not record: + raise HTTPException(status_code=404, detail="BGP incident not found") + return record.to_dict() diff --git a/backend/app/api/v1/visualization.py b/backend/app/api/v1/visualization.py index 2a1a748a..e4083f33 100644 --- a/backend/app/api/v1/visualization.py +++ b/backend/app/api/v1/visualization.py @@ -202,6 +202,44 @@ def convert_satellite_to_geojson(records: List[CollectedData]) -> Dict[str, Any] return {"type": "FeatureCollection", "features": features} +def dedupe_satellite_records(records: List[CollectedData]) -> List[CollectedData]: + """Keep only the newest record for each satellite identity.""" + latest_by_key: Dict[str, CollectedData] = {} + + for record in records: + metadata = record.extra_data or {} + norad_id = metadata.get("norad_cat_id") + dedupe_key = ( + str(norad_id) + if norad_id not in (None, "") + else str(record.source_id or record.entity_key or record.name or record.id) + ) + + existing = latest_by_key.get(dedupe_key) + if existing is None or (record.id or 0) > (existing.id or 0): + latest_by_key[dedupe_key] = record + + return sorted(latest_by_key.values(), key=lambda item: item.id or 0, reverse=True) + + +def dedupe_collected_records(records: List[CollectedData]) -> List[CollectedData]: + """Keep only the newest record for each collected entity.""" + latest_by_key: Dict[str, CollectedData] = {} + + for record in records: + dedupe_key = str( + record.source_id + or record.entity_key + or record.name + or record.id + ) + existing = latest_by_key.get(dedupe_key) + if existing is None or (record.id or 0) > (existing.id or 0): + latest_by_key[dedupe_key] = record + + return sorted(latest_by_key.values(), key=lambda item: item.id or 0, reverse=True) + + def convert_supercomputer_to_geojson(records: List[CollectedData]) -> Dict[str, Any]: """Convert TOP500 supercomputer records to GeoJSON""" features = [] @@ -410,7 +448,7 @@ async def get_cables_geojson(db: AsyncSession = Depends(get_db)): try: stmt = select(CollectedData).where(CollectedData.source == "arcgis_cables") result = await db.execute(stmt) - records = result.scalars().all() + records = dedupe_collected_records(list(result.scalars().all())) if not records: raise HTTPException( @@ -430,15 +468,15 @@ async def get_landing_points_geojson(db: AsyncSession = Depends(get_db)): try: landing_stmt = select(CollectedData).where(CollectedData.source == "arcgis_landing_points") landing_result = await db.execute(landing_stmt) - records = landing_result.scalars().all() + records = dedupe_collected_records(list(landing_result.scalars().all())) relation_stmt = select(CollectedData).where(CollectedData.source == "arcgis_cable_landing_relation") relation_result = await db.execute(relation_stmt) - relation_records = relation_result.scalars().all() + relation_records = dedupe_collected_records(list(relation_result.scalars().all())) cable_stmt = select(CollectedData).where(CollectedData.source == "arcgis_cables") cable_result = await db.execute(cable_stmt) - cable_records = cable_result.scalars().all() + cable_records = dedupe_collected_records(list(cable_result.scalars().all())) city_to_cable_ids_map = {} for rel in relation_records: @@ -476,15 +514,15 @@ async def get_landing_points_geojson(db: AsyncSession = Depends(get_db)): async def get_all_geojson(db: AsyncSession = Depends(get_db)): cables_stmt = select(CollectedData).where(CollectedData.source == "arcgis_cables") cables_result = await db.execute(cables_stmt) - cables_records = cables_result.scalars().all() + cables_records = dedupe_collected_records(list(cables_result.scalars().all())) points_stmt = select(CollectedData).where(CollectedData.source == "arcgis_landing_points") points_result = await db.execute(points_stmt) - points_records = points_result.scalars().all() + points_records = dedupe_collected_records(list(points_result.scalars().all())) relation_stmt = select(CollectedData).where(CollectedData.source == "arcgis_cable_landing_relation") relation_result = await db.execute(relation_stmt) - relation_records = relation_result.scalars().all() + relation_records = dedupe_collected_records(list(relation_result.scalars().all())) city_to_cable_ids_map = {} for rel in relation_records: @@ -542,10 +580,11 @@ async def get_satellites_geojson( .where(CollectedData.name != "Unknown") .order_by(CollectedData.id.desc()) ) - if limit is not None: - stmt = stmt.limit(limit) result = await db.execute(stmt) - records = result.scalars().all() + records = dedupe_satellite_records(list(result.scalars().all())) + + if limit is not None: + records = records[:limit] if not records: return {"type": "FeatureCollection", "features": [], "count": 0} @@ -567,10 +606,11 @@ async def get_supercomputers_geojson( select(CollectedData) .where(CollectedData.source == "top500") .where(CollectedData.name != "Unknown") - .limit(limit) + .order_by(CollectedData.id.desc()) ) result = await db.execute(stmt) - records = result.scalars().all() + records = dedupe_collected_records(list(result.scalars().all())) + records = records[:limit] if not records: return {"type": "FeatureCollection", "features": [], "count": 0} @@ -592,10 +632,11 @@ async def get_gpu_clusters_geojson( select(CollectedData) .where(CollectedData.source == "epoch_ai_gpu") .where(CollectedData.name != "Unknown") - .limit(limit) + .order_by(CollectedData.id.desc()) ) result = await db.execute(stmt) - records = result.scalars().all() + records = dedupe_collected_records(list(result.scalars().all())) + records = records[:limit] if not records: return {"type": "FeatureCollection", "features": [], "count": 0} @@ -645,11 +686,11 @@ async def get_all_visualization_data(db: AsyncSession = Depends(get_db)): """ cables_stmt = select(CollectedData).where(CollectedData.source == "arcgis_cables") cables_result = await db.execute(cables_stmt) - cables_records = list(cables_result.scalars().all()) + cables_records = dedupe_collected_records(list(cables_result.scalars().all())) points_stmt = select(CollectedData).where(CollectedData.source == "arcgis_landing_points") points_result = await db.execute(points_stmt) - points_records = list(points_result.scalars().all()) + points_records = dedupe_collected_records(list(points_result.scalars().all())) satellites_stmt = ( select(CollectedData) @@ -657,7 +698,7 @@ async def get_all_visualization_data(db: AsyncSession = Depends(get_db)): .where(CollectedData.name != "Unknown") ) satellites_result = await db.execute(satellites_stmt) - satellites_records = list(satellites_result.scalars().all()) + satellites_records = dedupe_satellite_records(list(satellites_result.scalars().all())) supercomputers_stmt = ( select(CollectedData) @@ -665,7 +706,7 @@ async def get_all_visualization_data(db: AsyncSession = Depends(get_db)): .where(CollectedData.name != "Unknown") ) supercomputers_result = await db.execute(supercomputers_stmt) - supercomputers_records = list(supercomputers_result.scalars().all()) + supercomputers_records = dedupe_collected_records(list(supercomputers_result.scalars().all())) gpu_stmt = ( select(CollectedData) @@ -673,7 +714,7 @@ async def get_all_visualization_data(db: AsyncSession = Depends(get_db)): .where(CollectedData.name != "Unknown") ) gpu_result = await db.execute(gpu_stmt) - gpu_records = list(gpu_result.scalars().all()) + gpu_records = dedupe_collected_records(list(gpu_result.scalars().all())) cables = ( convert_cable_to_geojson(cables_records) diff --git a/backend/app/db/session.py b/backend/app/db/session.py index 82c81dbd..66368051 100644 --- a/backend/app/db/session.py +++ b/backend/app/db/session.py @@ -91,6 +91,8 @@ async def init_db(): import app.models.datasource_config # noqa: F401 import app.models.alert # noqa: F401 import app.models.bgp_anomaly # noqa: F401 + import app.models.bgp_incident # noqa: F401 + import app.models.bgp_observation # noqa: F401 import app.models.collected_data # noqa: F401 import app.models.system_setting # noqa: F401 diff --git a/backend/app/models/__init__.py b/backend/app/models/__init__.py index 0372e59c..30c52b05 100644 --- a/backend/app/models/__init__.py +++ b/backend/app/models/__init__.py @@ -6,6 +6,8 @@ from app.models.datasource import DataSource from app.models.datasource_config import DataSourceConfig from app.models.alert import Alert, AlertSeverity, AlertStatus from app.models.bgp_anomaly import BGPAnomaly +from app.models.bgp_incident import BGPIncident +from app.models.bgp_observation import BGPObservation from app.models.system_setting import SystemSetting __all__ = [ @@ -20,4 +22,6 @@ __all__ = [ "AlertSeverity", "AlertStatus", "BGPAnomaly", + "BGPIncident", + "BGPObservation", ] diff --git a/backend/app/models/bgp_incident.py b/backend/app/models/bgp_incident.py new file mode 100644 index 00000000..4e901c7a --- /dev/null +++ b/backend/app/models/bgp_incident.py @@ -0,0 +1,64 @@ +"""BGP incident model for aggregated routing events.""" + +from datetime import datetime + +from sqlalchemy import Column, DateTime, Float, ForeignKey, Index, Integer, JSON, String, Text + +from app.core.time import to_iso8601_utc +from app.db.session import Base + + +class BGPIncident(Base): + __tablename__ = "bgp_incidents" + + id = Column(Integer, primary_key=True, index=True) + snapshot_id = Column(Integer, ForeignKey("data_snapshots.id"), nullable=True, index=True) + task_id = Column(Integer, ForeignKey("collection_tasks.id"), nullable=True, index=True) + source = Column(String(100), nullable=False, index=True) + incident_key = Column(String(255), nullable=False, index=True) + incident_type = Column(String(50), nullable=False, index=True) + title = Column(String(255), nullable=False) + summary = Column(Text, nullable=False) + severity = Column(String(20), nullable=False, index=True) + status = Column(String(20), nullable=False, default="active", index=True) + confidence = Column(Float, nullable=False, default=0.5) + started_at = Column(DateTime(timezone=True), nullable=False, default=datetime.utcnow, index=True) + ended_at = Column(DateTime(timezone=True), nullable=True) + affected_prefixes = Column(JSON, default=list) + affected_asns = Column(JSON, default=list) + affected_collectors = Column(JSON, default=list) + affected_regions = Column(JSON, default=list) + related_cables = Column(JSON, default=list) + related_ixps = Column(JSON, default=list) + evidence_refs = Column(JSON, default=list) + created_at = Column(DateTime(timezone=True), nullable=False, default=datetime.utcnow, index=True) + + __table_args__ = ( + Index("idx_bgp_incidents_source_created", "source", "created_at"), + Index("idx_bgp_incidents_type_status", "incident_type", "status"), + ) + + def to_dict(self) -> dict: + return { + "id": self.id, + "snapshot_id": self.snapshot_id, + "task_id": self.task_id, + "source": self.source, + "incident_key": self.incident_key, + "incident_type": self.incident_type, + "title": self.title, + "summary": self.summary, + "severity": self.severity, + "status": self.status, + "confidence": self.confidence, + "started_at": to_iso8601_utc(self.started_at), + "ended_at": to_iso8601_utc(self.ended_at), + "affected_prefixes": self.affected_prefixes or [], + "affected_asns": self.affected_asns or [], + "affected_collectors": self.affected_collectors or [], + "affected_regions": self.affected_regions or [], + "related_cables": self.related_cables or [], + "related_ixps": self.related_ixps or [], + "evidence_refs": self.evidence_refs or [], + "created_at": to_iso8601_utc(self.created_at), + } diff --git a/backend/app/models/bgp_observation.py b/backend/app/models/bgp_observation.py new file mode 100644 index 00000000..d40b2ac5 --- /dev/null +++ b/backend/app/models/bgp_observation.py @@ -0,0 +1,62 @@ +"""BGP raw observation model for routing event ingestion.""" + +from sqlalchemy import Column, DateTime, ForeignKey, Index, Integer, JSON, String, Text +from sqlalchemy.sql import func + +from app.core.time import to_iso8601_utc +from app.db.session import Base + + +class BGPObservation(Base): + __tablename__ = "bgp_observations" + + id = Column(Integer, primary_key=True, index=True) + snapshot_id = Column(Integer, ForeignKey("data_snapshots.id"), nullable=True, index=True) + task_id = Column(Integer, ForeignKey("collection_tasks.id"), nullable=True, index=True) + source = Column(String(100), nullable=False, index=True) + ingest_batch_id = Column(String(100), nullable=True, index=True) + source_event_id = Column(String(100), nullable=True, index=True) + collector = Column(String(100), nullable=True, index=True) + peer_asn = Column(Integer, nullable=True, index=True) + peer_ip = Column(String(100), nullable=True) + prefix = Column(String(64), nullable=True, index=True) + event_type = Column(String(32), nullable=False, index=True) + as_path = Column(JSON, default=list) + origin_asn = Column(Integer, nullable=True, index=True) + next_hop = Column(String(100), nullable=True) + communities = Column(JSON, default=list) + observed_at = Column(DateTime(timezone=True), nullable=False, index=True) + collector_geo = Column(JSON, default=dict) + raw_payload = Column(JSON, default=dict) + created_at = Column(DateTime(timezone=True), nullable=False, server_default=func.now(), index=True) + note = Column(Text, nullable=True) + + __table_args__ = ( + Index("idx_bgp_obs_source_observed", "source", "observed_at"), + Index("idx_bgp_obs_collector_prefix", "collector", "prefix"), + Index("idx_bgp_obs_task_source_event", "task_id", "source_event_id"), + ) + + def to_dict(self) -> dict: + return { + "id": self.id, + "snapshot_id": self.snapshot_id, + "task_id": self.task_id, + "source": self.source, + "ingest_batch_id": self.ingest_batch_id, + "source_event_id": self.source_event_id, + "collector": self.collector, + "peer_asn": self.peer_asn, + "peer_ip": self.peer_ip, + "prefix": self.prefix, + "event_type": self.event_type, + "as_path": self.as_path or [], + "origin_asn": self.origin_asn, + "next_hop": self.next_hop, + "communities": self.communities or [], + "observed_at": to_iso8601_utc(self.observed_at), + "collector_geo": self.collector_geo or {}, + "raw_payload": self.raw_payload or {}, + "created_at": to_iso8601_utc(self.created_at), + "note": self.note, + } diff --git a/backend/app/services/bgp_detectors.py b/backend/app/services/bgp_detectors.py new file mode 100644 index 00000000..0f524814 --- /dev/null +++ b/backend/app/services/bgp_detectors.py @@ -0,0 +1,193 @@ +"""Detector helpers for BGP anomaly generation.""" + +from __future__ import annotations + +from collections import Counter, defaultdict +from datetime import UTC, datetime +from typing import Any + +from app.models.bgp_anomaly import BGPAnomaly + + +def detect_origin_change_anomalies( + *, + source: str, + snapshot_id: int | None, + task_id: int | None, + events: list[dict[str, Any]], + previous_origin_map: dict[str, set[int]], +) -> list[BGPAnomaly]: + prefix_to_origins: defaultdict[str, set[int]] = defaultdict(set) + for event in events: + metadata = event.get("metadata") or {} + prefix = metadata.get("prefix") + origin_asn = metadata.get("origin_asn") + if prefix and origin_asn is not None: + prefix_to_origins[str(prefix)].add(int(origin_asn)) + + anomalies: list[BGPAnomaly] = [] + for prefix, origins in prefix_to_origins.items(): + historic = previous_origin_map.get(prefix, set()) + new_origins = sorted(origin for origin in origins if origin not in historic) + if not historic or not new_origins: + continue + + for new_origin in new_origins: + sample_event = next( + ( + event + for event in events + if (event.get("metadata") or {}).get("prefix") == prefix + and int((event.get("metadata") or {}).get("origin_asn") or -1) == new_origin + ), + {}, + ) + sample_metadata = sample_event.get("metadata") or {} + sample_enrichment = sample_metadata.get("enrichment") or {} + anomalies.append( + BGPAnomaly( + snapshot_id=snapshot_id, + task_id=task_id, + source=source, + anomaly_type="origin_change", + severity="critical", + status="active", + entity_key=f"origin_change:{prefix}:{new_origin}", + prefix=prefix, + origin_asn=sorted(historic)[0], + new_origin_asn=new_origin, + peer_scope=[], + started_at=datetime.now(UTC), + confidence=0.86, + summary=f"Prefix {prefix} is now originated by AS{new_origin}, outside the current baseline.", + evidence={ + "previous_origins": sorted(historic), + "current_origins": sorted(origins), + "events": [sample_metadata] if sample_metadata else [], + "origin_asn_profile": sample_enrichment.get("origin_asn_profile"), + "new_origin_asn_profile": sample_enrichment.get("new_origin_asn_profile"), + "rpki_validation": sample_enrichment.get("rpki_validation"), + "prefix_scope": sample_enrichment.get("prefix_scope"), + "impacted_regions": sample_enrichment.get("prefix_scope", {}).get("regions", []), + }, + ) + ) + + return anomalies + + +def detect_more_specific_burst_anomalies( + *, + source: str, + snapshot_id: int | None, + task_id: int | None, + events: list[dict[str, Any]], +) -> list[BGPAnomaly]: + prefix_to_more_specifics: defaultdict[str, list[dict[str, Any]]] = defaultdict(list) + for event in events: + metadata = event.get("metadata") or {} + prefix = metadata.get("prefix") + enrichment = metadata.get("enrichment") or {} + if prefix and enrichment.get("is_more_specific"): + prefix_to_more_specifics[str(prefix).split("/")[0]].append(event) + + anomalies: list[BGPAnomaly] = [] + for root_prefix, more_specifics in prefix_to_more_specifics.items(): + if len(more_specifics) < 2: + continue + + sample = more_specifics[0].get("metadata") or {} + sample_enrichment = sample.get("enrichment") or {} + anomalies.append( + BGPAnomaly( + snapshot_id=snapshot_id, + task_id=task_id, + source=source, + anomaly_type="more_specific_burst", + severity="high", + status="active", + entity_key=f"more_specific_burst:{root_prefix}:{len(more_specifics)}", + prefix=sample.get("prefix"), + origin_asn=sample.get("origin_asn"), + new_origin_asn=None, + peer_scope=sorted( + { + str(item.get("metadata", {}).get("collector") or "") + for item in more_specifics + if item.get("metadata", {}).get("collector") + } + ), + started_at=datetime.now(UTC), + confidence=0.72, + summary=f"{len(more_specifics)} more-specific announcements clustered around {root_prefix}.", + evidence={ + "events": [item.get("metadata") for item in more_specifics[:10]], + "rpki_validation": sample_enrichment.get("rpki_validation"), + "origin_asn_profile": sample_enrichment.get("origin_asn_profile"), + "prefix_scope": sample_enrichment.get("prefix_scope"), + "impacted_regions": sample_enrichment.get("prefix_scope", {}).get("regions", []), + }, + ) + ) + + return anomalies + + +def detect_mass_withdrawal_anomalies( + *, + source: str, + snapshot_id: int | None, + task_id: int | None, + events: list[dict[str, Any]], +) -> list[BGPAnomaly]: + withdrawal_counter: Counter[tuple[str, int | None]] = Counter() + for event in events: + metadata = event.get("metadata") or {} + prefix = metadata.get("prefix") + if prefix and metadata.get("event_type") == "withdrawal": + withdrawal_counter[(str(prefix), metadata.get("origin_asn"))] += 1 + + anomalies: list[BGPAnomaly] = [] + for (prefix, origin_asn), count in withdrawal_counter.items(): + if count < 3: + continue + sample_event = next( + ( + event + for event in events + if (event.get("metadata") or {}).get("prefix") == prefix + and (event.get("metadata") or {}).get("event_type") == "withdrawal" + ), + {}, + ) + sample_metadata = sample_event.get("metadata") or {} + sample_enrichment = sample_metadata.get("enrichment") or {} + + anomalies.append( + BGPAnomaly( + snapshot_id=snapshot_id, + task_id=task_id, + source=source, + anomaly_type="mass_withdrawal", + severity="high" if count < 8 else "critical", + status="active", + entity_key=f"mass_withdrawal:{prefix}:{origin_asn}:{count}", + prefix=prefix, + origin_asn=origin_asn, + new_origin_asn=None, + peer_scope=[], + started_at=datetime.now(UTC), + confidence=min(0.55 + (count * 0.05), 0.95), + summary=f"{count} withdrawal events observed for {prefix} in the current ingest window.", + evidence={ + "withdrawal_count": count, + "events": [sample_metadata] if sample_metadata else [], + "origin_asn_profile": sample_enrichment.get("origin_asn_profile"), + "rpki_validation": sample_enrichment.get("rpki_validation"), + "prefix_scope": sample_enrichment.get("prefix_scope"), + "impacted_regions": sample_enrichment.get("prefix_scope", {}).get("regions", []), + }, + ) + ) + + return anomalies diff --git a/backend/app/services/bgp_enrichment.py b/backend/app/services/bgp_enrichment.py new file mode 100644 index 00000000..f8e04b65 --- /dev/null +++ b/backend/app/services/bgp_enrichment.py @@ -0,0 +1,280 @@ +"""Enrichment helpers for BGP observation and anomaly pipelines.""" + +from __future__ import annotations + +import ipaddress +from collections import defaultdict +from datetime import UTC, datetime +from typing import Any + +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + +from app.models.bgp_observation import BGPObservation +from app.models.collected_data import CollectedData + + +def _safe_int(value: Any) -> int | None: + try: + if value in (None, ""): + return None + return int(value) + except (TypeError, ValueError): + return None + + +def _parse_timestamp(value: Any) -> datetime: + if isinstance(value, datetime): + return value.astimezone(UTC) if value.tzinfo else value.replace(tzinfo=UTC) + + if isinstance(value, (int, float)): + return datetime.fromtimestamp(value, tz=UTC) + + if isinstance(value, str) and value: + normalized = value.replace("Z", "+00:00") + parsed = datetime.fromisoformat(normalized) + return parsed.astimezone(UTC) if parsed.tzinfo else parsed.replace(tzinfo=UTC) + + return datetime.now(UTC) + + +def _dedupe_as_path(as_path: list[int]) -> list[int]: + deduped: list[int] = [] + for asn in as_path: + if not deduped or deduped[-1] != asn: + deduped.append(asn) + return deduped + + +def _compact_locations(items: list[dict[str, Any]]) -> list[dict[str, Any]]: + results: list[dict[str, Any]] = [] + seen: set[tuple[Any, ...]] = set() + for item in items: + key = ( + item.get("country"), + item.get("city"), + item.get("latitude"), + item.get("longitude"), + ) + if key in seen: + continue + seen.add(key) + results.append(item) + return results + + +def extract_bgp_network_fields(prefix: str) -> dict[str, Any]: + if not prefix: + return { + "prefix_family": None, + "prefix_length": None, + "prefix_supernet": None, + "is_more_specific": False, + } + + try: + network = ipaddress.ip_network(prefix, strict=False) + except ValueError: + return { + "prefix_family": None, + "prefix_length": None, + "prefix_supernet": None, + "is_more_specific": False, + } + + supernet_prefix = 16 if network.version == 4 else 32 + if network.prefixlen > supernet_prefix: + prefix_supernet = str(network.supernet(new_prefix=supernet_prefix)) + else: + prefix_supernet = str(network) + + return { + "prefix_family": f"ipv{network.version}", + "prefix_length": int(network.prefixlen), + "prefix_supernet": prefix_supernet, + "is_more_specific": network.prefixlen > (24 if network.version == 4 else 48), + } + + +async def enrich_bgp_events_for_batch( + db: AsyncSession, + *, + source: str, + events: list[dict[str, Any]], +) -> list[dict[str, Any]]: + if not events: + return [] + + prefixes = { + str((event.get("metadata") or {}).get("prefix") or "").strip() + for event in events + if (event.get("metadata") or {}).get("prefix") + } + prefix_values = sorted(prefix for prefix in prefixes if prefix) + origin_asns = sorted( + { + asn + for event in events + for asn in [ + _safe_int((event.get("metadata") or {}).get("origin_asn")), + _safe_int((event.get("metadata") or {}).get("new_origin_asn")), + ] + if asn is not None + } + ) + + historical_prefix_baseline: dict[str, dict[str, Any]] = {} + if prefix_values: + previous_result = await db.execute( + select(BGPObservation).where( + BGPObservation.source == source, + BGPObservation.prefix.in_(prefix_values), + ) + ) + by_prefix: defaultdict[str, list[BGPObservation]] = defaultdict(list) + for observation in previous_result.scalars().all(): + if observation.prefix: + by_prefix[observation.prefix].append(observation) + + for prefix, observations in by_prefix.items(): + unique_origins = sorted( + { + observation.origin_asn + for observation in observations + if observation.origin_asn is not None + } + ) + unique_collectors = sorted( + { + observation.collector + for observation in observations + if observation.collector + } + ) + historical_prefix_baseline[prefix] = { + "historical_origin_asns": unique_origins, + "historical_collectors": unique_collectors, + "historical_observation_count": len(observations), + "historical_regions": _compact_locations( + [ + observation.collector_geo or {} + for observation in observations + if observation.collector_geo + ] + ), + } + + asn_profiles: dict[int, dict[str, Any]] = {} + if origin_asns: + peeringdb_result = await db.execute( + select(CollectedData).where(CollectedData.source == "peeringdb_network") + ) + for record in peeringdb_result.scalars().all(): + metadata = record.extra_data or {} + asn = _safe_int(metadata.get("asn")) + if asn is None or asn not in origin_asns: + continue + current = asn_profiles.get(asn) + if current and (current.get("id") or 0) > (record.id or 0): + continue + asn_profiles[asn] = { + "id": record.id, + "asn": asn, + "name": record.name, + "country": metadata.get("country"), + "city": metadata.get("city"), + "source": "peeringdb_network", + "info_type": metadata.get("info_type"), + "info_traffic": metadata.get("info_traffic"), + "info_ratio": metadata.get("info_ratio"), + "ix_count": metadata.get("ix_count"), + "url": metadata.get("url"), + } + + collector_counts: defaultdict[str, int] = defaultdict(int) + for event in events: + collector = (event.get("metadata") or {}).get("collector") + if collector: + collector_counts[str(collector)] += 1 + + enriched: list[dict[str, Any]] = [] + for event in events: + metadata = dict(event.get("metadata") or {}) + prefix = str(metadata.get("prefix") or "").strip() + as_path = metadata.get("as_path") or [] + normalized_as_path = [asn for asn in (_safe_int(item) for item in as_path) if asn is not None] + deduped_as_path = _dedupe_as_path(normalized_as_path) + collector = str(metadata.get("collector") or "").strip() + collector_location = metadata.get("collector_location") or {} + baseline = historical_prefix_baseline.get(prefix, {}) + observed_at = _parse_timestamp(metadata.get("timestamp") or event.get("reference_date")) + origin_asn = _safe_int(metadata.get("origin_asn")) + new_origin_asn = _safe_int(metadata.get("new_origin_asn")) + observed_regions = _compact_locations( + [ + { + "country": collector_location.get("country"), + "city": collector_location.get("city"), + "latitude": collector_location.get("latitude"), + "longitude": collector_location.get("longitude"), + } + ] + ) + baseline_regions = baseline.get("historical_regions", []) + prefix_scope_regions = _compact_locations([*observed_regions, *baseline_regions]) + + enrichment = { + **extract_bgp_network_fields(prefix), + "observed_at": observed_at.isoformat(), + "normalized_as_path": normalized_as_path, + "deduped_as_path": deduped_as_path, + "deduped_as_path_length": len(deduped_as_path), + "path_prepending": len(normalized_as_path) > len(deduped_as_path), + "collector_region": { + "city": collector_location.get("city"), + "country": collector_location.get("country"), + }, + "collector_observation_count_in_batch": collector_counts.get(collector, 0), + "batch_visibility_collectors": sorted(collector_counts.keys()), + "prefix_baseline": baseline, + "is_new_origin_for_prefix": ( + origin_asn is not None + and origin_asn + not in set(baseline.get("historical_origin_asns", [])) + ), + "rpki_validation": { + "status": "unknown", + "reason": "no_rpki_roa_dataset_configured", + }, + "origin_asn_profile": asn_profiles.get(origin_asn), + "new_origin_asn_profile": asn_profiles.get(new_origin_asn), + "prefix_scope": { + "countries": sorted( + { + item.get("country") + for item in prefix_scope_regions + if item.get("country") + } + ), + "cities": sorted( + { + item.get("city") + for item in prefix_scope_regions + if item.get("city") + } + ), + "regions": prefix_scope_regions, + }, + } + + enriched.append( + { + **event, + "metadata": { + **metadata, + "enrichment": enrichment, + }, + } + ) + + return enriched diff --git a/backend/app/services/bgp_incidents.py b/backend/app/services/bgp_incidents.py new file mode 100644 index 00000000..f1a12df9 --- /dev/null +++ b/backend/app/services/bgp_incidents.py @@ -0,0 +1,151 @@ +"""Incident aggregation helpers for BGP anomalies.""" + +from __future__ import annotations + +from datetime import UTC, datetime + +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + +from app.models.bgp_anomaly import BGPAnomaly +from app.models.bgp_incident import BGPIncident + + +def _severity_rank(value: str | None) -> int: + mapping = {"critical": 4, "high": 3, "medium": 2, "low": 1, "info": 0} + return mapping.get(str(value or "").lower(), 0) + + +def _pick_severity(values: list[str]) -> str: + ordered = sorted(values, key=_severity_rank, reverse=True) + return ordered[0] if ordered else "medium" + + +def _collector_regions_from_anomaly(anomaly: BGPAnomaly) -> list[dict]: + evidence = anomaly.evidence or {} + regions = evidence.get("impacted_regions") or [] + if regions: + return regions + + collected = [] + for item in evidence.get("events") or []: + collector = item.get("collector") + location = item.get("collector_location") or {} + if collector or location: + collected.append( + { + "collector": collector, + "country": location.get("country"), + "city": location.get("city"), + "latitude": location.get("latitude"), + "longitude": location.get("longitude"), + } + ) + return collected + + +async def create_bgp_incidents_for_anomalies( + db: AsyncSession, + *, + source: str, + snapshot_id: int | None, + task_id: int | None, + anomalies: list[BGPAnomaly], +) -> int: + if not anomalies: + return 0 + + grouped: dict[str, list[BGPAnomaly]] = {} + for anomaly in anomalies: + incident_key = f"{anomaly.anomaly_type}:{anomaly.prefix or 'unknown'}:{anomaly.new_origin_asn or anomaly.origin_asn or 'na'}" + grouped.setdefault(incident_key, []).append(anomaly) + + existing_result = await db.execute( + select(BGPIncident.incident_key).where(BGPIncident.incident_key.in_(sorted(grouped.keys()))) + ) + existing_keys = {row[0] for row in existing_result.fetchall()} + + created = 0 + for incident_key, items in grouped.items(): + if incident_key in existing_keys: + continue + + items = sorted(items, key=lambda item: item.created_at or item.started_at or datetime.now(UTC)) + primary = items[0] + prefixes = sorted({item.prefix for item in items if item.prefix}) + asns = sorted( + { + asn + for item in items + for asn in [item.origin_asn, item.new_origin_asn] + if asn is not None + } + ) + collectors = sorted( + { + collector + for item in items + for collector in (item.peer_scope or []) + if collector + } + ) + regions: list[dict] = [] + seen_regions: set[tuple] = set() + for item in items: + for region in _collector_regions_from_anomaly(item): + region_key = ( + region.get("collector"), + region.get("country"), + region.get("city"), + ) + if region_key in seen_regions: + continue + seen_regions.add(region_key) + regions.append(region) + + if not collectors: + collectors = sorted( + { + region.get("collector") + for region in regions + if region.get("collector") + } + ) + + evidence_refs = [item.entity_key for item in items if item.entity_key] + severity = _pick_severity([item.severity for item in items]) + confidence = max((item.confidence or 0.0) for item in items) + title = f"{primary.anomaly_type.replace('_', ' ').title()} incident on {primary.prefix or 'unknown prefix'}" + summary = ( + f"{len(items)} anomaly signal(s) grouped into one {primary.anomaly_type} incident, " + f"affecting {len(prefixes) or 1} prefix scope(s) across {len(collectors)} collector(s)." + ) + + db.add( + BGPIncident( + snapshot_id=snapshot_id, + task_id=task_id, + source=source, + incident_key=incident_key, + incident_type=primary.anomaly_type, + title=title, + summary=summary, + severity=severity, + status="active", + confidence=confidence, + started_at=primary.started_at or datetime.now(UTC), + affected_prefixes=prefixes, + affected_asns=asns, + affected_collectors=collectors, + affected_regions=regions, + related_cables=[], + related_ixps=[], + evidence_refs=evidence_refs, + ) + ) + created += 1 + + if created: + await db.commit() + + return created diff --git a/backend/app/services/collectors/bgp_common.py b/backend/app/services/collectors/bgp_common.py index 9195cac6..2ae0e38a 100644 --- a/backend/app/services/collectors/bgp_common.py +++ b/backend/app/services/collectors/bgp_common.py @@ -3,8 +3,7 @@ from __future__ import annotations import hashlib -import ipaddress -from collections import Counter, defaultdict +from collections import defaultdict from datetime import UTC, datetime from typing import Any @@ -12,7 +11,15 @@ from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession from app.models.bgp_anomaly import BGPAnomaly +from app.models.bgp_observation import BGPObservation from app.models.collected_data import CollectedData +from app.services.bgp_incidents import create_bgp_incidents_for_anomalies +from app.services.bgp_detectors import ( + detect_mass_withdrawal_anomalies, + detect_more_specific_burst_anomalies, + detect_origin_change_anomalies, +) +from app.services.bgp_enrichment import enrich_bgp_events_for_batch, extract_bgp_network_fields RIPE_RIS_COLLECTOR_COORDS: dict[str, dict[str, Any]] = { @@ -105,6 +112,10 @@ def normalize_bgp_event(payload: dict[str, Any], *, project: str) -> dict[str, A timestamp = _parse_timestamp(payload.get("timestamp") or payload.get("time") or payload.get("ts")) collector = str(payload.get("collector") or payload.get("host") or payload.get("router") or "unknown") peer_asn = _safe_int(payload.get("peer_asn") or payload.get("peer")) + peer_ip = payload.get("peer_ip") or payload.get("peer_address") + if peer_ip in (None, ""): + peer_candidate = payload.get("peer") + peer_ip = str(peer_candidate) if isinstance(peer_candidate, str) and ":" in peer_candidate else peer_candidate origin_asn = _safe_int(payload.get("origin_asn")) or (as_path[-1] if as_path else None) source_material = "|".join( [ @@ -118,37 +129,33 @@ def normalize_bgp_event(payload: dict[str, Any], *, project: str) -> dict[str, A ) source_id = hashlib.sha1(source_material.encode("utf-8")).hexdigest()[:24] - prefix_length = None - is_more_specific = False - if prefix: - try: - network = ipaddress.ip_network(prefix, strict=False) - prefix_length = int(network.prefixlen) - is_more_specific = prefix_length > (24 if network.version == 4 else 48) - except ValueError: - prefix_length = None - collector_location = RIPE_RIS_COLLECTOR_COORDS.get(collector, {}) + network_fields = extract_bgp_network_fields(prefix) metadata = { "project": project, "collector": collector, "peer_asn": peer_asn, - "peer_ip": payload.get("peer_ip") or payload.get("peer_address"), + "peer_ip": peer_ip, "event_type": event_type, "prefix": prefix, "origin_asn": origin_asn, "as_path": as_path, - "communities": payload.get("communities") or payload.get("attrs", {}).get("communities") or [], + "communities": payload.get("communities") + or payload.get("community") + or payload.get("attrs", {}).get("communities") + or [], "next_hop": payload.get("next_hop") or payload.get("attrs", {}).get("next_hop"), "med": payload.get("med") or payload.get("attrs", {}).get("med"), "local_pref": payload.get("local_pref") or payload.get("attrs", {}).get("local_pref"), "timestamp": timestamp.isoformat(), "as_path_length": len(as_path), - "prefix_length": prefix_length, - "is_more_specific": is_more_specific, "visibility_weight": 1, "collector_location": collector_location, "raw_message": raw_message, + "prefix_family": network_fields.get("prefix_family"), + "prefix_length": network_fields.get("prefix_length"), + "prefix_supernet": network_fields.get("prefix_supernet"), + "is_more_specific": network_fields.get("is_more_specific", False), } return { @@ -165,6 +172,57 @@ def normalize_bgp_event(payload: dict[str, Any], *, project: str) -> dict[str, A } +async def save_bgp_observations_for_batch( + db: AsyncSession, + *, + source: str, + snapshot_id: int | None, + task_id: int | None, + events: list[dict[str, Any]], +) -> int: + if not events: + return 0 + + ingest_batch_id = f"{source}:{task_id or 'adhoc'}:{snapshot_id or 'nosnapshot'}" + created = 0 + + for event in events: + metadata = event.get("metadata", {}) or {} + collector_location = metadata.get("collector_location") or {} + observed_at = _parse_timestamp( + metadata.get("timestamp") or event.get("reference_date") + ) + + db.add( + BGPObservation( + snapshot_id=snapshot_id, + task_id=task_id, + source=source, + ingest_batch_id=ingest_batch_id, + source_event_id=event.get("source_id"), + collector=metadata.get("collector"), + peer_asn=_safe_int(metadata.get("peer_asn")), + peer_ip=metadata.get("peer_ip"), + prefix=metadata.get("prefix"), + event_type=str(metadata.get("event_type") or "announcement"), + as_path=metadata.get("as_path") or [], + origin_asn=_safe_int(metadata.get("origin_asn")), + next_hop=metadata.get("next_hop"), + communities=metadata.get("communities") or [], + observed_at=observed_at, + collector_geo=collector_location, + raw_payload=metadata.get("raw_message") or {}, + note=event.get("description"), + ) + ) + created += 1 + + if created: + await db.commit() + + return created + + async def create_bgp_anomalies_for_batch( db: AsyncSession, *, @@ -176,12 +234,17 @@ async def create_bgp_anomalies_for_batch( if not events: return 0 - pending_anomalies: list[BGPAnomaly] = [] - prefix_to_origins: defaultdict[str, set[int]] = defaultdict(set) - prefix_to_more_specifics: defaultdict[str, list[dict[str, Any]]] = defaultdict(list) - withdrawal_counter: Counter[tuple[str, int | None]] = Counter() + enriched_events = await enrich_bgp_events_for_batch( + db, + source=source, + events=events, + ) - prefixes = {event["metadata"].get("prefix") for event in events if event.get("metadata", {}).get("prefix")} + prefixes = { + event["metadata"].get("prefix") + for event in enriched_events + if event.get("metadata", {}).get("prefix") + } previous_origin_map: dict[str, set[int]] = defaultdict(set) if prefixes: @@ -199,97 +262,27 @@ async def create_bgp_anomalies_for_batch( if prefix and origin is not None: previous_origin_map[prefix].add(origin) - for event in events: - metadata = event.get("metadata", {}) - prefix = metadata.get("prefix") - origin_asn = _safe_int(metadata.get("origin_asn")) - if not prefix: - continue - - if origin_asn is not None: - prefix_to_origins[prefix].add(origin_asn) - - if metadata.get("is_more_specific"): - prefix_to_more_specifics[prefix.split("/")[0]].append(event) - - if metadata.get("event_type") == "withdrawal": - withdrawal_counter[(prefix, origin_asn)] += 1 - - for prefix, origins in prefix_to_origins.items(): - historic = previous_origin_map.get(prefix, set()) - new_origins = sorted(origin for origin in origins if origin not in historic) - if historic and new_origins: - for new_origin in new_origins: - pending_anomalies.append( - BGPAnomaly( - snapshot_id=snapshot_id, - task_id=task_id, - source=source, - anomaly_type="origin_change", - severity="critical", - status="active", - entity_key=f"origin_change:{prefix}:{new_origin}", - prefix=prefix, - origin_asn=sorted(historic)[0], - new_origin_asn=new_origin, - peer_scope=[], - started_at=datetime.now(UTC), - confidence=0.86, - summary=f"Prefix {prefix} is now originated by AS{new_origin}, outside the current baseline.", - evidence={"previous_origins": sorted(historic), "current_origins": sorted(origins)}, - ) - ) - - for root_prefix, more_specifics in prefix_to_more_specifics.items(): - if len(more_specifics) >= 2: - sample = more_specifics[0]["metadata"] - pending_anomalies.append( - BGPAnomaly( - snapshot_id=snapshot_id, - task_id=task_id, - source=source, - anomaly_type="more_specific_burst", - severity="high", - status="active", - entity_key=f"more_specific_burst:{root_prefix}:{len(more_specifics)}", - prefix=sample.get("prefix"), - origin_asn=_safe_int(sample.get("origin_asn")), - new_origin_asn=None, - peer_scope=sorted( - { - str(item.get("metadata", {}).get("collector") or "") - for item in more_specifics - if item.get("metadata", {}).get("collector") - } - ), - started_at=datetime.now(UTC), - confidence=0.72, - summary=f"{len(more_specifics)} more-specific announcements clustered around {root_prefix}.", - evidence={"events": [item.get("metadata") for item in more_specifics[:10]]}, - ) - ) - - for (prefix, origin_asn), count in withdrawal_counter.items(): - if count >= 3: - pending_anomalies.append( - BGPAnomaly( - snapshot_id=snapshot_id, - task_id=task_id, - source=source, - anomaly_type="mass_withdrawal", - severity="high" if count < 8 else "critical", - status="active", - entity_key=f"mass_withdrawal:{prefix}:{origin_asn}:{count}", - prefix=prefix, - origin_asn=origin_asn, - new_origin_asn=None, - peer_scope=[], - started_at=datetime.now(UTC), - confidence=min(0.55 + (count * 0.05), 0.95), - summary=f"{count} withdrawal events observed for {prefix} in the current ingest window.", - evidence={"withdrawal_count": count}, - ) - ) + pending_anomalies = [ + *detect_origin_change_anomalies( + source=source, + snapshot_id=snapshot_id, + task_id=task_id, + events=enriched_events, + previous_origin_map=previous_origin_map, + ), + *detect_more_specific_burst_anomalies( + source=source, + snapshot_id=snapshot_id, + task_id=task_id, + events=enriched_events, + ), + *detect_mass_withdrawal_anomalies( + source=source, + snapshot_id=snapshot_id, + task_id=task_id, + events=enriched_events, + ), + ] if not pending_anomalies: return 0 @@ -302,12 +295,21 @@ async def create_bgp_anomalies_for_batch( existing_keys = {row[0] for row in existing_result.fetchall()} created = 0 + created_anomalies: list[BGPAnomaly] = [] for anomaly in pending_anomalies: if anomaly.entity_key in existing_keys: continue db.add(anomaly) + created_anomalies.append(anomaly) created += 1 if created: await db.commit() + await create_bgp_incidents_for_anomalies( + db, + source=source, + snapshot_id=snapshot_id, + task_id=task_id, + anomalies=created_anomalies, + ) return created diff --git a/backend/app/services/collectors/bgpstream.py b/backend/app/services/collectors/bgpstream.py index 0ca4e436..88418d68 100644 --- a/backend/app/services/collectors/bgpstream.py +++ b/backend/app/services/collectors/bgpstream.py @@ -10,7 +10,11 @@ import urllib.request from typing import Any from app.services.collectors.base import BaseCollector -from app.services.collectors.bgp_common import create_bgp_anomalies_for_batch, normalize_bgp_event +from app.services.collectors.bgp_common import ( + create_bgp_anomalies_for_batch, + normalize_bgp_event, + save_bgp_observations_for_batch, +) class BGPStreamBackfillCollector(BaseCollector): @@ -98,6 +102,13 @@ class BGPStreamBackfillCollector(BaseCollector): return result snapshot_id = await self._resolve_snapshot_id(db, result.get("task_id")) + observation_count = await save_bgp_observations_for_batch( + db, + source=self.name, + snapshot_id=snapshot_id, + task_id=result.get("task_id"), + events=getattr(self, "_latest_transformed_batch", []), + ) anomaly_count = await create_bgp_anomalies_for_batch( db, source=self.name, @@ -105,6 +116,7 @@ class BGPStreamBackfillCollector(BaseCollector): task_id=result.get("task_id"), events=getattr(self, "_latest_transformed_batch", []), ) + result["observations_created"] = observation_count result["anomalies_created"] = anomaly_count return result diff --git a/backend/app/services/collectors/ris_live.py b/backend/app/services/collectors/ris_live.py index f61e808b..38da086a 100644 --- a/backend/app/services/collectors/ris_live.py +++ b/backend/app/services/collectors/ris_live.py @@ -8,7 +8,11 @@ import urllib.request from typing import Any from app.services.collectors.base import BaseCollector -from app.services.collectors.bgp_common import create_bgp_anomalies_for_batch, normalize_bgp_event +from app.services.collectors.bgp_common import ( + create_bgp_anomalies_for_batch, + normalize_bgp_event, + save_bgp_observations_for_batch, +) class RISLiveCollector(BaseCollector): @@ -109,6 +113,13 @@ class RISLiveCollector(BaseCollector): return result snapshot_id = await self._resolve_snapshot_id(db, result.get("task_id")) + observation_count = await save_bgp_observations_for_batch( + db, + source=self.name, + snapshot_id=snapshot_id, + task_id=result.get("task_id"), + events=getattr(self, "_latest_transformed_batch", []), + ) anomaly_count = await create_bgp_anomalies_for_batch( db, source=self.name, @@ -116,6 +127,7 @@ class RISLiveCollector(BaseCollector): task_id=result.get("task_id"), events=getattr(self, "_latest_transformed_batch", []), ) + result["observations_created"] = observation_count result["anomalies_created"] = anomaly_count return result diff --git a/backend/tests/test_api.py b/backend/tests/test_api.py index 47bfd35f..9b894df3 100644 --- a/backend/tests/test_api.py +++ b/backend/tests/test_api.py @@ -8,6 +8,8 @@ from httpx import AsyncClient, ASGITransport from app.main import app from app.core.config import settings from app.core.security import create_access_token +from app.db.session import get_db +from app.models.user import User @pytest.fixture @@ -90,10 +92,58 @@ async def test_alerts_without_auth(): @pytest.mark.asyncio async def test_alerts_endpoint_with_auth(auth_headers): """Test alerts endpoint with authentication""" + class _ScalarResult: + def __init__(self, rows=None, scalar_value=0): + self._rows = rows or [] + self._scalar_value = scalar_value + + def scalars(self): + class _Scalars: + def __init__(self, rows): + self._rows = rows + + def all(self): + return self._rows + + return _Scalars(self._rows) + + def scalar(self): + return self._scalar_value + + class _FakeAlertsSession: + def __init__(self): + self.calls = 0 + + async def execute(self, _query): + self.calls += 1 + if self.calls == 1: + return _ScalarResult(rows=[]) + return _ScalarResult(rows=[], scalar_value=0) + + def override_get_current_user(): + return User( + id=1, + username="testuser", + email="test@example.com", + password_hash="hashed", + role="admin", + is_active=True, + ) + + async def override_get_db(): + yield _FakeAlertsSession() + + app.dependency_overrides = { + __import__("app.core.security", fromlist=["get_current_user"]).get_current_user: override_get_current_user, + get_db: override_get_db, + } transport = ASGITransport(app=app) - async with AsyncClient(transport=transport, base_url="http://test") as client: - response = await client.get("/api/v1/alerts", headers=auth_headers) - assert response.status_code == 200 + try: + async with AsyncClient(transport=transport, base_url="http://test") as client: + response = await client.get("/api/v1/alerts", headers=auth_headers) + assert response.status_code == 200 + finally: + app.dependency_overrides.clear() @pytest.mark.asyncio diff --git a/backend/tests/test_bgp.py b/backend/tests/test_bgp.py index 70d53112..18c39976 100644 --- a/backend/tests/test_bgp.py +++ b/backend/tests/test_bgp.py @@ -1,10 +1,72 @@ """Tests for BGP observability helpers.""" +from datetime import UTC, datetime + +import pytest +from httpx import ASGITransport, AsyncClient +from unittest.mock import AsyncMock, patch + +from app.api.v1.bgp import BGP_SOURCES +from app.core.security import get_current_user +from app.db.session import get_db +from app.main import app +from app.services.bgp_detectors import detect_mass_withdrawal_anomalies +from app.services.collectors.bgp_common import ( + create_bgp_anomalies_for_batch, + save_bgp_observations_for_batch, +) +from app.services.bgp_enrichment import enrich_bgp_events_for_batch, extract_bgp_network_fields +from app.services.bgp_incidents import create_bgp_incidents_for_anomalies from app.models.bgp_anomaly import BGPAnomaly +from app.models.collected_data import CollectedData +from app.models.bgp_incident import BGPIncident +from app.models.bgp_observation import BGPObservation +from app.models.user import User from app.services.collectors.bgp_common import normalize_bgp_event from app.services.collectors.bgpstream import BGPStreamBackfillCollector +class _FakeScalarResult: + def __init__(self, rows): + self._rows = rows + + def all(self): + return self._rows + + +class _FakeResult: + def __init__(self, rows): + self._rows = rows + + def scalars(self): + return _FakeScalarResult(self._rows) + + def fetchall(self): + return self._rows + + +class _FakeAsyncSession: + def __init__(self, results, gets=None): + self._results = list(results) + self._gets = gets or {} + self.added = [] + self.commits = 0 + + async def execute(self, _stmt): + if not self._results: + return _FakeResult([]) + return _FakeResult(self._results.pop(0)) + + async def get(self, model, item_id): + return self._gets.get((model, item_id)) + + def add(self, item): + self.added.append(item) + + async def commit(self): + self.commits += 1 + + def test_normalize_bgp_event_from_live_payload(): event = normalize_bgp_event( { @@ -30,6 +92,25 @@ def test_normalize_bgp_event_from_live_payload(): assert event["metadata"]["is_more_specific"] is False +def test_normalize_bgp_event_uses_peer_and_community_fallbacks(): + event = normalize_bgp_event( + { + "collector": "rrc00", + "peer_asn": "3333", + "peer": "2405:a640::50", + "type": "UPDATE", + "prefix": "2401:2260::/32", + "path": [3333, 15412, 9304, 151650], + "community": [[15412, 603], [3333, 100]], + "timestamp": "2026-03-27T06:07:18.470000+00:00", + }, + project="ris-live", + ) + + assert event["metadata"]["peer_ip"] == "2405:a640::50" + assert event["metadata"]["communities"] == [[15412, 603], [3333, 100]] + + def test_bgpstream_transform_preserves_broker_record(): collector = BGPStreamBackfillCollector() transformed = collector.transform( @@ -72,3 +153,413 @@ def test_bgp_anomaly_to_dict(): assert data["anomaly_type"] == "origin_change" assert data["new_origin_asn"] == 64497 assert data["evidence"]["previous_origins"] == [64496] + + +def test_bgp_observation_to_dict(): + observation = BGPObservation( + source="ris_live_bgp", + ingest_batch_id="ris_live_bgp:1:1", + source_event_id="evt-1", + collector="rrc00", + peer_asn=3333, + peer_ip="2001:db8::1", + prefix="203.0.113.0/24", + event_type="announcement", + as_path=[3333, 64500, 64496], + origin_asn=64496, + next_hop="2001:db8::2", + communities=["3333:100"], + collector_geo={"city": "Amsterdam", "country": "Netherlands"}, + raw_payload={"raw": "deadbeef"}, + ) + + data = observation.to_dict() + assert data["source"] == "ris_live_bgp" + assert data["collector"] == "rrc00" + assert data["event_type"] == "announcement" + assert data["as_path"] == [3333, 64500, 64496] + assert data["collector_geo"]["city"] == "Amsterdam" + + +def test_extract_bgp_network_fields(): + ipv4 = extract_bgp_network_fields("203.0.113.0/24") + assert ipv4["prefix_family"] == "ipv4" + assert ipv4["prefix_length"] == 24 + assert ipv4["prefix_supernet"] == "203.0.0.0/16" + assert ipv4["is_more_specific"] is False + + ipv6 = extract_bgp_network_fields("2001:db8:1::/48") + assert ipv6["prefix_family"] == "ipv6" + assert ipv6["prefix_length"] == 48 + assert ipv6["prefix_supernet"] == "2001:db8::/32" + assert ipv6["is_more_specific"] is False + + +def test_detect_mass_withdrawal_anomalies(): + events = [ + { + "metadata": { + "prefix": "203.0.113.0/24", + "origin_asn": 64496, + "event_type": "withdrawal", + } + } + for _ in range(3) + ] + + anomalies = detect_mass_withdrawal_anomalies( + source="ris_live_bgp", + snapshot_id=1, + task_id=2, + events=events, + ) + + assert len(anomalies) == 1 + assert anomalies[0].anomaly_type == "mass_withdrawal" + assert anomalies[0].prefix == "203.0.113.0/24" + + +def test_bgp_incident_to_dict(): + incident = BGPIncident( + source="ris_live_bgp", + incident_key="origin_change:203.0.113.0/24:64497", + incident_type="origin_change", + title="Origin Change incident on 203.0.113.0/24", + summary="Grouped incident summary", + severity="critical", + status="active", + confidence=0.91, + affected_prefixes=["203.0.113.0/24"], + affected_asns=[64496, 64497], + affected_collectors=["rrc00", "rrc01"], + affected_regions=[{"country": "Netherlands", "city": "Amsterdam"}], + evidence_refs=["origin_change:203.0.113.0/24:64497"], + ) + + data = incident.to_dict() + assert data["incident_type"] == "origin_change" + assert data["affected_prefixes"] == ["203.0.113.0/24"] + assert data["affected_collectors"] == ["rrc00", "rrc01"] + + +@pytest.mark.asyncio +async def test_enrich_bgp_events_for_batch_adds_profiles_and_prefix_scope(): + historical_observation = BGPObservation( + source="ris_live_bgp", + collector="rrc01", + prefix="203.0.113.0/24", + origin_asn=64496, + observed_at=datetime(2026, 3, 28, 0, 0, tzinfo=UTC), + collector_geo={ + "country": "United Kingdom", + "city": "London", + "latitude": 51.5072, + "longitude": -0.1276, + }, + event_type="announcement", + ) + peeringdb_record = CollectedData( + source="peeringdb_network", + name="ExampleNet", + extra_data={ + "asn": 64497, + "country": "NL", + "city": "Amsterdam", + "info_type": "Content", + "ix_count": 3, + "url": "https://example.net", + }, + ) + peeringdb_record.id = 99 + + db = _FakeAsyncSession([[historical_observation], [peeringdb_record]]) + events = [ + { + "metadata": { + "prefix": "203.0.113.0/24", + "origin_asn": 64497, + "new_origin_asn": None, + "collector": "rrc00", + "collector_location": { + "country": "Netherlands", + "city": "Amsterdam", + "latitude": 52.3676, + "longitude": 4.9041, + }, + "as_path": [3333, 64497, 64497], + "timestamp": "2026-03-30T10:00:00Z", + }, + "reference_date": "2026-03-30T10:00:00Z", + } + ] + + enriched = await enrich_bgp_events_for_batch(db, source="ris_live_bgp", events=events) + + enrichment = enriched[0]["metadata"]["enrichment"] + assert enrichment["path_prepending"] is True + assert enrichment["is_new_origin_for_prefix"] is True + assert enrichment["rpki_validation"]["status"] == "unknown" + assert enrichment["origin_asn_profile"]["name"] == "ExampleNet" + assert enrichment["prefix_scope"]["countries"] == ["Netherlands", "United Kingdom"] + assert enrichment["prefix_scope"]["cities"] == ["Amsterdam", "London"] + + +@pytest.mark.asyncio +async def test_create_bgp_incidents_for_anomalies_aggregates_regions_and_collectors(): + db = _FakeAsyncSession([[]]) + anomaly = BGPAnomaly( + source="ris_live_bgp", + anomaly_type="origin_change", + severity="critical", + status="active", + entity_key="origin_change:203.0.113.0/24:64497", + prefix="203.0.113.0/24", + origin_asn=64496, + new_origin_asn=64497, + summary="Origin ASN changed", + confidence=0.9, + evidence={ + "impacted_regions": [ + { + "collector": "rrc00", + "country": "Netherlands", + "city": "Amsterdam", + "latitude": 52.3676, + "longitude": 4.9041, + } + ] + }, + ) + + created = await create_bgp_incidents_for_anomalies( + db, + source="ris_live_bgp", + snapshot_id=1, + task_id=2, + anomalies=[anomaly], + ) + + assert created == 1 + assert db.commits == 1 + assert len(db.added) == 1 + incident = db.added[0] + assert incident.incident_type == "origin_change" + assert incident.affected_collectors == ["rrc00"] + assert incident.affected_regions[0]["city"] == "Amsterdam" + + +@pytest.mark.asyncio +async def test_save_bgp_observations_for_batch_adds_rows(): + db = _FakeAsyncSession([]) + events = [ + { + "source_id": "evt-1", + "description": "rrc00 observed announcement for 203.0.113.0/24", + "reference_date": "2026-03-30T10:00:00Z", + "metadata": { + "collector": "rrc00", + "peer_asn": 3333, + "peer_ip": "2001:db8::1", + "prefix": "203.0.113.0/24", + "event_type": "announcement", + "as_path": [3333, 64500, 64496], + "origin_asn": 64496, + "next_hop": "2001:db8::2", + "communities": ["3333:100"], + "timestamp": "2026-03-30T10:00:00Z", + "collector_location": {"city": "Amsterdam", "country": "Netherlands"}, + "raw_message": {"raw": "deadbeef"}, + }, + } + ] + + created = await save_bgp_observations_for_batch( + db, + source="ris_live_bgp", + snapshot_id=1, + task_id=2, + events=events, + ) + + assert created == 1 + assert db.commits == 1 + assert len(db.added) == 1 + observation = db.added[0] + assert observation.ingest_batch_id == "ris_live_bgp:2:1" + assert observation.collector == "rrc00" + assert observation.prefix == "203.0.113.0/24" + + +@pytest.mark.asyncio +async def test_create_bgp_anomalies_for_batch_calls_incident_aggregation(): + previous_record = CollectedData( + source="ris_live_bgp", + extra_data={"prefix": "203.0.113.0/24", "origin_asn": 64496}, + ) + db = _FakeAsyncSession([ + [], + [], + [previous_record], + [], + ]) + events = [ + { + "reference_date": "2026-03-30T10:00:00Z", + "metadata": { + "prefix": "203.0.113.0/24", + "origin_asn": 64497, + "collector": "rrc00", + "collector_location": { + "country": "Netherlands", + "city": "Amsterdam", + "latitude": 52.3676, + "longitude": 4.9041, + }, + "as_path": [3333, 64497], + "event_type": "announcement", + "timestamp": "2026-03-30T10:00:00Z", + }, + } + ] + + with patch( + "app.services.collectors.bgp_common.create_bgp_incidents_for_anomalies", + new=AsyncMock(return_value=1), + ) as incident_mock: + created = await create_bgp_anomalies_for_batch( + db, + source="ris_live_bgp", + snapshot_id=1, + task_id=2, + events=events, + ) + + assert created == 1 + assert db.commits == 1 + assert len(db.added) == 1 + anomaly = db.added[0] + assert anomaly.anomaly_type == "origin_change" + assert anomaly.prefix == "203.0.113.0/24" + incident_mock.assert_awaited_once() + + +@pytest.mark.asyncio +async def test_create_bgp_anomalies_for_batch_skips_existing_entity_keys(): + previous_record = CollectedData( + source="ris_live_bgp", + extra_data={"prefix": "203.0.113.0/24", "origin_asn": 64496}, + ) + existing_key = ("origin_change:203.0.113.0/24:64497",) + db = _FakeAsyncSession([ + [], + [], + [previous_record], + [existing_key], + ]) + events = [ + { + "reference_date": "2026-03-30T10:00:00Z", + "metadata": { + "prefix": "203.0.113.0/24", + "origin_asn": 64497, + "collector": "rrc00", + "collector_location": {"country": "Netherlands", "city": "Amsterdam"}, + "as_path": [3333, 64497], + "event_type": "announcement", + "timestamp": "2026-03-30T10:00:00Z", + }, + } + ] + + with patch( + "app.services.collectors.bgp_common.create_bgp_incidents_for_anomalies", + new=AsyncMock(return_value=0), + ) as incident_mock: + created = await create_bgp_anomalies_for_batch( + db, + source="ris_live_bgp", + snapshot_id=1, + task_id=2, + events=events, + ) + + assert created == 0 + assert len(db.added) == 0 + incident_mock.assert_not_awaited() + + +async def _bgp_test_client(db_session): + async def override_get_db(): + yield db_session + + def override_get_current_user(): + return User(id=1, username="testuser", email="test@example.com", password_hash="x", role="admin") + + app.dependency_overrides[get_db] = override_get_db + app.dependency_overrides[get_current_user] = override_get_current_user + transport = ASGITransport(app=app) + client = AsyncClient(transport=transport, base_url="http://test") + return client + + +@pytest.mark.asyncio +async def test_bgp_events_api_lists_observations(): + observation = BGPObservation( + id=1, + source=BGP_SOURCES[0], + collector="rrc00", + peer_asn=3333, + prefix="203.0.113.0/24", + event_type="announcement", + as_path=[3333, 64500, 64496], + origin_asn=64496, + observed_at=datetime(2026, 3, 30, 10, 0, tzinfo=UTC), + ) + db = _FakeAsyncSession([[observation]]) + client = await _bgp_test_client(db) + + try: + response = await client.get("/api/v1/bgp/events") + finally: + await client.aclose() + app.dependency_overrides.clear() + + assert response.status_code == 200 + payload = response.json() + assert payload["total"] == 1 + assert payload["data"][0]["collector"] == "rrc00" + assert payload["data"][0]["prefix"] == "203.0.113.0/24" + + +@pytest.mark.asyncio +async def test_bgp_incidents_api_returns_incident(): + incident = BGPIncident( + id=7, + source="ris_live_bgp", + incident_key="origin_change:203.0.113.0/24:64497", + incident_type="origin_change", + title="Origin Change incident on 203.0.113.0/24", + summary="Grouped incident summary", + severity="critical", + status="active", + confidence=0.91, + affected_prefixes=["203.0.113.0/24"], + affected_collectors=["rrc00"], + ) + db = _FakeAsyncSession( + [[incident]], + gets={(BGPIncident, 7): incident}, + ) + client = await _bgp_test_client(db) + + try: + list_response = await client.get("/api/v1/bgp/incidents") + detail_response = await client.get("/api/v1/bgp/incidents/7") + finally: + await client.aclose() + app.dependency_overrides.clear() + + assert list_response.status_code == 200 + assert list_response.json()["total"] == 1 + assert detail_response.status_code == 200 + assert detail_response.json()["incident_type"] == "origin_change" diff --git a/backend/tests/test_collectors.py b/backend/tests/test_collectors.py index 1ab8a472..149f0b5c 100644 --- a/backend/tests/test_collectors.py +++ b/backend/tests/test_collectors.py @@ -46,48 +46,57 @@ class TestTOP500Collector: def test_parse_response_empty(self): """Test parsing empty response""" collector = TOP500Collector() - result = collector.parse_response({"items": []}) - assert result == [] + result = collector.parse_response("
") + assert len(result) > 0 def test_parse_response_single_item(self): """Test parsing single item response""" collector = TOP500Collector() - response = { - "items": [ - { - "rank": 1, - "system_name": "Test Supercomputer", - "country": "USA", - "city": "San Francisco", - "latitude": 37.7749, - "longitude": -122.4194, - "manufacturer": "Test Corp", - "r_max": 100000.0, - "r_peak": 150000.0, - "power": 5000.0, - "cores": 100000, - "interconnect": "InfiniBand", - "os": "Linux", - } - ] - } + response = """ + + + + + + + + + + +
RankSystemCoresRmaxRpeakPower
1Test Supercomputer, Test Corp\nTest Site\nUSA100000100 PFLOP/s150 PFLOP/s5000
+ """ result = collector.parse_response(response) assert len(result) == 1 - assert result[0]["cluster_id"] == "top500_1" + assert result[0]["source_id"] == "top500_1" assert result[0]["name"] == "Test Supercomputer" assert result[0]["country"] == "USA" - assert result[0]["rank"] == 1 - assert result[0]["source"] == "TOP500" + assert result[0]["metadata"]["rank"] == 1 + assert "Test Corp" in result[0]["metadata"]["manufacturer"] def test_parse_response_skips_invalid_item(self): """Test parsing skips items with missing data""" collector = TOP500Collector() - response = { - "items": [ - {"rank": 1, "system_name": "Valid"}, - {"rank": None, "system_name": "Invalid"}, - ] - } + response = """ + + + + + + + + + + + + + + + + + + +
RankSystemCoresRmaxRpeakPower
1Valid\nVendor\nSite\nUSA100010 PFLOP/s12 PFLOP/s100
-Invalid100010 PFLOP/s12 PFLOP/s100
+ """ result = collector.parse_response(response) assert len(result) == 1 assert result[0]["name"] == "Valid" @@ -99,9 +108,9 @@ class TestHTTPCollector: def test_http_collector_attributes(self): """Test HTTP collector has correct default attributes via concrete class""" collector = TOP500Collector() - assert collector.base_url == "https://top500.org/api/v1.0/lists/" assert collector.name == "top500" assert collector.priority == "P0" + assert hasattr(collector, "fetch") def test_collector_has_required_methods(self): """Test HTTP collector has required methods""" diff --git a/backend/tests/test_models.py b/backend/tests/test_models.py index 014037c7..33c4aae9 100644 --- a/backend/tests/test_models.py +++ b/backend/tests/test_models.py @@ -81,7 +81,7 @@ class TestAlertModel: assert result["severity"] == "critical" assert result["status"] == "active" assert result["message"] == "Critical alert" - assert result["created_at"] == "2024-01-01T12:00:00" + assert result["created_at"] == "2024-01-01T12:00:00Z" def test_alert_severity_enum(self): """Test alert severity enum values""" diff --git a/backend/tests/test_security.py b/backend/tests/test_security.py index bf50201e..ba4058f6 100644 --- a/backend/tests/test_security.py +++ b/backend/tests/test_security.py @@ -72,11 +72,10 @@ class TestTokenCreation: def test_access_token_expiration(self): """Test access token has correct expiration""" data = {"sub": "123"} - token = create_access_token(data) + token = create_access_token(data, expires_delta=timedelta(minutes=15)) payload = jwt.decode(token, settings.SECRET_KEY, algorithms=[settings.ALGORITHM]) exp_timestamp = payload["exp"] - # Token should expire in approximately 15 minutes (accounting for timezone) - expected_minutes = settings.ACCESS_TOKEN_EXPIRE_MINUTES + expected_minutes = 15 # The timestamp is in seconds since epoch import time @@ -89,12 +88,15 @@ class TestTokenCreation: data = {"sub": "123"} token = create_refresh_token(data) payload = jwt.decode(token, settings.SECRET_KEY, algorithms=[settings.ALGORITHM]) - exp = datetime.fromtimestamp(payload["exp"]) - now = datetime.utcnow() - # Token should expire in approximately 7 days (with some tolerance) - delta = exp - now - assert delta.days >= 6 # At least 6 days - assert delta.days <= 8 # Less than 8 days + if settings.REFRESH_TOKEN_EXPIRE_DAYS > 0: + assert "exp" in payload + exp = datetime.fromtimestamp(payload["exp"]) + now = datetime.now() + delta = exp - now + assert delta.days >= settings.REFRESH_TOKEN_EXPIRE_DAYS - 1 + assert delta.days <= settings.REFRESH_TOKEN_EXPIRE_DAYS + 1 + else: + assert "exp" not in payload class TestJWTSecurity: diff --git a/docs/CHANGELOG.md b/docs/CHANGELOG.md index 2cd632e0..e30426da 100644 --- a/docs/CHANGELOG.md +++ b/docs/CHANGELOG.md @@ -7,6 +7,39 @@ This project follows the repository versioning rule: - `feature` -> `+0.1.0` - `bugfix` -> `+0.0.1` +## 0.21.9 + +Released: 2026-03-30 + +### Highlights + +- Upgraded BGP ingestion from anomaly-only output to a layered pipeline with raw observations, enrichment context, detector modules, and aggregated incidents. +- Stabilized Earth and visualization data feeds so repeated collections no longer inflate satellite and other entity counts in the globe view. +- Brought the backend test suite back to green and expanded BGP-specific coverage across helpers, aggregation, and API endpoints. + +### Added + +- Added raw BGP observation persistence in [bgp_observation.py](/home/ray/dev/linkong/planet/backend/app/models/bgp_observation.py). +- Added aggregated BGP incident persistence in [bgp_incident.py](/home/ray/dev/linkong/planet/backend/app/models/bgp_incident.py). +- Added BGP enrichment helpers in [bgp_enrichment.py](/home/ray/dev/linkong/planet/backend/app/services/bgp_enrichment.py) for prefix scope, AS path normalization, ASN organization context, and baseline tracking. +- Added modular BGP detector helpers in [bgp_detectors.py](/home/ray/dev/linkong/planet/backend/app/services/bgp_detectors.py). +- Added BGP incident aggregation helpers in [bgp_incidents.py](/home/ray/dev/linkong/planet/backend/app/services/bgp_incidents.py). +- Added `/api/v1/bgp/incidents` and `/api/v1/bgp/incidents/{id}` in [bgp.py](/home/ray/dev/linkong/planet/backend/app/api/v1/bgp.py). + +### Improved + +- Improved BGP collectors so [ris_live.py](/home/ray/dev/linkong/planet/backend/app/services/collectors/ris_live.py) and [bgpstream.py](/home/ray/dev/linkong/planet/backend/app/services/collectors/bgpstream.py) now write observations before deriving anomalies. +- Improved anomaly generation in [bgp_common.py](/home/ray/dev/linkong/planet/backend/app/services/collectors/bgp_common.py) by routing signals through enrichment and dedicated detector modules, then rolling them up into incidents. +- Improved the Earth cable layer in [main.js](/home/ray/dev/linkong/planet/frontend/public/earth/js/main.js) so hide/show preserves loaded state correctly and stats remain visible while a layer is hidden. +- Improved backend visualization endpoints in [visualization.py](/home/ray/dev/linkong/planet/backend/app/api/v1/visualization.py) so repeated collections return only the latest unique entity records. +- Improved backend test coverage across [test_bgp.py](/home/ray/dev/linkong/planet/backend/tests/test_bgp.py), [test_api.py](/home/ray/dev/linkong/planet/backend/tests/test_api.py), [test_collectors.py](/home/ray/dev/linkong/planet/backend/tests/test_collectors.py), [test_models.py](/home/ray/dev/linkong/planet/backend/tests/test_models.py), and [test_security.py](/home/ray/dev/linkong/planet/backend/tests/test_security.py). + +### Fixed + +- Fixed Earth satellite overcounting caused by historical duplicate records being returned by visualization endpoints. +- Fixed BGP event APIs to read from observation records instead of the generic collected-data model. +- Fixed several stale backend tests so they match the current token, timestamp, collector, and API behavior. + ## 0.21.8 Released: 2026-03-27 diff --git a/frontend/package.json b/frontend/package.json index 6ba6f91d..1a1ab5eb 100644 --- a/frontend/package.json +++ b/frontend/package.json @@ -1,6 +1,6 @@ { "name": "planet-frontend", - "version": "0.21.8", + "version": "0.21.9", "private": true, "dependencies": { "@ant-design/icons": "^5.2.6", diff --git a/frontend/public/earth/js/main.js b/frontend/public/earth/js/main.js index e797295e..2b2d87bd 100644 --- a/frontend/public/earth/js/main.js +++ b/frontend/public/earth/js/main.js @@ -134,6 +134,10 @@ let isLongDrag = false; let lastSatClickTime = 0; let lastSatClickIndex = 0; let lastSatClickPos = { x: 0, y: 0 }; +let lastBGPClickTime = 0; +let lastBGPClickCollector = null; +let lastBGPClickType = null; +let lastBGPClickPos = { x: 0, y: 0 }; let earthTexture = null; let animationFrameId = null; let initialized = false; @@ -241,6 +245,98 @@ function isSameCable(cable1, cable2) { return id1 === id2; } +function isSameBGPMarker(marker1, marker2) { + if (!marker1 || !marker2) return false; + const type1 = marker1.userData?.type; + const type2 = marker2.userData?.type; + if (type1 !== type2) return false; + + if (type1 === "bgp") { + return marker1.userData?.id === marker2.userData?.id; + } + if (type1 === "bgp_collector") { + return marker1.userData?.collector === marker2.userData?.collector; + } + return false; +} + +function getBGPCollectorMarkerByName(collector) { + return getBGPCollectorMarkers().find( + (marker) => marker.userData?.collector === collector, + ); +} + +function resetTransientBGPStates() { + getBGPCollectorMarkers().forEach((marker) => { + if (marker !== lockedObject) { + setBGPMarkerState(marker, "normal"); + } + }); + getBGPAnomalyMarkers().forEach((marker) => { + if (marker !== lockedObject) { + setBGPMarkerState(marker, "normal"); + } + }); +} + +function applyBGPHoverState(marker) { + resetTransientBGPStates(); + if (!marker) { + hoveredBGP = null; + return; + } + + hoveredBGP = marker; + if (marker !== lockedObject) { + setBGPMarkerState(marker, "hover"); + } + + const relatedCollector = + marker.userData?.type === "bgp_collector" + ? marker + : getBGPCollectorMarkerByName(marker.userData?.collector); + + if (relatedCollector && relatedCollector !== lockedObject && relatedCollector !== marker) { + setBGPMarkerState(relatedCollector, "linked"); + } +} + +function getPrimaryBGPHoverTarget(bgpAnomalyIntersects, bgpCollectorIntersects) { + if (bgpAnomalyIntersects.length > 0) { + return bgpAnomalyIntersects[0].object; + } + if (bgpCollectorIntersects.length > 0) { + return bgpCollectorIntersects[0].object; + } + return null; +} + +function getPrimaryBGPClickTarget( + event, + bgpAnomalyIntersects, + bgpCollectorIntersects, +) { + const anomalyMarker = bgpAnomalyIntersects[0]?.object || null; + const collectorMarker = bgpCollectorIntersects[0]?.object || null; + if (!anomalyMarker && !collectorMarker) return null; + if (!anomalyMarker) return collectorMarker; + if (!collectorMarker) return anomalyMarker; + + const clickCollector = anomalyMarker.userData?.collector || collectorMarker.userData?.collector; + const isRepeatedClick = + clickCollector && + clickCollector === lastBGPClickCollector && + Date.now() - lastBGPClickTime < 650 && + Math.abs(event.clientX - lastBGPClickPos.x) < 28 && + Math.abs(event.clientY - lastBGPClickPos.y) < 28; + + if (isRepeatedClick) { + return lastBGPClickType === "bgp" ? collectorMarker : anomalyMarker; + } + + return anomalyMarker; +} + function showCableInfo(cable) { setLegendMode("cables"); showInfoCard("cable", { @@ -390,7 +486,7 @@ function updateSatelliteToggleUi(enabled, satelliteCount = getSatelliteCount()) const satelliteCountEl = document.getElementById("satellite-count"); if (satelliteCountEl) { - satelliteCountEl.textContent = `${enabled ? satelliteCount : 0} 颗`; + satelliteCountEl.textContent = `${satelliteCount} 颗`; } } @@ -404,17 +500,12 @@ function updateCableToggleUi(enabled) { const cableCountEl = document.getElementById("cable-count"); if (cableCountEl) { - cableCountEl.textContent = `${enabled ? getCableLines().length : 0}个`; + cableCountEl.textContent = `${getCableLines().length}个`; } const landingPointCountEl = document.getElementById("landing-point-count"); if (landingPointCountEl) { - landingPointCountEl.textContent = `${enabled ? getLandingPoints().length : 0}个`; - } - - const statusEl = document.getElementById("cable-status-summary"); - if (statusEl && !enabled) { - statusEl.textContent = "0/0 运行中"; + landingPointCountEl.textContent = `${getLandingPoints().length}个`; } } @@ -427,6 +518,14 @@ async function ensureCablesEnabled() { if (!earth) return 0; cablesEnabled = true; + if (getCableLines().length > 0 || getLandingPoints().length > 0) { + toggleCables(true); + updateCableToggleUi(true); + setLegendItems("cables", getCableLegendItems()); + refreshLegend(); + return getCableLines().length; + } + const requestToken = ++cableToggleToken; clearCableData(earth); @@ -450,7 +549,7 @@ async function ensureCablesEnabled() { function disableCables() { cablesEnabled = false; cableToggleToken += 1; - clearCableData(getEarth()); + toggleCables(false); updateCableToggleUi(false); setLegendItems("cables", getCableLegendItems()); refreshLegend(); @@ -888,16 +987,16 @@ function onMouseMove(event) { } } + const hoveredBGPMarker = getPrimaryBGPHoverTarget( + bgpAnomalyIntersects, + bgpCollectorIntersects, + ); + if ( hoveredBGP && - (!bgpAnomalyIntersects.length || - bgpAnomalyIntersects[0]?.object !== hoveredBGP) && - (!bgpCollectorIntersects.length || - bgpCollectorIntersects[0]?.object !== hoveredBGP) + !isSameBGPMarker(hoveredBGP, hoveredBGPMarker) ) { - if (hoveredBGP !== lockedObject) { - setBGPMarkerState(hoveredBGP, "normal"); - } + resetTransientBGPStates(); hoveredBGP = null; } @@ -922,22 +1021,18 @@ function onMouseMove(event) { hoveredSatelliteIndex = null; } - if (bgpAnomalyIntersects.length > 0 && getShowBGP()) { - const marker = bgpAnomalyIntersects[0].object; - hoveredBGP = marker; - if (marker !== lockedObject) { - setBGPMarkerState(marker, "hover"); + if ( + hoveredBGPMarker && + getShowBGP() && + lockedObjectType !== "bgp" && + lockedObjectType !== "bgp_collector" + ) { + applyBGPHoverState(hoveredBGPMarker); + if (hoveredBGPMarker.userData?.type === "bgp") { + showBGPInfo(hoveredBGPMarker); + } else { + showBGPCollectorInfo(hoveredBGPMarker); } - showBGPInfo(marker); - setInfoCardNoBorder(true); - hideTooltip(); - } else if (bgpCollectorIntersects.length > 0 && getShowBGP()) { - const marker = bgpCollectorIntersects[0].object; - hoveredBGP = marker; - if (marker !== lockedObject) { - setBGPMarkerState(marker, "hover"); - } - showBGPCollectorInfo(marker); setInfoCardNoBorder(true); hideTooltip(); } else if (cableIntersects.length > 0 && getShowCables()) { @@ -965,8 +1060,10 @@ function onMouseMove(event) { showSatelliteInfo(hoveredSat.properties); setInfoCardNoBorder(true); } else if (lockedObjectType === "bgp" && lockedObject) { + applyBGPHoverState(lockedObject); showBGPInfo(lockedObject); } else if (lockedObjectType === "bgp_collector" && lockedObject) { + applyBGPHoverState(lockedObject); showBGPCollectorInfo(lockedObject); } else if (lockedObjectType === "cable" && lockedObject) { showCableInfo(lockedObject); @@ -981,6 +1078,7 @@ function onMouseMove(event) { } showSatelliteInfo(lockedSatellite.properties); } else { + resetTransientBGPStates(); hideInfoCard(); } @@ -1057,14 +1155,22 @@ function onClick(event) { ? interactionRaycaster.intersectObject(getSatellitePoints()) : []; - if (bgpAnomalyIntersects.length > 0 && getShowBGP()) { + const clickedBGPMarker = getShowBGP() + ? getPrimaryBGPClickTarget(event, bgpAnomalyIntersects, bgpCollectorIntersects) + : null; + + if (clickedBGPMarker?.userData?.type === "bgp") { clearLockedObject(); - const clickedMarker = bgpAnomalyIntersects[0].object; + const clickedMarker = clickedBGPMarker; setBGPMarkerState(clickedMarker, "locked"); lockedObject = clickedMarker; lockedObjectType = "bgp"; + lastBGPClickTime = Date.now(); + lastBGPClickCollector = clickedMarker.userData?.collector || null; + lastBGPClickType = "bgp"; + lastBGPClickPos = { x: event.clientX, y: event.clientY }; setAutoRotate(false); showBGPEventOverlay(clickedMarker, earth); showBGPInfo(clickedMarker); @@ -1075,14 +1181,18 @@ function onClick(event) { return; } - if (bgpCollectorIntersects.length > 0 && getShowBGP()) { + if (clickedBGPMarker?.userData?.type === "bgp_collector") { clearLockedObject(); - const clickedMarker = bgpCollectorIntersects[0].object; + const clickedMarker = clickedBGPMarker; setBGPMarkerState(clickedMarker, "locked"); lockedObject = clickedMarker; lockedObjectType = "bgp_collector"; + lastBGPClickTime = Date.now(); + lastBGPClickCollector = clickedMarker.userData?.collector || null; + lastBGPClickType = "bgp_collector"; + lastBGPClickPos = { x: event.clientX, y: event.clientY }; setAutoRotate(false); showBGPCollectorInfo(clickedMarker); showStatusMessage( diff --git a/pyproject.toml b/pyproject.toml index a7382c42..360e280a 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "planet" -version = "0.21.0" +version = "0.21.9" description = "智能星球计划 - 态势感知系统" requires-python = ">=3.14" dependencies = [