From 8f67d8d108dac71163a4ad07277084a4cc2a1d1b Mon Sep 17 00:00:00 2001 From: opencode Date: Sat, 9 May 2026 11:32:07 +0800 Subject: [PATCH] feat: add ue_scene websocket integration --- backend/app/api/v1/websocket.py | 83 ++- backend/app/core/websocket/broadcaster.py | 32 +- backend/app/core/websocket/manager.py | 109 ++- backend/app/core/websocket/ue_scene.py | 562 +++++++++++++++ backend/tests/test_websocket_scene.py | 150 ++++ docs/ue-ndisplay-websocket-protocol.md | 352 +++++++++ scripts/test_ue_scene_ws.ps1 | 134 ++++ .../BP_GlobeController_WebSocket_Guide.md | 264 +++++++ unreal/Source/Planet/Planet.Build.cs | 22 + .../Planet/Private/UeSceneWebSocketClient.cpp | 677 ++++++++++++++++++ .../Planet/Public/UeSceneWebSocketClient.h | 159 ++++ unreal/UE57_WebSocket_Integration.md | 189 +++++ 12 files changed, 2684 insertions(+), 49 deletions(-) create mode 100644 backend/app/core/websocket/ue_scene.py create mode 100644 backend/tests/test_websocket_scene.py create mode 100644 docs/ue-ndisplay-websocket-protocol.md create mode 100644 scripts/test_ue_scene_ws.ps1 create mode 100644 unreal/Content/Blueprints/BP_GlobeController_WebSocket_Guide.md create mode 100644 unreal/Source/Planet/Planet.Build.cs create mode 100644 unreal/Source/Planet/Private/UeSceneWebSocketClient.cpp create mode 100644 unreal/Source/Planet/Public/UeSceneWebSocketClient.h create mode 100644 unreal/UE57_WebSocket_Integration.md diff --git a/backend/app/api/v1/websocket.py b/backend/app/api/v1/websocket.py index 23ccb3fe..e8fbbc02 100644 --- a/backend/app/api/v1/websocket.py +++ b/backend/app/api/v1/websocket.py @@ -1,7 +1,6 @@ """WebSocket API endpoints""" import asyncio -import json import logging from datetime import UTC, datetime from typing import Optional @@ -12,9 +11,20 @@ from jose import jwt, JWTError from app.core.config import settings from app.core.time import to_iso8601_utc from app.core.websocket.manager import manager +from app.core.websocket.ue_scene import expand_scene_payload_for_transport, ue_scene_state_store +from app.db.session import async_session_factory logger = logging.getLogger(__name__) router = APIRouter() +SUPPORTED_CHANNELS = [ + "gpu_clusters", + "submarine_cables", + "ixp_nodes", + "alerts", + "dashboard", + "datasource_tasks", + "ue_scene", +] async def authenticate_token(token: str) -> Optional[dict]: @@ -50,18 +60,12 @@ async def websocket_endpoint( await websocket.send_json( { "type": "connection_established", + "timestamp": to_iso8601_utc(datetime.now(UTC)), "data": { "connection_id": f"conn_{user_id}", "server_version": settings.VERSION, "heartbeat_interval": 30, - "supported_channels": [ - "gpu_clusters", - "submarine_cables", - "ixp_nodes", - "alerts", - "dashboard", - "datasource_tasks", - ], + "supported_channels": SUPPORTED_CHANNELS, }, } ) @@ -74,26 +78,79 @@ async def websocket_endpoint( await websocket.send_json( { "type": "heartbeat", + "timestamp": to_iso8601_utc(datetime.now(UTC)), "data": {"action": "pong", "timestamp": to_iso8601_utc(datetime.now(UTC))}, } ) elif data.get("type") == "subscribe": - channels = data.get("data", {}).get("channels", []) + requested_channels = data.get("data", {}).get("channels", []) + channels = manager.subscribe(websocket, requested_channels) await websocket.send_json( { "type": "subscription_confirmed", + "timestamp": to_iso8601_utc(datetime.now(UTC)), "data": {"action": "subscribe", "channels": channels}, } ) + elif data.get("type") == "sync_request": + sync_data = data.get("data", {}) + channel = sync_data.get("channel") + if channel == "ue_scene": + async with async_session_factory() as session: + scene_payloads = await ue_scene_state_store.get_sync_payloads( + session, + last_sequence=sync_data.get("last_sequence"), + reason=sync_data.get("reason"), + ) + for scene_payload in scene_payloads: + for transport_payload in expand_scene_payload_for_transport(scene_payload): + await websocket.send_json( + { + "type": "data_frame", + "channel": "ue_scene", + "timestamp": to_iso8601_utc(datetime.now(UTC)), + "data": transport_payload, + } + ) + else: + await websocket.send_json( + { + "type": "error", + "timestamp": to_iso8601_utc(datetime.now(UTC)), + "data": { + "message": f"Unsupported sync channel: {channel}", + }, + } + ) elif data.get("type") == "control_frame": await websocket.send_json( - {"type": "control_acknowledged", "data": {"received": True}} + { + "type": "control_acknowledged", + "timestamp": to_iso8601_utc(datetime.now(UTC)), + "data": { + "target": data.get("data", {}).get("target"), + "command": data.get("data", {}).get("command"), + "accepted": True, + }, + } ) else: - await websocket.send_json({"type": "ack", "data": {"received": True}}) + await websocket.send_json( + { + "type": "ack", + "timestamp": to_iso8601_utc(datetime.now(UTC)), + "data": {"received": True}, + } + ) except asyncio.TimeoutError: - await websocket.send_json({"type": "heartbeat", "data": {"action": "ping"}}) + await websocket.send_json( + { + "type": "heartbeat", + "timestamp": to_iso8601_utc(datetime.now(UTC)), + "data": {"action": "ping"}, + } + ) except WebSocketDisconnect: pass diff --git a/backend/app/core/websocket/broadcaster.py b/backend/app/core/websocket/broadcaster.py index 0fbc2523..3a4ce08e 100644 --- a/backend/app/core/websocket/broadcaster.py +++ b/backend/app/core/websocket/broadcaster.py @@ -1,11 +1,12 @@ -"""Data broadcaster for WebSocket connections""" +"""Data broadcaster for WebSocket connections.""" import asyncio from datetime import UTC, datetime -from typing import Dict, Any, Optional +from typing import Any, Dict from app.core.time import to_iso8601_utc from app.core.websocket.manager import manager +from app.db.session import async_session_factory @@ -45,6 +46,30 @@ class DataBroadcaster: pass await asyncio.sleep(interval) + async def broadcast_ue_scene(self, interval: int = 5) -> None: + """Broadcast UE scene updates for nDisplay primary nodes.""" + from app.core.websocket.ue_scene import expand_scene_payload_for_transport, ue_scene_state_store + + while self.running: + try: + async with async_session_factory() as session: + scene_payloads = await ue_scene_state_store.get_broadcast_payloads(session) + + for scene_data in scene_payloads: + for transport_payload in expand_scene_payload_for_transport(scene_data): + await manager.broadcast( + { + "type": "data_frame", + "channel": "ue_scene", + "timestamp": to_iso8601_utc(datetime.now(UTC)), + "data": transport_payload, + }, + channel="ue_scene", + ) + except Exception: + pass + await asyncio.sleep(interval) + async def broadcast_alert(self, alert: Dict[str, Any]): """Broadcast an alert to all connected clients""" await manager.broadcast( @@ -75,7 +100,7 @@ class DataBroadcaster: "timestamp": to_iso8601_utc(datetime.now(UTC)), "payload": data, }, - channel=channel if channel in manager.active_connections else "all", + channel=channel, ) async def broadcast_datasource_task_update(self, data: Dict[str, Any]): @@ -95,6 +120,7 @@ class DataBroadcaster: if not self.running: self.running = True self.tasks["dashboard"] = asyncio.create_task(self.broadcast_stats(5)) + self.tasks["ue_scene"] = asyncio.create_task(self.broadcast_ue_scene(5)) def stop(self): """Stop all broadcasters""" diff --git a/backend/app/core/websocket/manager.py b/backend/app/core/websocket/manager.py index a4994119..36e87d90 100644 --- a/backend/app/core/websocket/manager.py +++ b/backend/app/core/websocket/manager.py @@ -1,9 +1,7 @@ -"""WebSocket Connection Manager""" +"""WebSocket connection manager with channel subscriptions.""" + +from typing import Dict, Optional, Set -import json -import asyncio -from typing import Dict, Set, Optional -from datetime import datetime from fastapi import WebSocket import redis.asyncio as redis @@ -11,17 +9,25 @@ from app.core.config import settings class ConnectionManager: - """Manages WebSocket connections""" + """Manage user connections and channel subscriptions.""" - def __init__(self): - self.active_connections: Dict[str, Set[WebSocket]] = {} # user_id -> connections + def __init__(self) -> None: + self.user_connections: Dict[str, Set[WebSocket]] = {} + self.socket_users: Dict[WebSocket, str] = {} + self.channel_connections: Dict[str, Set[WebSocket]] = {} + self.socket_channels: Dict[WebSocket, Set[str]] = {} self.redis_client: Optional[redis.Redis] = None - async def connect(self, websocket: WebSocket, user_id: str): + @property + def active_connections(self) -> Dict[str, Set[WebSocket]]: + """Compatibility alias for existing callers.""" + return self.user_connections + + async def connect(self, websocket: WebSocket, user_id: str) -> None: await websocket.accept() - if user_id not in self.active_connections: - self.active_connections[user_id] = set() - self.active_connections[user_id].add(websocket) + self.user_connections.setdefault(user_id, set()).add(websocket) + self.socket_users[websocket] = user_id + self.socket_channels.setdefault(websocket, set()) if self.redis_client is None: redis_url = settings.REDIS_URL @@ -35,32 +41,69 @@ class ConnectionManager: decode_responses=True, ) - def disconnect(self, websocket: WebSocket, user_id: str): - if user_id in self.active_connections: - self.active_connections[user_id].discard(websocket) - if not self.active_connections[user_id]: - del self.active_connections[user_id] + def disconnect(self, websocket: WebSocket, user_id: str) -> None: + self.unsubscribe(websocket, list(self.socket_channels.get(websocket, set()))) - async def send_personal_message(self, message: dict, user_id: str): - if user_id in self.active_connections: - for connection in self.active_connections[user_id]: - try: - await connection.send_json(message) - except Exception: - pass + if user_id in self.user_connections: + self.user_connections[user_id].discard(websocket) + if not self.user_connections[user_id]: + del self.user_connections[user_id] - async def broadcast(self, message: dict, channel: str = "all"): + self.socket_users.pop(websocket, None) + self.socket_channels.pop(websocket, None) + + def subscribe(self, websocket: WebSocket, channels: list[str]) -> list[str]: + subscribed_channels: list[str] = [] + socket_channel_set = self.socket_channels.setdefault(websocket, set()) + + for channel in channels: + normalized_channel = channel.strip() + if not normalized_channel: + continue + self.channel_connections.setdefault(normalized_channel, set()).add(websocket) + socket_channel_set.add(normalized_channel) + subscribed_channels.append(normalized_channel) + + return subscribed_channels + + def unsubscribe(self, websocket: WebSocket, channels: list[str]) -> None: + socket_channel_set = self.socket_channels.setdefault(websocket, set()) + + for channel in channels: + normalized_channel = channel.strip() + if not normalized_channel: + continue + if normalized_channel in self.channel_connections: + self.channel_connections[normalized_channel].discard(websocket) + if not self.channel_connections[normalized_channel]: + del self.channel_connections[normalized_channel] + socket_channel_set.discard(normalized_channel) + + async def send_personal_message(self, message: dict, user_id: str) -> None: + for connection in list(self.user_connections.get(user_id, set())): + try: + await connection.send_json(message) + except Exception: + self.disconnect(connection, user_id) + + async def broadcast(self, message: dict, channel: str = "all") -> None: if channel == "all": - for user_id in self.active_connections: - await self.send_personal_message(message, user_id) + targets = list(self.socket_users.keys()) else: - await self.send_personal_message(message, channel) + targets = list(self.channel_connections.get(channel, set())) - async def close_all(self): - for user_id in self.active_connections: - for connection in self.active_connections[user_id]: - await connection.close() - self.active_connections.clear() + for connection in targets: + try: + await connection.send_json(message) + except Exception: + user_id = self.socket_users.get(connection) + if user_id is not None: + self.disconnect(connection, user_id) + + async def close_all(self) -> None: + for websocket, user_id in list(self.socket_users.items()): + await websocket.close() + self.disconnect(websocket, user_id) manager = ConnectionManager() diff --git a/backend/app/core/websocket/ue_scene.py b/backend/app/core/websocket/ue_scene.py new file mode 100644 index 00000000..a3bb6621 --- /dev/null +++ b/backend/app/core/websocket/ue_scene.py @@ -0,0 +1,562 @@ +"""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() diff --git a/backend/tests/test_websocket_scene.py b/backend/tests/test_websocket_scene.py new file mode 100644 index 00000000..e0ce53fa --- /dev/null +++ b/backend/tests/test_websocket_scene.py @@ -0,0 +1,150 @@ +import pytest + +from app.core.websocket.manager import ConnectionManager +from app.core.websocket.ue_scene import ( + UeSceneStateStore, + _item_revision, + build_incremental_changes, +) + + +class _FakeSession: + pass + + +class _FakeWebSocket: + def __init__(self): + self.accepted = False + self.messages = [] + + async def accept(self): + self.accepted = True + + async def send_json(self, message): + self.messages.append(message) + + async def close(self): + return None + + +def _scene_item(item_id: str, title: str, *, lat: float = 0.0, lng: float = 0.0): + item = { + "id": item_id, + "entity_type": "gpu_cluster", + "geo": {"lat": lat, "lng": lng, "alt": 0.0}, + "visual": {"style": "pulse_marker", "size": 1.0, "color": "#FF8C42"}, + "metrics": {"name": title}, + "labels": {"title": title, "subtitle": ""}, + "status": {"health": "normal", "alert_level": "none"}, + } + item["revision"] = _item_revision(item) + return item + + +def _scene_state(state_hash: int, items_by_layer: dict[str, dict[str, dict]]): + base_layers = { + "satellites": {"revision": 0, "items": {}}, + "supercomputers": {"revision": 0, "items": {}}, + "gpu_clusters": {"revision": 0, "items": {}}, + "submarine_cables": {"revision": 0, "items": {}}, + "landing_points": {"revision": 0, "items": {}}, + "alerts": {"revision": 0, "items": {}}, + } + for layer_name, items in items_by_layer.items(): + base_layers[layer_name]["items"] = items + base_layers[layer_name]["revision"] = len(items) + return { + "generated_at": "2026-04-17T00:00:00Z", + "state_hash": state_hash, + "total_records": sum(len(items) for items in items_by_layer.values()), + "layers": base_layers, + } + + +def test_build_incremental_changes_detects_add_update_remove(): + previous_state = _scene_state( + 1, + { + "gpu_clusters": { + "gpu:a": _scene_item("gpu:a", "A"), + "gpu:b": _scene_item("gpu:b", "B"), + } + }, + ) + current_state = _scene_state( + 2, + { + "gpu_clusters": { + "gpu:b": _scene_item("gpu:b", "B Updated"), + "gpu:c": _scene_item("gpu:c", "C"), + } + }, + ) + + changes = build_incremental_changes(previous_state, current_state) + + assert changes["gpu_clusters"]["added"][0]["id"] == "gpu:c" + assert changes["gpu_clusters"]["updated"][0]["id"] == "gpu:b" + assert changes["gpu_clusters"]["removed"] == ["gpu:a"] + + +@pytest.mark.asyncio +async def test_state_store_replays_incremental_history(monkeypatch): + state_list = [ + _scene_state(1, {"gpu_clusters": {"gpu:a": _scene_item("gpu:a", "A")}}), + _scene_state( + 2, + {"gpu_clusters": {"gpu:a": _scene_item("gpu:a", "A"), "gpu:b": _scene_item("gpu:b", "B")}}, + ), + _scene_state( + 3, + {"gpu_clusters": {"gpu:a": _scene_item("gpu:a", "A Updated"), "gpu:b": _scene_item("gpu:b", "B")}}, + ), + ] + state_index = {"value": 0} + + async def _fake_build_visualization_scene_state(_db): + current_index = state_index["value"] + if current_index < len(state_list) - 1: + state_index["value"] += 1 + return state_list[current_index] + + monkeypatch.setattr( + "app.core.websocket.ue_scene.build_visualization_scene_state", + _fake_build_visualization_scene_state, + ) + + store = UeSceneStateStore(history_limit=5) + session = _FakeSession() + + first_payloads = await store.get_broadcast_payloads(session) + second_payloads = await store.get_broadcast_payloads(session) + third_payloads = await store.get_broadcast_payloads(session) + replay_payloads = await store.get_sync_payloads( + session, + last_sequence=1, + reason="sequence_gap", + ) + + assert first_payloads[0]["update_type"] == "full" + assert second_payloads[0]["update_type"] == "incremental" + assert third_payloads[0]["sequence"] == 3 + assert [payload["sequence"] for payload in replay_payloads] == [2, 3] + + +@pytest.mark.asyncio +async def test_connection_manager_broadcasts_to_subscribed_channel_only(): + manager = ConnectionManager() + ws_dashboard = _FakeWebSocket() + ws_scene = _FakeWebSocket() + + await manager.connect(ws_dashboard, "user-dashboard") + await manager.connect(ws_scene, "user-scene") + + manager.subscribe(ws_dashboard, ["dashboard"]) + manager.subscribe(ws_scene, ["ue_scene"]) + + await manager.broadcast({"type": "data_frame", "channel": "ue_scene"}, channel="ue_scene") + + assert ws_dashboard.messages == [] + assert ws_scene.messages == [{"type": "data_frame", "channel": "ue_scene"}] diff --git a/docs/ue-ndisplay-websocket-protocol.md b/docs/ue-ndisplay-websocket-protocol.md new file mode 100644 index 00000000..9b89f61b --- /dev/null +++ b/docs/ue-ndisplay-websocket-protocol.md @@ -0,0 +1,352 @@ +# UE nDisplay WebSocket 协议草案 + +## 目的 + +本协议定义后端到 UE `nDisplay` 集群运行时之间的场景同步契约。 + +设计目标: + +- 由 UE 主节点单点接入后端 +- 使用全量快照重建完整场景 +- 使用增量更新降低同步成本 +- 支持控制帧和显示配置 +- 支持序列号、重放和重同步 + +## 适用范围 + +该协议建立在通用 WebSocket 协议之上,参考 [websocket_protocol.md](/D:/work/planet/planet/websocket_protocol.md:1)。 + +新增的核心内容是 `ue_scene` 频道。 + +## 连接模型 + +### 客户端 + +- UE 主节点连接 `ws://:8000/ws?token=` +- UE 渲染节点不直接连接后端场景通道 + +### 订阅 + +UE 主节点建议订阅: + +- `ue_scene` +- `alerts` +- `dashboard` + +## 协议规则 + +- 每次场景更新都带单调递增的 `sequence` +- 增量更新必须包含 `base_sequence` +- UE 如果发现序列断裂,应先发 `sync_request` +- 服务端优先返回可重放的增量历史;无法重放时返回全量快照 +- 后端始终发送地理坐标,UE 本地负责转成球体世界坐标 +- 控制消息与场景消息分离 + +## 消息类型 + +### `connection_established` + +服务端 -> 客户端 + +```json +{ + "type": "connection_established", + "timestamp": "2026-04-17T02:30:00.000Z", + "data": { + "connection_id": "conn_primary_node_01", + "server_version": "0.19.0", + "heartbeat_interval": 30, + "supported_channels": [ + "ue_scene", + "alerts", + "dashboard", + "datasource_tasks" + ] + } +} +``` + +### `subscribe` + +客户端 -> 服务端 + +```json +{ + "type": "subscribe", + "timestamp": "2026-04-17T02:30:01.000Z", + "data": { + "channels": ["ue_scene", "alerts"] + } +} +``` + +### `sync_request` + +客户端 -> 服务端 + +```json +{ + "type": "sync_request", + "timestamp": "2026-04-17T02:30:05.000Z", + "data": { + "channel": "ue_scene", + "reason": "initial_connect", + "last_sequence": null + } +} +``` + +`reason` 常见值: + +- `initial_connect` +- `sequence_gap` +- `manual_resync` +- `profile_changed` + +### `data_frame` 全量快照 + +服务端 -> 客户端 + +```json +{ + "type": "data_frame", + "channel": "ue_scene", + "timestamp": "2026-04-17T02:30:06.000Z", + "data": { + "update_type": "full", + "sequence": 12, + "cluster_time": "2026-04-17T02:30:06.000Z", + "display_profile": { + "profile_id": "polarized-wall-a", + "stereo_mode": "polarized", + "screen_width_m": 3.0, + "screen_height_m": 2.0, + "target_refresh_hz": 120 + }, + "camera_state": { + "mode": "auto_cruise", + "path_id": "global_overview", + "fov": 42.0 + }, + "payload": { + "meta": { + "generated_at": "2026-04-17T02:30:06.000Z", + "total_records": 20800, + "state_hash": 123456789 + }, + "layers": { + "satellites": { + "revision": 111, + "items": [] + }, + "supercomputers": { + "revision": 222, + "items": [] + }, + "gpu_clusters": { + "revision": 333, + "items": [] + }, + "submarine_cables": { + "revision": 444, + "items": [] + }, + "landing_points": { + "revision": 555, + "items": [] + }, + "alerts": { + "revision": 0, + "items": [] + } + } + } + } +} +``` + +### `data_frame` 增量更新 + +服务端 -> 客户端 + +```json +{ + "type": "data_frame", + "channel": "ue_scene", + "timestamp": "2026-04-17T02:31:00.000Z", + "data": { + "update_type": "incremental", + "sequence": 13, + "base_sequence": 12, + "cluster_time": "2026-04-17T02:31:00.000Z", + "meta": { + "generated_at": "2026-04-17T02:31:00.000Z", + "total_records": 20805, + "state_hash": 123456790 + }, + "changes": { + "gpu_clusters": { + "added": [], + "updated": [], + "removed": [] + }, + "submarine_cables": { + "added": [], + "updated": [], + "removed": [] + } + } + } +} +``` + +### `control_frame` + +客户端 -> 服务端 + +```json +{ + "type": "control_frame", + "timestamp": "2026-04-17T02:32:00.000Z", + "data": { + "target": "ue_scene", + "command": "set_camera_mode", + "arguments": { + "mode": "manual", + "camera_id": "operator_cam_1" + } + } +} +``` + +常见命令: + +- `set_camera_mode` +- `set_camera_path` +- `focus_entity` +- `set_layer_visibility` +- `set_alert_filter` +- `set_display_profile` +- `pause_auto_cruise` +- `resume_auto_cruise` + +### `control_acknowledged` + +服务端 -> 客户端 + +```json +{ + "type": "control_acknowledged", + "timestamp": "2026-04-17T02:32:00.050Z", + "data": { + "target": "ue_scene", + "command": "set_camera_mode", + "accepted": true + } +} +``` + +### `heartbeat` + +双向 + +```json +{ + "type": "heartbeat", + "timestamp": "2026-04-17T02:32:30.000Z", + "data": { + "action": "ping" + } +} +``` + +服务端返回: + +```json +{ + "type": "heartbeat", + "timestamp": "2026-04-17T02:32:30.020Z", + "data": { + "action": "pong" + } +} +``` + +## 场景实体结构 + +每个场景实体建议包含: + +- `id` +- `entity_type` +- `revision` +- `geo` +- `visual` +- `metrics` +- `labels` +- `status` + +示例: + +```json +{ + "id": "gpu_clusters:123", + "entity_type": "gpu_cluster", + "revision": 987654321, + "geo": { + "lat": 35.9327, + "lng": -84.3107, + "alt": 0.0 + }, + "visual": { + "style": "pulse_marker", + "size": 1.0, + "color": "#FF8C42" + }, + "metrics": { + "name": "Frontier" + }, + "labels": { + "title": "Frontier", + "subtitle": "Oak Ridge, US" + }, + "status": { + "health": "normal", + "alert_level": "none" + } +} +``` + +## 同步策略 + +UE 主节点建议这样处理: + +1. 保存本地最后一次成功应用的 `sequence` +2. 收到增量帧时检查 `base_sequence` +3. 若 `base_sequence` 不等于本地序列,则立即发 `sync_request` +4. 如果服务端能重放缺失增量,则按顺序依次应用 +5. 如果服务端无法重放,则接收并应用新的全量快照 +6. 场景更新应在主节点原子提交,再交给集群同步显示 + +## 服务端重同步规则 + +服务端建议行为: + +- `initial_connect`:直接返回全量快照 +- `manual_resync`:直接返回全量快照 +- `profile_changed`:直接返回全量快照 +- `sequence_gap`:优先尝试从历史缓存中回放增量 +- 如果历史缓存不完整:回退到全量快照 +- 如果客户端已是最新序列:返回空变化增量或保持静默 + +## 集群职责约束 + +- 只有 UE 主节点接收后端场景同步消息 +- 渲染节点不直接重放后端消息 +- 服务端表达的是场景意图,而不是每个节点的底层渲染状态 + +## 推荐的后端扩展 + +- 为 `ue_scene` 增加更多真实可视化层 +- 将告警/BGP 事件并入统一场景层 +- 为控制命令增加真正的状态持久化与回执 +- 为主节点断线重连补状态恢复策略 diff --git a/scripts/test_ue_scene_ws.ps1 b/scripts/test_ue_scene_ws.ps1 new file mode 100644 index 00000000..db00cab2 --- /dev/null +++ b/scripts/test_ue_scene_ws.ps1 @@ -0,0 +1,134 @@ +param( + [string]$BaseUrl = "http://localhost:8000", + [string]$Username, + [string]$Password, + [int]$ReadMessageCount = 5 +) + +Set-StrictMode -Version Latest +$ErrorActionPreference = "Stop" + +if ([string]::IsNullOrWhiteSpace($Username) -or [string]::IsNullOrWhiteSpace($Password)) { + throw "Missing username/password. Pass -Username and -Password." +} + +Add-Type -AssemblyName System.Net.Http + +function Read-WebSocketMessage { + param( + [System.Net.WebSockets.ClientWebSocket]$Socket + ) + + $buffer = New-Object byte[] 8192 + $segment = [System.ArraySegment[byte]]::new($buffer) + $builder = New-Object System.Text.StringBuilder + + do { + $result = $Socket.ReceiveAsync($segment, [Threading.CancellationToken]::None).GetAwaiter().GetResult() + if ($result.MessageType -eq [System.Net.WebSockets.WebSocketMessageType]::Close) { + return $null + } + + $builder.Append([System.Text.Encoding]::UTF8.GetString($buffer, 0, $result.Count)) | Out-Null + } while (-not $result.EndOfMessage) + + return $builder.ToString() +} + +$loginUri = "$BaseUrl/api/v1/auth/login" +$loginBody = "username=$([uri]::EscapeDataString($Username))&password=$([uri]::EscapeDataString($Password))" + +Write-Host "Requesting access token from $loginUri ..." +$tokenResponse = Invoke-RestMethod ` + -Method Post ` + -Uri $loginUri ` + -ContentType "application/x-www-form-urlencoded" ` + -Body $loginBody + +if (-not $tokenResponse.access_token) { + throw "Login succeeded but access_token missing." +} + +$accessToken = [string]$tokenResponse.access_token +Write-Host "Token acquired. Expires in: $($tokenResponse.expires_in)" + +$wsBase = $BaseUrl -replace "^http", "ws" +$wsUri = "$wsBase/ws?token=$accessToken" + +$socket = [System.Net.WebSockets.ClientWebSocket]::new() +$socket.Options.KeepAliveInterval = [TimeSpan]::FromSeconds(20) + +Write-Host "Connecting WebSocket: $wsUri" +$socket.ConnectAsync([Uri]$wsUri, [Threading.CancellationToken]::None).GetAwaiter().GetResult() +Write-Host "WebSocket connected." + +try { + $initial = Read-WebSocketMessage -Socket $socket + if ($null -ne $initial) { + Write-Host "" + Write-Host "=== connection_established ===" + Write-Host $initial + } + + $subscribeFrame = @{ + type = "subscribe" + data = @{ + channels = @("ue_scene") + } + } | ConvertTo-Json -Depth 8 -Compress + + $subscribeBytes = [System.Text.Encoding]::UTF8.GetBytes($subscribeFrame) + $subscribeSegment = [System.ArraySegment[byte]]::new($subscribeBytes) + $socket.SendAsync( + $subscribeSegment, + [System.Net.WebSockets.WebSocketMessageType]::Text, + $true, + [Threading.CancellationToken]::None + ).GetAwaiter().GetResult() + + $subAck = Read-WebSocketMessage -Socket $socket + if ($null -ne $subAck) { + Write-Host "" + Write-Host "=== subscription_confirmed ===" + Write-Host $subAck + } + + $syncFrame = @{ + type = "sync_request" + data = @{ + channel = "ue_scene" + reason = "initial_connect" + last_sequence = $null + } + } | ConvertTo-Json -Depth 8 -Compress + + $syncBytes = [System.Text.Encoding]::UTF8.GetBytes($syncFrame) + $syncSegment = [System.ArraySegment[byte]]::new($syncBytes) + $socket.SendAsync( + $syncSegment, + [System.Net.WebSockets.WebSocketMessageType]::Text, + $true, + [Threading.CancellationToken]::None + ).GetAwaiter().GetResult() + + for ($i = 1; $i -le $ReadMessageCount; $i++) { + $message = Read-WebSocketMessage -Socket $socket + if ($null -eq $message) { + break + } + + Write-Host "" + Write-Host "=== message $i ===" + Write-Host $message + } +} +finally { + if ($socket.State -eq [System.Net.WebSockets.WebSocketState]::Open) { + $socket.CloseAsync( + [System.Net.WebSockets.WebSocketCloseStatus]::NormalClosure, + "done", + [Threading.CancellationToken]::None + ).GetAwaiter().GetResult() + } + $socket.Dispose() +} diff --git a/unreal/Content/Blueprints/BP_GlobeController_WebSocket_Guide.md b/unreal/Content/Blueprints/BP_GlobeController_WebSocket_Guide.md new file mode 100644 index 00000000..2cbd78d4 --- /dev/null +++ b/unreal/Content/Blueprints/BP_GlobeController_WebSocket_Guide.md @@ -0,0 +1,264 @@ +# BP_GlobeController 接 `ue_scene` 数据接口指南 + +## 这份文档解决什么问题 + +这份文档专门讲: + +1. 后端 `ue_scene` 已经准备好之后,如何先在本机验证接口 +2. UE 展示端接接口时,消息应该怎么走 +3. `BP_GlobeController` 后续该预留哪些函数和状态 + +适合当前阶段: + +- 你已经有 `Earth_Globe` +- 你已经有一个能展示地球的 Demo +- 现在要开始把后端的 `ue_scene` 数据接进来 + +## 当前后端入口 + +后端 WebSocket 入口在: + +- [backend/app/api/v1/websocket.py](/D:/work/planet/planet/backend/app/api/v1/websocket.py:1) + +`ue_scene` 数据构建在: + +- [backend/app/core/websocket/ue_scene.py](/D:/work/planet/planet/backend/app/core/websocket/ue_scene.py:1) + +协议草案在: + +- [docs/ue-ndisplay-websocket-protocol.md](/D:/work/planet/planet/docs/ue-ndisplay-websocket-protocol.md:1) + +## 先不要直接进 UE + +建议先做一件事: + +先在本机验证这条链路是通的: + +1. 登录后端拿 `access_token` +2. 用这个 token 连接 `/ws` +3. 订阅 `ue_scene` +4. 发送 `sync_request` +5. 确认后端回了 `full` 快照 + +这一步跑通后,再把同样的消息流程搬到 UE。 + +## 本机验证脚本 + +仓库里已经补了测试脚本: + +- [scripts/test_ue_scene_ws.ps1](/D:/work/planet/planet/scripts/test_ue_scene_ws.ps1) + +### 用法 + +```powershell +powershell -ExecutionPolicy Bypass -File scripts/test_ue_scene_ws.ps1 ` + -BaseUrl http://localhost:8000 ` + -Username 你的用户名 ` + -Password 你的密码 ` + -ReadMessageCount 3 +``` + +### 预期输出 + +你应该能看到这些消息: + +1. `connection_established` +2. `subscription_confirmed` +3. 至少一条 `data_frame` + +其中第一条 `data_frame` 里: + +- `channel = ue_scene` +- `update_type = full` + +如果这一步失败,不要先去 UE 里找问题。 + +## 连接顺序 + +UE 主节点后续接接口时,顺序就是下面这样: + +1. 登录拿 token +2. 连接 `ws://:8000/ws?token=` +3. 收到 `connection_established` +4. 发 `subscribe` +5. 发 `sync_request` +6. 收 `data_frame` +7. 应用 `full` +8. 后续继续应用 `incremental` + +## WebSocket 地址 + +```text +ws://localhost:8000/ws?token= +``` + +## `subscribe` 消息 + +```json +{ + "type": "subscribe", + "data": { + "channels": ["ue_scene"] + } +} +``` + +## `sync_request` 消息 + +```json +{ + "type": "sync_request", + "data": { + "channel": "ue_scene", + "reason": "initial_connect", + "last_sequence": null + } +} +``` + +## 你在 UE 里要接收的重点字段 + +### 1. `data.update_type` + +可能是: + +- `full` +- `incremental` + +### 2. `data.sequence` + +每次场景更新都会递增。 + +UE 主节点要保存自己最近一次成功应用的 `sequence`。 + +### 3. `data.payload.layers` + +`full` 快照时,场景实体都在这里。 + +重点层包括: + +- `satellites` +- `supercomputers` +- `gpu_clusters` +- `submarine_cables` +- `landing_points` +- `alerts` + +### 4. `data.changes` + +`incremental` 更新时,变化都在这里。 + +每一层会包含: + +- `added` +- `updated` +- `removed` + +## `BP_GlobeController` 建议先预留的变量 + +### `LastSequence` + +- 类型:整数 +- 用途:记录最近一次成功应用的序列号 + +### `SceneConnected` + +- 类型:布尔 +- 用途:当前是否已连上后端 WebSocket + +### `PendingFullSync` + +- 类型:布尔 +- 用途:是否正在等待全量快照 + +### `CurrentSceneHash` + +- 类型:字符串 +- 用途:记录当前已应用场景的 `state_hash` + +## `BP_GlobeController` 建议先预留的函数 + +### `ConnectWebSocket` + +负责: + +- 建立连接 +- 进入已连接状态 + +### `SendSubscribe` + +负责: + +- 订阅 `ue_scene` + +### `SendSyncRequest` + +负责: + +- 请求 `full` 快照 + +### `HandleMessage` + +负责: + +- 解析服务端返回的原始 JSON + +### `ApplyFullScene` + +负责: + +- 清理旧场景数据 +- 按 `layers` 重建当前场景 + +### `ApplyIncrementalScene` + +负责: + +- 处理 `added` +- 处理 `updated` +- 处理 `removed` + +### `GeoToWorld` + +负责: + +- 把 `lat/lng/alt` 转成球面世界坐标 + +## 第一个可落地目标 + +不要一上来就把所有层都接进来。 + +建议第一阶段只做: + +1. 连接成功 +2. 收到 `full` +3. 先只解析 `gpu_clusters` +4. 在地球上生成几个测试点位 + +这样最容易判断问题在哪: + +- 是连接没通 +- 还是 JSON 没解析对 +- 还是地理坐标映射有问题 + +## 当前阶段建议 + +对于你现在这个 Demo,建议下一步按这个顺序: + +1. 先跑 [scripts/test_ue_scene_ws.ps1](/D:/work/planet/planet/scripts/test_ue_scene_ws.ps1:1) +2. 确认后端确实回 `ue_scene` +3. 再在 UE 里接登录和 WebSocket +4. 先只落 `ApplyFullScene` +5. 最后再接 `ApplyIncrementalScene` + +## 一个现实提醒 + +UE 原生 `WebSockets` 插件更适合 C++ 层接入,不是最友好的纯蓝图体验。 + +所以如果你下一步要我继续带着落地,我建议走这条路线: + +1. 先把协议和测试跑通 +2. 再补一个 UE 侧最小 C++ WebSocket 接收层 +3. 最后把消息转给 `BP_GlobeController` + +这样会比纯蓝图硬接稳定很多。 diff --git a/unreal/Source/Planet/Planet.Build.cs b/unreal/Source/Planet/Planet.Build.cs new file mode 100644 index 00000000..5cd2eb72 --- /dev/null +++ b/unreal/Source/Planet/Planet.Build.cs @@ -0,0 +1,22 @@ +using UnrealBuildTool; + +public class Planet : ModuleRules +{ + public Planet(ReadOnlyTargetRules Target) : base(Target) + { + PCHUsage = PCHUsageMode.UseExplicitOrSharedPCHs; + + PublicDependencyModuleNames.AddRange( + new[] + { + "Core", + "CoreUObject", + "Engine", + "HTTP", + "Json", + "JsonUtilities", + "WebSockets", + } + ); + } +} diff --git a/unreal/Source/Planet/Private/UeSceneWebSocketClient.cpp b/unreal/Source/Planet/Private/UeSceneWebSocketClient.cpp new file mode 100644 index 00000000..f9aced0e --- /dev/null +++ b/unreal/Source/Planet/Private/UeSceneWebSocketClient.cpp @@ -0,0 +1,677 @@ +#include "UeSceneWebSocketClient.h" + +#include "Dom/JsonObject.h" +#include "Dom/JsonValue.h" +#include "HttpModule.h" +#include "JsonObjectConverter.h" +#include "Serialization/JsonReader.h" +#include "Serialization/JsonSerializer.h" +#include "WebSocketsModule.h" + +namespace +{ +constexpr double GeoPointMergeTolerance = 1e-4; + +bool TryParseGeoPoint(const TSharedPtr& PointObject, FVector& OutGeoPoint) +{ + if (!PointObject.IsValid()) + { + return false; + } + + double Latitude = 0.0; + double Longitude = 0.0; + double Altitude = 0.0; + if (!PointObject->TryGetNumberField(TEXT("lat"), Latitude) + || !PointObject->TryGetNumberField(TEXT("lng"), Longitude)) + { + return false; + } + + PointObject->TryGetNumberField(TEXT("alt"), Altitude); + OutGeoPoint = FVector(Longitude, Latitude, Altitude); + return true; +} + +FString ResolveCableName(const TSharedPtr& ItemObject) +{ + FString CableName; + + const TSharedPtr* LabelsObject = nullptr; + if (ItemObject->TryGetObjectField(TEXT("labels"), LabelsObject) + && LabelsObject != nullptr + && LabelsObject->IsValid()) + { + (*LabelsObject)->TryGetStringField(TEXT("title"), CableName); + } + + if (CableName.IsEmpty()) + { + ItemObject->TryGetStringField(TEXT("id"), CableName); + } + + return CableName; +} + +bool AreGeoPointsConnected(const FVector& A, const FVector& B) +{ + return FMath::Abs(A.X - B.X) <= GeoPointMergeTolerance + && FMath::Abs(A.Y - B.Y) <= GeoPointMergeTolerance + && FMath::Abs(A.Z - B.Z) <= GeoPointMergeTolerance; +} + +void AppendMergedSegment( + TArray>& MergedSegments, + const TArray& IncomingSegment +) +{ + if (IncomingSegment.Num() < 2) + { + return; + } + + if (MergedSegments.Num() == 0) + { + MergedSegments.Add(IncomingSegment); + return; + } + + TArray& LastMergedSegment = MergedSegments.Last(); + if (LastMergedSegment.Num() == 0) + { + LastMergedSegment = IncomingSegment; + return; + } + + if (AreGeoPointsConnected(LastMergedSegment.Last(), IncomingSegment[0])) + { + for (int32 Index = 1; Index < IncomingSegment.Num(); ++Index) + { + LastMergedSegment.Add(IncomingSegment[Index]); + } + return; + } + + MergedSegments.Add(IncomingSegment); +} +} + +UUeSceneWebSocketClient::UUeSceneWebSocketClient() +{ + PrimaryComponentTick.bCanEverTick = false; +} + +void UUeSceneWebSocketClient::BeginPlay() +{ + Super::BeginPlay(); +} + +void UUeSceneWebSocketClient::EndPlay(const EEndPlayReason::Type EndPlayReason) +{ + DisconnectFromBackend(); + Super::EndPlay(EndPlayReason); +} + +void UUeSceneWebSocketClient::ConnectToBackend() +{ + if (Username.IsEmpty() || Password.IsEmpty()) + { + EmitError(TEXT("Username / Password 不能为空。")); + return; + } + + RequestAccessToken(); +} + +void UUeSceneWebSocketClient::DisconnectFromBackend() +{ + if (Socket.IsValid()) + { + Socket->Close(); + Socket.Reset(); + } + + PendingFullScene.Reset(); + OnConnectedChanged.Broadcast(false); +} + +void UUeSceneWebSocketClient::SendSubscribe() +{ + TSharedRef Root = MakeShared(); + Root->SetStringField(TEXT("type"), TEXT("subscribe")); + + TSharedRef Data = MakeShared(); + TArray> Channels; + for (const FString& Channel : DefaultChannels) + { + Channels.Add(MakeShared(Channel)); + } + Data->SetArrayField(TEXT("channels"), Channels); + Root->SetObjectField(TEXT("data"), Data); + + FString Payload; + TSharedRef> Writer = TJsonWriterFactory<>::Create(&Payload); + FJsonSerializer::Serialize(Root, Writer); + SendJsonString(Payload); +} + +void UUeSceneWebSocketClient::SendInitialSyncRequest() +{ + SendSyncRequest(TEXT("initial_connect"), INDEX_NONE); +} + +void UUeSceneWebSocketClient::SendSyncRequest(const FString& Reason, int32 InLastSequence) +{ + TSharedRef Root = MakeShared(); + Root->SetStringField(TEXT("type"), TEXT("sync_request")); + + TSharedRef Data = MakeShared(); + Data->SetStringField(TEXT("channel"), TEXT("ue_scene")); + Data->SetStringField(TEXT("reason"), Reason); + if (InLastSequence >= 0) + { + Data->SetNumberField(TEXT("last_sequence"), InLastSequence); + } + Root->SetObjectField(TEXT("data"), Data); + + FString Payload; + TSharedRef> Writer = TJsonWriterFactory<>::Create(&Payload); + FJsonSerializer::Serialize(Root, Writer); + SendJsonString(Payload); +} + +void UUeSceneWebSocketClient::RequestAccessToken() +{ + TSharedRef Request = FHttpModule::Get().CreateRequest(); + Request->SetURL(BuildHttpUrl(TEXT("/api/v1/auth/login"))); + Request->SetVerb(TEXT("POST")); + Request->SetHeader(TEXT("Content-Type"), TEXT("application/x-www-form-urlencoded")); + Request->SetContentAsString( + FString::Printf(TEXT("username=%s&password=%s"), *Username, *Password) + ); + Request->OnProcessRequestComplete().BindUObject( + this, &UUeSceneWebSocketClient::HandleLoginResponse + ); + Request->ProcessRequest(); +} + +void UUeSceneWebSocketClient::HandleLoginResponse( + FHttpRequestPtr Request, + FHttpResponsePtr Response, + bool bSucceeded +) +{ + if (!bSucceeded || !Response.IsValid()) + { + EmitError(TEXT("登录请求失败。")); + return; + } + + if (Response->GetResponseCode() < 200 || Response->GetResponseCode() >= 300) + { + EmitError(FString::Printf(TEXT("登录失败,HTTP %d"), Response->GetResponseCode())); + return; + } + + TSharedPtr Root; + TSharedRef> Reader = TJsonReaderFactory<>::Create(Response->GetContentAsString()); + if (!FJsonSerializer::Deserialize(Reader, Root) || !Root.IsValid()) + { + EmitError(TEXT("登录响应 JSON 解析失败。")); + return; + } + + AccessToken = Root->GetStringField(TEXT("access_token")); + if (AccessToken.IsEmpty()) + { + EmitError(TEXT("登录响应里没有 access_token。")); + return; + } + + OpenWebSocket(); +} + +void UUeSceneWebSocketClient::OpenWebSocket() +{ + if (!FModuleManager::Get().IsModuleLoaded(TEXT("WebSockets"))) + { + FModuleManager::Get().LoadModule(TEXT("WebSockets")); + } + + Socket = FWebSocketsModule::Get().CreateWebSocket(BuildWsUrl()); + BindSocketEvents(); + Socket->Connect(); +} + +void UUeSceneWebSocketClient::BindSocketEvents() +{ + if (!Socket.IsValid()) + { + return; + } + + Socket->OnConnected().AddLambda([this]() + { + OnConnectedChanged.Broadcast(true); + }); + + Socket->OnConnectionError().AddLambda([this](const FString& Error) + { + EmitError(FString::Printf(TEXT("WebSocket 连接失败: %s"), *Error)); + }); + + Socket->OnClosed().AddLambda([this](int32 StatusCode, const FString& Reason, bool bWasClean) + { + OnConnectedChanged.Broadcast(false); + }); + + Socket->OnMessage().AddLambda([this](const FString& Message) + { + OnRawMessage.Broadcast(Message); + TryBroadcastSceneFrame(Message); + }); +} + +void UUeSceneWebSocketClient::EmitError(const FString& ErrorMessage) +{ + OnError.Broadcast(ErrorMessage); +} + +void UUeSceneWebSocketClient::TryBroadcastSceneFrame(const FString& Message) +{ + TSharedPtr Root; + TSharedRef> Reader = TJsonReaderFactory<>::Create(Message); + if (!FJsonSerializer::Deserialize(Reader, Root) || !Root.IsValid()) + { + return; + } + + FString Type; + if (!Root->TryGetStringField(TEXT("type"), Type) || Type != TEXT("data_frame")) + { + return; + } + + FString Channel; + if (!Root->TryGetStringField(TEXT("channel"), Channel) || Channel != TEXT("ue_scene")) + { + return; + } + + const TSharedPtr* DataObject = nullptr; + if (!Root->TryGetObjectField(TEXT("data"), DataObject) || DataObject == nullptr || !DataObject->IsValid()) + { + return; + } + TSharedPtr EffectiveDataObject = *DataObject; + + FString UpdateType = TEXT("unknown"); + EffectiveDataObject->TryGetStringField(TEXT("update_type"), UpdateType); + + int32 Sequence = 0; + EffectiveDataObject->TryGetNumberField(TEXT("sequence"), Sequence); + + int32 TotalRecords = 0; + const TSharedPtr* PayloadObject = nullptr; + if (EffectiveDataObject->TryGetObjectField(TEXT("payload"), PayloadObject) && PayloadObject != nullptr && PayloadObject->IsValid()) + { + const TSharedPtr* MetaObject = nullptr; + if ((*PayloadObject)->TryGetObjectField(TEXT("meta"), MetaObject) && MetaObject != nullptr && MetaObject->IsValid()) + { + (*MetaObject)->TryGetNumberField(TEXT("total_records"), TotalRecords); + } + } + + if (UpdateType == TEXT("full")) + { + TSharedPtr AssembledDataObject; + if (!TryAssembleFullScene(EffectiveDataObject, Sequence, TotalRecords, AssembledDataObject)) + { + return; + } + if (AssembledDataObject.IsValid()) + { + EffectiveDataObject = AssembledDataObject; + } + } + + LastSequence = Sequence; + OnSceneFrame.Broadcast(UpdateType, Sequence, TotalRecords); + BroadcastScenePayload(UpdateType, Sequence, EffectiveDataObject); +} + +bool UUeSceneWebSocketClient::TryAssembleFullScene( + const TSharedPtr& DataObject, + int32 Sequence, + int32 TotalRecords, + TSharedPtr& OutAssembledDataObject +) +{ + OutAssembledDataObject = nullptr; + + const TSharedPtr* PayloadObject = nullptr; + if (!DataObject->TryGetObjectField(TEXT("payload"), PayloadObject) || PayloadObject == nullptr || !PayloadObject->IsValid()) + { + EmitError(TEXT("Full scene payload is missing.")); + return false; + } + + const TSharedPtr* MetaObject = nullptr; + if (!(*PayloadObject)->TryGetObjectField(TEXT("meta"), MetaObject) || MetaObject == nullptr || !MetaObject->IsValid()) + { + EmitError(TEXT("Full scene meta is missing.")); + return false; + } + + bool bChunked = false; + (*MetaObject)->TryGetBoolField(TEXT("chunked"), bChunked); + if (!bChunked) + { + OutAssembledDataObject = DataObject; + PendingFullScene.Reset(); + return true; + } + + int32 TransportChunkIndex = 0; + int32 TransportChunkCount = 0; + int32 LayerChunkIndex = 0; + int32 LayerChunkCount = 0; + FString LayerName; + (*MetaObject)->TryGetNumberField(TEXT("transport_chunk_index"), TransportChunkIndex); + (*MetaObject)->TryGetNumberField(TEXT("transport_chunk_count"), TransportChunkCount); + (*MetaObject)->TryGetNumberField(TEXT("layer_chunk_index"), LayerChunkIndex); + (*MetaObject)->TryGetNumberField(TEXT("layer_chunk_count"), LayerChunkCount); + (*MetaObject)->TryGetStringField(TEXT("layer_name"), LayerName); + + if (TransportChunkIndex <= 0 || TransportChunkCount <= 0 || LayerChunkIndex <= 0 || LayerChunkCount <= 0 || LayerName.IsEmpty()) + { + EmitError(TEXT("Chunked full scene metadata is incomplete.")); + PendingFullScene.Reset(); + return false; + } + + if (PendingFullScene.Sequence != Sequence) + { + PendingFullScene.Reset(); + PendingFullScene.Sequence = Sequence; + PendingFullScene.ExpectedTransportChunkCount = TransportChunkCount; + PendingFullScene.TotalRecords = TotalRecords; + PendingFullScene.DisplayProfile = DataObject->GetObjectField(TEXT("display_profile")); + PendingFullScene.CameraState = DataObject->GetObjectField(TEXT("camera_state")); + PendingFullScene.Meta = MakeShared(*MetaObject->Get()); + DataObject->TryGetStringField(TEXT("cluster_time"), PendingFullScene.ClusterTime); + } + + if (PendingFullScene.Sequence != Sequence || PendingFullScene.ExpectedTransportChunkCount != TransportChunkCount) + { + EmitError(TEXT("Chunked full scene sequence mismatch.")); + PendingFullScene.Reset(); + return false; + } + + if (PendingFullScene.ReceivedTransportChunkIndices.Contains(TransportChunkIndex)) + { + return false; + } + + PendingFullScene.ReceivedTransportChunkIndices.Add(TransportChunkIndex); + + const TSharedPtr* LayersObject = nullptr; + if (!(*PayloadObject)->TryGetObjectField(TEXT("layers"), LayersObject) || LayersObject == nullptr || !LayersObject->IsValid()) + { + EmitError(TEXT("Chunked full scene layers are missing.")); + PendingFullScene.Reset(); + return false; + } + + const TSharedPtr* LayerDataObject = nullptr; + if (!(*LayersObject)->TryGetObjectField(LayerName, LayerDataObject) || LayerDataObject == nullptr || !LayerDataObject->IsValid()) + { + EmitError(TEXT("Chunked full scene layer payload is missing.")); + PendingFullScene.Reset(); + return false; + } + + FLayerChunkAccumulator& LayerAccumulator = PendingFullScene.Layers.FindOrAdd(LayerName); + if (LayerAccumulator.ExpectedChunkCount == 0) + { + LayerAccumulator.ExpectedChunkCount = LayerChunkCount; + LayerDataObject->Get()->TryGetNumberField(TEXT("revision"), LayerAccumulator.Revision); + } + else if (LayerAccumulator.ExpectedChunkCount != LayerChunkCount) + { + EmitError(TEXT("Layer chunk count changed while assembling full scene.")); + PendingFullScene.Reset(); + return false; + } + + if (LayerAccumulator.ReceivedChunkIndices.Contains(LayerChunkIndex)) + { + return false; + } + + const TArray>* LayerItems = nullptr; + if (!LayerDataObject->Get()->TryGetArrayField(TEXT("items"), LayerItems) || LayerItems == nullptr) + { + EmitError(TEXT("Chunked full scene items are missing.")); + PendingFullScene.Reset(); + return false; + } + + LayerAccumulator.ReceivedChunkIndices.Add(LayerChunkIndex); + LayerAccumulator.Items.Append(*LayerItems); + + if (PendingFullScene.ReceivedTransportChunkIndices.Num() != PendingFullScene.ExpectedTransportChunkCount) + { + return false; + } + + for (const TPair& LayerPair : PendingFullScene.Layers) + { + if (LayerPair.Value.ReceivedChunkIndices.Num() != LayerPair.Value.ExpectedChunkCount) + { + return false; + } + } + + TSharedRef AssembledDataObject = MakeShared(); + AssembledDataObject->SetStringField(TEXT("update_type"), TEXT("full")); + AssembledDataObject->SetNumberField(TEXT("sequence"), Sequence); + AssembledDataObject->SetStringField(TEXT("cluster_time"), PendingFullScene.ClusterTime); + if (PendingFullScene.DisplayProfile.IsValid()) + { + AssembledDataObject->SetObjectField(TEXT("display_profile"), PendingFullScene.DisplayProfile); + } + if (PendingFullScene.CameraState.IsValid()) + { + AssembledDataObject->SetObjectField(TEXT("camera_state"), PendingFullScene.CameraState); + } + + TSharedRef AssembledPayloadObject = MakeShared(); + TSharedRef AssembledMetaObject = PendingFullScene.Meta.IsValid() + ? MakeShared(*PendingFullScene.Meta.Get()) + : MakeShared(); + AssembledMetaObject->SetBoolField(TEXT("chunked"), false); + AssembledMetaObject->RemoveField(TEXT("layer_name")); + AssembledMetaObject->RemoveField(TEXT("layer_chunk_index")); + AssembledMetaObject->RemoveField(TEXT("layer_chunk_count")); + AssembledMetaObject->RemoveField(TEXT("transport_chunk_index")); + AssembledMetaObject->RemoveField(TEXT("transport_chunk_count")); + AssembledPayloadObject->SetObjectField(TEXT("meta"), AssembledMetaObject); + + TSharedRef AssembledLayersObject = MakeShared(); + for (TPair& LayerPair : PendingFullScene.Layers) + { + TSharedRef LayerObject = MakeShared(); + LayerObject->SetNumberField(TEXT("revision"), LayerPair.Value.Revision); + LayerObject->SetArrayField(TEXT("items"), LayerPair.Value.Items); + AssembledLayersObject->SetObjectField(LayerPair.Key, LayerObject); + } + AssembledPayloadObject->SetObjectField(TEXT("layers"), AssembledLayersObject); + AssembledDataObject->SetObjectField(TEXT("payload"), AssembledPayloadObject); + + OutAssembledDataObject = AssembledDataObject; + PendingFullScene.Reset(); + return true; +} + +void UUeSceneWebSocketClient::BroadcastScenePayload( + const FString& UpdateType, + int32 Sequence, + const TSharedPtr& DataObject +) +{ + FString PayloadJson; + TSharedRef> Writer = TJsonWriterFactory<>::Create(&PayloadJson); + if (FJsonSerializer::Serialize(DataObject.ToSharedRef(), Writer)) + { + OnScenePayload.Broadcast(UpdateType, Sequence, PayloadJson); + } + + BroadcastLegacySubmarineCableEvents(DataObject); +} + +void UUeSceneWebSocketClient::BroadcastLegacySubmarineCableEvents( + const TSharedPtr& DataObject +) +{ + if (!DataObject.IsValid()) + { + return; + } + + const TSharedPtr* PayloadObject = nullptr; + if (!DataObject->TryGetObjectField(TEXT("payload"), PayloadObject) + || PayloadObject == nullptr + || !PayloadObject->IsValid()) + { + return; + } + + const TSharedPtr* LayersObject = nullptr; + if (!(*PayloadObject)->TryGetObjectField(TEXT("layers"), LayersObject) + || LayersObject == nullptr + || !LayersObject->IsValid()) + { + return; + } + + const TSharedPtr* CableLayerObject = nullptr; + if (!(*LayersObject)->TryGetObjectField(TEXT("submarine_cables"), CableLayerObject) + || CableLayerObject == nullptr + || !CableLayerObject->IsValid()) + { + return; + } + + const TArray>* CableItems = nullptr; + if (!(*CableLayerObject)->TryGetArrayField(TEXT("items"), CableItems) || CableItems == nullptr) + { + return; + } + + for (const TSharedPtr& CableItemValue : *CableItems) + { + const TSharedPtr CableItemObject = CableItemValue.IsValid() + ? CableItemValue->AsObject() + : nullptr; + if (!CableItemObject.IsValid()) + { + continue; + } + + const FString CableName = ResolveCableName(CableItemObject); + + const TSharedPtr* GeoObject = nullptr; + if (!CableItemObject->TryGetObjectField(TEXT("geo"), GeoObject) + || GeoObject == nullptr + || !GeoObject->IsValid()) + { + continue; + } + + const TArray>* SegmentArray = nullptr; + if ((*GeoObject)->TryGetArrayField(TEXT("segments"), SegmentArray) && SegmentArray != nullptr) + { + TArray> MergedSegments; + for (const TSharedPtr& SegmentValue : *SegmentArray) + { + const TArray>* SegmentPoints = nullptr; + if (!SegmentValue.IsValid() || !SegmentValue->TryGetArray(SegmentPoints) || SegmentPoints == nullptr) + { + continue; + } + + TArray GeoPoints; + GeoPoints.Reserve(SegmentPoints->Num()); + + for (const TSharedPtr& PointValue : *SegmentPoints) + { + FVector GeoPoint = FVector::ZeroVector; + if (PointValue.IsValid() && TryParseGeoPoint(PointValue->AsObject(), GeoPoint)) + { + GeoPoints.Add(GeoPoint); + } + } + + AppendMergedSegment(MergedSegments, GeoPoints); + } + + for (const TArray& MergedSegment : MergedSegments) + { + if (MergedSegment.Num() >= 2) + { + OnSubmarineCableReceived.Broadcast(CableName, MergedSegment); + } + } + + continue; + } + + const TArray>* PathArray = nullptr; + if (!(*GeoObject)->TryGetArrayField(TEXT("path"), PathArray) || PathArray == nullptr) + { + continue; + } + + TArray GeoPoints; + GeoPoints.Reserve(PathArray->Num()); + for (const TSharedPtr& PointValue : *PathArray) + { + FVector GeoPoint = FVector::ZeroVector; + if (PointValue.IsValid() && TryParseGeoPoint(PointValue->AsObject(), GeoPoint)) + { + GeoPoints.Add(GeoPoint); + } + } + + if (GeoPoints.Num() >= 2) + { + OnSubmarineCableReceived.Broadcast(CableName, GeoPoints); + } + } +} + +FString UUeSceneWebSocketClient::BuildWsUrl() const +{ + FString WsBase = BackendBaseUrl; + WsBase = WsBase.Replace(TEXT("http://"), TEXT("ws://")); + WsBase = WsBase.Replace(TEXT("https://"), TEXT("wss://")); + return FString::Printf(TEXT("%s/ws?token=%s"), *WsBase, *AccessToken); +} + +FString UUeSceneWebSocketClient::BuildHttpUrl(const FString& Path) const +{ + return FString::Printf(TEXT("%s%s"), *BackendBaseUrl, *Path); +} + +void UUeSceneWebSocketClient::SendJsonString(const FString& JsonPayload) +{ + if (!Socket.IsValid() || !Socket->IsConnected()) + { + EmitError(TEXT("WebSocket 尚未连接。")); + return; + } + + Socket->Send(JsonPayload); +} diff --git a/unreal/Source/Planet/Public/UeSceneWebSocketClient.h b/unreal/Source/Planet/Public/UeSceneWebSocketClient.h new file mode 100644 index 00000000..c25871c9 --- /dev/null +++ b/unreal/Source/Planet/Public/UeSceneWebSocketClient.h @@ -0,0 +1,159 @@ +#pragma once + +#include "CoreMinimal.h" +#include "Components/ActorComponent.h" +#include "Interfaces/IHttpRequest.h" +#include "Interfaces/IHttpResponse.h" +#include "IWebSocket.h" +#include "UeSceneWebSocketClient.generated.h" + +DECLARE_DYNAMIC_MULTICAST_DELEGATE_OneParam(FUeSceneRawMessageEvent, const FString&, Message); +DECLARE_DYNAMIC_MULTICAST_DELEGATE_OneParam(FUeSceneConnectedEvent, bool, bConnected); +DECLARE_DYNAMIC_MULTICAST_DELEGATE_OneParam(FUeSceneErrorEvent, const FString&, ErrorMessage); +DECLARE_DYNAMIC_MULTICAST_DELEGATE_ThreeParams( + FUeSceneFrameEvent, + const FString&, + UpdateType, + int32, + Sequence, + int32, + TotalRecords +); +DECLARE_DYNAMIC_MULTICAST_DELEGATE_ThreeParams( + FUeScenePayloadEvent, + const FString&, + UpdateType, + int32, + Sequence, + const FString&, + PayloadJson +); +DECLARE_DYNAMIC_MULTICAST_DELEGATE_TwoParams( + FUeSceneSubmarineCableEvent, + const FString&, + CableName, + const TArray&, + GeoPoints +); + +UCLASS(ClassGroup=(Planet), meta=(BlueprintSpawnableComponent)) +class PLANET_API UUeSceneWebSocketClient : public UActorComponent +{ + GENERATED_BODY() + +public: + UUeSceneWebSocketClient(); + + UPROPERTY(EditAnywhere, BlueprintReadWrite, Category="Planet|Backend") + FString BackendBaseUrl = TEXT("http://127.0.0.1:8000"); + + UPROPERTY(EditAnywhere, BlueprintReadWrite, Category="Planet|Backend") + FString Username; + + UPROPERTY(EditAnywhere, BlueprintReadWrite, Category="Planet|Backend") + FString Password; + + UPROPERTY(EditAnywhere, BlueprintReadWrite, Category="Planet|Backend") + TArray DefaultChannels = {TEXT("ue_scene")}; + + UPROPERTY(EditAnywhere, BlueprintReadWrite, Category="Planet|Backend") + int32 LastSequence = INDEX_NONE; + + UPROPERTY(BlueprintAssignable, Category="Planet|Backend") + FUeSceneConnectedEvent OnConnectedChanged; + + UPROPERTY(BlueprintAssignable, Category="Planet|Backend") + FUeSceneRawMessageEvent OnRawMessage; + + UPROPERTY(BlueprintAssignable, Category="Planet|Backend") + FUeSceneErrorEvent OnError; + + UPROPERTY(BlueprintAssignable, Category="Planet|Backend") + FUeSceneFrameEvent OnSceneFrame; + + UPROPERTY(BlueprintAssignable, Category="Planet|Backend") + FUeScenePayloadEvent OnScenePayload; + + UPROPERTY(BlueprintAssignable, Category="Planet|Backend|Legacy") + FUeSceneSubmarineCableEvent OnSubmarineCableReceived; + + UFUNCTION(BlueprintCallable, Category="Planet|Backend") + void ConnectToBackend(); + + UFUNCTION(BlueprintCallable, Category="Planet|Backend") + void DisconnectFromBackend(); + + UFUNCTION(BlueprintCallable, Category="Planet|Backend") + void SendSubscribe(); + + UFUNCTION(BlueprintCallable, Category="Planet|Backend") + void SendInitialSyncRequest(); + + UFUNCTION(BlueprintCallable, Category="Planet|Backend") + void SendSyncRequest(const FString& Reason, int32 InLastSequence); + +protected: + virtual void BeginPlay() override; + virtual void EndPlay(const EEndPlayReason::Type EndPlayReason) override; + +private: + struct FLayerChunkAccumulator + { + int32 Revision = 0; + int32 ExpectedChunkCount = 0; + TSet ReceivedChunkIndices; + TArray> Items; + }; + + struct FFullSceneAssemblyState + { + int32 Sequence = INDEX_NONE; + int32 ExpectedTransportChunkCount = 0; + int32 TotalRecords = 0; + FString ClusterTime; + TSharedPtr DisplayProfile; + TSharedPtr CameraState; + TSharedPtr Meta; + TMap Layers; + TSet ReceivedTransportChunkIndices; + + void Reset() + { + Sequence = INDEX_NONE; + ExpectedTransportChunkCount = 0; + TotalRecords = 0; + ClusterTime.Reset(); + DisplayProfile.Reset(); + CameraState.Reset(); + Meta.Reset(); + Layers.Reset(); + ReceivedTransportChunkIndices.Reset(); + } + }; + + FString AccessToken; + TSharedPtr Socket; + FFullSceneAssemblyState PendingFullScene; + + void RequestAccessToken(); + void HandleLoginResponse(FHttpRequestPtr Request, FHttpResponsePtr Response, bool bSucceeded); + void OpenWebSocket(); + void BindSocketEvents(); + void EmitError(const FString& ErrorMessage); + void TryBroadcastSceneFrame(const FString& Message); + bool TryAssembleFullScene( + const TSharedPtr& DataObject, + int32 Sequence, + int32 TotalRecords, + TSharedPtr& OutAssembledDataObject + ); + void BroadcastScenePayload( + const FString& UpdateType, + int32 Sequence, + const TSharedPtr& DataObject + ); + void BroadcastLegacySubmarineCableEvents(const TSharedPtr& DataObject); + FString BuildWsUrl() const; + FString BuildHttpUrl(const FString& Path) const; + void SendJsonString(const FString& JsonPayload); +}; diff --git a/unreal/UE57_WebSocket_Integration.md b/unreal/UE57_WebSocket_Integration.md new file mode 100644 index 00000000..d57001b9 --- /dev/null +++ b/unreal/UE57_WebSocket_Integration.md @@ -0,0 +1,189 @@ +# UE5.7 连接 `ue_scene` 数据接口实操 + +## 目标 + +把当前后端的 `ue_scene` WebSocket 数据接到 UE5.7 工程里。 + +当前后端入口: + +- [backend/app/api/v1/websocket.py](/D:/work/planet/planet/backend/app/api/v1/websocket.py:1) +- [backend/app/core/websocket/ue_scene.py](/D:/work/planet/planet/backend/app/core/websocket/ue_scene.py:1) + +本次补的 UE 侧最小 C++ 文件: + +- [Planet.Build.cs](/D:/work/planet/planet/unreal/Source/Planet/Planet.Build.cs:1) +- [UeSceneWebSocketClient.h](/D:/work/planet/planet/unreal/Source/Planet/Public/UeSceneWebSocketClient.h:1) +- [UeSceneWebSocketClient.cpp](/D:/work/planet/planet/unreal/Source/Planet/Private/UeSceneWebSocketClient.cpp:1) + +## 这一步先做什么 + +先做最小闭环: + +1. UE 登录后端拿 token +2. UE 连接 `/ws` +3. UE 订阅 `ue_scene` +4. UE 发送 `sync_request` +5. UE 收到原始 JSON + +先不急着做完整场景生成。 + +## 第 1 步:把代码放进你的 UE5.7 工程 + +你的真实 UE 工程里应该有: + +```text +YourProject/ + Source/ + YourProject/ +``` + +把这里的示例代码内容按模块名改进去: + +- `Planet.Build.cs` +- `UeSceneWebSocketClient.h` +- `UeSceneWebSocketClient.cpp` + +如果你的工程模块名不是 `Planet`,要把: + +- `PLANET_API` + +改成你自己的模块导出宏。 + +## 第 2 步:启用插件 + +在 UE5.7 编辑器里启用: + +- `WebSockets` +- `HTTP` +- `Json Blueprint Utilities`(可选) + +启用后重启编辑器。 + +## 第 3 步:重新生成并编译工程 + +如果你是 Visual Studio 工作流: + +1. 右键 `.uproject` +2. 选择“生成 Visual Studio 项目文件” +3. 打开 `.sln` +4. 编译 `Development Editor` + +## 第 4 步:把组件挂到 `BP_GlobeController` + +编译成功后: + +1. 打开 `BP_GlobeController` +2. 点击 `添加` +3. 搜索: + - `Ue Scene Web Socket Client` +4. 把这个组件挂进去 + +## 第 5 步:填写后端参数 + +选中这个组件,在 `细节` 面板里填写: + +- `Backend Base Url`:`http://127.0.0.1:8000` +- `Username`:你的后端用户名 +- `Password`:你的后端密码 + +默认频道已经是: + +- `ue_scene` + +## 第 6 步:在蓝图里接事件 + +在 `BP_GlobeController` 里接这几个事件: + +### `On Connected Changed` + +连接成功后: + +1. 调 `Send Subscribe` +2. 再调 `Send Initial Sync Request` + +### `On Raw Message` + +先不要急着解析全部结构。 + +第一步只做: + +1. 把原始 JSON 打印出来 +2. 或保存到字符串变量 + +这样先确认 UE 确实收到后端消息。 + +### `On Error` + +把错误字符串打印到屏幕和日志里。 + +## 蓝图推荐连法 + +### 开始游戏时 + +1. `开始游戏时` +2. 调组件的 `Connect To Backend` + +### 连接成功后 + +1. `On Connected Changed` +2. `分支` +3. 如果 `bConnected = true` +4. 调 `Send Subscribe` +5. 调 `Send Initial Sync Request` + +### 收消息后 + +1. `On Raw Message` +2. `打印字符串` + +## 第 7 步:先验证消息类型 + +你应该先在屏幕或输出日志里确认这些消息已经进来了: + +1. `connection_established` +2. `subscription_confirmed` +3. `data_frame` + +其中 `data_frame` 的第一条应是: + +- `channel = ue_scene` +- `update_type = full` + +## 第 8 步:再接 JSON 解析 + +当前这一步不建议纯蓝图手搓整个 JSON 结构。 + +更稳的是: + +1. 先在 C++ 里加一个简单解析层 +2. 只抽出 `type` +3. 只抽出 `data.update_type` +4. 只抽出 `payload.layers.gpu_clusters.items` + +先把 `gpu_clusters` 画出来,再扩到别的层。 + +## 当前阶段建议 + +先完成这 3 件事就够了: + +1. UE 能成功连接后端 +2. UE 能收到 `full` 快照 +3. UE 能把原始 JSON 打印出来 + +做到这一步,我们再继续补: + +- JSON 结构体 +- `ApplyFullScene` +- `ApplyIncrementalScene` +- `GeoToWorld` + +## 一句最重要的提醒 + +现在先别在 UE 里同时做: + +- WebSocket 连接 +- 全量 JSON 解析 +- 点位生成 +- 海缆生成 + +一步一步来,先确认“连通”,再做“可视化”。这样最不容易卡死。