Files
planet/backend/app/core/websocket/ue_scene.py
2026-05-09 15:51:46 +08:00

563 lines
19 KiB
Python

"""UE scene state built from visualization aggregate output."""
import asyncio
import json
import zlib
from datetime import UTC, datetime
from typing import TYPE_CHECKING, Any, Dict, List
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from app.core.time import to_iso8601_utc
from app.models.bgp_anomaly import BGPAnomaly
from app.models.bgp_incident import BGPIncident
if TYPE_CHECKING:
from app.models.collected_data import CollectedData
DEFAULT_LAYER_ORDER = (
"satellites",
"supercomputers",
"gpu_clusters",
"submarine_cables",
"landing_points",
"bgp_anomalies",
"bgp_incidents",
"bgp_collectors",
"alerts",
)
FULL_RESYNC_REASONS = {"initial_connect", "manual_resync", "profile_changed"}
UE_SCENE_LAYER_CHUNK_ITEM_LIMITS = {
"satellites": 500,
"supercomputers": 200,
"gpu_clusters": 100,
"submarine_cables": 25,
"landing_points": 250,
"bgp_anomalies": 100,
"bgp_incidents": 100,
"bgp_collectors": 100,
"alerts": 100,
}
def _default_display_profile() -> Dict[str, Any]:
return {
"profile_id": "polarized-wall-a",
"stereo_mode": "polarized",
"screen_width_m": 3.0,
"screen_height_m": 2.0,
"target_refresh_hz": 120,
}
def _default_camera_state() -> Dict[str, Any]:
return {
"mode": "auto_cruise",
"path_id": "global_overview",
"fov": 42.0,
}
def _stable_json(value: Any) -> str:
return json.dumps(value, ensure_ascii=True, sort_keys=True, separators=(",", ":"))
def _entity_type_for_layer(layer_name: str) -> str:
return {
"satellites": "satellite",
"supercomputers": "supercomputer",
"gpu_clusters": "gpu_cluster",
"submarine_cables": "submarine_cable",
"landing_points": "landing_point",
"bgp_anomalies": "bgp_anomaly",
"bgp_incidents": "bgp_incident",
"bgp_collectors": "bgp_collector",
"alerts": "alert",
}[layer_name]
def _default_visual_for_layer(layer_name: str) -> Dict[str, Any]:
visuals = {
"satellites": {"style": "satellite_marker", "size": 0.7, "color": "#9BDBFF"},
"supercomputers": {"style": "supercomputer_marker", "size": 1.2, "color": "#FF6B6B"},
"gpu_clusters": {"style": "pulse_marker", "size": 1.0, "color": "#FF8C42"},
"submarine_cables": {"style": "cable_arc", "width": 2.0, "color": "#4ECDC4"},
"landing_points": {"style": "landing_point_marker", "size": 0.8, "color": "#45B7D1"},
"bgp_anomalies": {"style": "anomaly_marker", "size": 1.0, "color": "#FFB703"},
"bgp_incidents": {"style": "incident_marker", "size": 1.2, "color": "#E63946"},
"bgp_collectors": {"style": "collector_marker", "size": 0.9, "color": "#7B9ACC"},
"alerts": {"style": "alert_marker", "size": 1.0, "color": "#FFD166"},
}
return visuals[layer_name]
def _empty_scene_layers() -> Dict[str, Dict[str, Any]]:
return {
layer_name: {"revision": 0, "items": {}}
for layer_name in DEFAULT_LAYER_ORDER
}
def _empty_changes() -> Dict[str, Dict[str, List[Any]]]:
return {
layer_name: {"added": [], "updated": [], "removed": []}
for layer_name in DEFAULT_LAYER_ORDER
}
def _point_geo(coordinates: List[Any]) -> Dict[str, Any]:
lng = coordinates[0] if len(coordinates) > 0 else None
lat = coordinates[1] if len(coordinates) > 1 else None
alt = coordinates[2] if len(coordinates) > 2 else 0.0
return {"lat": lat, "lng": lng, "alt": alt}
def _path_geo(geometry_type: str, coordinates: Any) -> Dict[str, Any]:
if geometry_type == "LineString":
return {
"path": [
{"lat": point[1], "lng": point[0], "alt": point[2] if len(point) > 2 else 0.0}
for point in coordinates
if isinstance(point, list) and len(point) >= 2
]
}
if geometry_type == "MultiLineString":
return {
"segments": [
[
{"lat": point[1], "lng": point[0], "alt": point[2] if len(point) > 2 else 0.0}
for point in line
if isinstance(point, list) and len(point) >= 2
]
for line in coordinates
if isinstance(line, list)
]
}
return {}
def _item_identifier(layer_name: str, feature: Dict[str, Any], index: int) -> str:
properties = feature.get("properties") or {}
feature_id = feature.get("id") or properties.get("id") or properties.get("source_id")
return f"{layer_name}:{feature_id or index}"
def _item_revision(item: Dict[str, Any]) -> int:
return zlib.crc32(_stable_json(item).encode("utf-8")) & 0xFFFFFFFF
def _layer_revision(items: Dict[str, Dict[str, Any]]) -> int:
if not items:
return 0
joined = "|".join(
f"{item_id}:{items[item_id]['revision']}"
for item_id in sorted(items)
)
return zlib.crc32(joined.encode("utf-8")) & 0xFFFFFFFF
def _serialize_feature(layer_name: str, feature: Dict[str, Any], index: int) -> Dict[str, Any]:
properties = dict(feature.get("properties") or {})
geometry = feature.get("geometry") or {}
geometry_type = geometry.get("type", "")
coordinates = geometry.get("coordinates") or []
item_id = _item_identifier(layer_name, feature, index)
geo: Dict[str, Any] = {}
if geometry_type == "Point":
geo = _point_geo(coordinates)
elif geometry_type in {"LineString", "MultiLineString"}:
geo = _path_geo(geometry_type, coordinates)
title = (
properties.get("name")
or properties.get("Name")
or properties.get("title")
or item_id
)
subtitle_parts = [
properties.get("city"),
properties.get("country"),
properties.get("region"),
]
item = {
"id": item_id,
"entity_type": _entity_type_for_layer(layer_name),
"geo": geo,
"visual": {
**_default_visual_for_layer(layer_name),
**({"color": properties["color"]} if properties.get("color") else {}),
},
"metrics": properties,
"labels": {
"title": title,
"subtitle": ", ".join([part for part in subtitle_parts if part]),
},
"status": {
"health": properties.get("status", "normal"),
"alert_level": properties.get("severity", "none"),
},
}
item["revision"] = _item_revision(item)
return item
async def build_visualization_scene_state(db: AsyncSession) -> Dict[str, Any]:
from app.api.v1.visualization import (
_build_landing_point_cable_maps,
_filter_known_records,
_load_current_collected_data_by_sources,
build_anomaly_geography_hints,
build_incident_geography_hints,
convert_bgp_anomalies_to_geojson,
convert_bgp_collectors_to_geojson,
convert_bgp_incidents_to_geojson,
convert_cable_to_geojson,
convert_gpu_cluster_to_geojson,
convert_landing_point_to_geojson,
convert_satellite_to_geojson,
convert_supercomputer_to_geojson,
)
from app.services.bgp_collectors import build_bgp_collector_coverage
records_by_source = await _load_current_collected_data_by_sources(
db,
[
"arcgis_cables",
"arcgis_landing_points",
"arcgis_cable_landing_relation",
"celestrak_tle",
"top500",
"epoch_ai_gpu",
],
)
cables_records = records_by_source.get("arcgis_cables", [])
landing_point_records = records_by_source.get("arcgis_landing_points", [])
relation_records = records_by_source.get("arcgis_cable_landing_relation", [])
satellites_records = _filter_known_records(records_by_source.get("celestrak_tle", []))
supercomputer_records = _filter_known_records(records_by_source.get("top500", []))
gpu_records = _filter_known_records(records_by_source.get("epoch_ai_gpu", []))
bgp_anomalies_result = await db.execute(
select(BGPAnomaly)
.where(BGPAnomaly.status == "active")
.order_by(BGPAnomaly.created_at.desc())
.limit(200)
)
bgp_anomalies = list(bgp_anomalies_result.scalars().all())
bgp_anomaly_geography_hints = await build_anomaly_geography_hints(db, bgp_anomalies)
bgp_incidents_result = await db.execute(
select(BGPIncident)
.where(BGPIncident.status == "active")
.order_by(BGPIncident.created_at.desc())
.limit(100)
)
bgp_incidents = list(bgp_incidents_result.scalars().all())
bgp_incident_geography_hints = await build_incident_geography_hints(db, bgp_incidents)
bgp_collector_coverage = await build_bgp_collector_coverage(
db,
source_filter=("ris_live_bgp", "bgpstream_bgp"),
)
bgp_coverage_by_collector = {
item["collector"]: item
for item in bgp_collector_coverage
if item.get("collector")
}
city_to_cable_ids_map, cable_id_to_name_map = _build_landing_point_cable_maps(
relation_records,
cables_records,
)
aggregate = {
"satellites": convert_satellite_to_geojson(satellites_records),
"supercomputers": convert_supercomputer_to_geojson(supercomputer_records),
"gpu_clusters": convert_gpu_cluster_to_geojson(gpu_records),
"submarine_cables": convert_cable_to_geojson(cables_records),
"landing_points": convert_landing_point_to_geojson(
landing_point_records,
city_to_cable_ids_map,
cable_id_to_name_map,
),
"bgp_anomalies": convert_bgp_anomalies_to_geojson(
bgp_anomalies,
bgp_anomaly_geography_hints,
),
"bgp_incidents": convert_bgp_incidents_to_geojson(
bgp_incidents,
bgp_incident_geography_hints,
),
"bgp_collectors": convert_bgp_collectors_to_geojson(bgp_coverage_by_collector),
"alerts": {"type": "FeatureCollection", "features": []},
}
layers = _empty_scene_layers()
total_records = 0
for layer_name in DEFAULT_LAYER_ORDER:
feature_collection = aggregate.get(layer_name) or {"features": []}
items = {
item["id"]: item
for index, feature in enumerate(feature_collection.get("features", []), start=1)
for item in [_serialize_feature(layer_name, feature, index)]
}
layers[layer_name]["items"] = items
layers[layer_name]["revision"] = _layer_revision(items)
total_records += len(items)
timestamp = to_iso8601_utc(datetime.now(UTC))
state_hash = zlib.crc32(
_stable_json(
{
layer_name: {
item_id: item["revision"]
for item_id, item in layers[layer_name]["items"].items()
}
for layer_name in DEFAULT_LAYER_ORDER
}
).encode("utf-8")
) & 0xFFFFFFFF
return {
"generated_at": timestamp,
"state_hash": state_hash,
"total_records": total_records,
"layers": layers,
}
def _changes_exist(changes: Dict[str, Dict[str, List[Any]]]) -> bool:
return any(
layer_changes["added"] or layer_changes["updated"] or layer_changes["removed"]
for layer_changes in changes.values()
)
def build_incremental_changes(
previous_state: Dict[str, Any],
current_state: Dict[str, Any],
) -> Dict[str, Dict[str, List[Any]]]:
changes = _empty_changes()
for layer_name in DEFAULT_LAYER_ORDER:
previous_items = previous_state["layers"][layer_name]["items"]
current_items = current_state["layers"][layer_name]["items"]
previous_ids = set(previous_items)
current_ids = set(current_items)
for added_id in sorted(current_ids - previous_ids):
changes[layer_name]["added"].append(current_items[added_id])
for removed_id in sorted(previous_ids - current_ids):
changes[layer_name]["removed"].append(removed_id)
for common_id in sorted(previous_ids & current_ids):
if _stable_json(previous_items[common_id]) != _stable_json(current_items[common_id]):
changes[layer_name]["updated"].append(current_items[common_id])
return changes
def _full_payload_from_state(state: Dict[str, Any], sequence: int) -> Dict[str, Any]:
return {
"update_type": "full",
"sequence": sequence,
"cluster_time": state["generated_at"],
"display_profile": _default_display_profile(),
"camera_state": _default_camera_state(),
"payload": {
"meta": {
"generated_at": state["generated_at"],
"total_records": state["total_records"],
"state_hash": state["state_hash"],
},
"layers": {
layer_name: {
"revision": layer_data["revision"],
"items": [
layer_data["items"][item_id]
for item_id in sorted(layer_data["items"])
],
}
for layer_name, layer_data in state["layers"].items()
},
},
}
def _incremental_payload(
*,
current_state: Dict[str, Any],
sequence: int,
base_sequence: int,
changes: Dict[str, Dict[str, List[Any]]],
) -> Dict[str, Any]:
return {
"update_type": "incremental",
"sequence": sequence,
"base_sequence": base_sequence,
"cluster_time": current_state["generated_at"],
"changes": changes,
"meta": {
"generated_at": current_state["generated_at"],
"total_records": current_state["total_records"],
"state_hash": current_state["state_hash"],
},
}
def _noop_incremental_payload(state: Dict[str, Any], sequence: int) -> Dict[str, Any]:
return _incremental_payload(
current_state=state,
sequence=sequence,
base_sequence=sequence,
changes=_empty_changes(),
)
def expand_scene_payload_for_transport(scene_payload: Dict[str, Any]) -> List[Dict[str, Any]]:
"""Split oversized scene payloads into smaller transport-safe chunks."""
update_type = scene_payload.get("update_type")
if update_type != "full":
return [scene_payload]
payload = scene_payload.get("payload") or {}
layers = payload.get("layers") or {}
if not layers:
return [scene_payload]
chunks: List[Dict[str, Any]] = []
for layer_name in DEFAULT_LAYER_ORDER:
layer_data = layers.get(layer_name) or {}
items = list(layer_data.get("items") or [])
if not items:
continue
chunk_size = UE_SCENE_LAYER_CHUNK_ITEM_LIMITS.get(layer_name, 100)
total_layer_chunks = max(1, (len(items) + chunk_size - 1) // chunk_size)
for chunk_index, offset in enumerate(range(0, len(items), chunk_size), start=1):
chunk_items = items[offset : offset + chunk_size]
chunks.append(
{
**scene_payload,
"payload": {
"meta": {
**(payload.get("meta") or {}),
"chunked": True,
"layer_name": layer_name,
"layer_chunk_index": chunk_index,
"layer_chunk_count": total_layer_chunks,
},
"layers": {
layer_name: {
"revision": layer_data.get("revision", 0),
"items": chunk_items,
}
},
},
}
)
if not chunks:
return [scene_payload]
total_chunks = len(chunks)
for transport_index, chunk in enumerate(chunks, start=1):
chunk_payload = chunk.setdefault("payload", {})
chunk_meta = chunk_payload.setdefault("meta", {})
chunk_meta["transport_chunk_index"] = transport_index
chunk_meta["transport_chunk_count"] = total_chunks
return chunks
class UeSceneStateStore:
"""Maintain cached ue_scene state and a bounded incremental replay history."""
def __init__(self, history_limit: int = 20) -> None:
self.history_limit = history_limit
self.sequence = 0
self.state: Dict[str, Any] | None = None
self.history: List[Dict[str, Any]] = []
self.lock = asyncio.Lock()
async def _refresh_locked(self, db: AsyncSession) -> List[Dict[str, Any]]:
current_state = await build_visualization_scene_state(db)
if self.state is None:
self.sequence = 1
self.state = current_state
return [_full_payload_from_state(self.state, self.sequence)]
if current_state["state_hash"] == self.state["state_hash"]:
self.state = current_state
return []
previous_sequence = self.sequence
previous_state = self.state
self.sequence += 1
self.state = current_state
changes = build_incremental_changes(previous_state, current_state)
incremental = _incremental_payload(
current_state=current_state,
sequence=self.sequence,
base_sequence=previous_sequence,
changes=changes,
)
self.history.append(incremental)
if len(self.history) > self.history_limit:
self.history = self.history[-self.history_limit :]
return [incremental]
async def get_broadcast_payloads(self, db: AsyncSession) -> List[Dict[str, Any]]:
async with self.lock:
payloads = await self._refresh_locked(db)
return [payload for payload in payloads if payload["update_type"] == "full" or _changes_exist(payload.get("changes", {}))]
async def get_sync_payloads(
self,
db: AsyncSession,
*,
last_sequence: int | None,
reason: str | None,
) -> List[Dict[str, Any]]:
async with self.lock:
await self._refresh_locked(db)
if self.state is None:
return []
normalized_reason = (reason or "").strip()
if last_sequence is None or normalized_reason in FULL_RESYNC_REASONS:
return [_full_payload_from_state(self.state, self.sequence)]
if last_sequence == self.sequence:
return [_noop_incremental_payload(self.state, self.sequence)]
if last_sequence > self.sequence:
return [_full_payload_from_state(self.state, self.sequence)]
replay_payloads = [
payload
for payload in self.history
if payload["base_sequence"] >= last_sequence
and payload["sequence"] > last_sequence
]
if replay_payloads:
expected_base = last_sequence
ordered_payloads: List[Dict[str, Any]] = []
for payload in replay_payloads:
if payload["base_sequence"] != expected_base:
return [_full_payload_from_state(self.state, self.sequence)]
ordered_payloads.append(payload)
expected_base = payload["sequence"]
if ordered_payloads and ordered_payloads[-1]["sequence"] == self.sequence:
return ordered_payloads
return [_full_payload_from_state(self.state, self.sequence)]
ue_scene_state_store = UeSceneStateStore()