From 29921ae5ae72900327912e2a6a0a719623ccc071 Mon Sep 17 00:00:00 2001 From: SXP-Simon Date: Tue, 23 Sep 2025 22:36:23 +0800 Subject: [PATCH 1/5] =?UTF-8?q?[feat]=20(=E8=87=AA=E5=8A=A8=E5=88=86?= =?UTF-8?q?=E6=9E=90=E5=99=A8=E7=9A=84=E7=BE=A4=E8=81=8A=E5=B9=B6=E5=8F=91?= =?UTF-8?q?=E5=A4=84=E7=90=86)=20=E5=B0=86=E5=8E=9F=E6=9C=AC=E7=9A=84?= =?UTF-8?q?=E4=B8=B2=E8=A1=8C=E5=88=86=E6=9E=90=E9=80=BB=E8=BE=91=E8=B0=83?= =?UTF-8?q?=E6=95=B4=E4=B8=BA=E5=B9=B6=E8=A1=8C=E4=BB=BB=E5=8A=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/scheduler/auto_scheduler.py | 157 ++++++++++++++++++++++---------- 1 file changed, 108 insertions(+), 49 deletions(-) diff --git a/src/scheduler/auto_scheduler.py b/src/scheduler/auto_scheduler.py index 585e606..51a8834 100644 --- a/src/scheduler/auto_scheduler.py +++ b/src/scheduler/auto_scheduler.py @@ -114,64 +114,123 @@ class AutoScheduler: await asyncio.sleep(300) async def _run_auto_analysis(self): - """执行自动分析""" + """执行自动分析 - 并发处理所有群聊""" try: - logger.info("开始执行自动群聊分析") + logger.info("开始执行自动群聊分析(并发模式)") - # 为每个启用的群执行分析 enabled_groups = self.config_manager.get_enabled_groups() + if not enabled_groups: + logger.info("没有启用的群聊需要分析") + return + + logger.info(f"将为 {len(enabled_groups)} 个群聊并发执行分析: {enabled_groups}") + + # 创建并发任务 - 为每个群聊创建独立的分析任务 + analysis_tasks = [] for group_id in enabled_groups: - try: - logger.info(f"为群 {group_id} 执行自动分析") - await self._perform_auto_analysis_for_group(group_id) - except Exception as e: - logger.error(f"群 {group_id} 自动分析失败: {e}") + task = asyncio.create_task( + self._perform_auto_analysis_for_group_with_timeout(group_id), + name=f"analysis_group_{group_id}" + ) + analysis_tasks.append(task) + + # 并发执行所有分析任务,使用 return_exceptions=True 确保单个任务失败不影响其他任务 + results = await asyncio.gather(*analysis_tasks, return_exceptions=True) + + # 统计执行结果 + success_count = 0 + error_count = 0 + + for i, result in enumerate(results): + group_id = enabled_groups[i] + if isinstance(result, Exception): + logger.error(f"群 {group_id} 分析任务异常: {result}") + error_count += 1 + else: + success_count += 1 + + logger.info(f"并发分析完成 - 成功: {success_count}, 失败: {error_count}, 总计: {len(enabled_groups)}") except Exception as e: - logger.error(f"自动分析执行失败: {e}") + logger.error(f"自动分析执行失败: {e}", exc_info=True) - async def _perform_auto_analysis_for_group(self, group_id: str): - """为指定群执行自动分析""" + async def _perform_auto_analysis_for_group_with_timeout(self, group_id: str): + """为指定群执行自动分析(带超时控制)""" try: - # 检查bot管理器状态 - if not self.bot_manager.is_ready_for_auto_analysis(): - status = self.bot_manager.get_status_info() - logger.warning(f"群 {group_id} 自动分析跳过:bot管理器未就绪 - {status}") - return - - logger.info(f"开始为群 {group_id} 执行自动分析") - - # 获取群聊消息 - analysis_days = self.config_manager.get_analysis_days() - bot_instance = self.bot_manager.get_bot_instance() - - messages = await self.message_handler.fetch_group_messages(bot_instance, group_id, analysis_days) - - if not messages: - logger.warning(f"群 {group_id} 未获取到足够的消息记录") - return - - # 检查消息数量 - min_threshold = self.config_manager.get_min_messages_threshold() - if len(messages) < min_threshold: - logger.warning(f"群 {group_id} 消息数量不足({len(messages)}条),跳过分析") - return - - logger.info(f"群 {group_id} 获取到 {len(messages)} 条消息,开始分析") - - # 进行分析 - 构造正确的 unified_msg_origin - platform_id = self._get_platform_id() - umo = f"{platform_id}:group:{group_id}" if platform_id else None - analysis_result = await self.analyzer.analyze_messages(messages, group_id, umo) - if not analysis_result: - logger.error(f"群 {group_id} 分析失败") - return - - # 生成并发送报告 - await self._send_analysis_report(group_id, analysis_result) - + # 为每个群聊设置独立的超时时间(20分钟) + async with asyncio.timeout(1200): + await self._perform_auto_analysis_for_group(group_id) + except asyncio.TimeoutError: + logger.error(f"群 {group_id} 分析超时(20分钟),跳过该群分析") except Exception as e: - logger.error(f"群 {group_id} 自动分析执行失败: {e}", exc_info=True) + logger.error(f"群 {group_id} 分析任务执行失败: {e}") + + async def _perform_auto_analysis_for_group(self, group_id: str): + """为指定群执行自动分析(核心逻辑)""" + # 为每个群聊使用独立的锁,避免全局锁导致串行化 + group_lock_key = f"analysis_{group_id}" + if not hasattr(self, '_group_locks'): + self._group_locks = {} + + if group_lock_key not in self._group_locks: + self._group_locks[group_lock_key] = asyncio.Lock() + + async with self._group_locks[group_lock_key]: + try: + start_time = asyncio.get_event_loop().time() + + # 检查bot管理器状态 + if not self.bot_manager.is_ready_for_auto_analysis(): + status = self.bot_manager.get_status_info() + logger.warning(f"群 {group_id} 自动分析跳过:bot管理器未就绪 - {status}") + return + + logger.info(f"开始为群 {group_id} 执行自动分析(并发任务)") + + # 获取群聊消息 + analysis_days = self.config_manager.get_analysis_days() + bot_instance = self.bot_manager.get_bot_instance() + + messages = await self.message_handler.fetch_group_messages(bot_instance, group_id, analysis_days) + + if not messages: + logger.warning(f"群 {group_id} 未获取到足够的消息记录") + return + + # 检查消息数量 + min_threshold = self.config_manager.get_min_messages_threshold() + if len(messages) < min_threshold: + logger.warning(f"群 {group_id} 消息数量不足({len(messages)}条),跳过分析") + return + + logger.info(f"群 {group_id} 获取到 {len(messages)} 条消息,开始分析") + + # 进行分析 - 构造正确的 unified_msg_origin + platform_id = self._get_platform_id() + umo = f"{platform_id}:group:{group_id}" if platform_id else None + analysis_result = await self.analyzer.analyze_messages(messages, group_id, umo) + if not analysis_result: + logger.error(f"群 {group_id} 分析失败") + return + + # 生成并发送报告 + await self._send_analysis_report(group_id, analysis_result) + + # 记录执行时间 + end_time = asyncio.get_event_loop().time() + execution_time = end_time - start_time + logger.info(f"群 {group_id} 分析完成,耗时: {execution_time:.2f}秒") + + except Exception as e: + logger.error(f"群 {group_id} 自动分析执行失败: {e}", exc_info=True) + + finally: + # 清理群聊锁资源(可选,防止内存泄漏) + if hasattr(self, '_group_locks') and len(self._group_locks) > 50: + old_locks = list(self._group_locks.keys())[:10] + for lock_key in old_locks: + if not self._group_locks[lock_key].locked(): + del self._group_locks[lock_key] async def _send_analysis_report(self, group_id: str, analysis_result: dict): """发送分析报告到群""" From a192081defec25efe8adc0e3afdc53188f52b331 Mon Sep 17 00:00:00 2001 From: SXP-Simon Date: Tue, 23 Sep 2025 22:54:49 +0800 Subject: [PATCH 2/5] =?UTF-8?q?[fix]=20(=5Freload=5Fconfig=5Fand=5Frestart?= =?UTF-8?q?=5Fscheduler)=20=E5=88=A0=E9=99=A4=E4=B8=8D=E9=9C=80=E8=A6=81?= =?UTF-8?q?=E7=9A=84=E5=87=BD=E6=95=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- main.py | 16 ---------------- 1 file changed, 16 deletions(-) diff --git a/main.py b/main.py index 73be5e4..de0ad8b 100644 --- a/main.py +++ b/main.py @@ -82,22 +82,6 @@ class QQGroupDailyAnalysis(Star): logger.error(f"延迟启动调度器失败: {e}") - - - async def _reload_config_and_restart_scheduler(self): - """重新加载配置并重启调度器""" - try: - # 重新加载配置 - config_manager.reload_config() - logger.info(f"重新加载配置: 自动分析={config_manager.get_enable_auto_analysis()}") - - # 重启调度器 - await auto_scheduler.restart_scheduler() - logger.info("配置重载和调度器重启完成") - - except Exception as e: - logger.error(f"重新加载配置失败: {e}") - @filter.command("群分析") @filter.permission_type(PermissionType.ADMIN) async def analyze_group_daily(self, event: AiocqhttpMessageEvent, days: Optional[int] = None): From b6f0d99a6104e421807f5564b65891ca6f0beabd Mon Sep 17 00:00:00 2001 From: SXP-Simon Date: Tue, 23 Sep 2025 22:55:57 +0800 Subject: [PATCH 3/5] =?UTF-8?q?[feat]=20(terminate)=20=E5=8F=82=E8=80=83?= =?UTF-8?q?=E6=8F=92=E4=BB=B6=E5=BC=80=E5=8F=91=E6=96=87=E6=A1=A3=EF=BC=8C?= =?UTF-8?q?=E5=AE=9E=E7=8E=B0=E6=8F=92=E4=BB=B6=E7=9A=84=E8=B5=84=E6=BA=90?= =?UTF-8?q?=E5=8D=B8=E8=BD=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- main.py | 39 +++++++++++++++++++++++++++++++++++++++ 1 file changed, 39 insertions(+) diff --git a/main.py b/main.py index de0ad8b..c1d0560 100644 --- a/main.py +++ b/main.py @@ -81,6 +81,45 @@ class QQGroupDailyAnalysis(Star): except Exception as e: logger.error(f"延迟启动调度器失败: {e}") + async def terminate(self): + """插件被卸载/停用时调用,清理资源""" + try: + logger.info("开始清理QQ群日常分析插件资源...") + + global auto_scheduler, bot_manager, message_analyzer, report_generator, config_manager + + # 停止自动调度器 + if auto_scheduler: + logger.info("正在停止自动调度器...") + await auto_scheduler.stop_scheduler() + logger.info("自动调度器已停止") + + # 清理bot管理器资源 + # if bot_manager: + # logger.info("正在清理bot管理器资源...") + # # 如果有其他需要清理的资源,可以在这里添加 + + # # 清理消息分析器资源 + # if message_analyzer: + # logger.info("正在清理消息分析器资源...") + # # 如果有其他需要清理的资源,可以在这里添加 + + # # 清理报告生成器资源 + # if report_generator: + # logger.info("正在清理报告生成器资源...") + # # 如果有其他需要清理的资源,可以在这里添加 + + # 重置全局变量 + auto_scheduler = None + bot_manager = None + message_analyzer = None + report_generator = None + config_manager = None + + logger.info("QQ群日常分析插件资源清理完成") + + except Exception as e: + logger.error(f"插件资源清理失败: {e}") @filter.command("群分析") @filter.permission_type(PermissionType.ADMIN) From 9602b31014fb5128d5857faf5ed0768256d5a14a Mon Sep 17 00:00:00 2001 From: SXP-Simon Date: Tue, 23 Sep 2025 23:26:12 +0800 Subject: [PATCH 4/5] =?UTF-8?q?[fix]=20(debug)=20=E7=BA=A0=E6=AD=A3?= =?UTF-8?q?=E6=97=A5=E5=BF=97=E7=AD=89=E7=BA=A7=E6=83=85=E5=86=B5=EF=BC=8C?= =?UTF-8?q?=E7=BA=A0=E6=AD=A3=E4=B8=8D=E8=83=BD=E9=80=82=E9=85=8D=E4=BD=8E?= =?UTF-8?q?=E7=89=88=E6=9C=AC=E7=9A=84=E7=9A=84=20async=20timeout=20?= =?UTF-8?q?=E6=96=B9=E6=B3=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- main.py | 2 +- src/analysis/llm_analyzer.py | 6 +++--- src/scheduler/auto_scheduler.py | 5 ++--- 3 files changed, 6 insertions(+), 7 deletions(-) diff --git a/main.py b/main.py index c1d0560..8652bc0 100644 --- a/main.py +++ b/main.py @@ -79,7 +79,7 @@ class QQGroupDailyAnalysis(Star): logger.info(f"Bot管理器状态: {status}") except Exception as e: - logger.error(f"延迟启动调度器失败: {e}") + logger.debug(f"延迟启动调度器失败,可能由于短时间内多次更新插件配置: {e}") async def terminate(self): """插件被卸载/停用时调用,清理资源""" diff --git a/src/analysis/llm_analyzer.py b/src/analysis/llm_analyzer.py index bd5bb34..3ad0d27 100644 --- a/src/analysis/llm_analyzer.py +++ b/src/analysis/llm_analyzer.py @@ -253,7 +253,7 @@ class LLMAnalyzer: else: logger.warning(f"话题分析响应中未找到JSON格式,响应内容: {result_text[:200]}...") except json.JSONDecodeError as e: - logger.error(f"话题分析JSON解析失败: {e}") + logger.debug(f"话题分析JSON解析失败: {e}") logger.debug(f"修复后的JSON: {json_text if 'json_text' in locals() else 'N/A'}") logger.debug(f"原始响应: {result_text}") @@ -467,7 +467,7 @@ class LLMAnalyzer: titles_data = json.loads(json_match.group()) return [UserTitle(**title) for title in titles_data], token_usage except Exception as e: - logger.error(f"用户称号分析JSON解析失败: {e}") + logger.debug(f"用户称号分析JSON解析失败: {e}") logger.debug(f"原始响应: {result_text}") return [], token_usage @@ -570,7 +570,7 @@ class LLMAnalyzer: quotes_data = json.loads(json_match.group()) return [GoldenQuote(**quote) for quote in quotes_data[:max_golden_quotes]], token_usage except Exception as e: - logger.error(f"金句分析JSON解析失败: {e}") + logger.debug(f"金句分析JSON解析失败: {e}") logger.debug(f"原始响应: {result_text}") return [], token_usage diff --git a/src/scheduler/auto_scheduler.py b/src/scheduler/auto_scheduler.py index 51a8834..1647f1f 100644 --- a/src/scheduler/auto_scheduler.py +++ b/src/scheduler/auto_scheduler.py @@ -157,9 +157,8 @@ class AutoScheduler: async def _perform_auto_analysis_for_group_with_timeout(self, group_id: str): """为指定群执行自动分析(带超时控制)""" try: - # 为每个群聊设置独立的超时时间(20分钟) - async with asyncio.timeout(1200): - await self._perform_auto_analysis_for_group(group_id) + # 为每个群聊设置独立的超时时间(20分钟)- 使用 asyncio.wait_for 兼容所有 Python 版本 + await asyncio.wait_for(self._perform_auto_analysis_for_group(group_id), timeout=1200) except asyncio.TimeoutError: logger.error(f"群 {group_id} 分析超时(20分钟),跳过该群分析") except Exception as e: From afe569d1e3238884ab2c9f0f24ee301754d23101 Mon Sep 17 00:00:00 2001 From: SXP-Simon Date: Tue, 23 Sep 2025 23:49:59 +0800 Subject: [PATCH 5/5] =?UTF-8?q?[fix]=20(umo)=20=E7=BA=A0=E6=AD=A3=E8=87=AA?= =?UTF-8?q?=E5=8A=A8=E5=88=86=E6=9E=90=E5=99=A8=E4=B8=AD=E4=BC=A0=E9=80=92?= =?UTF-8?q?=E7=9A=84=20umo?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/scheduler/auto_scheduler.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/scheduler/auto_scheduler.py b/src/scheduler/auto_scheduler.py index 1647f1f..95d6fed 100644 --- a/src/scheduler/auto_scheduler.py +++ b/src/scheduler/auto_scheduler.py @@ -206,7 +206,7 @@ class AutoScheduler: # 进行分析 - 构造正确的 unified_msg_origin platform_id = self._get_platform_id() - umo = f"{platform_id}:group:{group_id}" if platform_id else None + umo = f"{platform_id}:GroupMessage:{group_id}" if platform_id else None analysis_result = await self.analyzer.analyze_messages(messages, group_id, umo) if not analysis_result: logger.error(f"群 {group_id} 分析失败")