diff --git a/main.py b/main.py index 8bc59da..07b0a44 100644 --- a/main.py +++ b/main.py @@ -467,20 +467,21 @@ class GroupDailyAnalysis(Star): ) if image_url: - # 优先使用适配器的 send_image (支持 Base64 转换) - success = await adapter.send_image(group_id, image_url) - if not success: - logger.warning( - "适配器发送图片失败,尝试使用 AstrBot 内置方式回退 (可能因跨容器路径问题失败)" - ) - yield event.image_result(image_url) + # 优先使用适配器的 send_image (由插件适配器统一处理 Base64 转换和路径问题) + # 不再使用 yield event.image_result 回退,防止适配器超时回复导致重复发送图片 + await adapter.send_image(group_id, image_url) # 上传到群文件/群相册 (属于附加功能,不影响消息发送) await self._try_upload_image(group_id, image_url, platform_id) elif html_content: - yield event.plain_result("⚠️ 图片生成暂不可用,已尝试加入队列。") + yield event.plain_result("⚠️ 群分析报告图片发送失败,自动重试中。") + # 使用带提示词的重试任务,确保排队发送时视觉一致 await self.retry_manager.add_task( - html_content, analysis_result, group_id, platform_id + html_content, + analysis_result, + group_id, + platform_id, + caption="📊 每日群聊分析报告已生成:", ) else: text_report = self.report_generator.generate_text_report( diff --git a/src/infrastructure/platform/adapters/onebot_adapter.py b/src/infrastructure/platform/adapters/onebot_adapter.py index e253984..6a0c034 100644 --- a/src/infrastructure/platform/adapters/onebot_adapter.py +++ b/src/infrastructure/platform/adapters/onebot_adapter.py @@ -515,9 +515,78 @@ class OneBotAdapter(PlatformAdapter): return True except Exception as e: + # 识别 OneBot 的“假失败”情况:如果由于图片过大导致超时,其实图片往往已在后台由 OneBot 自动重传并最终会成功。 + error_str = str(e).lower() + # 判定为“疑似成功”的特征:超时、1200、网络错误 + is_potential_success = ( + "timeout" in error_str or "1200" in error_str or "网络错误" in error_str + ) + + if is_potential_success: + logger.warning( + f"OneBot 发送群 {group_id} 图片出现疑似超时 ({e})。 " + "这通常是因为图片较大导致上传缓慢。如果群内稍后出现了图片,请忽略随后可能的重试提示。" + ) + return ( + False # 返回 False,由上层 RetryManager 接管(带 20s 延迟观察期) + ) + logger.error(f"OneBot 图片发送最终失败: {e}") return False + async def was_image_sent_recently(self, group_id: str, seconds: int = 60) -> bool: + """ + [真相检查] 检查最近 X 秒内,机器人是否已经向该群发送过图片。 + 用于判断之前的“超时/1200”错误是否其实已经在后台发送成功。 + """ + try: + # 1. 获取最近的消息历史 (OneBot 标准 API) + history = await self.bot.call_action( + "get_group_msg_history", + group_id=int(group_id), + count=50, # 只检查最近 50 条消息,足够覆盖大多数情况 + ) + + if not history or "messages" not in history: + # 某些 OneBot 实现返回值结构不同 + messages = history if isinstance(history, list) else [] + else: + messages = history["messages"] + + # 2. 逆序检查 + import time + + now = time.time() + self_id = str(getattr(self.bot, "self_id", "")) + + for msg in reversed(messages): + msg_time = msg.get("time", 0) + # 只检查约定时间范围内的消息 + if now - msg_time > seconds: + break + + # 检查发送者是否是机器人自己 + user_id = str( + msg.get("user_id", msg.get("sender", {}).get("user_id", "")) + ) + if user_id != self_id: + continue + + # 检查消息内容是否包含图片 + raw_message = msg.get("message", []) + # 适配字符串形式或列表形式的消息 + msg_str = str(raw_message) + if "[CQ:image" in msg_str or '"type": "image"' in msg_str: + logger.debug( + f"自检发现群 {group_id} 已有成功发送的图片回显,无需重试。" + ) + return True + + return False + except Exception as e: + logger.debug(f"回显自检失败 (可能不支持 get_group_msg_history): {e}") + return False + async def send_file( self, group_id: str, @@ -1016,6 +1085,7 @@ class OneBotAdapter(PlatformAdapter): raise e1 # 重新尝试两个可能的 API 名 + params = {} try: params = { "group_id": int(group_id), diff --git a/src/infrastructure/reporting/dispatcher.py b/src/infrastructure/reporting/dispatcher.py index 174360b..9166e12 100644 --- a/src/infrastructure/reporting/dispatcher.py +++ b/src/infrastructure/reporting/dispatcher.py @@ -115,7 +115,11 @@ class ReportDispatcher: if platform_id: await self.retry_manager.add_task( - html_content, analysis_result, group_id, platform_id + html_content, + analysis_result, + group_id, + platform_id, + caption="📊 每日群聊分析报告已生成:", ) return True # 已加入队列视作处理成功 (不在此处报错) else: diff --git a/src/infrastructure/scheduler/auto_scheduler.py b/src/infrastructure/scheduler/auto_scheduler.py index c54a2dc..90f7115 100644 --- a/src/infrastructure/scheduler/auto_scheduler.py +++ b/src/infrastructure/scheduler/auto_scheduler.py @@ -258,59 +258,45 @@ class AutoScheduler: async def _get_enabled_targets(self) -> set[tuple[str, str]]: """ 获取所有启用分析的群聊目标。 - - 根据群组列表模式(白名单/黑名单/无限制)过滤群聊, - 返回去重后的 (group_id, platform_id) 集合。 - - Returns: - set[tuple[str, str]]: 启用分析的 (群ID, 平台ID) 集合 + 确保每个群组 ID 在一次调度任务中只出现一次(多平台/多适配器去重)。 """ group_list_mode = self.config_manager.get_group_list_mode() - # 使用 set 存储 (group_id, platform_id) 元组,避免重复 - enabled_targets = set() + # 使用字典记录 group_id -> platform_id 实现去重 + # 优先级:先发现先处理 + targets_map: dict[str, str] = {} # 1. 通过 API 获取所有群组(自动发现) - logger.info(f"自动分析使用 {group_list_mode} 模式,正在获取群列表...") + logger.info(f"自动分析使用 {group_list_mode} 模式,正在扫描所有平台的群列表...") all_groups = await self._get_all_groups() - logger.info(f"共获取到 {len(all_groups)} 个群组") for platform_id, group_id in all_groups: - # 构造 UMO 进行权限检查 - umo = f"{platform_id}:GroupMessage:{group_id}" - if self.config_manager.is_group_allowed(umo): - enabled_targets.add((str(group_id), str(platform_id))) + group_id_str = str(group_id) + if group_id_str in targets_map: + continue - # 2. 白名单模式下,额外检查配置中的 UMO - # 解决 get_group_list 失败但配置了明确 UMO 的情况 + # 权限检查 + umo = f"{platform_id}:GroupMessage:{group_id_str}" + if self.config_manager.is_group_allowed(umo): + targets_map[group_id_str] = platform_id + + # 2. 白名单模式补全 if group_list_mode == "whitelist": whitelist_config = self.config_manager.get_group_list() - logger.info( - f"正在检查白名单配置中的额外 UMO ({len(whitelist_config)} 条)..." - ) - for item in whitelist_config: item = str(item).strip() - # 如果是 UMO 格式 (例: platform_id:GroupMessage:group_id) if ":" in item: parts = item.split(":") if len(parts) >= 3: p_id = parts[0] g_id = parts[-1] + if g_id not in targets_map: + if self.bot_manager.get_bot_instance(p_id): + targets_map[g_id] = p_id - # 检查该平台是否存在 - if self.bot_manager.get_bot_instance(p_id): - enabled_targets.add((str(g_id), str(p_id))) - logger.debug(f"添加白名单 UMO 目标: {item}") - else: - logger.warning( - f"白名单 UMO {item} 对应的平台 {p_id} 不存在或未加载" - ) - - logger.info( - f"根据 {group_list_mode} 过滤及合并后,共有 {len(enabled_targets)} 个群聊需要分析" - ) - + # 转换为集合形式返回 + enabled_targets = set(targets_map.items()) + logger.info(f"扫码完成:共有 {len(enabled_targets)} 个群聊目标将执行分析任务") return enabled_targets # ================================================================ @@ -375,7 +361,7 @@ class AutoScheduler: logger.error(f"自动分析执行失败: {e}", exc_info=True) async def _perform_auto_analysis_for_group_with_timeout( - self, group_id: str, target_platform_id: str = None + self, group_id: str, target_platform_id: str | None = None ): """为指定群执行自动分析(带超时控制)""" try: @@ -390,7 +376,7 @@ class AutoScheduler: logger.error(f"群 {group_id} 分析任务执行失败: {e}") async def _perform_auto_analysis_for_group( - self, group_id: str, target_platform_id: str = None + self, group_id: str, target_platform_id: str | None = None ): """为指定群执行自动分析(业务逻辑委派给 AnalysisApplicationService)""" # 为每个群聊使用独立的锁 @@ -534,7 +520,7 @@ class AutoScheduler: logger.error(f"增量分析执行失败: {e}", exc_info=True) async def _perform_incremental_analysis_for_group_with_timeout( - self, group_id: str, target_platform_id: str = None + self, group_id: str, target_platform_id: str | None = None ): """为指定群执行增量分析(带超时控制,10分钟)""" try: @@ -553,7 +539,7 @@ class AutoScheduler: return {"success": False, "reason": str(e)} async def _perform_incremental_analysis_for_group( - self, group_id: str, target_platform_id: str = None + self, group_id: str, target_platform_id: str | None = None ): """为指定群执行增量分析(业务逻辑委派给 AnalysisApplicationService)""" # 为每个群聊使用独立的锁 @@ -677,7 +663,7 @@ class AutoScheduler: logger.error(f"增量最终报告执行失败: {e}", exc_info=True) async def _perform_incremental_final_report_for_group_with_timeout( - self, group_id: str, target_platform_id: str = None + self, group_id: str, target_platform_id: str | None = None ): """为指定群生成增量最终报告(带超时控制,20分钟)""" try: @@ -696,7 +682,7 @@ class AutoScheduler: return {"success": False, "reason": str(e)} async def _perform_incremental_final_report_for_group( - self, group_id: str, target_platform_id: str = None + self, group_id: str, target_platform_id: str | None = None ): """为指定群生成增量最终报告(业务逻辑委派给 AnalysisApplicationService)""" # 为每个群聊使用独立的锁 diff --git a/src/infrastructure/scheduler/retry.py b/src/infrastructure/scheduler/retry.py index c86f17f..55b8b93 100644 --- a/src/infrastructure/scheduler/retry.py +++ b/src/infrastructure/scheduler/retry.py @@ -18,6 +18,7 @@ class RetryTask: analysis_result: dict # 保存原始分析结果,用于文本回退 group_id: str platform_id: str # 需要保存 platform_id 以便找回 Bot + caption: str = "" # 保存原始消息提示词 retry_count: int = 0 max_retries: int = 3 created_at: float = 0.0 @@ -46,6 +47,7 @@ class RetryManager: self.running = False self.worker_task = None self._dlq = [] # 死信队列 (Failures) + self._active_groups = set() # 正在处理中的群,防止重试地狱 async def start(self): """启动重试工作进程""" @@ -75,7 +77,12 @@ class RetryManager: logger.info("[RetryManager] 图片重试管理器已停止") async def add_task( - self, html_content: str, analysis_result: dict, group_id: str, platform_id: str + self, + html_content: str, + analysis_result: dict, + group_id: str, + platform_id: str, + caption: str = "", ): """添加重试任务""" if not self.running: @@ -84,59 +91,104 @@ class RetryManager: ) await self.start() + # 核心去重:如果该群已经在重试流程中,不再重复加入队列 + if group_id in self._active_groups: + logger.debug(f"[RetryManager] 群 {group_id} 已在重试观察期,跳过任务添加") + return + task = RetryTask( html_content=html_content, analysis_result=analysis_result, group_id=group_id, platform_id=platform_id, + caption=caption, created_at=time.time(), ) await self.queue.put(task) logger.info(f"[RetryManager] 已添加群 {group_id} 的重试任务") async def _worker(self): - """工作进程循环""" + """工作进程主循环:仅负责分发任务到协程,不阻塞""" while self.running: try: task: RetryTask = await self.queue.get() - # 延迟策略:指数回退 (5s, 10s, 20s...) + 随机波动 (1~5s) - jitter = random.uniform(1, 5) - delay = 5 * (2**task.retry_count) + jitter - - logger.info( - f"[RetryManager] 处理群 {task.group_id} 的重试任务 (第 {task.retry_count + 1} 次尝试)" - ) - - success = await self._process_task(task) - - if success: - logger.info(f"[RetryManager] 群 {task.group_id} 重试成功") + # 如果该群已经在处理中,不再重复处理(双重保险) + if task.group_id in self._active_groups: self.queue.task_done() - else: - task.retry_count += 1 - if task.retry_count < task.max_retries: - logger.warning( - f"[RetryManager] 群 {task.group_id} 重试失败,{delay}秒后再次尝试" - ) - asyncio.create_task(self._requeue_after_delay(task, delay)) - self.queue.task_done() - else: - logger.error( - f"[RetryManager] 群 {task.group_id} 超过最大重试次数,移入死信队列并尝试文本回退" - ) - self._dlq.append(task) - self.queue.task_done() - # 尝试发送文本回退 - await self._send_fallback_text(task) + continue + + # 启动非阻塞的延迟执行协程 + asyncio.create_task(self._run_task_with_delay(task)) + self.queue.task_done() except asyncio.CancelledError: break except Exception as e: - logger.error(f"[RetryManager] Worker 异常: {e}", exc_info=True) + logger.error(f"[RetryManager] Worker 调度异常: {e}") await asyncio.sleep(1) + async def _run_task_with_delay(self, task: RetryTask): + """异步执行带延迟的单体重试任务""" + # 锁定该群,防止其他重试任务进入 + if task.group_id in self._active_groups: + return + self._active_groups.add(task.group_id) + + try: + # 1. 延迟策略:指数回落 + 观察期 (应对 OneBot 假超时) + jitter = random.uniform(2, 8) + delay = 20 * (2**task.retry_count) + jitter + + if task.retry_count == 0: + logger.info( + f"[RetryManager] 群 {task.group_id} 启动 {delay:.1f}s 重试观察期..." + ) + + await asyncio.sleep(delay) + + if not self.running: + return + + # 【核心判定逻辑】:睡醒后先别急着发,去群里看看那张“疑似失败”的图是不是其实已经出来了 + adapter = self.bot_manager.get_adapter(task.platform_id) + if adapter and hasattr(adapter, "was_image_sent_recently"): + # 检查过去 3 分钟内的消息回显 (覆盖初发和之前的重试) + if await adapter.was_image_sent_recently(task.group_id, seconds=180): + logger.info( + f"[RetryManager] 根据消息回显判断,群 {task.group_id} 的图片报告已成功送达。取消后续重试。" + ) + return + + # 2. 执行渲染与发送 + success = await self._process_task(task) + + if success: + logger.info(f"[RetryManager] 群 {task.group_id} 重试发送成功") + else: + # 3. 失败后续处理 + task.retry_count += 1 + if task.retry_count < task.max_retries: + # 释放锁,以便下次取到时能重新进入观察期 + self._active_groups.discard(task.group_id) + await self.queue.put(task) + logger.warning( + f"[RetryManager] 群 {task.group_id} 本轮重试失败,准备进入第 {task.retry_count + 1} 轮..." + ) + else: + logger.error( + f"[RetryManager] 群 {task.group_id} 已达最大重试次数,执行文本回退" + ) + await self._send_fallback_text(task) + except Exception as e: + logger.error(f"[RetryManager] 重试协程执行异常 (群 {task.group_id}): {e}") + finally: + # 确保最终释放该群的重试锁 + if task.group_id in self._active_groups: + self._active_groups.discard(task.group_id) + async def _requeue_after_delay(self, task: RetryTask, delay: float): + # 这是一个遗留辅助方法,新逻辑已在 _run_task_with_delay 中处理 await asyncio.sleep(delay) await self.queue.put(task) @@ -238,7 +290,9 @@ class RetryManager: # 适配器内部通常应处理好 bytes/base64 的发送。 # 这里我们尝试直接传 image_file_str (base64://) try: - success = await adapter.send_image(task.group_id, image_file_str) + success = await adapter.send_image( + task.group_id, image_file_str, caption=task.caption + ) return success except Exception as e: logger.error(f"[RetryManager] 适配器发送图片异常: {e}") @@ -248,9 +302,6 @@ class RetryManager: logger.error(f"[RetryManager] 处理任务时发生意外错误: {e}", exc_info=True) return False - except Exception: - pass - async def _send_fallback_text(self, task: RetryTask): """发送文本回退报告(业务逻辑委派给适配器)""" if not self.report_generator: