feat: UDP data forwarding to multiple downstream devices

- DataForwarder class in listener: sends raw packets to all enabled targets
- REST API: GET/PUT /api/forward for managing targets
- Settings page: add/remove/edit forward targets with host/port/name/enabled
- Status endpoint includes forward_active count
- Config persisted in data/config.json via forward_targets array
This commit is contained in:
2026-07-25 18:07:48 +08:00
parent 8d7cf09bb5
commit 858a5b4b93
4 changed files with 139 additions and 131 deletions
+19
View File
@@ -31,6 +31,7 @@ async def get_status():
"packet_count": telemetry_listener.packet_count,
"last_packet_time": telemetry_listener.last_packet_time,
"selected_game_id": cfg.get("selected_game_id"),
"forward_active": telemetry_listener.forwarder.active,
"latest_data": latest.to_dict() if latest else None,
}
@@ -310,3 +311,21 @@ async def api_export_game_plugin(plugin_id: str):
async def api_reload_game_plugins():
game_plugin_manager.reload()
return {"ok": True, "count": len(game_plugin_manager.list_all())}
# ---- Data Forwarding ----
@router.get("/forward")
async def api_list_forward_targets():
cfg = get_config()
return cfg.get("forward_targets", [])
@router.put("/forward")
async def api_update_forward_targets(data: dict[str, Any]):
targets = data.get("targets", [])
cfg = get_config()
cfg["forward_targets"] = targets
from config.settings import save_config
save_config(cfg)
telemetry_listener.forwarder.reload_targets()
return {"ok": True, "active": telemetry_listener.forwarder.active}
+39
View File
@@ -1,6 +1,7 @@
from __future__ import annotations
import asyncio
import socket
import time
from typing import Any, Callable
@@ -13,6 +14,39 @@ from utils.logger import get_logger
logger = get_logger(__name__)
class DataForwarder:
def __init__(self):
self._sock: socket.socket | None = None
self._targets: list[dict[str, Any]] = []
def reload_targets(self):
cfg = get_config()
self._targets = [
t for t in cfg.get("forward_targets", [])
if t.get("enabled", True) and t.get("host") and t.get("port")
]
if self._targets and not self._sock:
self._sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
self._sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
elif not self._targets and self._sock:
self._sock.close()
self._sock = None
logger.info("Forwarder reloaded: %d targets", len(self._targets))
def forward(self, data: bytes):
if not self._sock or not self._targets:
return
for t in self._targets:
try:
self._sock.sendto(data, (t["host"], t["port"]))
except Exception as e:
logger.error("Forward to %s:%d failed: %s", t["host"], t["port"], e)
@property
def active(self) -> int:
return len(self._targets)
class TelemetryListener:
def __init__(self):
self._transport: asyncio.DatagramTransport | None = None
@@ -23,6 +57,7 @@ class TelemetryListener:
self._latest_data: TelemetryData | None = None
self._last_packet_time: float = 0.0
self._packet_count: int = 0
self.forwarder = DataForwarder()
@property
def is_running(self) -> bool:
@@ -71,6 +106,8 @@ class TelemetryListener:
)
self._running = True
self.forwarder.reload_targets()
if parser_key and parser_key in PARSER_MAP:
self._parser = self._get_parser(parser_key)
logger.info("Telemetry listener started on %s:%d [parser=%s]", host, port, parser_key)
@@ -101,6 +138,8 @@ class TelemetryListener:
self._last_packet_time = time.time()
self._packet_count += 1
self.forwarder.forward(data)
td = None
if self._parser:
td = self._parser.parse(data, addr)