feat: TCP listener + protocol field in manifest, dual UDP/TCP support
This commit is contained in:
@@ -8,6 +8,7 @@ from typing import Any, Callable
|
||||
from config.settings import get_config
|
||||
from server.telemetry.parsers import PARSER_MAP, BaseParser
|
||||
from server.telemetry.data import TelemetryData
|
||||
from server.telemetry.tcp_listener import TCPListener
|
||||
from server.game_manager import game_plugin_manager
|
||||
from utils.logger import get_logger
|
||||
|
||||
@@ -50,7 +51,9 @@ class DataForwarder:
|
||||
class TelemetryListener:
|
||||
def __init__(self):
|
||||
self._transport: asyncio.DatagramTransport | None = None
|
||||
self._tcp: TCPListener | None = None
|
||||
self._running = False
|
||||
self._current_protocol = "udp"
|
||||
self._callbacks: list[Callable[[TelemetryData], None]] = []
|
||||
self._parser: BaseParser | None = None
|
||||
self._parser_cache: dict[str, BaseParser] = {}
|
||||
@@ -90,6 +93,28 @@ class TelemetryListener:
|
||||
host = cfg.get("telemetry_host", "0.0.0.0")
|
||||
port = cfg.get("telemetry_port", 20777)
|
||||
|
||||
selected_game = cfg.get("selected_game_id")
|
||||
if selected_game:
|
||||
self.set_parser_for_game(selected_game)
|
||||
gp = game_plugin_manager.get(selected_game)
|
||||
if gp and gp.manifest_path:
|
||||
import json as _json
|
||||
try:
|
||||
with open(f"{gp.manifest_path}/manifest.json", "r") as f:
|
||||
manifest = _json.load(f)
|
||||
self._current_protocol = manifest.get("protocol", "udp")
|
||||
except Exception:
|
||||
self._current_protocol = "udp"
|
||||
|
||||
if self._current_protocol == "tcp":
|
||||
self._tcp = TCPListener(self._handle_packet)
|
||||
ok = await self._tcp.start(host, port)
|
||||
if ok:
|
||||
self._running = True
|
||||
self.forwarder.reload_targets()
|
||||
logger.info("TCP listener started on %s:%d [game=%s]", host, port, selected_game or "none")
|
||||
return ok
|
||||
else:
|
||||
loop = asyncio.get_event_loop()
|
||||
try:
|
||||
self._transport, _ = await loop.create_datagram_endpoint(
|
||||
@@ -98,16 +123,10 @@ class TelemetryListener:
|
||||
)
|
||||
self._running = True
|
||||
self.forwarder.reload_targets()
|
||||
|
||||
selected_game = cfg.get("selected_game_id")
|
||||
if selected_game:
|
||||
self.set_parser_for_game(selected_game)
|
||||
|
||||
logger.info("Telemetry listener started on %s:%d [game=%s]",
|
||||
host, port, selected_game or "none")
|
||||
logger.info("UDP listener started on %s:%d [game=%s]", host, port, selected_game or "none")
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error("Failed to start telemetry listener: %s", e)
|
||||
logger.error("Failed to start UDP listener: %s", e)
|
||||
return False
|
||||
|
||||
def stop(self):
|
||||
@@ -116,6 +135,9 @@ class TelemetryListener:
|
||||
if self._transport:
|
||||
self._transport.close()
|
||||
self._transport = None
|
||||
if self._tcp:
|
||||
asyncio.ensure_future(self._tcp.stop())
|
||||
self._tcp = None
|
||||
if self.forwarder._sock:
|
||||
self.forwarder._sock.close()
|
||||
self.forwarder._sock = None
|
||||
|
||||
@@ -0,0 +1,59 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import socket
|
||||
import struct
|
||||
import time
|
||||
from typing import Any, Callable
|
||||
|
||||
from utils.logger import get_logger
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
|
||||
class TCPListener:
|
||||
def __init__(self, handle_packet):
|
||||
self._server: asyncio.Server | None = None
|
||||
self._handle_packet = handle_packet
|
||||
self._running = False
|
||||
|
||||
@property
|
||||
def is_running(self) -> bool:
|
||||
return self._running
|
||||
|
||||
async def start(self, host: str, port: int) -> bool:
|
||||
if self._running:
|
||||
await self.stop()
|
||||
try:
|
||||
self._server = await asyncio.start_server(
|
||||
self._on_client_connected, host, port
|
||||
)
|
||||
self._running = True
|
||||
logger.info("TCP listener started on %s:%d", host, port)
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error("Failed to start TCP listener: %s", e)
|
||||
return False
|
||||
|
||||
async def stop(self):
|
||||
self._running = False
|
||||
if self._server:
|
||||
self._server.close()
|
||||
await self._server.wait_closed()
|
||||
self._server = None
|
||||
logger.info("TCP listener stopped")
|
||||
|
||||
async def _on_client_connected(self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter):
|
||||
addr = writer.get_extra_info('peername')
|
||||
logger.debug("TCP client connected: %s", addr)
|
||||
try:
|
||||
while self._running:
|
||||
data = await reader.read(65536)
|
||||
if not data:
|
||||
break
|
||||
self._handle_packet(data, addr)
|
||||
except Exception as e:
|
||||
logger.debug("TCP client disconnected: %s - %s", addr, e)
|
||||
finally:
|
||||
writer.close()
|
||||
await writer.wait_closed()
|
||||
Reference in New Issue
Block a user