""" 消息路由器模块 处理客户端间的消息传递 """ import time import threading import json from typing import Dict, Any, List, Optional class MessageRouter: """ 消息路由器 处理客户端间的消息传递 """ def __init__(self, plugin): """ 初始化消息路由器 Args: plugin: 实时通信插件实例 """ self.plugin = plugin self.enabled = False self.initialized = False # 消息路由配置 self.message_config = { "enable_message_filtering": True, "max_message_size": 65536, # 64KB "enable_message_queue": True, "message_queue_size": 10000, "enable_broadcast_optimization": True, "enable_message_compression": True, "enable_rate_limiting": True, "max_messages_per_second_per_client": 100, "enable_message_logging": False } # 消息路由状态 self.router_state = { "is_processing": False, "pending_messages": 0, "active_routes": 0, "last_message_time": 0.0 } # 消息路由统计 self.message_stats = { "messages_routed": 0, "messages_broadcast": 0, "messages_filtered": 0, "messages_dropped": 0, "bytes_routed": 0, "routing_errors": 0 } # 消息处理器存储 self.message_handlers = {} # 消息队列 self.message_queue = [] self.message_queue_lock = threading.RLock() # 路由表 self.routing_table = {} # 回调函数 self.message_callbacks = { "message_routed": [], "message_broadcast": [], "message_filtered": [], "message_error": [] } # 时间戳记录 self.last_queue_process = 0.0 self.last_stats_reset = 0.0 print("✓ 消息路由器已创建") def initialize(self) -> bool: """ 初始化消息路由器 Returns: 是否初始化成功 """ try: print("正在初始化消息路由器...") # 注册默认消息处理器 self._register_default_handlers() self.initialized = True print("✓ 消息路由器初始化完成") return True except Exception as e: print(f"✗ 消息路由器初始化失败: {e}") self.message_stats["routing_errors"] += 1 import traceback traceback.print_exc() return False def _register_default_handlers(self): """注册默认消息处理器""" try: # 注册系统消息处理器 self.register_message_handler("ping", self._handle_ping_message) self.register_message_handler("pong", self._handle_pong_message) self.register_message_handler("chat", self._handle_chat_message) self.register_message_handler("join_room", self._handle_join_room_message) self.register_message_handler("leave_room", self._handle_leave_room_message) self.register_message_handler("broadcast", self._handle_broadcast_message) print("✓ 默认消息处理器已注册") except Exception as e: print(f"✗ 默认消息处理器注册失败: {e}") self.message_stats["routing_errors"] += 1 def enable(self) -> bool: """ 启用消息路由器 Returns: 是否启用成功 """ try: if not self.initialized: print("✗ 消息路由器未初始化") return False self.enabled = True print("✓ 消息路由器已启用") return True except Exception as e: print(f"✗ 消息路由器启用失败: {e}") self.message_stats["routing_errors"] += 1 import traceback traceback.print_exc() return False def disable(self): """禁用消息路由器""" try: self.enabled = False # 清空消息队列 with self.message_queue_lock: self.message_queue.clear() print("✓ 消息路由器已禁用") except Exception as e: print(f"✗ 消息路由器禁用失败: {e}") self.message_stats["routing_errors"] += 1 import traceback traceback.print_exc() def finalize(self): """清理消息路由器资源""" try: # 禁用消息路由器 if self.enabled: self.disable() # 清理回调和处理器 self.message_callbacks.clear() self.message_handlers.clear() self.initialized = False print("✓ 消息路由器资源已清理") except Exception as e: print(f"✗ 消息路由器资源清理失败: {e}") import traceback traceback.print_exc() def update(self, dt: float): """ 更新消息路由器状态 Args: dt: 时间增量(秒) """ try: if not self.enabled: return current_time = time.time() # 处理消息队列 if self.message_config["enable_message_queue"]: self._process_message_queue() self.router_state["last_message_time"] = current_time except Exception as e: print(f"✗ 消息路由器更新失败: {e}") self.message_stats["routing_errors"] += 1 import traceback traceback.print_exc() def _process_message_queue(self): """处理消息队列""" try: if not self.message_queue: return with self.message_queue_lock: messages_to_process = self.message_queue[:] self.message_queue.clear() for message_data in messages_to_process: try: client_id = message_data.get("client_id") message = message_data.get("message") # 处理消息 self._route_message_internal(client_id, message) except Exception as e: print(f"✗ 队列消息处理失败: {e}") self.message_stats["routing_errors"] += 1 self.router_state["pending_messages"] = len(self.message_queue) except Exception as e: print(f"✗ 消息队列处理失败: {e}") self.message_stats["routing_errors"] += 1 def register_message_handler(self, message_type: str, handler: callable): """ 注册消息处理器 Args: message_type: 消息类型 handler: 处理器函数 """ try: self.message_handlers[message_type] = handler print(f"✓ 消息处理器已注册: {message_type}") except Exception as e: print(f"✗ 消息处理器注册失败: {e}") self.message_stats["routing_errors"] += 1 def unregister_message_handler(self, message_type: str): """ 注销消息处理器 Args: message_type: 消息类型 """ try: if message_type in self.message_handlers: del self.message_handlers[message_type] print(f"✓ 消息处理器已注销: {message_type}") else: print(f"✗ 消息处理器不存在: {message_type}") except Exception as e: print(f"✗ 消息处理器注销失败: {e}") self.message_stats["routing_errors"] += 1 def route_message(self, client_id: str, message: Dict[str, Any]) -> bool: """ 路由消息 Args: client_id: 客户端ID message: 消息数据 Returns: 是否路由成功 """ try: if not self.enabled: return False # 检查消息大小 if self.message_config["enable_message_filtering"]: message_size = len(json.dumps(message, ensure_ascii=False).encode('utf-8')) if message_size > self.message_config["max_message_size"]: print(f"✗ 消息过大: {message_size} bytes") self.message_stats["messages_dropped"] += 1 return False # 如果启用消息队列,将消息加入队列 if self.message_config["enable_message_queue"]: with self.message_queue_lock: if len(self.message_queue) >= self.message_config["message_queue_size"]: print("✗ 消息队列已满,丢弃消息") self.message_stats["messages_dropped"] += 1 return False self.message_queue.append({ "client_id": client_id, "message": message, "timestamp": time.time() }) self.router_state["pending_messages"] = len(self.message_queue) return True else: # 直接处理消息 return self._route_message_internal(client_id, message) except Exception as e: print(f"✗ 消息路由失败: {e}") self.message_stats["routing_errors"] += 1 return False def _route_message_internal(self, client_id: str, message: Dict[str, Any]) -> bool: """ 内部消息路由处理 Args: client_id: 客户端ID message: 消息数据 Returns: 是否路由成功 """ try: # 更新统计信息 self.message_stats["messages_routed"] += 1 message_size = len(json.dumps(message, ensure_ascii=False).encode('utf-8')) self.message_stats["bytes_routed"] += message_size # 触发消息路由回调 self._trigger_message_callback("message_routed", { "client_id": client_id, "message": message, "timestamp": time.time() }) # 获取消息类型 message_type = message.get("type", "unknown") # 查找对应的消息处理器 if message_type in self.message_handlers: try: # 调用消息处理器 handler_result = self.message_handlers[message_type](client_id, message) return handler_result if handler_result is not None else True except Exception as e: print(f"✗ 消息处理器执行失败: {message_type} - {e}") self.message_stats["routing_errors"] += 1 return False else: # 默认处理:广播给所有客户端(除了发送者) return self._handle_default_message(client_id, message) except Exception as e: print(f"✗ 内部消息路由处理失败: {e}") self.message_stats["routing_errors"] += 1 return False def _handle_default_message(self, client_id: str, message: Dict[str, Any]) -> bool: """ 处理默认消息(广播给其他客户端) Args: client_id: 发送客户端ID message: 消息数据 Returns: 是否处理成功 """ try: # 广播消息给所有其他客户端 if self.plugin.websocket_server: broadcast_success = self.plugin.websocket_server.broadcast_message( message, exclude_client_id=client_id) if broadcast_success: self.message_stats["messages_broadcast"] += 1 return True else: return False else: return False except Exception as e: print(f"✗ 默认消息处理失败: {e}") self.message_stats["routing_errors"] += 1 return False # 默认消息处理器 def _handle_ping_message(self, client_id: str, message: Dict[str, Any]) -> bool: """ 处理ping消息 Args: client_id: 客户端ID message: 消息数据 Returns: 是否处理成功 """ try: # 更新客户端心跳 if self.plugin.client_manager: self.plugin.client_manager.update_client_heartbeat(client_id) # 回复pong消息 pong_message = { "type": "pong", "timestamp": time.time(), "ping_timestamp": message.get("timestamp", 0) } if self.plugin.websocket_server: return self.plugin.websocket_server.send_message_to_client(client_id, pong_message) else: return False except Exception as e: print(f"✗ ping消息处理失败: {e}") self.message_stats["routing_errors"] += 1 return False def _handle_pong_message(self, client_id: str, message: Dict[str, Any]) -> bool: """ 处理pong消息 Args: client_id: 客户端ID message: 消息数据 Returns: 是否处理成功 """ try: # 更新客户端心跳 if self.plugin.client_manager: self.plugin.client_manager.update_client_heartbeat(client_id) return True except Exception as e: print(f"✗ pong消息处理失败: {e}") self.message_stats["routing_errors"] += 1 return False def _handle_chat_message(self, client_id: str, message: Dict[str, Any]) -> bool: """ 处理聊天消息 Args: client_id: 客户端ID message: 消息数据 Returns: 是否处理成功 """ try: # 获取客户端信息 client_data = self.plugin.client_manager.get_client(client_id) if self.plugin.client_manager else None if not client_data: return False # 添加发送者信息 message["sender_id"] = client_id message["sender_name"] = client_data.get("username", "Unknown") message["sender_time"] = time.time() # 广播聊天消息 if self.plugin.websocket_server: broadcast_success = self.plugin.websocket_server.broadcast_message(message) if broadcast_success: self.message_stats["messages_broadcast"] += 1 return broadcast_success else: return False except Exception as e: print(f"✗ 聊天消息处理失败: {e}") self.message_stats["routing_errors"] += 1 return False def _handle_join_room_message(self, client_id: str, message: Dict[str, Any]) -> bool: """ 处理加入房间消息 Args: client_id: 客户端ID message: 消息数据 Returns: 是否处理成功 """ try: room_id = message.get("room_id") if not room_id: return False # 使用房间管理器处理加入房间请求 if self.plugin.room_manager: join_success = self.plugin.room_manager.add_client_to_room(room_id, client_id) if join_success: # 通知客户端加入成功 response_message = { "type": "room_joined", "room_id": room_id, "timestamp": time.time() } if self.plugin.websocket_server: self.plugin.websocket_server.send_message_to_client(client_id, response_message) return join_success else: return False except Exception as e: print(f"✗ 加入房间消息处理失败: {e}") self.message_stats["routing_errors"] += 1 return False def _handle_leave_room_message(self, client_id: str, message: Dict[str, Any]) -> bool: """ 处理离开房间消息 Args: client_id: 客户端ID message: 消息数据 Returns: 是否处理成功 """ try: room_id = message.get("room_id") if not room_id: return False # 使用房间管理器处理离开房间请求 if self.plugin.room_manager: leave_success = self.plugin.room_manager.remove_client_from_room(room_id, client_id) if leave_success: # 通知客户端离开成功 response_message = { "type": "room_left", "room_id": room_id, "timestamp": time.time() } if self.plugin.websocket_server: self.plugin.websocket_server.send_message_to_client(client_id, response_message) return leave_success else: return False except Exception as e: print(f"✗ 离开房间消息处理失败: {e}") self.message_stats["routing_errors"] += 1 return False def _handle_broadcast_message(self, client_id: str, message: Dict[str, Any]) -> bool: """ 处理广播消息 Args: client_id: 客户端ID message: 消息数据 Returns: 是否处理成功 """ try: # 广播消息给所有客户端 if self.plugin.websocket_server: broadcast_success = self.plugin.websocket_server.broadcast_message(message) if broadcast_success: self.message_stats["messages_broadcast"] += 1 return broadcast_success else: return False except Exception as e: print(f"✗ 广播消息处理失败: {e}") self.message_stats["routing_errors"] += 1 return False def send_direct_message(self, from_client_id: str, to_client_id: str, message: Dict[str, Any]) -> bool: """ 发送直接消息 Args: from_client_id: 发送客户端ID to_client_id: 接收客户端ID message: 消息数据 Returns: 是否发送成功 """ try: if not self.enabled: return False # 添加发送者信息 message["sender_id"] = from_client_id message["timestamp"] = time.time() # 发送消息到目标客户端 if self.plugin.websocket_server: send_success = self.plugin.websocket_server.send_message_to_client(to_client_id, message) if send_success: self.message_stats["messages_routed"] += 1 return send_success else: return False except Exception as e: print(f"✗ 直接消息发送失败: {e}") self.message_stats["routing_errors"] += 1 return False def broadcast_to_room(self, room_id: str, message: Dict[str, Any], exclude_client_id: str = None) -> bool: """ 广播消息到房间 Args: room_id: 房间ID message: 消息数据 exclude_client_id: 要排除的客户端ID Returns: 是否广播成功 """ try: if not self.enabled: return False # 使用房间管理器获取房间内的客户端 if self.plugin.room_manager: room_clients = self.plugin.room_manager.get_room_clients(room_id) if not room_clients: return False success_count = 0 for client_id in room_clients: # 跳过排除的客户端 if client_id == exclude_client_id: continue # 发送消息 if self.plugin.websocket_server: send_success = self.plugin.websocket_server.send_message_to_client(client_id, message) if send_success: success_count += 1 return success_count > 0 else: return False except Exception as e: print(f"✗ 房间广播失败: {e}") self.message_stats["routing_errors"] += 1 return False def get_message_stats(self) -> Dict[str, Any]: """ 获取消息统计信息 Returns: 消息统计字典 """ return { "state": self.router_state.copy(), "stats": self.message_stats.copy(), "config": self.message_config.copy() } def reset_stats(self): """重置消息统计信息""" try: self.message_stats = { "messages_routed": 0, "messages_broadcast": 0, "messages_filtered": 0, "messages_dropped": 0, "bytes_routed": 0, "routing_errors": 0 } print("✓ 消息统计信息已重置") except Exception as e: print(f"✗ 消息统计信息重置失败: {e}") def set_message_config(self, config: Dict[str, Any]) -> bool: """ 设置消息配置 Args: config: 消息配置字典 Returns: 是否设置成功 """ try: self.message_config.update(config) print(f"✓ 消息配置已更新: {self.message_config}") return True except Exception as e: print(f"✗ 消息配置设置失败: {e}") return False def get_message_config(self) -> Dict[str, Any]: """ 获取消息配置 Returns: 消息配置字典 """ return self.message_config.copy() def _trigger_message_callback(self, callback_type: str, data: Dict[str, Any]): """ 触发消息回调 Args: callback_type: 回调类型 data: 回调数据 """ try: if callback_type in self.message_callbacks: for callback in self.message_callbacks[callback_type]: try: callback(data) except Exception as e: print(f"✗ 消息回调执行失败: {callback_type} - {e}") except Exception as e: print(f"✗ 消息回调触发失败: {e}") def register_message_callback(self, callback_type: str, callback: callable): """ 注册消息回调 Args: callback_type: 回调类型 callback: 回调函数 """ try: if callback_type in self.message_callbacks: self.message_callbacks[callback_type].append(callback) print(f"✓ 消息回调已注册: {callback_type}") else: print(f"✗ 无效的回调类型: {callback_type}") except Exception as e: print(f"✗ 消息回调注册失败: {e}") def unregister_message_callback(self, callback_type: str, callback: callable): """ 注销消息回调 Args: callback_type: 回调类型 callback: 回调函数 """ try: if callback_type in self.message_callbacks: if callback in self.message_callbacks[callback_type]: self.message_callbacks[callback_type].remove(callback) print(f"✓ 消息回调已注销: {callback_type}") else: print(f"✗ 无效的回调类型: {callback_type}") except Exception as e: print(f"✗ 消息回调注销失败: {e}")