From f715cb85e6c991a63e2cc453f66556abaf4dba03 Mon Sep 17 00:00:00 2001 From: AskaEth Date: Mon, 27 Jul 2026 19:04:16 +0800 Subject: [PATCH] feat: TCP listener + protocol field in manifest, dual UDP/TCP support --- server/telemetry/listener.py | 58 +++++++++++++++++++++---------- server/telemetry/tcp_listener.py | 59 ++++++++++++++++++++++++++++++++ 2 files changed, 99 insertions(+), 18 deletions(-) create mode 100644 server/telemetry/tcp_listener.py diff --git a/server/telemetry/listener.py b/server/telemetry/listener.py index 0022025..b9ffb09 100644 --- a/server/telemetry/listener.py +++ b/server/telemetry/listener.py @@ -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,25 +93,41 @@ class TelemetryListener: host = cfg.get("telemetry_host", "0.0.0.0") port = cfg.get("telemetry_port", 20777) - loop = asyncio.get_event_loop() - try: - self._transport, _ = await loop.create_datagram_endpoint( - lambda: _TelemetryProtocol(self), - local_addr=(host, port), - ) - self._running = True - self.forwarder.reload_targets() + 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" - 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") - return True - except Exception as e: - logger.error("Failed to start telemetry listener: %s", e) - return False + 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( + lambda: _TelemetryProtocol(self), + local_addr=(host, port), + ) + self._running = True + self.forwarder.reload_targets() + 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 UDP listener: %s", e) + return False def stop(self): self._running = False @@ -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 diff --git a/server/telemetry/tcp_listener.py b/server/telemetry/tcp_listener.py new file mode 100644 index 0000000..170fea4 --- /dev/null +++ b/server/telemetry/tcp_listener.py @@ -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()