#!/usr/bin/env python3 # -*- coding: utf-8 -*- import logging import asyncio from typing import Dict, Any from aiohttp import web import json # 导入命令装饰器 try: from core.plugin_command_decorator import plugin_command, command except ImportError: # 回退方案 def plugin_command(name=None, description=None, permissions=None): def decorator(func): return func return decorator command = plugin_command # 导入网络桥接类 - 修正路径 try: from bridges.plugin_network_bridge import PluginNetworkBridge except ImportError: # 如果导入失败,创建一个虚拟类 class PluginNetworkBridge: def __init__(self, plugin_name, internet_service, plugin_bridge): self.plugin_name = plugin_name logger.warning(f"PluginNetworkBridge 不可用,插件 {plugin_name} 将以无网络模式运行") async def register_http_route(self, *args, **kwargs): logger.warning("网络功能不可用,跳过HTTP路由注册") async def register_websocket(self, *args, **kwargs): logger.warning("网络功能不可用,跳过WebSocket注册") async def broadcast_websocket(self, *args, **kwargs): logger.warning("网络功能不可用,无法广播消息") def get_network_info(self): return { 'plugin_name': self.plugin_name, 'registered_routes': [], 'websocket_handlers': [], 'base_url': '网络服务不可用' } async def setup_data_transfer(self, *args, **kwargs): logger.warning("网络功能不可用,跳过数据传输设置") logger = logging.getLogger(__name__) class Plugin: """示例插件 - 展示命令注册和网络交互""" def __init__(self, plugin_name: str, config: Dict, bridge): self.plugin_name = plugin_name self.config = config self.bridge = bridge self.network_bridge = None self.is_running = False logger.debug(f"示例插件初始化: {plugin_name}") async def initialize(self): """初始化插件 - 安全版本""" try: logger.info(f"初始化示例插件: {self.plugin_name}") # 安全地获取网络服务 internet_service = None try: internet_service = self.bridge.service_manager.get_service("internet") logger.debug(f"网络服务获取: {internet_service is not None}") except (ValueError, AttributeError) as e: logger.warning(f"网络服务不可用: {str(e)}") except Exception as e: logger.error(f"获取网络服务时出错: {str(e)}") # 只有在网络服务可用时才设置网络功能 if internet_service: try: # 创建网络桥接 self.network_bridge = PluginNetworkBridge( self.plugin_name, internet_service, self.bridge ) # 注册网络路由 await self._setup_network_routes() logger.info(f"插件网络功能初始化完成: {self.plugin_name}") except Exception as e: logger.error(f"设置网络功能时出错: {str(e)}") logger.info("插件将以无网络模式运行") else: logger.info(f"插件 {self.plugin_name} 将以无网络模式运行") # 创建虚拟网络桥接以便命令能正常工作 self.network_bridge = PluginNetworkBridge(self.plugin_name, None, self.bridge) # 注册事件处理器(不依赖网络服务) self.bridge.subscribe_plugin( self.plugin_name, "event.framework.start", self._handle_framework_start ) self.is_running = True logger.debug(f"示例插件初始化完成: {self.plugin_name}") except Exception as e: logger.error(f"初始化示例插件时出错: {str(e)}", exc_info=True) raise async def _setup_network_routes(self): """设置网络路由 - 安全版本""" try: if not self.network_bridge: logger.warning("网络桥接不可用,跳过路由设置") return # 注册HTTP API端点 await self.network_bridge.register_http_route( "/api/info", self._handle_api_info, methods=["GET"], require_auth=False ) await self.network_bridge.register_http_route( "/api/echo", self._handle_api_echo, methods=["POST"], require_auth=True ) # 注册WebSocket端点 await self.network_bridge.register_websocket( "/chat", self._handle_websocket_chat ) # 设置跨端数据传输 await self.network_bridge.setup_data_transfer( self._handle_cross_platform_data ) logger.info(f"示例插件网络路由设置完成: {self.plugin_name}") except Exception as e: logger.error(f"设置网络路由时出错: {str(e)}", exc_info=True) # 不抛出异常,让插件继续运行 async def register_delayed_routes(self, internet_service): """延迟注册网络路由(在网络服务启动后调用)""" try: logger.info(f"为插件 {self.plugin_name} 延迟注册网络路由") # 重新创建网络桥接,使用真实的网络服务 if internet_service: try: # 重新初始化网络桥接 self.network_bridge = PluginNetworkBridge( self.plugin_name, internet_service, self.bridge ) # 重新设置网络路由 await self._setup_network_routes() logger.info(f"插件 {self.plugin_name} 网络功能重新初始化完成") except Exception as e: logger.error(f"重新初始化网络桥接时出错: {str(e)}") logger.info(f"插件 {self.plugin_name} 将继续使用无网络模式") else: logger.warning(f"网络服务不可用,插件 {self.plugin_name} 保持无网络模式") except Exception as e: logger.error(f"延迟注册网络路由时出错: {str(e)}") async def _handle_api_info(self, request): """处理API信息请求""" try: info = { "plugin_name": self.plugin_name, "version": self.config.get('version', '1.0.0'), "description": self.config.get('description', '示例插件'), "network_info": self.network_bridge.get_network_info() if self.network_bridge else None, "timestamp": asyncio.get_event_loop().time() } return web.json_response(info) except Exception as e: logger.error(f"处理API信息请求时出错: {str(e)}") return web.json_response({"error": str(e)}, status=500) async def _handle_api_echo(self, request): """处理API回显请求""" try: data = await request.json() response = { "plugin_name": self.plugin_name, "echo": data, "timestamp": asyncio.get_event_loop().time() } return web.json_response(response) except Exception as e: logger.error(f"处理API回显请求时出错: {str(e)}") return web.json_response({"error": str(e)}, status=400) async def _handle_websocket_chat(self, ws, request): """处理WebSocket聊天""" try: logger.info(f"WebSocket聊天连接建立: {self.plugin_name}") async for msg in ws: if msg.type == web.WSMsgType.TEXT: try: data = json.loads(msg.data) # 处理不同类型的消息 if data.get('type') == 'message': # 广播消息给所有客户端 if self.network_bridge: await self.network_bridge.broadcast_websocket({ "type": "message", "from": data.get('user', 'anonymous'), "content": data.get('content', ''), "timestamp": asyncio.get_event_loop().time() }) except json.JSONDecodeError: logger.warning(f"收到无效的JSON消息: {msg.data}") elif msg.type == web.WSMsgType.ERROR: logger.error(f"WebSocket错误: {ws.exception()}") except Exception as e: logger.error(f"WebSocket聊天处理出错: {str(e)}") finally: logger.info(f"WebSocket聊天连接关闭: {self.plugin_name}") async def _handle_cross_platform_data(self, event_type: str, data: Dict): """处理跨端数据""" try: logger.info(f"收到跨端数据: {event_type}") # 在这里处理来自其他平台的数据 if event_type == "network.data.receive": # 广播到WebSocket if self.network_bridge: await self.network_bridge.broadcast_websocket({ "type": "cross_platform", "source": data.get('source', 'unknown'), "data": data.get('data', {}), "timestamp": asyncio.get_event_loop().time() }) except Exception as e: logger.error(f"处理跨端数据时出错: {str(e)}") async def _handle_framework_start(self, event_type: str, data: Dict): """处理框架启动事件""" try: logger.info(f"框架启动事件: {event_type}") # 发送欢迎消息 if self.network_bridge: await self.network_bridge.broadcast_websocket({ "type": "system", "message": f"插件 {self.plugin_name} 已启动,框架已就绪", "timestamp": asyncio.get_event_loop().time() }) except Exception as e: logger.error(f"处理框架启动事件时出错: {str(e)}") # 确保所有网络相关方法都检查 network_bridge @plugin_command(name="chat_broadcast", description="向所有聊天客户端广播消息", permissions=["plugin.example.chat.broadcast"]) async def cmd_chat_broadcast(self, *args): """向所有聊天客户端广播消息""" try: if not args: return "❌ 请提供要广播的消息内容" message = " ".join(args) if self.network_bridge: await self.network_bridge.broadcast_websocket({ "type": "broadcast", "from": "system", "content": message, "timestamp": asyncio.get_event_loop().time() }) return f"✅ 已广播消息: {message}" else: return "❌ 网络服务不可用,无法广播消息" except Exception as e: logger.error(f"广播消息时出错: {str(e)}") return f"❌ 广播失败: {str(e)}" @plugin_command(name="network_info", description="显示插件网络信息") async def cmd_network_info(self, *args): """显示插件的网络配置信息""" try: if not self.network_bridge: result = ["🌐 **插件网络信息:**"] result.append("❌ 网络服务不可用") result.append("\n💡 **网络服务状态:**") # 尝试获取网络服务状态 try: internet_service = self.bridge.service_manager.get_service("internet") if internet_service: result.append(" ✅ 网络服务已注册") health_info = await internet_service.check_service_health() result.append(f" 🔄 服务运行: {'✅ 是' if health_info.get('is_running') else '❌ 否'}") result.append(f" 🌐 HTTP活跃: {'✅ 是' if health_info.get('http_active') else '❌ 否'}") else: result.append(" ❌ 网络服务未注册") except: result.append(" ❓ 无法获取网络服务状态") result.append("\n🔧 **建议:**") result.append(" - 检查网络服务启动日志") result.append(" - 使用 'services' 命令查看服务状态") result.append(" - 使用 'netdiag' 命令进行网络诊断") return "\n".join(result) info = self.network_bridge.get_network_info() result = ["🌐 **插件网络信息:**"] result.append(f" 插件名称: {info['plugin_name']}") result.append(f" 基础URL: {info['base_url']}") result.append(f" HTTP路由数: {len(info['registered_routes'])}") result.append(f" WebSocket处理器数: {len(info['websocket_handlers'])}") if info['registered_routes']: result.append("\n📡 **注册的HTTP路由:**") for route in info['registered_routes']: result.append(f" {route['path']} [{','.join(route['methods'])}]") if info['websocket_handlers']: result.append("\n🔗 **注册的WebSocket:**") for ws in info['websocket_handlers']: result.append(f" {ws['path']}") return "\n".join(result) except Exception as e: logger.error(f"获取网络信息时出错: {str(e)}") return f"❌ 获取网络信息失败: {str(e)}" # ... 其余方法保持不变 ... async def shutdown(self): """关闭插件""" try: logger.info(f"关闭示例插件: {self.plugin_name}") self.is_running = False # 只有在网络桥接可用时才发送关闭通知 if self.network_bridge and hasattr(self.network_bridge, 'broadcast_websocket'): try: await self.network_bridge.broadcast_websocket({ "type": "system", "message": f"插件 {self.plugin_name} 正在关闭", "timestamp": asyncio.get_event_loop().time() }) except Exception as e: logger.warning(f"发送关闭通知失败: {str(e)}") # 清理资源 self.bridge.cleanup_plugin_subscriptions(self.plugin_name) logger.debug(f"示例插件关闭完成: {self.plugin_name}") except Exception as e: logger.error(f"关闭示例插件时出错: {str(e)}", exc_info=True)