fix: stabilize earth bgp geography and rendering

This commit is contained in:
linkong
2026-04-02 15:36:20 +08:00
parent 07e4f519a1
commit e5fec8ba3d
26 changed files with 1788 additions and 137 deletions

View File

@@ -20,6 +20,7 @@ async def build_bgp_collector_coverage(
source_filter: tuple[str, ...] | None = None,
) -> list[dict[str, Any]]:
now = datetime.now(UTC)
recent_15m_threshold = now - timedelta(minutes=15)
recent_24h_threshold = now - timedelta(hours=24)
recent_7d_threshold = now - timedelta(days=7)
@@ -52,8 +53,10 @@ async def build_bgp_collector_coverage(
"event_types": defaultdict(int),
"countries": set(),
"cities": set(),
"recent_15m_observation_count": 0,
"recent_24h_observation_count": 0,
"recent_7d_observation_count": 0,
"recent_15m_prefixes": set(),
"recent_24h_prefixes": set(),
"recent_7d_prefixes": set(),
"latest_observed_at": None,
@@ -78,6 +81,10 @@ async def build_bgp_collector_coverage(
if observed_at.tzinfo
else observed_at.replace(tzinfo=UTC)
)
if aware_observed_at >= recent_15m_threshold:
coverage["recent_15m_observation_count"] += 1
if record.prefix:
coverage["recent_15m_prefixes"].add(record.prefix)
if aware_observed_at >= recent_24h_threshold:
coverage["recent_24h_observation_count"] += 1
if record.prefix:
@@ -116,8 +123,10 @@ async def build_bgp_collector_coverage(
"event_types": defaultdict(int),
"countries": {location.get("country")} if location.get("country") else set(),
"cities": {location.get("city")} if location.get("city") else set(),
"recent_15m_observation_count": 0,
"recent_24h_observation_count": 0,
"recent_7d_observation_count": 0,
"recent_15m_prefixes": set(),
"recent_24h_prefixes": set(),
"recent_7d_prefixes": set(),
"latest_observed_at": None,
@@ -142,8 +151,10 @@ async def build_bgp_collector_coverage(
"prefix_count": len(item["prefixes"]),
"origin_asn_count": len(item["origin_asns"]),
"peer_asn_count": len(item["peer_asns"]),
"recent_15m_observation_count": item["recent_15m_observation_count"],
"recent_24h_observation_count": item["recent_24h_observation_count"],
"recent_7d_observation_count": item["recent_7d_observation_count"],
"recent_15m_prefix_count": len(item["recent_15m_prefixes"]),
"recent_24h_prefix_count": len(item["recent_24h_prefixes"]),
"recent_7d_prefix_count": len(item["recent_7d_prefixes"]),
"top_event_types": [

View File

@@ -55,6 +55,11 @@ def _unique_peers(events: list[dict[str, Any]]) -> list[int]:
return sorted(peers)
def _path_signature(metadata: dict[str, Any]) -> tuple[int, ...]:
path = metadata.get("as_path") or []
return tuple(int(asn) for asn in path if asn is not None)
def detect_origin_change_anomalies(
*,
source: str,
@@ -100,6 +105,7 @@ def detect_origin_change_anomalies(
)
sample_metadata = sample_event.get("metadata") or {}
sample_enrichment = sample_metadata.get("enrichment") or {}
sample_prefix_geography = sample_enrichment.get("prefix_geography") or {}
anomaly_type = "origin_change"
severity = "critical"
confidence = 0.86
@@ -140,8 +146,10 @@ def detect_origin_change_anomalies(
"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_geography": sample_prefix_geography,
"prefix_scope": sample_enrichment.get("prefix_scope"),
"impacted_regions": related_regions
"impacted_regions": sample_prefix_geography.get("regions")
or related_regions
or sample_enrichment.get("prefix_scope", {}).get("regions", []),
},
)
@@ -180,6 +188,7 @@ def detect_more_specific_burst_anomalies(
sample = more_specifics[0].get("metadata") or {}
sample_enrichment = sample.get("enrichment") or {}
sample_prefix_geography = sample_enrichment.get("prefix_geography") or {}
event_count = len(more_specifics)
anomalies.append(
BGPAnomaly(
@@ -205,8 +214,10 @@ def detect_more_specific_burst_anomalies(
"unique_prefixes": unique_prefixes,
"rpki_validation": sample_enrichment.get("rpki_validation"),
"origin_asn_profile": sample_enrichment.get("origin_asn_profile"),
"prefix_geography": sample_prefix_geography,
"prefix_scope": sample_enrichment.get("prefix_scope"),
"impacted_regions": _iter_event_regions(more_specifics)
"impacted_regions": sample_prefix_geography.get("regions")
or _iter_event_regions(more_specifics)
or sample_enrichment.get("prefix_scope", {}).get("regions", []),
},
)
@@ -242,6 +253,7 @@ def detect_mass_withdrawal_anomalies(
sample_event = related_events[0] if related_events else {}
sample_metadata = sample_event.get("metadata") or {}
sample_enrichment = sample_metadata.get("enrichment") or {}
sample_prefix_geography = sample_enrichment.get("prefix_geography") or {}
severity = "medium"
if count >= 4 or len(related_collectors) >= 3:
severity = "high"
@@ -277,8 +289,175 @@ def detect_mass_withdrawal_anomalies(
],
"origin_asn_profile": sample_enrichment.get("origin_asn_profile"),
"rpki_validation": sample_enrichment.get("rpki_validation"),
"prefix_geography": sample_prefix_geography,
"prefix_scope": sample_enrichment.get("prefix_scope"),
"impacted_regions": _iter_event_regions(related_events)
"impacted_regions": sample_prefix_geography.get("regions")
or _iter_event_regions(related_events)
or sample_enrichment.get("prefix_scope", {}).get("regions", []),
},
)
)
return anomalies
def detect_route_leak_anomalies(
*,
source: str,
snapshot_id: int | None,
task_id: int | None,
events: list[dict[str, Any]],
) -> list[BGPAnomaly]:
events_by_prefix: defaultdict[str, list[dict[str, Any]]] = defaultdict(list)
for event in events:
metadata = event.get("metadata") or {}
prefix = metadata.get("prefix")
if prefix and metadata.get("event_type") == "announcement":
events_by_prefix[str(prefix)].append(event)
anomalies: list[BGPAnomaly] = []
for prefix, related_events in events_by_prefix.items():
related_collectors = _unique_collectors(related_events)
if len(related_collectors) < 2:
continue
path_signatures = Counter()
max_path_length = 0
for event in related_events:
metadata = event.get("metadata") or {}
signature = _path_signature(metadata)
if signature:
path_signatures[signature] += 1
max_path_length = max(max_path_length, len(signature))
if len(path_signatures) < 2:
continue
dominant_length = len(path_signatures.most_common(1)[0][0])
if max_path_length < max(dominant_length + 2, 5):
continue
sample_event = max(
related_events,
key=lambda event: len(_path_signature((event.get("metadata") or {}))),
)
sample_metadata = sample_event.get("metadata") or {}
sample_enrichment = sample_metadata.get("enrichment") or {}
sample_prefix_geography = sample_enrichment.get("prefix_geography") or {}
peer_scope = related_collectors
path_lengths = sorted({len(signature) for signature in path_signatures if signature})
anomalies.append(
BGPAnomaly(
snapshot_id=snapshot_id,
task_id=task_id,
source=source,
anomaly_type="route_leak_candidate",
severity="high" if max_path_length >= dominant_length + 3 else "medium",
status="active",
entity_key=f"route_leak_candidate:{prefix}:{max_path_length}:{len(related_collectors)}",
prefix=prefix,
origin_asn=sample_metadata.get("origin_asn"),
new_origin_asn=None,
peer_scope=peer_scope,
started_at=datetime.now(UTC),
confidence=min(0.58 + (0.05 * min(len(related_collectors), 4)) + (0.03 * min(max_path_length - dominant_length, 4)), 0.88),
summary=(
f"Prefix {prefix} shows divergent long AS paths across "
f"{len(related_collectors)} collectors, suggesting a possible route leak."
),
evidence={
"path_lengths": path_lengths,
"dominant_path_length": dominant_length,
"max_path_length": max_path_length,
"path_signatures": [
{"path": list(signature), "count": count}
for signature, count in path_signatures.most_common(5)
],
"events": [(item.get("metadata") or {}) for item in related_events[:10]],
"origin_asn_profile": sample_enrichment.get("origin_asn_profile"),
"rpki_validation": sample_enrichment.get("rpki_validation"),
"prefix_geography": sample_prefix_geography,
"prefix_scope": sample_enrichment.get("prefix_scope"),
"impacted_regions": sample_prefix_geography.get("regions")
or _iter_event_regions(related_events)
or sample_enrichment.get("prefix_scope", {}).get("regions", []),
},
)
)
return anomalies
def detect_path_flap_anomalies(
*,
source: str,
snapshot_id: int | None,
task_id: int | None,
events: list[dict[str, Any]],
) -> list[BGPAnomaly]:
events_by_prefix: defaultdict[str, list[dict[str, Any]]] = defaultdict(list)
for event in events:
metadata = event.get("metadata") or {}
prefix = metadata.get("prefix")
if prefix:
events_by_prefix[str(prefix)].append(event)
anomalies: list[BGPAnomaly] = []
for prefix, related_events in events_by_prefix.items():
ordered = sorted(
related_events,
key=lambda event: str((event.get("metadata") or {}).get("timestamp") or ""),
)
event_types = [str((item.get("metadata") or {}).get("event_type") or "") for item in ordered]
transitions = sum(1 for index in range(1, len(event_types)) if event_types[index] != event_types[index - 1])
distinct_paths = {
_path_signature(item.get("metadata") or {})
for item in ordered
if _path_signature(item.get("metadata") or {})
}
related_collectors = _unique_collectors(ordered)
if transitions < 3 and len(distinct_paths) < 3:
continue
sample_metadata = (ordered[0].get("metadata") or {}) if ordered else {}
sample_enrichment = sample_metadata.get("enrichment") or {}
sample_prefix_geography = sample_enrichment.get("prefix_geography") or {}
severity = "medium"
if transitions >= 5 or len(distinct_paths) >= 4:
severity = "high"
anomalies.append(
BGPAnomaly(
snapshot_id=snapshot_id,
task_id=task_id,
source=source,
anomaly_type="path_flap",
severity=severity,
status="active",
entity_key=f"path_flap:{prefix}:{transitions}:{len(distinct_paths)}",
prefix=prefix,
origin_asn=sample_metadata.get("origin_asn"),
new_origin_asn=None,
peer_scope=related_collectors,
started_at=datetime.now(UTC),
confidence=min(0.54 + (0.05 * min(transitions, 5)) + (0.03 * min(len(distinct_paths), 4)), 0.9),
summary=(
f"Prefix {prefix} shows repeated state/path changes "
f"({transitions} transitions, {len(distinct_paths)} distinct paths) in the current window."
),
evidence={
"transitions": transitions,
"event_types": event_types[:12],
"distinct_paths": [list(path) for path in list(distinct_paths)[:6]],
"events": [(item.get("metadata") or {}) for item in ordered[:10]],
"origin_asn_profile": sample_enrichment.get("origin_asn_profile"),
"rpki_validation": sample_enrichment.get("rpki_validation"),
"prefix_geography": sample_prefix_geography,
"prefix_scope": sample_enrichment.get("prefix_scope"),
"impacted_regions": sample_prefix_geography.get("regions")
or _iter_event_regions(ordered)
or sample_enrichment.get("prefix_scope", {}).get("regions", []),
},
)

View File

@@ -7,9 +7,10 @@ from collections import defaultdict
from datetime import UTC, datetime
from typing import Any
from sqlalchemy import select
from sqlalchemy import select, text
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.countries import get_country_centroid, normalize_country
from app.models.bgp_observation import BGPObservation
from app.models.collected_data import CollectedData
@@ -96,6 +97,83 @@ def extract_bgp_network_fields(prefix: str) -> dict[str, Any]:
}
async def _lookup_prefix_geography(
db: AsyncSession,
prefix_values: list[str],
) -> dict[str, dict[str, Any]]:
results: dict[str, dict[str, Any]] = {}
for prefix in prefix_values:
try:
network = ipaddress.ip_network(prefix, strict=False)
except ValueError:
continue
family = f"ipv{network.version}"
range_start = str(network.network_address)
range_end = str(network.broadcast_address)
result = await db.execute(
text(
"""
SELECT metadata
FROM collected_data
WHERE source = 'iptoasn_prefix_geo'
AND COALESCE(is_current, TRUE) = TRUE
AND metadata->>'family' = :family
AND CAST(metadata->>'range_start' AS inet) <= CAST(:range_start AS inet)
AND CAST(metadata->>'range_end' AS inet) >= CAST(:range_end AS inet)
ORDER BY id DESC
LIMIT 1
"""
),
{
"family": family,
"range_start": range_start,
"range_end": range_end,
},
)
row = result.fetchone()
if not row:
continue
if isinstance(row, dict):
payload = row.get("metadata") or row.get("extra_data")
elif hasattr(row, "_mapping"):
payload = row._mapping.get("metadata") or row._mapping.get("extra_data")
else:
payload = row[0]
if not isinstance(payload, dict):
continue
country = normalize_country(payload.get("country") or payload.get("country_code"))
prefix_hint = payload.get("prefix") or prefix
asn = _safe_int(payload.get("asn"))
as_name = payload.get("as_name")
centroid = get_country_centroid(country)
regions = []
if country:
regions.append(
{
"country": country,
"city": None,
"latitude": centroid.get("latitude") if centroid else None,
"longitude": centroid.get("longitude") if centroid else None,
}
)
results[prefix] = {
"prefix": prefix_hint,
"country": country,
"asn": asn,
"as_name": as_name,
"source": payload.get("source_dataset") or "iptoasn_combined",
"confidence": "country_range",
"regions": regions,
}
return results
async def enrich_bgp_events_for_batch(
db: AsyncSession,
*,
@@ -165,6 +243,7 @@ async def enrich_bgp_events_for_batch(
}
asn_profiles: dict[int, dict[str, Any]] = {}
prefix_geographies = await _lookup_prefix_geography(db, prefix_values) if prefix_values else {}
if origin_asns:
peeringdb_result = await db.execute(
select(CollectedData).where(CollectedData.source == "peeringdb_network")
@@ -207,21 +286,12 @@ async def enrich_bgp_events_for_batch(
collector = str(metadata.get("collector") or "").strip()
collector_location = metadata.get("collector_location") or {}
baseline = historical_prefix_baseline.get(prefix, {})
prefix_geography = prefix_geographies.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])
prefix_scope_regions = _compact_locations([*baseline_regions])
enrichment = {
**extract_bgp_network_fields(prefix),
@@ -248,6 +318,7 @@ async def enrich_bgp_events_for_batch(
},
"origin_asn_profile": asn_profiles.get(origin_asn),
"new_origin_asn_profile": asn_profiles.get(new_origin_asn),
"prefix_geography": prefix_geography,
"prefix_scope": {
"countries": sorted(
{

View File

@@ -209,15 +209,14 @@ async def create_bgp_incidents_for_anomalies(
grouped.setdefault(incident_key, []).append(anomaly)
existing_result = await db.execute(
select(BGPIncident.incident_key).where(BGPIncident.incident_key.in_(sorted(grouped.keys())))
select(BGPIncident).where(BGPIncident.incident_key.in_(sorted(grouped.keys())))
)
existing_keys = {row[0] for row in existing_result.fetchall()}
existing_incidents = {
incident.incident_key: incident for incident in existing_result.scalars().all()
}
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})
@@ -270,6 +269,28 @@ async def create_bgp_incidents_for_anomalies(
)
related_infrastructure = await infer_related_infrastructure(db, regions)
existing = existing_incidents.get(incident_key)
if existing is not None:
existing.snapshot_id = snapshot_id
existing.task_id = task_id
existing.source = source
existing.incident_type = primary.anomaly_type
existing.title = title
existing.summary = summary
existing.severity = severity
existing.status = "active"
existing.confidence = confidence
existing.started_at = primary.started_at or existing.started_at or datetime.now(UTC)
existing.ended_at = None
existing.affected_prefixes = prefixes
existing.affected_asns = asns
existing.affected_collectors = collectors
existing.affected_regions = regions
existing.related_cables = related_infrastructure["related_cables"]
existing.related_ixps = related_infrastructure["related_ixps"]
existing.evidence_refs = evidence_refs
continue
db.add(
BGPIncident(
snapshot_id=snapshot_id,
@@ -294,7 +315,7 @@ async def create_bgp_incidents_for_anomalies(
)
created += 1
if created:
if created or existing_incidents:
await db.commit()
return created

View File

@@ -32,6 +32,7 @@ from app.services.collectors.spacetrack import SpaceTrackTLECollector
from app.services.collectors.celestrak import CelesTrakTLECollector
from app.services.collectors.ris_live import RISLiveCollector
from app.services.collectors.bgpstream import BGPStreamBackfillCollector
from app.services.collectors.iptoasn import IPtoASNPrefixGeoCollector
collector_registry.register(TOP500Collector())
collector_registry.register(EpochAIGPUCollector())
@@ -55,3 +56,4 @@ collector_registry.register(SpaceTrackTLECollector())
collector_registry.register(CelesTrakTLECollector())
collector_registry.register(RISLiveCollector())
collector_registry.register(BGPStreamBackfillCollector())
collector_registry.register(IPtoASNPrefixGeoCollector())

View File

@@ -18,6 +18,8 @@ from app.services.bgp_detectors import (
detect_mass_withdrawal_anomalies,
detect_more_specific_burst_anomalies,
detect_origin_change_anomalies,
detect_path_flap_anomalies,
detect_route_leak_anomalies,
)
from app.services.bgp_enrichment import enrich_bgp_events_for_batch, extract_bgp_network_fields
@@ -282,6 +284,18 @@ async def create_bgp_anomalies_for_batch(
task_id=task_id,
events=enriched_events,
),
*detect_route_leak_anomalies(
source=source,
snapshot_id=snapshot_id,
task_id=task_id,
events=enriched_events,
),
*detect_path_flap_anomalies(
source=source,
snapshot_id=snapshot_id,
task_id=task_id,
events=enriched_events,
),
]
if not pending_anomalies:
@@ -302,16 +316,29 @@ async def create_bgp_anomalies_for_batch(
created = 0
created_anomalies: list[BGPAnomaly] = []
refreshed_anomalies: list[BGPAnomaly] = []
existing_map = {item.entity_key: item for item in existing_anomalies if item.entity_key}
for anomaly in pending_anomalies:
if anomaly.entity_key in existing_keys:
existing = existing_map.get(anomaly.entity_key)
if existing is not None:
existing.severity = anomaly.severity
existing.status = anomaly.status
existing.summary = anomaly.summary
existing.confidence = anomaly.confidence
existing.peer_scope = anomaly.peer_scope
existing.evidence = anomaly.evidence
existing.new_origin_asn = anomaly.new_origin_asn
existing.origin_asn = anomaly.origin_asn
refreshed_anomalies.append(existing)
continue
db.add(anomaly)
created_anomalies.append(anomaly)
created += 1
if created:
if created or refreshed_anomalies:
await db.commit()
incident_seed_anomalies = [*created_anomalies, *existing_anomalies]
incident_seed_anomalies = [*created_anomalies, *refreshed_anomalies]
if incident_seed_anomalies:
await create_bgp_incidents_for_anomalies(
db,

View File

@@ -0,0 +1,111 @@
"""IPtoASN prefix geography collector.
Downloads the public combined IPv4+IPv6 TSV database and stores coarse
prefix-to-country/ASN geography hints for BGP enrichment.
"""
from __future__ import annotations
import gzip
from datetime import UTC, datetime
from ipaddress import summarize_address_range, ip_address
from typing import Any
import httpx
from app.services.collectors.base import BaseCollector
class IPtoASNPrefixGeoCollector(BaseCollector):
name = "iptoasn_prefix_geo"
priority = "P1"
module = "L3"
frequency_hours = 24
data_type = "prefix_geography"
fail_on_empty = True
async def fetch(self) -> list[dict[str, Any]]:
if not self._resolved_url:
raise RuntimeError("IPtoASN combined URL is not configured")
async with httpx.AsyncClient(timeout=180.0, follow_redirects=True) as client:
response = await client.get(
self._resolved_url,
headers={
"User-Agent": "Planet-Intelligence-System/1.0 (Python/collector)",
"Accept": "application/gzip,application/octet-stream,*/*",
},
)
response.raise_for_status()
body = gzip.decompress(response.content).decode("utf-8", errors="replace")
rows: list[dict[str, Any]] = []
for raw_line in body.splitlines():
line = raw_line.strip()
if not line or line.startswith("#"):
continue
parts = line.split("\t")
if len(parts) < 5:
continue
range_start, range_end, asn, country_code, as_name = parts[:5]
rows.append(
{
"range_start": range_start,
"range_end": range_end,
"asn": asn,
"country_code": country_code,
"as_name": as_name,
}
)
return rows
def transform(self, raw_data: list[dict[str, Any]]) -> list[dict[str, Any]]:
reference_date = datetime.now(UTC).isoformat()
transformed: list[dict[str, Any]] = []
for item in raw_data:
try:
start_ip = ip_address(str(item["range_start"]))
end_ip = ip_address(str(item["range_end"]))
except ValueError:
continue
if start_ip.version != end_ip.version:
continue
summarized = list(summarize_address_range(start_ip, end_ip))
primary_prefix = str(summarized[0]) if summarized else f"{start_ip}/{32 if start_ip.version == 4 else 128}"
family = f"ipv{start_ip.version}"
asn_value = item.get("asn")
try:
normalized_asn = int(str(asn_value))
except (TypeError, ValueError):
normalized_asn = None
transformed.append(
{
"source_id": f"{family}:{item['range_start']}-{item['range_end']}",
"name": primary_prefix,
"title": f"{primary_prefix} {item.get('country_code', '').strip()}".strip(),
"country": item.get("country_code"),
"city": "",
"latitude": None,
"longitude": None,
"metadata": {
"family": family,
"range_start": item["range_start"],
"range_end": item["range_end"],
"prefix": primary_prefix,
"prefixes": [str(prefix) for prefix in summarized[:8]],
"range_prefix_count": len(summarized),
"country_code": item.get("country_code"),
"asn": normalized_asn,
"as_name": item.get("as_name"),
"source_dataset": "iptoasn_combined",
},
"reference_date": reference_date,
}
)
return transformed