refactor(core): 重构核心模块结构并添加开发文档
将核心模块按功能重新组织为更清晰的结构,包括 managers、handlers 和 utils 目录 添加完整的开发文档,涵盖快速开始、项目结构、核心概念和插件开发指南 更新所有相关模块的导入路径以匹配新的结构 将单例模式实现提取到单独的 singleton.py 文件
This commit is contained in:
0
core/managers/__init__.py
Normal file
0
core/managers/__init__.py
Normal file
157
core/managers/admin_manager.py
Normal file
157
core/managers/admin_manager.py
Normal file
@@ -0,0 +1,157 @@
|
||||
"""
|
||||
管理员管理器模块
|
||||
|
||||
该模块负责管理机器人的管理员列表。
|
||||
它实现了文件和 Redis 缓存之间的数据同步,并提供了一套清晰的 API
|
||||
供其他模块调用。
|
||||
"""
|
||||
import json
|
||||
import os
|
||||
from typing import Set
|
||||
|
||||
from ..utils.logger import logger
|
||||
from ..utils.singleton import Singleton
|
||||
from .redis_manager import redis_manager
|
||||
|
||||
|
||||
class AdminManager(Singleton):
|
||||
"""
|
||||
管理员管理器类
|
||||
|
||||
负责加载、缓存和管理管理员列表。
|
||||
使用单例模式,确保全局只有一个实例。
|
||||
"""
|
||||
_REDIS_KEY = "neobot:admins" # 用于存储管理员集合的 Redis 键
|
||||
def __init__(self):
|
||||
"""
|
||||
初始化 AdminManager
|
||||
"""
|
||||
super().__init__()
|
||||
if not self._initialized:
|
||||
return
|
||||
|
||||
# 管理员数据文件路径
|
||||
self.data_file = os.path.join(
|
||||
os.path.dirname(os.path.abspath(__file__)),
|
||||
"..",
|
||||
"data",
|
||||
"admin.json"
|
||||
)
|
||||
|
||||
self._admins: Set[int] = set()
|
||||
logger.info("管理员管理器初始化完成")
|
||||
|
||||
async def initialize(self):
|
||||
"""
|
||||
异步初始化,加载数据并同步到 Redis
|
||||
"""
|
||||
await self._load_from_file()
|
||||
await self._sync_to_redis()
|
||||
logger.info("管理员数据加载并同步到 Redis 完成")
|
||||
|
||||
async def _load_from_file(self):
|
||||
"""
|
||||
从 admin.json 加载管理员列表
|
||||
"""
|
||||
try:
|
||||
if os.path.exists(self.data_file):
|
||||
with open(self.data_file, "r", encoding="utf-8") as f:
|
||||
data = json.load(f)
|
||||
admins = data.get("admins", [])
|
||||
self._admins = set(int(admin_id) for admin_id in admins)
|
||||
logger.debug(f"从 {self.data_file} 加载了 {len(self._admins)} 位管理员")
|
||||
else:
|
||||
# 如果文件不存在,创建一个空的
|
||||
self._admins = set()
|
||||
await self._save_to_file()
|
||||
except (json.JSONDecodeError, ValueError) as e:
|
||||
logger.error(f"加载或解析 admin.json 失败: {e}")
|
||||
self._admins = set()
|
||||
|
||||
async def _save_to_file(self):
|
||||
"""
|
||||
将当前管理员列表保存回 admin.json
|
||||
"""
|
||||
try:
|
||||
# 确保目录存在
|
||||
os.makedirs(os.path.dirname(self.data_file), exist_ok=True)
|
||||
# 将 set 转换为 list 以便 JSON 序列化
|
||||
admin_list = [str(admin_id) for admin_id in self._admins]
|
||||
with open(self.data_file, "w", encoding="utf-8") as f:
|
||||
json.dump({"admins": admin_list}, f, indent=2, ensure_ascii=False)
|
||||
logger.debug(f"管理员列表已保存到 {self.data_file}")
|
||||
except Exception as e:
|
||||
logger.error(f"保存 admin.json 失败: {e}")
|
||||
|
||||
async def _sync_to_redis(self):
|
||||
"""
|
||||
将内存中的管理员集合同步到 Redis
|
||||
"""
|
||||
from core.managers.redis_manager import redis_manager
|
||||
try:
|
||||
# 首先清空旧的集合
|
||||
await redis_manager.redis.delete(self._REDIS_KEY)
|
||||
if self._admins:
|
||||
# 将所有管理员ID添加到集合中
|
||||
await redis_manager.redis.sadd(self._REDIS_KEY, *self._admins)
|
||||
logger.debug(f"已将 {len(self._admins)} 位管理员同步到 Redis")
|
||||
except Exception as e:
|
||||
logger.error(f"同步管理员到 Redis 失败: {e}")
|
||||
|
||||
async def is_admin(self, user_id: int) -> bool:
|
||||
"""
|
||||
检查用户是否为管理员(从 Redis 缓存读取)
|
||||
"""
|
||||
|
||||
try:
|
||||
return await redis_manager.redis.sismember(self._REDIS_KEY, user_id)
|
||||
except Exception as e:
|
||||
logger.error(f"从 Redis 检查管理员权限失败: {e}")
|
||||
# Redis 失败时,回退到内存检查
|
||||
return user_id in self._admins
|
||||
|
||||
async def add_admin(self, user_id: int) -> bool:
|
||||
"""
|
||||
添加管理员,并同步到文件和 Redis
|
||||
"""
|
||||
from .redis_manager import redis_manager
|
||||
if user_id in self._admins:
|
||||
return False # 用户已经是管理员
|
||||
|
||||
self._admins.add(user_id)
|
||||
await self._save_to_file()
|
||||
try:
|
||||
await redis_manager.redis.sadd(self._REDIS_KEY, user_id)
|
||||
logger.info(f"已添加新管理员 {user_id} 并更新缓存")
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error(f"添加管理员 {user_id} 到 Redis 失败: {e}")
|
||||
return False
|
||||
|
||||
async def remove_admin(self, user_id: int) -> bool:
|
||||
"""
|
||||
移除管理员,并同步到文件和 Redis
|
||||
"""
|
||||
from .redis_manager import redis_manager
|
||||
if user_id not in self._admins:
|
||||
return False # 用户不是管理员
|
||||
|
||||
self._admins.remove(user_id)
|
||||
await self._save_to_file()
|
||||
try:
|
||||
await redis_manager.redis.srem(self._REDIS_KEY, user_id)
|
||||
logger.info(f"已移除管理员 {user_id} 并更新缓存")
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error(f"从 Redis 移除管理员 {user_id} 失败: {e}")
|
||||
return False
|
||||
|
||||
async def get_all_admins(self) -> Set[int]:
|
||||
"""
|
||||
获取所有管理员的集合
|
||||
"""
|
||||
return self._admins.copy()
|
||||
|
||||
|
||||
# 全局 AdminManager 实例
|
||||
admin_manager = AdminManager()
|
||||
143
core/managers/command_manager.py
Normal file
143
core/managers/command_manager.py
Normal file
@@ -0,0 +1,143 @@
|
||||
"""
|
||||
命令与事件管理器模块
|
||||
|
||||
该模块定义了 `CommandManager` 类,它是整个机器人框架事件处理的核心。
|
||||
它通过装饰器模式,为插件提供了注册消息指令、通知事件处理器和
|
||||
请求事件处理器的能力。
|
||||
"""
|
||||
from typing import Any, Callable, Dict, Optional, Tuple
|
||||
|
||||
from ..config_loader import global_config
|
||||
from ..handlers.event_handler import MessageHandler, NoticeHandler, RequestHandler
|
||||
|
||||
|
||||
# 从配置中获取命令前缀
|
||||
comm_prefixes = global_config.bot.get("command", ("/",))
|
||||
|
||||
|
||||
class CommandManager:
|
||||
"""
|
||||
命令管理器,负责注册和分发所有类型的事件。
|
||||
|
||||
这是一个单例对象(`matcher`),在整个应用中共享。
|
||||
它将不同类型的事件处理委托给专门的处理器类。
|
||||
"""
|
||||
|
||||
def __init__(self, prefixes: Tuple[str, ...]):
|
||||
"""
|
||||
初始化命令管理器。
|
||||
|
||||
Args:
|
||||
prefixes (Tuple[str, ...]): 一个包含所有合法命令前缀的元组。
|
||||
"""
|
||||
self.plugins: Dict[str, Dict[str, Any]] = {}
|
||||
|
||||
# 初始化专门的事件处理器
|
||||
self.message_handler = MessageHandler(prefixes)
|
||||
self.notice_handler = NoticeHandler()
|
||||
self.request_handler = RequestHandler()
|
||||
|
||||
# 将处理器映射到事件类型
|
||||
self.handler_map = {
|
||||
"message": self.message_handler,
|
||||
"notice": self.notice_handler,
|
||||
"request": self.request_handler,
|
||||
}
|
||||
|
||||
# 注册内置的 /help 命令
|
||||
self._register_internal_commands()
|
||||
|
||||
def _register_internal_commands(self):
|
||||
"""
|
||||
注册框架内置的命令
|
||||
"""
|
||||
# Help 命令
|
||||
self.message_handler.command("help")(self._help_command)
|
||||
self.plugins["core.help"] = {
|
||||
"name": "帮助",
|
||||
"description": "显示所有可用指令的帮助信息",
|
||||
"usage": "/help",
|
||||
}
|
||||
|
||||
# --- 装饰器代理 ---
|
||||
|
||||
def on_message(self) -> Callable:
|
||||
"""
|
||||
装饰器:注册一个通用的消息处理器。
|
||||
"""
|
||||
return self.message_handler.on_message()
|
||||
|
||||
def command(
|
||||
self,
|
||||
*names: str,
|
||||
permission: Optional[Any] = None,
|
||||
override_permission_check: bool = False
|
||||
) -> Callable:
|
||||
"""
|
||||
装饰器:注册一个消息指令处理器。
|
||||
"""
|
||||
return self.message_handler.command(
|
||||
*names,
|
||||
permission=permission,
|
||||
override_permission_check=override_permission_check
|
||||
)
|
||||
|
||||
def on_notice(self, notice_type: Optional[str] = None) -> Callable:
|
||||
"""
|
||||
装饰器:注册一个通知事件处理器。
|
||||
"""
|
||||
return self.notice_handler.register(notice_type=notice_type)
|
||||
|
||||
def on_request(self, request_type: Optional[str] = None) -> Callable:
|
||||
"""
|
||||
装饰器:注册一个请求事件处理器。
|
||||
"""
|
||||
return self.request_handler.register(request_type=request_type)
|
||||
|
||||
# --- 事件处理 ---
|
||||
|
||||
async def handle_event(self, bot, event):
|
||||
"""
|
||||
统一的事件分发入口。
|
||||
|
||||
根据事件的 `post_type` 将其分发给对应的处理器。
|
||||
"""
|
||||
if event.post_type == 'message' and global_config.bot.get('ignore_self_message', False):
|
||||
if hasattr(event, 'user_id') and hasattr(event, 'self_id') and event.user_id == event.self_id:
|
||||
return
|
||||
|
||||
handler = self.handler_map.get(event.post_type)
|
||||
if handler:
|
||||
await handler.handle(bot, event)
|
||||
|
||||
# --- 内置命令实现 ---
|
||||
|
||||
async def _help_command(self, bot, event):
|
||||
"""
|
||||
内置的 `/help` 命令的实现。
|
||||
"""
|
||||
help_text = "--- 可用指令列表 ---\n"
|
||||
|
||||
for plugin_name, meta in self.plugins.items():
|
||||
name = meta.get("name", "未命名插件")
|
||||
description = meta.get("description", "暂无描述")
|
||||
usage = meta.get("usage", "暂无用法说明")
|
||||
|
||||
help_text += f"\n{name}:\n"
|
||||
help_text += f" 功能: {description}\n"
|
||||
help_text += f" 用法: {usage}\n"
|
||||
|
||||
await bot.send(event, help_text.strip())
|
||||
|
||||
|
||||
# --- 全局单例 ---
|
||||
|
||||
# 确保前缀配置是元组格式
|
||||
if isinstance(comm_prefixes, list):
|
||||
comm_prefixes = tuple(comm_prefixes)
|
||||
elif isinstance(comm_prefixes, str):
|
||||
comm_prefixes = (comm_prefixes,)
|
||||
|
||||
# 实例化全局唯一的命令管理器
|
||||
matcher = CommandManager(prefixes=comm_prefixes)
|
||||
|
||||
264
core/managers/permission_manager.py
Normal file
264
core/managers/permission_manager.py
Normal file
@@ -0,0 +1,264 @@
|
||||
"""
|
||||
权限管理器模块
|
||||
|
||||
该模块负责管理用户权限,支持 admin、op、user 三个权限级别。
|
||||
权限数据存储在 `permissions.json` 文件中,格式为:
|
||||
{
|
||||
"users": {
|
||||
"123456": "admin",
|
||||
"789012": "op",
|
||||
"345678": "user"
|
||||
}
|
||||
}
|
||||
"""
|
||||
import json
|
||||
import os
|
||||
from functools import total_ordering
|
||||
from typing import Dict
|
||||
|
||||
from ..utils.logger import logger
|
||||
from ..utils.singleton import Singleton
|
||||
from .admin_manager import admin_manager
|
||||
|
||||
|
||||
@total_ordering
|
||||
class Permission:
|
||||
"""
|
||||
权限封装类
|
||||
|
||||
封装了权限的名称和等级,并提供了比较方法。
|
||||
使用 @total_ordering 装饰器可以自动生成所有的比较运算符。
|
||||
"""
|
||||
def __init__(self, name: str, level: int):
|
||||
"""
|
||||
初始化权限对象
|
||||
|
||||
Args:
|
||||
name (str): 权限名称 (e.g., "admin", "op")
|
||||
level (int): 权限等级,数字越大权限越高
|
||||
"""
|
||||
self.name = name
|
||||
self.level = level
|
||||
|
||||
def __eq__(self, other):
|
||||
"""
|
||||
判断权限是否相等
|
||||
"""
|
||||
if not isinstance(other, Permission):
|
||||
return NotImplemented
|
||||
return self.level == other.level
|
||||
|
||||
def __lt__(self, other):
|
||||
"""
|
||||
判断权限是否小于另一个权限
|
||||
"""
|
||||
if not isinstance(other, Permission):
|
||||
return NotImplemented
|
||||
return self.level < other.level
|
||||
|
||||
def __str__(self) -> str:
|
||||
"""
|
||||
返回权限的字符串表示(即权限名称)
|
||||
"""
|
||||
return self.name
|
||||
|
||||
|
||||
# 定义全局权限常量
|
||||
ADMIN = Permission("admin", 3)
|
||||
OP = Permission("op", 2)
|
||||
USER = Permission("user", 1)
|
||||
|
||||
# 用于从字符串名称查找权限对象的字典
|
||||
_PERMISSIONS: Dict[str, Permission] = {
|
||||
p.name: p for p in [ADMIN, OP, USER]
|
||||
}
|
||||
|
||||
|
||||
class PermissionManager(Singleton):
|
||||
"""
|
||||
权限管理器类
|
||||
|
||||
负责加载、保存和查询用户权限数据。
|
||||
使用单例模式,确保全局只有一个权限管理器实例。
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
"""
|
||||
初始化权限管理器
|
||||
|
||||
如果已经初始化过,则直接返回。
|
||||
"""
|
||||
super().__init__()
|
||||
if not self._initialized:
|
||||
return
|
||||
|
||||
# 权限数据文件路径
|
||||
self.data_file = os.path.join(
|
||||
os.path.dirname(os.path.abspath(__file__)),
|
||||
"..",
|
||||
"data",
|
||||
"permissions.json"
|
||||
)
|
||||
|
||||
# 确保数据目录存在
|
||||
data_dir = os.path.dirname(self.data_file)
|
||||
os.makedirs(data_dir, exist_ok=True)
|
||||
|
||||
# 权限数据存储结构:{"users": {"user_id": "level_name"}}
|
||||
self._data: Dict[str, Dict[str, str]] = {"users": {}}
|
||||
|
||||
# 加载现有数据
|
||||
self.load()
|
||||
|
||||
logger.info("权限管理器初始化完成")
|
||||
|
||||
def load(self) -> None:
|
||||
"""
|
||||
从文件加载权限数据
|
||||
|
||||
如果文件不存在,则创建空文件并初始化默认数据结构。
|
||||
"""
|
||||
try:
|
||||
if os.path.exists(self.data_file):
|
||||
with open(self.data_file, "r", encoding="utf-8") as f:
|
||||
data = json.load(f)
|
||||
# 兼容旧格式
|
||||
if "users" in data:
|
||||
self._data["users"] = data["users"]
|
||||
else:
|
||||
self._data["users"] = {}
|
||||
logger.debug(f"权限数据已从 {self.data_file} 加载")
|
||||
else:
|
||||
# 文件不存在,创建空文件
|
||||
self.save()
|
||||
logger.debug(f"创建空的权限数据文件: {self.data_file}")
|
||||
except json.JSONDecodeError as e:
|
||||
logger.error(f"权限数据文件格式错误: {e}")
|
||||
# 文件损坏,重置为空数据
|
||||
self._data["users"] = {}
|
||||
self.save()
|
||||
except Exception as e:
|
||||
logger.error(f"加载权限数据失败: {e}")
|
||||
self._data["users"] = {}
|
||||
|
||||
def save(self) -> None:
|
||||
"""
|
||||
将权限数据保存到文件
|
||||
"""
|
||||
try:
|
||||
with open(self.data_file, "w", encoding="utf-8") as f:
|
||||
json.dump(self._data, f, indent=2, ensure_ascii=False)
|
||||
logger.debug(f"权限数据已保存到 {self.data_file}")
|
||||
except Exception as e:
|
||||
logger.error(f"保存权限数据失败: {e}")
|
||||
|
||||
async def get_user_permission(self, user_id: int) -> Permission:
|
||||
"""
|
||||
获取指定用户的权限对象
|
||||
|
||||
Args:
|
||||
user_id (int): 用户 QQ 号
|
||||
|
||||
Returns:
|
||||
Permission: 用户的权限对象,如果用户不存在则返回默认级别 USER
|
||||
"""
|
||||
# 首先,通过 AdminManager 检查是否为管理员
|
||||
if await admin_manager.is_admin(user_id):
|
||||
return ADMIN
|
||||
|
||||
# 如果不是管理员,则从 permissions.json 中查找
|
||||
user_id_str = str(user_id)
|
||||
level_name = self._data["users"].get(user_id_str, USER.name)
|
||||
return _PERMISSIONS.get(level_name, USER)
|
||||
|
||||
def set_user_permission(self, user_id: int, permission: Permission) -> None:
|
||||
"""
|
||||
设置指定用户的权限级别
|
||||
|
||||
Args:
|
||||
user_id (int): 用户 QQ 号
|
||||
permission (Permission): 权限对象
|
||||
|
||||
Raises:
|
||||
ValueError: 如果权限对象无效
|
||||
"""
|
||||
if not isinstance(permission, Permission) or permission.name not in _PERMISSIONS:
|
||||
raise ValueError(f"无效的权限对象: {permission}")
|
||||
|
||||
user_id_str = str(user_id)
|
||||
self._data["users"][user_id_str] = permission.name
|
||||
self.save()
|
||||
logger.info(f"设置用户 {user_id} 的权限级别为 {permission.name}")
|
||||
|
||||
def remove_user(self, user_id: int) -> None:
|
||||
"""
|
||||
移除指定用户的权限设置,恢复为默认级别
|
||||
|
||||
Args:
|
||||
user_id (int): 用户 QQ 号
|
||||
"""
|
||||
user_id_str = str(user_id)
|
||||
if user_id_str in self._data["users"]:
|
||||
del self._data["users"][user_id_str]
|
||||
self.save()
|
||||
logger.info(f"移除用户 {user_id} 的权限设置")
|
||||
|
||||
async def check_permission(self, user_id: int, required_permission: Permission) -> bool:
|
||||
"""
|
||||
检查用户是否具有指定权限级别
|
||||
|
||||
Args:
|
||||
user_id (int): 用户 QQ 号
|
||||
required_permission (Permission): 所需的权限对象
|
||||
|
||||
Returns:
|
||||
bool: 如果用户权限 >= 所需权限,返回 True,否则返回 False
|
||||
"""
|
||||
# 如果传入的是字符串,先转换为 Permission 对象
|
||||
if isinstance(required_permission, str):
|
||||
required_permission = _PERMISSIONS.get(required_permission.lower())
|
||||
if not required_permission:
|
||||
# 如果是无效的权限字符串,默认拒绝
|
||||
logger.warning(f"检测到无效的权限检查字符串: {required_permission}")
|
||||
return False
|
||||
|
||||
user_permission = await self.get_user_permission(user_id)
|
||||
return user_permission >= required_permission
|
||||
|
||||
def get_all_users(self) -> Dict[str, str]:
|
||||
"""
|
||||
获取所有设置了权限的用户及其级别名称
|
||||
|
||||
Returns:
|
||||
Dict[str, str]: 用户ID到权限级别名称的映射
|
||||
"""
|
||||
return self._data["users"].copy()
|
||||
|
||||
def clear_all(self) -> None:
|
||||
"""
|
||||
清空所有权限设置
|
||||
"""
|
||||
self._data["users"].clear()
|
||||
self.save()
|
||||
logger.info("已清空所有权限设置")
|
||||
|
||||
|
||||
# 全局权限管理器实例
|
||||
permission_manager = PermissionManager()
|
||||
|
||||
def require_admin(func):
|
||||
"""
|
||||
一个装饰器,用于限制命令只能由管理员执行。
|
||||
"""
|
||||
from functools import wraps
|
||||
from models.events.message import MessageEvent
|
||||
|
||||
@wraps(func)
|
||||
async def wrapper(event: MessageEvent, *args, **kwargs):
|
||||
user_id = event.user_id
|
||||
if await permission_manager.check_permission(user_id, ADMIN):
|
||||
return await func(event, *args, **kwargs)
|
||||
else:
|
||||
await event.reply("抱歉,您没有权限执行此命令。")
|
||||
return None
|
||||
return wrapper
|
||||
126
core/managers/plugin_manager.py
Normal file
126
core/managers/plugin_manager.py
Normal file
@@ -0,0 +1,126 @@
|
||||
"""
|
||||
插件管理器模块
|
||||
|
||||
负责扫描、加载和管理 `base_plugins` 目录下的所有插件。
|
||||
"""
|
||||
|
||||
import importlib
|
||||
import json
|
||||
import os
|
||||
import pkgutil
|
||||
import sys
|
||||
|
||||
from .command_manager import matcher
|
||||
from ..utils.exceptions import SyncHandlerError
|
||||
from ..utils.logger import logger
|
||||
from ..utils.executor import run_in_thread_pool
|
||||
|
||||
|
||||
def load_all_plugins():
|
||||
"""
|
||||
扫描并加载 `plugins` 目录下的所有插件。
|
||||
|
||||
该函数会遍历 `plugins` 目录下的所有模块:
|
||||
1. 如果模块已加载,则执行 reload 操作(用于热重载)。
|
||||
2. 如果模块未加载,则执行 import 操作。
|
||||
|
||||
加载过程中会提取插件元数据 `__plugin_meta__` 并注册到 CommandManager。
|
||||
"""
|
||||
plugin_dir = os.path.join(
|
||||
os.path.dirname(os.path.abspath(__file__)), "..", "plugins"
|
||||
)
|
||||
package_name = "plugins"
|
||||
|
||||
logger.info(f"正在从 {package_name} 加载插件...")
|
||||
|
||||
for loader, module_name, is_pkg in pkgutil.iter_modules([plugin_dir]):
|
||||
full_module_name = f"{package_name}.{module_name}"
|
||||
|
||||
try:
|
||||
if full_module_name in sys.modules:
|
||||
module = importlib.reload(sys.modules[full_module_name])
|
||||
action = "重载"
|
||||
else:
|
||||
module = importlib.import_module(full_module_name)
|
||||
action = "加载"
|
||||
|
||||
# 提取插件元数据
|
||||
if hasattr(module, "__plugin_meta__"):
|
||||
meta = getattr(module, "__plugin_meta__")
|
||||
matcher.plugins[full_module_name] = meta
|
||||
|
||||
type_str = "包" if is_pkg else "文件"
|
||||
logger.success(f" [{type_str}] 成功{action}: {module_name}")
|
||||
except SyncHandlerError as e:
|
||||
logger.error(f" 插件 {module_name} 加载失败: {e} (跳过此插件)")
|
||||
except Exception as e:
|
||||
print(
|
||||
f" {action if 'action' in locals() else '加载'}插件 {module_name} 失败: {e}"
|
||||
)
|
||||
|
||||
|
||||
class PluginDataManager:
|
||||
"""
|
||||
用于管理插件产生的数据文件的类
|
||||
"""
|
||||
|
||||
def __init__(self, plugin_name: str):
|
||||
"""
|
||||
初始化插件数据管理器
|
||||
|
||||
:param plugin_name: 插件名称
|
||||
"""
|
||||
self.plugin_name = plugin_name
|
||||
self.data_file = os.path.join(
|
||||
os.path.dirname(os.path.abspath(__file__)),
|
||||
"..",
|
||||
"plugins",
|
||||
"data",
|
||||
self.plugin_name + ".json",
|
||||
)
|
||||
self.data = {}
|
||||
|
||||
async def load(self):
|
||||
"""读取配置文件"""
|
||||
if not os.path.exists(self.data_file):
|
||||
await self.set(self.plugin_name, [])
|
||||
try:
|
||||
with open(self.data_file, "r", encoding="utf-8") as f:
|
||||
self.data = await run_in_thread_pool(json.load, f)
|
||||
except json.JSONDecodeError:
|
||||
self.data = {}
|
||||
|
||||
async def save(self):
|
||||
"""保存配置到文件"""
|
||||
with open(self.data_file, "w", encoding="utf-8") as f:
|
||||
await run_in_thread_pool(json.dump, self.data, f, indent=2, ensure_ascii=False)
|
||||
|
||||
def get(self, key, default=None):
|
||||
"""获取配置项"""
|
||||
return self.data.get(key, default)
|
||||
|
||||
async def set(self, key, value):
|
||||
"""设置配置项"""
|
||||
self.data[key] = value
|
||||
await self.save()
|
||||
|
||||
async def add(self, key, value):
|
||||
"""添加配置项"""
|
||||
if key not in self.data:
|
||||
self.data[key] = []
|
||||
self.data[key].append(value)
|
||||
await self.save()
|
||||
|
||||
async def remove(self, key):
|
||||
"""删除配置项"""
|
||||
if key in self.data:
|
||||
del self.data[key]
|
||||
await self.save()
|
||||
|
||||
async def clear(self):
|
||||
"""清空所有配置"""
|
||||
self.data.clear()
|
||||
await self.save()
|
||||
|
||||
def get_all(self):
|
||||
return self.data.copy()
|
||||
58
core/managers/redis_manager.py
Normal file
58
core/managers/redis_manager.py
Normal file
@@ -0,0 +1,58 @@
|
||||
import redis.asyncio as redis
|
||||
from ..config_loader import global_config as config
|
||||
from ..utils.logger import logger
|
||||
|
||||
class RedisManager:
|
||||
"""
|
||||
Redis 连接管理器(异步单例)
|
||||
"""
|
||||
_instance = None
|
||||
_redis = None
|
||||
|
||||
def __new__(cls):
|
||||
if cls._instance is None:
|
||||
cls._instance = super().__new__(cls)
|
||||
return cls._instance
|
||||
|
||||
async def initialize(self):
|
||||
"""
|
||||
异步初始化 Redis 连接并进行健康检查
|
||||
"""
|
||||
if self._redis is None:
|
||||
try:
|
||||
host = config.redis['host']
|
||||
port = config.redis['port']
|
||||
db = config.redis['db']
|
||||
password = config.redis.get('password')
|
||||
|
||||
logger.info(f"正在尝试连接 Redis: {host}:{port}, DB: {db}")
|
||||
|
||||
self._redis = redis.Redis(
|
||||
host=host,
|
||||
port=port,
|
||||
db=db,
|
||||
password=password,
|
||||
decode_responses=True
|
||||
)
|
||||
if await self._redis.ping():
|
||||
logger.success("Redis 连接成功!")
|
||||
else:
|
||||
logger.error("Redis 连接失败: PING 命令无响应")
|
||||
except redis.exceptions.ConnectionError as e:
|
||||
logger.error(f"Redis 连接失败: {e}")
|
||||
self._redis = None
|
||||
except Exception as e:
|
||||
logger.exception(f"Redis 初始化时发生未知错误: {e}")
|
||||
self._redis = None
|
||||
|
||||
@property
|
||||
def redis(self):
|
||||
"""
|
||||
获取 Redis 连接实例
|
||||
"""
|
||||
if self._redis is None:
|
||||
raise ConnectionError("Redis 未初始化或连接失败,请先调用 initialize()")
|
||||
return self._redis
|
||||
|
||||
# 全局 Redis 管理器实例
|
||||
redis_manager = RedisManager()
|
||||
Reference in New Issue
Block a user