EG/plugins/user/realtime_communication/messaging/message_router.py
2025-10-30 15:01:29 +08:00

771 lines
26 KiB
Python

"""
消息路由器模块
处理客户端间的消息传递
"""
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}")