From 858a5b4b934bae04ab2b018152490c77ec4e3884 Mon Sep 17 00:00:00 2001 From: AskaEth Date: Sat, 25 Jul 2026 18:07:48 +0800 Subject: [PATCH] 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 --- config/settings.py | 1 + server/api.py | 19 ++++ server/telemetry/listener.py | 39 +++++++ static/js/pages/settings.js | 211 +++++++++++++---------------------- 4 files changed, 139 insertions(+), 131 deletions(-) diff --git a/config/settings.py b/config/settings.py index 0fa0394..6bc0a2c 100644 --- a/config/settings.py +++ b/config/settings.py @@ -20,6 +20,7 @@ DEFAULT_CONFIG: dict[str, Any] = { "telemetry_host": "0.0.0.0", "server_host": "0.0.0.0", "server_port": 9527, + "forward_targets": [], } diff --git a/server/api.py b/server/api.py index bc3170d..7c5575e 100644 --- a/server/api.py +++ b/server/api.py @@ -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} diff --git a/server/telemetry/listener.py b/server/telemetry/listener.py index f43ecf5..695c5e0 100644 --- a/server/telemetry/listener.py +++ b/server/telemetry/listener.py @@ -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) diff --git a/static/js/pages/settings.js b/static/js/pages/settings.js index 512c157..d6778de 100644 --- a/static/js/pages/settings.js +++ b/static/js/pages/settings.js @@ -8,10 +8,20 @@ const PageSettings = { async render() { const cfg = await API.getConfig(); const games = await API.getGames(); - this._el.innerHTML = this._template(cfg, games); + const forwards = await this._getForwards(); + this._el.innerHTML = this._template(cfg, games, forwards); this._bindEvents(); }, + async _getForwards() { + try { const res = await fetch('/api/forward'); return await res.json(); } + catch (e) { return []; } + }, + + async _saveForwards(targets) { + await fetch('/api/forward', { method: 'PUT', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ targets }) }); + }, + _bindEvents() { document.getElementById('settings-telemetry-port')?.addEventListener('change', async (e) => { await API.updateConfig({ telemetry_port: parseInt(e.target.value) || 20777 }); @@ -19,166 +29,105 @@ const PageSettings = { }); document.getElementById('settings-restart-telemetry')?.addEventListener('click', async () => { - await API.stopTelemetry(); - await API.startTelemetry(); + await API.stopTelemetry(); await API.startTelemetry(); Toast.show('遥测监听已重启', 'success'); }); + document.getElementById('btn-add-forward')?.addEventListener('click', () => this._addForward()); + document.getElementById('forward-list')?.addEventListener('click', async (e) => { + if (e.target.closest('.btn-del-forward')) { + e.target.closest('.forward-row').remove(); + this._collectAndSave(); + } + }); + document.getElementById('btn-save-forwards')?.addEventListener('click', () => this._collectAndSave()); + document.getElementById('btn-import-plugin')?.addEventListener('click', () => { - const input = document.createElement('input'); - input.type = 'file'; - input.accept = '.tsp'; + const input = document.createElement('input'); input.type = 'file'; input.accept = '.tsp'; input.onchange = async (e) => { - const file = e.target.files[0]; - if (!file) return; - const formData = new FormData(); - formData.append('file', file); - try { - const res = await fetch('/api/games/install', { method: 'POST', body: formData }); - const data = await res.json(); - if (data.ok) { Toast.show('游戏插件安装成功', 'success'); this.render(); } - else Toast.show('安装失败', 'error'); - } catch (err) { Toast.show('安装失败', 'error'); } - }; - input.click(); + const file = e.target.files[0]; if (!file) return; + const fd = new FormData(); fd.append('file', file); + const res = await fetch('/api/games/install', { method: 'POST', body: fd }); + const data = await res.json(); + if (data.ok) { Toast.show('插件安装成功', 'success'); this.render(); } + else Toast.show('安装失败', 'error'); + }; input.click(); }); document.getElementById('btn-reload-plugins')?.addEventListener('click', async () => { - await API.reloadGamePlugins(); - Toast.show('插件已重新加载', 'success'); - this.render(); + await API.reloadGamePlugins(); Toast.show('插件已重新加载', 'success'); this.render(); }); document.getElementById('settings-game-list')?.addEventListener('click', async (e) => { const exportBtn = e.target.closest('.btn-export-plugin'); const removeBtn = e.target.closest('.btn-remove-plugin'); if (exportBtn) { - const pluginId = exportBtn.dataset.pluginId; - const a = document.createElement('a'); - a.href = `/api/games/${pluginId}/export`; - a.download = `${pluginId}.tsp`; - a.click(); - Toast.show('正在下载 .tsp 文件...', 'info'); + const pid = exportBtn.dataset.pluginId; + const a = document.createElement('a'); a.href = '/api/games/' + pid + '/export'; a.download = pid + '.tsp'; a.click(); } if (removeBtn) { - const pluginId = removeBtn.dataset.pluginId; - if (confirm('确定移除这个游戏插件?')) { - await API.removeGamePlugin(pluginId); - Toast.show('插件已移除', 'info'); - this.render(); - } + if (confirm('移除?')) { await API.removeGamePlugin(removeBtn.dataset.pluginId); Toast.show('已移除', 'info'); this.render(); } } }); document.getElementById('btn-import-dashboard')?.addEventListener('click', () => { - const input = document.createElement('input'); - input.type = 'file'; - input.accept = '.tsd'; + const input = document.createElement('input'); input.type = 'file'; input.accept = '.tsd'; input.onchange = async (e) => { - const file = e.target.files[0]; - if (!file) return; - const formData = new FormData(); - formData.append('file', file); - try { - const res = await fetch('/api/dashboards/import', { method: 'POST', body: formData }); - const data = await res.json(); - if (data.ok) Toast.show('仪表盘导入成功', 'success'); - else Toast.show('导入失败', 'error'); - } catch (err) { Toast.show('导入失败', 'error'); } - }; - input.click(); + const file = e.target.files[0]; if (!file) return; + const fd = new FormData(); fd.append('file', file); + const res = await fetch('/api/dashboards/import', { method: 'POST', body: fd }); + const data = await res.json(); + Toast.show(data.ok ? '导入成功' : '导入失败', data.ok ? 'success' : 'error'); + }; input.click(); }); }, - _gamePluginListHtml(games) { - if (!games || games.length === 0) { - return '
暂无游戏插件
'; - } - return games.map(g => ` -暂无转发目标,点击下方添加
'; + return forwards.map(f => '').join(''); + }, -暂无游戏插件
'; + return games.map(g => 'manifest.json + parser.py 的文件夹manifest.json 定义元信息,parser.py 需有 get_parser() 函数game_id() 和 parse(data, addr) 方法.tsp 即可导入TelemetryData 字段参考见 server/telemetry/data.py
- 将收到的游戏原始 UDP 数据包完整转发到下游物理外设