import asyncio import websockets from config.logger import setup_logging from core.connection import ConnectionHandler from core.utils.util import initialize_modules TAG = __name__ class WebSocketServer: def __init__(self, config: dict): self.config = config self.logger = setup_logging() self.config_lock = asyncio.Lock() modules = initialize_modules( self.logger, self.config, "VAD" in self.config["selected_module"], "ASR" in self.config["selected_module"], "LLM" in self.config["selected_module"], "TTS" in self.config["selected_module"], "Memory" in self.config["selected_module"], "Intent" in self.config["selected_module"], ) self._vad = modules["vad"] if "vad" in modules else None self._asr = modules["asr"] if "asr" in modules else None self._tts = modules["tts"] if "tts" in modules else None self._llm = modules["llm"] if "llm" in modules else None self._intent = modules["intent"] if "intent" in modules else None self._memory = modules["memory"] if "memory" in modules else None self.active_connections = set() async def start(self): server_config = self.config["server"] host = server_config.get("ip", "0.0.0.0") port = int(server_config.get("port", 8000)) async with websockets.serve( self._handle_connection, host, port, process_request=self._http_response ): await asyncio.Future() async def _handle_connection(self, websocket): """处理新连接,每次创建独立的ConnectionHandler""" # 创建ConnectionHandler时传入当前server实例 handler = ConnectionHandler( self.config, self._vad, self._asr, self._llm, self._tts, self._memory, self._intent, self, # 传入当前 WebSocketServer 实例 ) self.active_connections.add(handler) try: await handler.handle_connection(websocket) finally: self.active_connections.discard(handler) async def _http_response(self, websocket, request_headers): # 检查是否为 WebSocket 升级请求 if request_headers.headers.get("connection", "").lower() == "upgrade": # 如果是 WebSocket 请求,返回 None 允许握手继续 return None else: # 如果是普通 HTTP 请求,返回 "server is running" return websocket.respond(200, "Server is running\n")