e6875f0b4b
- 13-service async plugin framework - Textual TUI with CLI fallback - Plugin hot-reload + permission system - Web management panel (aiohttp) - Bridge-based inter-module communication - 10 regression tests Fixes applied: - PBKDF2-SHA256 auth (was plain SHA256) - Auth bypass removed (was allow-all on fail) - Bare excepts replaced with logged errors - CatFramework/DreamSu -> SenSu naming unified - ServiceManager: health checks + startup_order - Env var credentials (SENSU_ADMIN_PASSWORD etc)
387 lines
16 KiB
Python
387 lines
16 KiB
Python
#!/usr/bin/env python3
|
|
# -*- coding: utf-8 -*-
|
|
|
|
import logging
|
|
import asyncio
|
|
from typing import Dict, Any
|
|
from aiohttp import web
|
|
import json
|
|
|
|
# 导入命令装饰器
|
|
try:
|
|
from fmfuncs.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)
|