feat: 添加分析编排器和日志过滤器,优化群聊分析功能

This commit is contained in:
SXP-Simon
2026-02-08 18:55:26 +08:00
parent b4d50c2571
commit 46356d1c19
+311 -4
View File
@@ -16,9 +16,8 @@ from astrbot.core.platform.sources.aiocqhttp.aiocqhttp_message_event import (
)
from astrbot.core.star.filter.permission import PermissionType
from .src.core.bot_manager import BotManager
# 导入重构后的模块
from .src.application.analysis_orchestrator import AnalysisOrchestrator, AnalysisConfig
from .src.infrastructure.platform.factory import PlatformAdapterFactory
from .src.core.config import ConfigManager
from .src.core.history_manager import HistoryManager
from .src.reports.generators import ReportGenerator
@@ -26,7 +25,7 @@ from .src.scheduler.auto_scheduler import AutoScheduler
from .src.scheduler.retry import RetryManager
from .src.utils.helpers import MessageAnalyzer
from .src.utils.pdf_utils import PDFInstaller
from .src.domain.value_objects.unified_message import UnifiedMessage
class QQGroupDailyAnalysis(Star):
def __init__(self, context: Context, config: AstrBotConfig):
@@ -55,6 +54,314 @@ class QQGroupDailyAnalysis(Star):
self.history_manager,
self.html_render, # 传入html_render函数
)
# 注册分析编排器缓存
self.orchestrators = {} # {platform_id: AnalysisOrchestrator}
# 注册日志过滤器
from .src.utils.trace_context import TraceLogFilter
logger.addFilter(TraceLogFilter())
logger.info("QQ群日常分析插件已初始化(模块化版本)")
def _get_orchestrator(self, platform_id: str, bot_instance: Any = None) -> Optional[AnalysisOrchestrator]:
"""获取或创建分析编排器"""
if platform_id in self.orchestrators:
return self.orchestrators[platform_id]
# 如果缓存中没有,尝试创建
if not bot_instance:
bot_instance = self.bot_manager.get_bot_instance(platform_id)
if not bot_instance:
return None
# 检测平台名称
platform_name = self.bot_manager._detect_platform_name(bot_instance)
if not platform_name:
return None
# 创建编排器
analysis_config = AnalysisConfig(
days=self.config_manager.get_analysis_days(),
min_messages_threshold=self.config_manager.get_min_messages_threshold(),
output_format=self.config_manager.get_output_format()
)
orchestrator = AnalysisOrchestrator.create_for_platform(
platform_name,
bot_instance,
config={"bot_qq_ids": self.config_manager.get_bot_qq_ids()},
analysis_config=analysis_config
)
if orchestrator:
self.orchestrators[platform_id] = orchestrator
return orchestrator
@filter.on_platform_loaded()
async def on_platform_loaded(self):
"""平台加载完成后初始化"""
try:
# 检查插件是否被启用 (Fix for empty plugin_set issue)
# 如果 plugin_set 为空列表,会导致所有插件不响应消息
if self.context:
# 获取配置对象
config = self.context.get_config()
plugin_set = config.get("plugin_set")
if isinstance(plugin_set, list) and not plugin_set:
logger.warning("检测到 plugin_set 为空,自动修正以启用插件")
config["plugin_set"].append("astrbot_plugin_qq_group_daily_analysis")
elif isinstance(plugin_set, list) and "*" not in plugin_set and "astrbot_plugin_qq_group_daily_analysis" not in plugin_set:
logger.warning("检测到当前插件未在 plugin_set 中,自动添加")
config["plugin_set"].append("astrbot_plugin_qq_group_daily_analysis")
# 初始化所有bot实例
discovered = await self.bot_manager.initialize_from_config()
if discovered:
platform_count = len(discovered)
logger.info(f"Bot管理器初始化成功,发现 {platform_count} 个适配器")
for platform_id, bot_instance in discovered.items():
logger.info(
f" - 平台 {platform_id}: {type(bot_instance).__name__}"
)
# 预先创建编排器
self._get_orchestrator(platform_id, bot_instance)
# 启动调度器
self.auto_scheduler.schedule_jobs(self.context)
else:
logger.warning("Bot管理器初始化失败,未发现任何适配器")
status = self.bot_manager.get_status_info()
logger.info(f"Bot管理器状态: {status}")
# 始终启动重试管理器,确保手动触发也能使用重试队列
await self.retry_manager.start()
except Exception as e:
logger.error(f"平台加载事件处理失败: {e}", exc_info=True)
async def terminate(self):
"""插件被卸载/停用时调用,清理资源"""
try:
logger.info("开始清理QQ群日常分析插件资源...")
# 停止自动调度器
if self.auto_scheduler:
logger.info("正在停止自动调度器...")
self.auto_scheduler.unschedule_jobs(self.context)
logger.info("自动调度器已停止")
if self.retry_manager:
await self.retry_manager.stop()
# 重置实例属性
self.auto_scheduler = None
self.bot_manager = None
self.message_analyzer = None
self.report_generator = None
self.config_manager = None
self.orchestrators = {}
logger.info("QQ群日常分析插件资源清理完成")
except Exception as e:
logger.error(f"插件资源清理失败: {e}")
@filter.command("群分析")
@filter.command("group_analysis")
@filter.permission_type(PermissionType.ADMIN)
async def analyze_group_daily(
self, event: AstrMessageEvent, days: int | None = None
):
"""
分析群聊日常活动
用法: /群分析 [天数]
"""
# 1. 获取 group_id 和 platform_id
group_id = None
platform_id = None
if hasattr(event, "message_obj"):
group_id = getattr(event.message_obj, "group_id", None)
# 尝试从 metadata 获取 platform_id
if hasattr(event, "platform") and isinstance(event.platform, str):
platform_id = event.platform
elif hasattr(event, "metadata") and hasattr(event.metadata, "id"):
platform_id = event.metadata.id
# 如果无法获取,尝试从 bot_manager 推断
if not platform_id and hasattr(event, "bot"):
platform_id = self.bot_manager._get_platform_id_from_instance(event.bot)
if not group_id:
yield event.plain_result("❌ 请在群聊中使用此命令")
return
# 更新bot实例(用于手动命令)
if hasattr(event, "bot"):
self.bot_manager.update_from_event(event)
# 2. 检查群组权限
if not self.config_manager.is_group_allowed(group_id):
yield event.plain_result("❌ 此群未启用日常分析功能")
return
# 3. 设置分析天数
analysis_days = (
days if days and 1 <= days <= 7 else self.config_manager.get_analysis_days()
)
yield event.plain_result(f"🔍 开始分析群聊近{analysis_days}天的活动,请稍候...")
logger.info(f"收到分析请求: group_id={group_id}, platform_id={platform_id}, days={analysis_days}")
try:
# 4. 获取编排器
orchestrator = self._get_orchestrator(platform_id)
if not orchestrator:
# 尝试使用 bot_manager 获取 bot 实例再创建
bot_instance = self.bot_manager.get_bot_instance(platform_id)
if bot_instance:
orchestrator = self._get_orchestrator(platform_id, bot_instance)
if not orchestrator:
yield event.plain_result(
f"❌ 未找到平台 {platform_id} 的分析编排器,请检查配置或联系开发者"
)
return
# 5. 获取群聊消息 (使用编排器,支持 DDD)
# 使用 fetch_messages_as_raw 保持向后兼容性,或者重构 message_analyzer 支持 UnifiedMessage
# 这里我们尝试重构为使用 UnifiedMessage,但为了稳健性,我们暂时获取 raw 格式
# 实际上,AnalysisOrchestrator 提供了 fetch_messages_as_raw 方法
messages = await orchestrator.fetch_messages_as_raw(
group_id=group_id,
days=analysis_days
)
if not messages:
yield event.plain_result(
"❌ 未找到足够的群聊记录,请确保群内有足够的消息历史"
)
return
# 检查消息数量是否足够分析
min_threshold = self.config_manager.get_min_messages_threshold()
if len(messages) < min_threshold:
yield event.plain_result(
f"❌ 消息数量不足({len(messages)}条),至少需要{min_threshold}条消息才能进行有效分析"
)
return
yield event.plain_result(
f"📊 已获取{len(messages)}条消息,正在进行智能分析..."
)
# 6. 进行分析
# 传递 unified_msg_origin 以获取正确的 LLM 提供商
analysis_result = await self.message_analyzer.analyze_messages(
messages, group_id, event.unified_msg_origin
)
if not analysis_result or not analysis_result.get("statistics"):
yield event.plain_result("❌ 分析过程中出现错误,请稍后重试")
return
# 7. 保存到历史记录
await self.history_manager.save_analysis(group_id, analysis_result)
# 8. 生成并发送报告
output_format = self.config_manager.get_output_format()
if output_format == "image":
# 生成图片报告
(image_url, html_content) = await self.report_generator.generate_image_report(
analysis_result, group_id, self.html_render
)
if image_url:
# 使用编排器发送图片
if await orchestrator.send_image(group_id, image_url):
logger.info(f"图片报告发送成功: {group_id}")
else:
# 发送失败,尝试 yield
yield event.image_result(image_url)
elif html_content:
# 生成失败但有HTML,加入重试队列
logger.warning("图片报告生成失败,加入重试队列")
yield event.plain_result(
"[AstrBot QQ群日常分析总结插件] ⚠️ 图片报告暂无法生成,已加入重试队列,稍后将自动重试发送。"
)
await self.retry_manager.add_task(
html_content, analysis_result, group_id, platform_id
)
else:
# 回退到文本报告
logger.warning("图片报告生成失败(无HTML),回退到文本报告")
text_report = self.report_generator.generate_text_report(analysis_result)
yield event.plain_result(
f"[AstrBot QQ群日常分析总结插件] ⚠️ 图片报告生成失败,以下是文本版本:\\n\\n{text_report}"
)
elif output_format == "pdf":
if not self.config_manager.playwright_available:
yield event.plain_result("❌ PDF 功能不可用,请使用 /安装PDF 命令安装依赖")
return
pdf_path = await self.report_generator.generate_pdf_report(
analysis_result, group_id
)
if pdf_path:
# 使用编排器发送文件
if await orchestrator.send_file(group_id, pdf_path):
pass # 发送成功
else:
# 回退 yield
from pathlib import Path
pdf_file = File(name=Path(pdf_path).name, file=pdf_path)
result = event.make_result()
result.chain.append(pdf_file)
yield result
else:
logger.warning("PDF 报告生成失败,回退到文本报告")
text_report = self.report_generator.generate_text_report(analysis_result)
yield event.plain_result(
f"\\n📝 以下是文本版本的分析报告:\\n\\n{text_report}"
)
else:
# 文本报告
text_report = self.report_generator.generate_text_report(analysis_result)
# 使用编排器发送文本
if not await orchestrator.send_text(group_id, text_report):
yield event.plain_result(text_report)
except Exception as e:
logger.error(f"群分析失败: {e}", exc_info=True)
yield event.plain_result(
f"❌ 分析失败: {str(e)}。请检查网络连接和LLM配置,或联系管理员"
)
self.report_generator = ReportGenerator(self.config_manager)
self.history_manager = HistoryManager(self)
self.retry_manager = RetryManager(
self.bot_manager, self.html_render, self.report_generator
)
self.auto_scheduler = AutoScheduler(
self.config_manager,
self.message_analyzer.message_handler,
self.message_analyzer,
self.report_generator,
self.bot_manager,
self.retry_manager,
self.history_manager,
self.html_render, # 传入html_render函数
)
# 注册日志过滤器
from .src.utils.trace_context import TraceLogFilter