Files
planet/backend/tests/test_bgp.py

1400 lines
44 KiB
Python

"""Tests for BGP observability helpers."""
from datetime import UTC, datetime, timedelta
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,
detect_origin_change_anomalies,
detect_path_flap_anomalies,
detect_route_leak_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,
infer_related_infrastructure,
)
from app.services.bgp_collectors import build_bgp_collector_coverage
from app.api.v1.visualization import convert_bgp_incidents_to_geojson, build_incident_geography_hints
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
from app.services.collectors.iptoasn import IPtoASNPrefixGeoCollector
from app.services.collectors.opengeofeed import OpenGeoFeedPrefixGeoCollector
from app.services.collectors.nro_delegated import NRODelegatedPrefixGeoCollector
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
def fetchone(self):
return self._rows[0] if self._rows else None
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, _params=None):
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(
{
"collector": "rrc00",
"peer_asn": "3333",
"peer_ip": "2001:db8::1",
"type": "UPDATE",
"event_type": "announcement",
"prefix": "203.0.113.0/24",
"path": ["3333", "64500", "64496"],
"communities": ["3333:100"],
"timestamp": "2026-03-26T08:00:00Z",
},
project="ris-live",
)
assert event["name"] == "203.0.113.0/24"
assert event["metadata"]["collector"] == "rrc00"
assert event["metadata"]["peer_asn"] == 3333
assert event["metadata"]["origin_asn"] == 64496
assert event["metadata"]["as_path_length"] == 3
assert event["metadata"]["prefix_length"] == 24
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(
[
{
"project": "routeviews",
"collector": "route-views.sg",
"filename": "rib.20260326.0800.gz",
"startTime": "2026-03-26T08:00:00Z",
"prefix": "198.51.100.0/24",
"origin_asn": 64512,
}
]
)
assert len(transformed) == 1
record = transformed[0]
assert record["name"] == "rib.20260326.0800.gz"
assert record["metadata"]["project"] == "bgpstream"
assert record["metadata"]["broker_record"]["filename"] == "rib.20260326.0800.gz"
def test_iptoasn_transform_creates_prefix_geography_records():
collector = IPtoASNPrefixGeoCollector()
transformed = collector.transform(
[
{
"range_start": "1.0.0.0",
"range_end": "1.0.0.255",
"asn": "13335",
"country_code": "AU",
"as_name": "CLOUDFLARENET",
}
]
)
assert len(transformed) == 1
record = transformed[0]
assert record["name"] == "1.0.0.0/24"
assert record["metadata"]["family"] == "ipv4"
assert record["metadata"]["country_code"] == "AU"
assert record["metadata"]["asn"] == 13335
assert record["metadata"]["source_dataset"] == "iptoasn_combined"
def test_iptoasn_build_dataset_urls_from_combined():
urls = IPtoASNPrefixGeoCollector._build_dataset_urls(
"https://iptoasn.com/data/ip2asn-combined.tsv.gz"
)
assert urls == [
"https://iptoasn.com/data/ip2asn-v4.tsv.gz",
"https://iptoasn.com/data/ip2asn-v6.tsv.gz",
]
def test_iptoasn_build_dataset_urls_passthrough_non_combined():
url = "https://example.com/custom.tsv.gz"
urls = IPtoASNPrefixGeoCollector._build_dataset_urls(url)
assert urls == [url]
def test_opengeofeed_transform_creates_prefix_geography_records():
collector = OpenGeoFeedPrefixGeoCollector()
transformed = collector.transform(
[
{
"prefix": "203.0.113.0/24",
"country_code": "GB",
"region": "GB-LND",
"city": "London",
"postal_code": "EC1A",
"extra_columns": ["source:example"],
}
]
)
assert len(transformed) == 1
record = transformed[0]
assert record["name"] == "203.0.113.0/24"
assert record["metadata"]["family"] == "ipv4"
assert record["metadata"]["range_start"] == "203.0.113.0"
assert record["metadata"]["range_end"] == "203.0.113.255"
assert record["metadata"]["country_code"] == "GB"
assert record["metadata"]["source_dataset"] == "opengeofeed_public"
assert record["metadata"]["confidence"] == "geofeed"
def test_nro_delegated_transform_creates_prefix_geography_records():
collector = NRODelegatedPrefixGeoCollector()
transformed = collector.transform(
[
{
"rir": "ripencc",
"country_code": "DE",
"type": "ipv4",
"start": "198.51.100.0",
"value": "256",
"allocated_date": "20250401",
"status": "allocated",
}
]
)
assert len(transformed) == 1
record = transformed[0]
assert record["name"] == "198.51.100.0/24"
assert record["metadata"]["family"] == "ipv4"
assert record["metadata"]["range_start"] == "198.51.100.0"
assert record["metadata"]["range_end"] == "198.51.100.255"
assert record["metadata"]["country_code"] == "DE"
assert record["metadata"]["source_dataset"] == "nro_delegated_stats"
assert record["metadata"]["confidence"] == "registry_allocated"
def test_bgp_anomaly_to_dict():
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={"previous_origins": [64496], "current_origins": [64497]},
)
data = anomaly.to_dict()
assert data["source"] == "ris_live_bgp"
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_detect_origin_change_anomalies_creates_conflict_without_baseline():
events = [
{
"metadata": {
"prefix": "203.0.113.0/24",
"origin_asn": 64496,
"collector": "rrc00",
"collector_location": {
"country": "Netherlands",
"city": "Amsterdam",
"latitude": 52.3676,
"longitude": 4.9041,
},
}
},
{
"metadata": {
"prefix": "203.0.113.0/24",
"origin_asn": 64497,
"collector": "rrc01",
"collector_location": {
"country": "United Kingdom",
"city": "London",
"latitude": 51.5072,
"longitude": -0.1276,
},
}
},
]
anomalies = detect_origin_change_anomalies(
source="ris_live_bgp",
snapshot_id=1,
task_id=2,
events=events,
previous_origin_map={},
)
assert len(anomalies) == 2
assert {item.anomaly_type for item in anomalies} == {"origin_conflict"}
assert anomalies[0].peer_scope == ["rrc00", "rrc01"]
def test_detect_mass_withdrawal_anomalies_accepts_cross_collector_pair():
events = [
{
"metadata": {
"prefix": "203.0.113.0/24",
"origin_asn": 64496,
"event_type": "withdrawal",
"collector": "rrc00",
"peer_asn": 3333,
"collector_location": {
"country": "Netherlands",
"city": "Amsterdam",
"latitude": 52.3676,
"longitude": 4.9041,
},
}
},
{
"metadata": {
"prefix": "203.0.113.0/24",
"origin_asn": 64496,
"event_type": "withdrawal",
"collector": "rrc01",
"peer_asn": 3334,
"collector_location": {
"country": "United Kingdom",
"city": "London",
"latitude": 51.5072,
"longitude": -0.1276,
},
}
},
]
anomalies = detect_mass_withdrawal_anomalies(
source="ris_live_bgp",
snapshot_id=1,
task_id=2,
events=events,
)
assert len(anomalies) == 1
assert anomalies[0].severity == "medium"
assert anomalies[0].evidence["collector_count"] == 2
def test_detect_route_leak_anomalies_creates_candidate_for_divergent_long_paths():
events = [
{
"metadata": {
"collector": "rrc00",
"event_type": "announcement",
"prefix": "203.0.113.0/24",
"origin_asn": 64496,
"as_path": [64500, 64496],
"collector_location": {"country": "NL", "city": "Amsterdam", "latitude": 52.3, "longitude": 4.9},
"enrichment": {"prefix_scope": {"regions": [{"country": "NL"}]}},
}
},
{
"metadata": {
"collector": "rrc01",
"event_type": "announcement",
"prefix": "203.0.113.0/24",
"origin_asn": 64496,
"as_path": [64510, 64520, 64530, 64540, 64496],
"collector_location": {"country": "GB", "city": "London", "latitude": 51.5, "longitude": -0.1},
"enrichment": {"prefix_scope": {"regions": [{"country": "GB"}]}},
}
},
]
anomalies = detect_route_leak_anomalies(
source="ris_live_bgp",
snapshot_id=1,
task_id=2,
events=events,
)
assert len(anomalies) == 1
assert anomalies[0].anomaly_type == "route_leak_candidate"
assert anomalies[0].evidence["max_path_length"] == 5
def test_detect_path_flap_anomalies_creates_signal_for_repeated_state_changes():
base_timestamp = datetime(2026, 3, 27, 0, 0, tzinfo=UTC)
events = []
for index, event_type in enumerate(["announcement", "withdrawal", "announcement", "withdrawal"]):
events.append(
{
"metadata": {
"collector": "rrc00",
"event_type": event_type,
"timestamp": (base_timestamp + timedelta(minutes=index)).isoformat(),
"prefix": "198.51.100.0/24",
"origin_asn": 64512,
"as_path": [64500 + index, 64512] if event_type == "announcement" else [],
"collector_location": {"country": "NL", "city": "Amsterdam", "latitude": 52.3, "longitude": 4.9},
"enrichment": {"prefix_scope": {"regions": [{"country": "NL"}]}},
}
}
)
anomalies = detect_path_flap_anomalies(
source="ris_live_bgp",
snapshot_id=1,
task_id=2,
events=events,
)
assert len(anomalies) == 1
assert anomalies[0].anomaly_type == "path_flap"
assert anomalies[0].evidence["transitions"] == 3
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"]
def test_convert_bgp_incidents_to_geojson_adds_estimated_geography():
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_collectors=["rrc00", "rrc01"],
affected_regions=[
{
"collector": "rrc00",
"country": "Netherlands",
"city": "Amsterdam",
"latitude": 52.3676,
"longitude": 4.9041,
},
{
"collector": "rrc01",
"country": "United Kingdom",
"city": "London",
"latitude": 51.5072,
"longitude": -0.1276,
},
],
)
payload = convert_bgp_incidents_to_geojson([incident])
feature = payload["features"][0]
assert feature["properties"]["geography_mode"] == "collector_centroid"
assert feature["properties"]["estimated_radius_km"] > 0
assert feature["properties"]["estimated_center"]["latitude"] != 0
def test_convert_bgp_incidents_to_geojson_prefers_prefix_scope_hint():
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_collectors=["rrc00"],
affected_regions=[
{
"collector": "rrc00",
"country": "Netherlands",
"city": "Amsterdam",
"latitude": 52.3676,
"longitude": 4.9041,
}
],
)
payload = convert_bgp_incidents_to_geojson(
[incident],
{
incident.incident_key: {
"geography_mode": "prefix_scope",
"regions": [
{
"country": "Japan",
"city": "Tokyo",
"latitude": 35.6764,
"longitude": 139.65,
}
],
}
},
)
feature = payload["features"][0]
assert feature["properties"]["geography_mode"] == "prefix_scope"
assert feature["geometry"]["coordinates"] == [139.65, 35.6764]
def test_convert_bgp_incidents_to_geojson_prefers_prefix_geography_hint():
incident = BGPIncident(
source="ris_live_bgp",
incident_key="origin_change:198.51.100.0/24:64512",
incident_type="origin_change",
title="Origin Change incident on 198.51.100.0/24",
summary="Grouped incident summary",
severity="critical",
status="active",
confidence=0.91,
affected_prefixes=["198.51.100.0/24"],
affected_collectors=["rrc00"],
affected_regions=[
{
"collector": "rrc00",
"country": "Netherlands",
"city": "Amsterdam",
"latitude": 52.3676,
"longitude": 4.9041,
}
],
)
payload = convert_bgp_incidents_to_geojson(
[incident],
{
incident.incident_key: {
"geography_mode": "prefix_geography",
"regions": [
{
"country": "日本",
"city": None,
"latitude": 35.6764,
"longitude": 139.65,
}
],
}
},
)
feature = payload["features"][0]
assert feature["properties"]["geography_mode"] == "prefix_geography"
assert feature["geometry"]["coordinates"] == [139.65, 35.6764]
@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
iptoasn_row = {
"extra_data": {
"family": "ipv4",
"range_start": "203.0.113.0",
"range_end": "203.0.113.255",
"prefix": "203.0.113.0/24",
"country": "英国",
"country_code": "GB",
"asn": 64497,
"as_name": "Example ASN",
"source_dataset": "iptoasn_combined",
}
}
db = _FakeAsyncSession([[historical_observation], [], [iptoasn_row], [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_geography"]["country"] == "英国"
assert enrichment["prefix_geography"]["source"] == "iptoasn_combined"
assert enrichment["prefix_scope"]["countries"] == ["United Kingdom"]
assert enrichment["prefix_scope"]["cities"] == ["London"]
@pytest.mark.asyncio
async def test_enrich_bgp_events_for_batch_prefers_opengeofeed_over_iptoasn():
opengeofeed_row = {
"extra_data": {
"family": "ipv4",
"range_start": "203.0.113.0",
"range_end": "203.0.113.255",
"prefix": "203.0.113.0/24",
"country_code": "GB",
"region": "GB-LND",
"city": "London",
"source_dataset": "opengeofeed_public",
"confidence": "geofeed",
}
}
db = _FakeAsyncSession([[], [opengeofeed_row], []])
events = [
{
"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],
"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)
geography = enriched[0]["metadata"]["enrichment"]["prefix_geography"]
assert geography["source"] == "opengeofeed_public"
assert geography["confidence"] == "geofeed"
assert geography["city"] == "London"
@pytest.mark.asyncio
async def test_enrich_bgp_events_for_batch_falls_back_to_nro_delegated():
nro_row = {
"extra_data": {
"family": "ipv4",
"range_start": "198.51.100.0",
"range_end": "198.51.100.255",
"prefix": "198.51.100.0/24",
"country_code": "DE",
"rir": "ripencc",
"source_dataset": "nro_delegated_stats",
"confidence": "registry_allocated",
}
}
db = _FakeAsyncSession([[], [], [nro_row], []])
events = [
{
"metadata": {
"prefix": "198.51.100.0/24",
"origin_asn": 64512,
"collector": "rrc00",
"collector_location": {
"country": "Netherlands",
"city": "Amsterdam",
"latitude": 52.3676,
"longitude": 4.9041,
},
"as_path": [3333, 64512],
"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)
geography = enriched[0]["metadata"]["enrichment"]["prefix_geography"]
assert geography["source"] == "nro_delegated_stats"
assert geography["confidence"] == "registry_allocated"
assert geography["country"] == "德国"
@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,
}
]
},
)
with patch(
"app.services.bgp_incidents.infer_related_infrastructure",
new=AsyncMock(return_value={"related_cables": [], "related_ixps": []}),
):
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_create_bgp_incidents_for_anomalies_refreshes_existing_incident():
existing = BGPIncident(
source="ris_live_bgp",
incident_key="origin_change:203.0.113.0/24:64497",
incident_type="origin_change",
title="Old title",
summary="Old summary",
severity="medium",
status="active",
confidence=0.4,
affected_prefixes=["203.0.113.0/24"],
affected_asns=[64496, 64497],
affected_collectors=["rrc00"],
affected_regions=[
{
"collector": "rrc00",
"country": "Netherlands",
"city": "Amsterdam",
"latitude": 52.3676,
"longitude": 4.9041,
}
],
related_cables=[],
related_ixps=[],
evidence_refs=["old-key"],
)
db = _FakeAsyncSession([[existing]])
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": None,
"country": "United States",
"city": None,
"latitude": 39.8283,
"longitude": -98.5795,
}
]
},
)
with patch(
"app.services.bgp_incidents.infer_related_infrastructure",
new=AsyncMock(return_value={"related_cables": [{"landing_point": "NYC"}], "related_ixps": []}),
):
created = await create_bgp_incidents_for_anomalies(
db,
source="ris_live_bgp",
snapshot_id=1,
task_id=2,
anomalies=[anomaly],
)
assert created == 0
assert db.commits == 1
assert len(db.added) == 0
assert existing.summary != "Old summary"
assert existing.severity == "critical"
assert existing.confidence == 0.9
assert existing.affected_regions[0]["country"] == "United States"
assert existing.related_cables == [{"landing_point": "NYC"}]
assert existing.evidence_refs == ["origin_change:203.0.113.0/24:64497"]
@pytest.mark.asyncio
async def test_build_incident_geography_hints_prefers_evidence_prefix_scope_when_no_cached_prefix_geo():
incident = BGPIncident(
source="ris_live_bgp",
incident_key="origin_change:93.175.153.0/24:16509",
incident_type="origin_change",
title="Origin Change incident on 93.175.153.0/24",
summary="summary",
severity="critical",
status="active",
affected_prefixes=["93.175.153.0/24"],
evidence_refs=["origin_change:93.175.153.0/24:16509"],
)
anomaly = BGPAnomaly(
source="ris_live_bgp",
anomaly_type="origin_change",
severity="critical",
status="active",
entity_key="origin_change:93.175.153.0/24:16509",
prefix="93.175.153.0/24",
origin_asn=12654,
new_origin_asn=16509,
summary="summary",
confidence=0.8,
evidence={
"prefix_scope": {
"regions": [
{
"country": "Netherlands",
"city": "Amsterdam",
"latitude": 52.3676,
"longitude": 4.9041,
}
]
}
},
)
db = _FakeAsyncSession([[anomaly]])
hints = await build_incident_geography_hints(db, [incident])
assert hints["origin_change:93.175.153.0/24:16509"]["geography_mode"] == "prefix_scope"
assert hints["origin_change:93.175.153.0/24:16509"]["regions"][0]["country"] == "Netherlands"
@pytest.mark.asyncio
async def test_infer_related_infrastructure_links_nearby_cables():
landing = CollectedData(
source="arcgis_landing_points",
name="Amsterdam Landing",
data_type="landing_point",
extra_data={
"city_id": 10,
"country": "Netherlands",
"city": "Amsterdam",
"latitude": 52.3676,
"longitude": 4.9041,
},
)
relation = CollectedData(
source="arcgis_cable_landing_relation",
name="rel-1",
data_type="landing_relation",
extra_data={"city_id": 10, "cable_id": 20},
)
cable = CollectedData(
source="arcgis_cables",
name="AEConnect-1",
data_type="cable",
extra_data={"cable_id": 20},
)
db = _FakeAsyncSession([[landing], [relation], [cable]])
result = await infer_related_infrastructure(
db,
[
{
"collector": "rrc00",
"country": "Netherlands",
"city": "Amsterdam",
"latitude": 52.36,
"longitude": 4.90,
}
],
)
assert len(result["related_cables"]) == 1
assert result["related_cables"][0]["landing_point"] == "Amsterdam Landing"
assert result["related_cables"][0]["cable_names"] == ["AEConnect-1"]
assert result["related_ixps"][0]["name"] == "Amsterdam, Netherlands"
@pytest.mark.asyncio
async def test_build_bgp_collector_coverage_summarizes_observations():
now = datetime.now(UTC)
obs_one = BGPObservation(
source="ris_live_bgp",
collector="rrc00",
prefix="203.0.113.0/24",
origin_asn=64496,
peer_asn=3333,
event_type="announcement",
observed_at=now,
collector_geo={"city": "Amsterdam", "country": "Netherlands"},
)
obs_two = BGPObservation(
source="ris_live_bgp",
collector="rrc00",
prefix="198.51.100.0/24",
origin_asn=64497,
peer_asn=3334,
event_type="withdrawal",
observed_at=now + timedelta(minutes=5),
collector_geo={"city": "Amsterdam", "country": "Netherlands"},
)
db = _FakeAsyncSession([[obs_one, obs_two]])
coverage = await build_bgp_collector_coverage(db, source_filter=BGP_SOURCES)
first = next(item for item in coverage if item["collector"] == "rrc00")
assert first["observation_count"] == 2
assert first["recent_15m_observation_count"] == 2
assert first["recent_24h_observation_count"] == 2
assert first["recent_7d_observation_count"] == 2
assert first["prefix_count"] == 2
assert first["recent_15m_prefix_count"] == 2
assert first["recent_24h_prefix_count"] == 2
assert first["origin_asn_count"] == 2
assert first["latest_event_type"] == "withdrawal"
assert first["baseline_scope"]["countries"] == ["Netherlands"]
assert first["baseline_scope"]["cities"] == ["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",)
existing_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,
)
db = _FakeAsyncSession([
[],
[],
[],
[previous_record],
[existing_key],
[existing_anomaly],
])
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_awaited_once()
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"
@pytest.mark.asyncio
async def test_bgp_incident_summary_api_returns_aggregates():
class _SummaryResult:
def __init__(self, scalar_value=None, rows=None):
self._scalar_value = scalar_value
self._rows = rows or []
def scalar(self):
return self._scalar_value
def fetchall(self):
return self._rows
class _SummarySession:
def __init__(self):
self.calls = 0
async def execute(self, _stmt):
self.calls += 1
if self.calls == 1:
return _SummaryResult(scalar_value=2)
if self.calls == 2:
return _SummaryResult(rows=[("origin_change", 2)])
if self.calls == 3:
return _SummaryResult(rows=[("critical", 1), ("high", 1)])
return _SummaryResult(rows=[("active", 2)])
db = _SummarySession()
client = await _bgp_test_client(db)
try:
response = await client.get("/api/v1/bgp/incidents/summary")
finally:
await client.aclose()
app.dependency_overrides.clear()
assert response.status_code == 200
payload = response.json()
assert payload["total"] == 2
assert payload["by_type"]["origin_change"] == 2
assert payload["by_severity"]["critical"] == 1
assert payload["by_status"]["active"] == 2
@pytest.mark.asyncio
async def test_bgp_event_summary_api_returns_aggregates():
observation_one = BGPObservation(
id=1,
source="ris_live_bgp",
collector="rrc00",
prefix="203.0.113.0/24",
event_type="announcement",
observed_at=datetime(2026, 3, 30, 10, 0, tzinfo=UTC),
)
observation_two = BGPObservation(
id=2,
source="ris_live_bgp",
collector="rrc01",
prefix="198.51.100.0/24",
event_type="withdrawal",
observed_at=datetime(2026, 3, 30, 10, 5, tzinfo=UTC),
)
db = _FakeAsyncSession([[observation_one, observation_two]])
client = await _bgp_test_client(db)
try:
response = await client.get("/api/v1/bgp/events/summary")
finally:
await client.aclose()
app.dependency_overrides.clear()
assert response.status_code == 200
payload = response.json()
assert payload["total"] == 2
assert payload["collector_count"] == 2
assert payload["prefix_count"] == 2
assert payload["by_type"]["announcement"] == 1
assert payload["by_type"]["withdrawal"] == 1
@pytest.mark.asyncio
async def test_bgp_collectors_api_returns_coverage():
now = datetime.now(UTC)
observation = BGPObservation(
id=1,
source="ris_live_bgp",
collector="rrc00",
peer_asn=3333,
prefix="203.0.113.0/24",
event_type="announcement",
origin_asn=64496,
observed_at=now,
collector_geo={"city": "Amsterdam", "country": "Netherlands"},
)
db = _FakeAsyncSession([[observation], [observation]])
client = await _bgp_test_client(db)
try:
list_response = await client.get("/api/v1/bgp/collectors")
summary_response = await client.get("/api/v1/bgp/collectors/summary")
finally:
await client.aclose()
app.dependency_overrides.clear()
assert list_response.status_code == 200
list_payload = list_response.json()
assert list_payload["total"] >= 1
target = next(item for item in list_payload["data"] if item["collector"] == "rrc00")
assert target["observation_count"] == 1
assert target["prefix_count"] == 1
assert summary_response.status_code == 200
summary_payload = summary_response.json()
assert summary_payload["active_collectors"] >= 1
assert summary_payload["observed_prefixes"] >= 1
assert summary_payload["recent_24h_events"] >= 1
assert summary_payload["recent_7d_events"] >= 1