diff --git a/main.py b/main.py index d6ad796..220785a 100644 --- a/main.py +++ b/main.py @@ -45,7 +45,6 @@ from .src.infrastructure.platform.template_preview import ( ) from .src.infrastructure.reporting.generators import ReportGenerator from .src.infrastructure.scheduler.auto_scheduler import AutoScheduler -from .src.infrastructure.scheduler.retry import RetryManager from .src.shared.trace_context import TraceContext, TraceLogFilter from .src.utils.logger import logger from .src.utils.pdf_utils import PDFInstaller @@ -72,7 +71,6 @@ class GroupDailyAnalysis(Star): template_command_service: TemplateCommandService telegram_template_preview_handler: TelegramTemplatePreviewHandler template_preview_router: TemplatePreviewRouter - retry_manager: RetryManager auto_scheduler: AutoScheduler message_sender: MessageSender @@ -145,18 +143,12 @@ class GroupDailyAnalysis(Star): handlers=[self.telegram_template_preview_handler] ) - # 调度与重试 - self.retry_manager = RetryManager( - self.bot_manager, self.html_render, self.report_generator - ) - self.message_sender = MessageSender( - self.bot_manager, self.config_manager, self.retry_manager - ) + # 调度与发送 + self.message_sender = MessageSender(self.bot_manager, self.config_manager) self.auto_scheduler = AutoScheduler( self.config_manager, self.analysis_service, self.bot_manager, - self.retry_manager, self.report_generator, self.html_render, plugin_instance=self, @@ -239,10 +231,6 @@ class GroupDailyAnalysis(Star): if self.auto_scheduler: self.auto_scheduler.schedule_jobs(self.context) - # 4. 始终启动重试管理器 - if self.retry_manager: - await self.retry_manager.start() - self._initialized = True self._discovery_run = True logger.info(f"插件任务注册完成 (来源: {source})") @@ -278,9 +266,6 @@ class GroupDailyAnalysis(Star): logger.debug("正在停止自动调度器...") self.auto_scheduler.unschedule_jobs(self.context) - if self.retry_manager: - await self.retry_manager.stop() - if self.template_preview_router: await self.template_preview_router.unregister_handlers() @@ -625,23 +610,16 @@ class GroupDailyAnalysis(Star): if image_url: caption = TraceContext.make_report_caption() - await adapter.send_image(group_id, image_url, caption=caption) - await self._try_upload_image(group_id, image_url, platform_id) - elif html_content: - yield event.plain_result("⚠️ 群分析报告图片发送失败,自动重试中。") - caption = TraceContext.make_report_caption() - await self.retry_manager.add_task( - html_content, - analysis_result, - group_id, - platform_id, - caption=caption, - ) - else: - text_report = self.report_generator.generate_text_report( - analysis_result - ) - yield event.plain_result(f"⚠️ 图片生成失败,回退文本:\n\n{text_report}") + sent = await adapter.send_image(group_id, image_url, caption=caption) + if sent: + await self._try_upload_image(group_id, image_url, platform_id) + return # 成功发送 + + # 如果图片生成或发送失败,直接回退到文本 + logger.warning(f"图片报告发送失败,正在发送文本回退报告。群: {group_id}") + text_report = self.report_generator.generate_text_report(analysis_result) + await adapter.send_text_report(group_id, text_report) + return elif output_format == "pdf": pdf_path = await self.report_generator.generate_pdf_report( @@ -694,8 +672,7 @@ class GroupDailyAnalysis(Star): else: text_report = self.report_generator.generate_text_report(analysis_result) - if not await adapter.send_text(group_id, text_report): - yield event.plain_result(text_report) + await adapter.send_text_report(group_id, text_report) @filter.command("设置格式", alias={"set_format"}) @filter.permission_type(PermissionType.ADMIN) diff --git a/src/infrastructure/messaging/message_sender.py b/src/infrastructure/messaging/message_sender.py index 1ed0817..dab4834 100644 --- a/src/infrastructure/messaging/message_sender.py +++ b/src/infrastructure/messaging/message_sender.py @@ -12,10 +12,9 @@ class MessageSender: 封装了 PlatformAdapter 的底层调用,提供更高层的发送接口 """ - def __init__(self, bot_manager, config_manager, retry_manager): + def __init__(self, bot_manager, config_manager): self.bot_manager = bot_manager self.config_manager = config_manager - self.retry_manager = retry_manager async def send_text( self, group_id: str, text: str, platform_id: str | None = None diff --git a/src/infrastructure/platform/adapters/discord_adapter.py b/src/infrastructure/platform/adapters/discord_adapter.py index 6c57fb8..1157e51 100644 --- a/src/infrastructure/platform/adapters/discord_adapter.py +++ b/src/infrastructure/platform/adapters/discord_adapter.py @@ -124,7 +124,7 @@ class DiscordAdapter(PlatformAdapter): list[UnifiedMessage]: 统一格式的消息对象列表 """ if not discord: - logger.error("Discord module (py-cord) not found. Cannot fetch messages.") + logger.error("未找到 Discord 模块 (py-cord),无法拉取历史消息。") return [] try: diff --git a/src/infrastructure/platform/adapters/lark_adapter.py b/src/infrastructure/platform/adapters/lark_adapter.py index 5223a71..092c26f 100644 --- a/src/infrastructure/platform/adapters/lark_adapter.py +++ b/src/infrastructure/platform/adapters/lark_adapter.py @@ -857,23 +857,6 @@ class LarkAdapter(PlatformAdapter): logger.error(f"飞书文件发送失败: {e}") return False - async def send_forward_msg(self, group_id: str, nodes: list[dict]) -> bool: - if not nodes: - return True - chunks: list[str] = ["📊 群分析报告摘要"] - for node in nodes: - data = node.get("data", node) - name = str(data.get("name", "AstrBot")) - content = data.get("content", "") - if isinstance(content, list): - text_parts = [] - for seg in content: - if isinstance(seg, dict) and seg.get("type") == "text": - text_parts.append(str(seg.get("data", {}).get("text", ""))) - content = "".join(text_parts) - chunks.append(f"[{name}] {content}") - return await self.send_text(group_id, "\n".join(chunks)) - async def get_group_info(self, group_id: str) -> UnifiedGroup | None: if not self._lark_client or not self._lark_client.im: return None diff --git a/src/infrastructure/platform/adapters/onebot_adapter.py b/src/infrastructure/platform/adapters/onebot_adapter.py index 2751e83..c2cfc53 100644 --- a/src/infrastructure/platform/adapters/onebot_adapter.py +++ b/src/infrastructure/platform/adapters/onebot_adapter.py @@ -22,7 +22,6 @@ from ....domain.value_objects.unified_message import ( MessageContentType, UnifiedMessage, ) -from ....shared.trace_context import REPORT_CAPTION_PATTERN from ....utils.logger import logger from ..base import PlatformAdapter @@ -494,7 +493,7 @@ class OneBotAdapter(PlatformAdapter): """ try: use_base64 = False - plugin = self.config.get("plugin_instance") if self.config else None + plugin: Any = self.config.get("plugin_instance") if self.config else None if plugin and hasattr(plugin, "config_manager"): use_base64 = plugin.config_manager.get_enable_base64_image() @@ -574,150 +573,16 @@ class OneBotAdapter(PlatformAdapter): logger.info(f"Base64 回退模式发送图片成功: 群 {group_id}") 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 + await self.bot.call_action( + "send_group_msg", + group_id=int(group_id), + message=message, ) + logger.info(f"Base64 回退模式发送图片成功: 群 {group_id}") + return True - if is_potential_success: - logger.warning( - f"OneBot 发送群 {group_id} 图片出现疑似超时 ({e})。 " - "进入多轮观察期,尝试通过历史回显核实..." - ) - - # Multi-stage observation. Some OneBot implementations commit - # history with delay after timeout-like errors. - observe_windows = (10, 20, 30) - for wait_seconds in observe_windows: - await asyncio.sleep(wait_seconds) - if await self.was_image_sent_recently( - group_id, seconds=420, token=caption - ): - logger.info( - f"[OneBot] [真相拦截] 群 {group_id} 在 {wait_seconds}s 观察后确认已送达,拦截重试。" - ) - return True - - 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, token: str | None = None - ) -> bool: - """ - [真相检查] 检查最近 X 秒内,机器人是否已经向该群发送过图片。 - 用于判断之前的“超时/1200”错误是否其实已经在后台发送成功。 - """ - try: - # 1. 获取最近的消息历史 (OneBot 标准 API) - try: - history = await self.bot.call_action( - "get_group_msg_history", - group_id=int(group_id), - count=100, # 增大扫描深度以应对高频群聊 - ) - except Exception as e: - logger.warning( - f"[OneBot] was_image_sent_recently: get_group_msg_history 失败 (可能 API 繁忙): {e}" - ) - return False # API 失败时,我们保持谨慎,但不阻止重试 - - if not history or "messages" not in history: - messages = history if isinstance(history, list) else [] - else: - messages = history["messages"] - - # 2. 逆序检查 - import time - - now = time.time() - # 1. 优先从内存缓存中获取机器人 ID - self_id = self.bot_self_ids[0] if self.bot_self_ids else "" - - if not self_id: - # 尝试从 bot 实例中获取多个可能的 ID 属性 - self_id = ( - str(getattr(self.bot, "self_id", "")) - or str(getattr(self.bot, "uin", "")) - or str(getattr(self.bot, "user_id", "")) - ) - - if not self_id: - # 最后的 API 兜底:尝试从 login_info 获取 - try: - login_info = await self.bot.call_action("get_login_info") - if login_info and "user_id" in login_info: - self_id = str(login_info["user_id"]) - # 更新缓存,下次无需重复请求 - if self_id not in self.bot_self_ids: - self.bot_self_ids.append(self_id) - logger.info(f"[OneBot] 成功通过 API 获取到机器人 ID: {self_id}") - except Exception as e: - logger.debug( - f"[OneBot] was_image_sent_recently: get_login_info API 调用失败: {e}" - ) - - if not self_id: - logger.warning( - "[OneBot] was_image_sent_recently: 无法确定机器人 ID,历史回显校验可能不准确" - ) - - # [优化] 从 Caption 中提取基于时间戳的去重 Token - search_token = None - if token: - match = REPORT_CAPTION_PATTERN.search(token) - if match: - search_token = match.group(0) # 例如 "| 03-12 17:33:20" - - 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 not in self.bot_self_ids: - # 如果内存中没有,尝试最后一次实时提取作为兜底 - if not self_id or user_id != self_id: - continue - - # 检查消息内容是否包含图片 - raw_message = msg.get("message", []) - # 适配字符串形式或列表形式的消息 - msg_str = str(raw_message) - - has_image = "[CQ:image" in msg_str or '"type": "image"' in msg_str - - if has_image: - if search_token: - # 精确匹配 TraceID - if search_token in msg_str: - logger.info( - f"[OneBot] [真相检查] 发现匹配 ID ({search_token}) 的历史图片。拦截重复发送。群: {group_id}" - ) - return True - else: - logger.debug( - f"[OneBot] [真相检查] 发现机器人发送的图片,但 ID 不匹配。跳过。群: {group_id}" - ) - else: - # 广义匹配(回退模式) - logger.info( - f"[OneBot] [真相检查] 发现近期发送过的图片回显 (广义匹配)。无需重试。群: {group_id}" - ) - return True - - return False except Exception as e: - logger.debug(f"回显自检失败: {e}") + logger.error(f"OneBot 图片发送最终失败: {e}") return False async def send_file( @@ -761,10 +626,13 @@ class OneBotAdapter(PlatformAdapter): file=file_b64, name=filename or os.path.basename(file_path), ) - logger.info(f"Base64 回退模式发送文件成功: {filename or file_path}") - return True + # ... 实现省略 ... + logger.info( + f"[OneBot] 文件发送成功(Base64 模式): {filename or file_path}" + ) + return True except Exception as e: - logger.error(f"OneBot 文件发送最终失败: {e}") + logger.error(f"[OneBot] 文件发送最终失败: {e}") return False async def send_forward_msg( @@ -774,13 +642,6 @@ class OneBotAdapter(PlatformAdapter): ) -> bool: """ 发送群合并转发消息。 - - Args: - group_id (str): 目标群号 - nodes (list[dict]): 转发节点列表 - - Returns: - bool: 是否发送成功 """ if not hasattr(self.bot, "call_action"): return False @@ -799,7 +660,7 @@ class OneBotAdapter(PlatformAdapter): ) return True except Exception as e: - logger.warning(f"OneBot 发送合并转发消息失败: {e}") + logger.warning(f"[OneBot] 发送合并转发消息失败: {e}") return False # ==================== IGroupInfoRepository 实现 ==================== diff --git a/src/infrastructure/platform/base.py b/src/infrastructure/platform/base.py index 4df2109..9ef06fb 100644 --- a/src/infrastructure/platform/base.py +++ b/src/infrastructure/platform/base.py @@ -4,6 +4,7 @@ from abc import ABC, abstractmethod from collections.abc import Mapping +from typing import Any from ...domain.repositories.avatar_repository import IAvatarRepository from ...domain.repositories.message_repository import ( @@ -25,20 +26,22 @@ class PlatformAdapter( 充当领域层与具体聊天平台(如 OneBot, Discord)之间的中转站。 Attributes: - bot (object): 平台对应的机器人 SDK 实例 + bot (Any): 平台对应的机器人 SDK 实例,显式标注为 Any 以支持动态属性调用 config (dict): 针对该平台的特定配置 """ + bot: Any + def __init__( self, - bot_instance: object, - config: Mapping[str, object] | None = None, + bot_instance: Any, + config: Mapping[str, Any] | None = None, ): """ 初始化平台适配器。 Args: - bot_instance (object): 后端机器人实例 + bot_instance (Any): 后端机器人实例 config (dict, optional): 平台特定配置项 """ self.bot = bot_instance @@ -46,11 +49,13 @@ class PlatformAdapter( self.bot_self_ids: list[str] = [] self._capabilities: PlatformCapabilities | None = None - def set_context(self, context: object) -> None: + def set_context(self, context: Any): """ - 可选的上下文注入钩子,供需要访问插件核心服务的适配器使用。 + 设置上下文对象(用于部分需要 ctx 的平台如 Telegram)。 + + Args: + context (Any): 上下文对象 """ - # 具体适配器可覆盖此方法 pass @property @@ -102,6 +107,52 @@ class PlatformAdapter( """ raise NotImplementedError + async def send_forward_msg( + self, + group_id: str, + nodes: list[dict], + ) -> bool: + """ + 发送合并转发消息(基类默认实现:转换为格式化文本分段发送)。 + 各适配器可覆盖此方法实现原生合并转发。 + """ + if not nodes: + return True + + # 万能回退:将节点重新组合成易读的长文本 + lines = [] + for node in nodes: + data = node.get("data", node) + name = data.get("name", "Daily Analysis") + content = data.get("content", "") + if content: + lines.append(f"【{name}】\n{content}") + + full_text = "\n\n".join(lines) + + # 处理超长文本分段(取大部分平台的安全阈值 1800 字符) + max_chunk_size = 1800 + if len(full_text) > max_chunk_size: + # 尝试在换行处拆分 + chunks = [] + curr = full_text + while len(curr) > max_chunk_size: + # 寻找最近的换行符 + split_idx = curr.rfind("\n", 0, max_chunk_size) + if split_idx == -1: + split_idx = max_chunk_size + chunks.append(curr[:split_idx].strip()) + curr = curr[split_idx:].strip() + if curr: + chunks.append(curr) + + for chunk in chunks: + if not await self.send_text(group_id, chunk): + return False + return True + else: + return await self.send_text(group_id, full_text) + async def set_reaction( self, group_id: str, message_id: str, emoji: str | int, is_add: bool = True ) -> bool: @@ -118,3 +169,43 @@ class PlatformAdapter( bool: 平台是否支持并成功执行 """ return False + + async def send_text_report(self, group_id: str, content: str) -> bool: + """ + 以最适合当前平台的方式发送长文本报告。 + 默认逻辑:将长文本切分为多个节点,然后调用 send_forward_msg。 + 各平台适配器通过实现 send_forward_msg 来决定最终呈现形式(合并转发、分段发送等)。 + """ + import re + + try: + # 1. 准备节点基础信息 + self_id = self.bot_self_ids[0] if self.bot_self_ids else "bot" + self_name = "分析报告" + # 2. 切分文本为逻辑段落(按标题、空行切分) + raw_content = str(content) + sections = re.split(r"\n+(?=[🎯📊💬🏆])|\n{2,}", raw_content.strip()) + nodes = [] + + for sec in sections: + if not sec.strip(): + continue + nodes.append( + { + "type": "node", + "data": { + "name": self_name, + "uin": self_id, + "content": sec.strip(), + }, + } + ) + + if not nodes: + return await self.send_text(group_id, raw_content) + + # 3. 尝试发送转发消息/长消息链 + return await self.send_forward_msg(group_id, nodes) + except Exception: + # 兜底:直接发送 + return await self.send_text(group_id, str(content)) diff --git a/src/infrastructure/reporting/dispatcher.py b/src/infrastructure/reporting/dispatcher.py index fb82b4c..872187b 100644 --- a/src/infrastructure/reporting/dispatcher.py +++ b/src/infrastructure/reporting/dispatcher.py @@ -15,11 +15,15 @@ class ReportDispatcher: 负责协调报告生成、格式选择、消息发送和失败重试 """ - def __init__(self, config_manager, report_generator, message_sender, retry_manager): + def __init__( + self, + config_manager, + report_generator, + message_sender, + ): self.config_manager = config_manager self.report_generator = report_generator self.message_sender = message_sender - self.retry_manager = retry_manager self._html_render_func: Callable | None = None def set_html_render(self, render_func: Callable): @@ -88,44 +92,21 @@ class ReportDispatcher: logger.error(f"[{trace_id}] Failed to generate image report: {e}") # image_url and html_content remain None - # 3. 发送图片 + # 4. 发送图片 if image_url: caption = TraceContext.make_report_caption() sent = await self.message_sender.send_image_smart( group_id, image_url, caption, platform_id ) if sent: - # 4. 发送成功后,尝试上传到群文件/群相册(静默处理) + # 5. 发送成功后,尝试上传到群文件/群相册(静默处理) await self._try_upload_image(group_id, image_url, platform_id) return True - # 5. 发送失败或生成失败的处理 -> 加入重试队列 - if html_content: - logger.warning( - f"[{trace_id}] Image dispatch failed, adding to retry queue..." - ) - # 尝试获取 platform_id 如果没有提供 - if not platform_id: - platforms = self.message_sender._get_available_platforms(group_id) - if platforms: - platform_id = platforms[0][0] # use first available - - if platform_id: - await self.retry_manager.add_task( - html_content, - analysis_result, - group_id, - platform_id, - caption=TraceContext.make_report_caption(), - ) - return True # 已加入队列视作处理成功 (不在此处报错) - else: - logger.error( - f"[{trace_id}] Cannot add to retry queue: No platform_id available." - ) - - # 6. 最终回退:文本报告 - logger.warning(f"[{trace_id}] Falling back to text report.") + # 6. 最终回退:如果图片发送失败(包括生成失败或发送接口报错),直接尝试发送文本报告 + logger.warning( + f"[{trace_id}] Image dispatch failed, falling back to text report." + ) return await self._dispatch_text(group_id, analysis_result, platform_id) async def _dispatch_pdf( @@ -198,13 +179,20 @@ class ReportDispatcher: async def _dispatch_text( self, group_id: str, analysis_result: dict[str, Any], platform_id: str | None ) -> bool: + """分发文本报告""" + logger.info(f"[分发器] 正在向群组 {group_id} 分发文本报告") + text_report = self.report_generator.generate_text_report(analysis_result) + adapter = self.message_sender.bot_manager.get_adapter(platform_id) + # 尝试通过适配器发送文本报告 + logger.info(f"[分发器] 正在尝试通过适配器发送文本报告。群: {group_id}") try: - text_report = self.report_generator.generate_text_report(analysis_result) + if adapter and await adapter.send_text_report(group_id, text_report): + return True return await self.message_sender.send_text( group_id, f"📊 每日群聊分析报告:\n\n{text_report}", platform_id ) except Exception as e: - logger.error(f"[{TraceContext.get()}] Failed to dispatch text report: {e}") + logger.error(f"[分发器] 发送文本报告最终失败。群: {group_id}, 错误: {e}") return False # ================================================================ diff --git a/src/infrastructure/scheduler/auto_scheduler.py b/src/infrastructure/scheduler/auto_scheduler.py index bb62820..0016c51 100644 --- a/src/infrastructure/scheduler/auto_scheduler.py +++ b/src/infrastructure/scheduler/auto_scheduler.py @@ -25,7 +25,6 @@ class AutoScheduler: config_manager, analysis_service, bot_manager, - retry_manager, report_generator=None, html_render_func=None, plugin_instance: Any | None = None, @@ -33,15 +32,14 @@ class AutoScheduler: self.config_manager = config_manager self.analysis_service = analysis_service self.bot_manager = bot_manager - self.retry_manager = retry_manager self.report_generator = report_generator self.html_render_func = html_render_func self.plugin_instance = plugin_instance # 初始化核心组件 - self.message_sender = MessageSender(bot_manager, config_manager, retry_manager) + self.message_sender = MessageSender(bot_manager, config_manager) self.report_dispatcher = ReportDispatcher( - config_manager, report_generator, self.message_sender, retry_manager + config_manager, report_generator, self.message_sender ) if html_render_func: self.report_dispatcher.set_html_render(html_render_func) diff --git a/src/infrastructure/scheduler/retry.py b/src/infrastructure/scheduler/retry.py deleted file mode 100644 index 8c51aa1..0000000 --- a/src/infrastructure/scheduler/retry.py +++ /dev/null @@ -1,436 +0,0 @@ -import asyncio -import base64 -import hashlib -import random -import time -from collections.abc import Callable -from dataclasses import dataclass - -import aiohttp - -from ...shared.trace_context import REPORT_CAPTION_PATTERN -from ...utils.logger import logger - - -@dataclass -class RetryTask: - """重试任务数据类""" - - html_content: str - analysis_result: dict # 保存原始分析结果,用于文本回退 - group_id: str - platform_id: str # 需要保存 platform_id 以便找回 Bot - caption: str = "" # 保存原始消息提示词 - retry_count: int = 0 - max_retries: int = 2 - created_at: float = 0.0 - task_key: str = "" - - def __post_init__(self): - if self.created_at == 0.0: - self.created_at = time.time() - - -class RetryManager: - """ - 重试管理器 - - 实现了一个简单的延迟队列 + 死信队列机制: - 1. 任务加入队列 - 2. Worker 取出任务,尝试执行 - 3. 失败则指数退避(延迟)后放回队列 - 4. 超过最大重试次数放入死信队列 - """ - - def __init__(self, bot_manager, html_render_func: Callable, report_generator=None): - self.bot_manager = bot_manager - self.html_render_func = html_render_func - self.report_generator = report_generator # 用于生成文本报告 - self.queue = asyncio.Queue() - self.running = False - self.worker_task = None - self._dlq = [] # 死信队列 (Failures) - self._active_groups = set() # 正在处理中的群,防止重试地狱 - self._active_task_keys = set() # 任务级去重锁,防止同一报告重复入队 - self._recent_task_key_ts: dict[str, float] = {} # 最近完成任务用于短时去重 - self._dedupe_ttl_seconds = 15 * 60 - - async def start(self): - """启动重试工作进程""" - if self.running: - return - self.running = True - self.worker_task = asyncio.create_task(self._worker()) - logger.info("[RetryManager] 图片重试管理器已启动") - - async def stop(self): - """停止重试工作进程""" - self.running = False - if self.worker_task: - self.worker_task.cancel() - try: - await self.worker_task - except asyncio.CancelledError: - pass - - # 检查剩余任务 - pending_count = self.queue.qsize() - if pending_count > 0: - logger.warning( - f"[RetryManager] 停止时仍有 {pending_count} 个任务在队列中 pending" - ) - self._active_groups.clear() - self._active_task_keys.clear() - - logger.info("[RetryManager] 图片重试管理器已停止") - - async def add_task( - self, - html_content: str, - analysis_result: dict, - group_id: str, - platform_id: str, - caption: str = "", - ): - """添加重试任务""" - if not self.running: - logger.warning( - "[RetryManager] 警告:添加任务时管理器未运行,正在尝试启动..." - ) - await self.start() - - self._cleanup_expired_task_keys() - task_key = self._build_task_key(group_id, platform_id, caption, html_content) - - # 任务级去重:同一份报告在观察窗口内只保留一个任务 - if task_key in self._active_task_keys: - logger.debug(f"[RetryManager] 任务 {task_key} 已在重试流程中,跳过重复入队") - return - - if task_key in self._recent_task_key_ts: - logger.debug( - f"[RetryManager] 任务 {task_key} 在去重窗口内已处理过,跳过重复入队" - ) - 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(), - task_key=task_key, - ) - self._active_task_keys.add(task_key) - self._recent_task_key_ts[task_key] = time.time() - await self.queue.put(task) - logger.info(f"[RetryManager] 已添加群 {group_id} 的重试任务 (key={task_key})") - - async def _worker(self): - """工作进程主循环:仅负责分发任务到协程,不阻塞""" - while self.running: - try: - task: RetryTask = await self.queue.get() - - # 【修复】去掉原本在这里的 group_id in self._active_groups 判断 - # 因为重新排队的任务本身就在 active_groups 中,会导致任务死在队列里被彻底丢弃 - - # 启动非阻塞的延迟执行协程 - 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}") - await asyncio.sleep(1) - - def _cleanup_expired_task_keys(self): - """清理过期的去重记录。""" - now = time.time() - expired_keys = [ - key - for key, ts in self._recent_task_key_ts.items() - if now - ts > self._dedupe_ttl_seconds - ] - for key in expired_keys: - self._recent_task_key_ts.pop(key, None) - - def _build_task_key( - self, group_id: str, platform_id: str, caption: str, html_content: str - ) -> str: - """ - 构造用于去重的稳定任务 key。 - - 优先级: - 1)优先使用 caption 中的 TraceID 时间戳 - 2)若无则取 HTML 内容的哈希前缀 - """ - token = None - if caption: - match = REPORT_CAPTION_PATTERN.search(caption) - if match: - token = match.group(0) - - if not token: - digest = hashlib.sha1(html_content[:2048].encode("utf-8")).hexdigest()[:16] - token = f"html:{digest}" - - return f"{platform_id}:{group_id}:{token}" - - def _mark_task_finished(self, task: RetryTask): - """释放任务级去重锁并刷新冷却时间。""" - if task.task_key: - self._active_task_keys.discard(task.task_key) - self._recent_task_key_ts[task.task_key] = time.time() - - async def _run_task_with_delay(self, task: RetryTask): - """异步执行带延迟的单体重试任务""" - # 锁定该群,防止其他“新”重试任务进入。 - # 如果是重试任务(retry_count > 0),它已经在队列循环中,之前已经释放过锁。 - if task.group_id in self._active_groups and task.retry_count == 0: - # 群正在被其它重试任务占用,重新排队避免任务被静默丢弃 - await asyncio.sleep(5) - if self.running: - await self.queue.put(task) - return - self._active_groups.add(task.group_id) - - try: - # 1. 策略计算:指数回落 + 抖动 - 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 重试观察期..." - ) - else: - logger.info( - f"[RetryManager] 群 {task.group_id} 准备第 {task.retry_count + 1} 轮重试,退避 {delay:.1f}s..." - ) - - await asyncio.sleep(delay) - - if not self.running: - return - - # 【真相检查 1】:睡醒后先核实群里图片是不是其实已经出来了 - adapter = self.bot_manager.get_adapter(task.platform_id) - if adapter and hasattr(adapter, "was_image_sent_recently"): - # 检查过去 5 分钟内的消息回显 (覆盖初发和之前的重试) - if await adapter.was_image_sent_recently( - task.group_id, seconds=300, token=task.caption - ): - logger.info( - f"[RetryManager] [拦截] 根据历史回显,群 {task.group_id} 的图片已成功送达。取消本次重试。" - ) - return - - # 2. 执行渲染与发送 - success = await self._process_task(task) - - if success: - logger.info(f"[RetryManager] 群 {task.group_id} 重试流程圆满完成") - self._mark_task_finished(task) - else: - # 3. 失败后续处理 - task.retry_count += 1 - if task.retry_count < task.max_retries: - # 将任务重新放回队列。 - # 注意:锁会在 finally 释放,这样下一个 worker 就能拉取到它并进入睡眠。 - await self.queue.put(task) - logger.warning( - f"[RetryManager] 群 {task.group_id} 本轮调用返回失败,已排期下一轮..." - ) - else: - logger.error( - f"[RetryManager] 群 {task.group_id} 已达最大重试次数,执行文本回退" - ) - await self._send_fallback_text(task) - self._mark_task_finished(task) - except Exception as e: - logger.error(f"[RetryManager] 重试协程发生意外: {e}", exc_info=True) - self._mark_task_finished(task) - 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) - - async def _process_task(self, task: RetryTask) -> bool: - """执行具体的渲染和发送逻辑""" - try: - # 1. 尝试渲染 - image_options = { - "full_page": True, - "type": "jpeg", - "quality": 85, - } - logger.debug(f"[RetryManager] 正在重新渲染群 {task.group_id} 的图片...") - - # 修改:return_url=False 获取二进制数据而不是 URL - # 这可以规避 NTQQ/NT 的“Timeout”假失败,因为下载本地/内网 URL 造成的网络等待会被跳过 - image_data = await self.html_render_func( - task.html_content, - {}, - False, # return_url=False,获取 bytes - image_options, - ) - - # 修复:某些实现即使 return_url=False 也仍返回 URL 字符串 - if isinstance(image_data, str): - if image_data.startswith(("http://", "https://")): - logger.warning( - f"[RetryManager] html_render 返回了 URL 而不是 bytes,尝试下载: {image_data}" - ) - async with aiohttp.ClientSession() as session: - async with session.get(image_data) as resp: - if resp.status == 200: - image_data = await resp.read() - else: - logger.error( - f"[RetryManager] 下载重试图片失败: {resp.status}" - ) - image_data = None - else: - # 本地文件路径 - try: - import os - - if os.path.exists(image_data): - with open(image_data, "rb") as f: - image_data = f.read() - - # 校验文件头 (防御性编程,避免发送错误文本) - if not image_data.startswith( - b"\xff\xd8" - ) and not image_data.startswith(b"\x89PNG"): - if len(image_data) < 1024 and ( - b"Error" in image_data or b"Exception" in image_data - ): - logger.error( - f"[RetryManager] 渲染器生成了错误文件而非图片: {image_data.decode('utf-8', errors='ignore')}" - ) - return False - else: - logger.error( - f"[RetryManager] 渲染器返回的路径不存在: {image_data}" - ) - image_data = None - except Exception as e: - logger.error(f"[RetryManager] 读取本地图片失败: {e}") - image_data = None - - if not image_data: - logger.warning( - f"[RetryManager] 重新渲染失败(返回空数据){task.group_id}" - ) - return False - - # 将 bytes 转换为 base64 字符串 - try: - base64_str = base64.b64encode(image_data).decode("utf-8") - image_file_str = f"base64://{base64_str}" - logger.debug( - f"[RetryManager] 图片转Base64成功,长度: {len(base64_str)}" - ) - except Exception as e: - logger.error(f"[RetryManager] Base64编码失败: {e}") - return False - - # 2. 获取适配器 (DDD 基础设施层) - adapter = self.bot_manager.get_adapter(task.platform_id) - if not adapter: - logger.error( - f"[RetryManager] 平台 {task.platform_id} 的适配器未找到,无法重试" - ) - return False - - # 3. 【临界检查 2】发送图片前最后一次复核 (针对渲染耗时极长产生的盲窗) - # 例如渲染 10s 期间图片出来了,这里可以最后贴身拦截一次 - if adapter and hasattr(adapter, "was_image_sent_recently"): - if await adapter.was_image_sent_recently( - task.group_id, seconds=120, token=task.caption - ): - logger.info( - f"[RetryManager] [临界拦截] 渲染完成后检测到群 {task.group_id} 已有报告。拦截重复发送。" - ) - return True - - # 4. 执行实际发送 - logger.info( - f"[RetryManager] 正在向群 {task.group_id} 发送回补图片 (Adapter: {type(adapter).__name__})..." - ) - - # 注意:某些适配器可能需要 URL,某些需要 Base64。 - # 适配器内部通常应处理好 bytes/base64 的发送。 - # 这里我们尝试直接传 image_file_str (base64://) - try: - 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}") - return False - - except Exception as e: - logger.error(f"[RetryManager] 处理任务时发生意外错误: {e}", exc_info=True) - return False - - async def _send_fallback_text(self, task: RetryTask): - """发送文本回退报告(业务逻辑委派给适配器)""" - if not self.report_generator: - logger.warning("[RetryManager] 未配置 ReportGenerator,无法发送文本回退") - return - - try: - logger.info(f"[RetryManager] 正在为群 {task.group_id} 生成文本回退报告...") - text_report = self.report_generator.generate_text_report( - task.analysis_result - ) - - # 2. 获取适配器 (DDD 基础设施层) - adapter = self.bot_manager.get_adapter(task.platform_id) - if not adapter: - logger.error( - f"[RetryManager] 无法获取适配器 {task.platform_id},放弃发送回退文本" - ) - return - - nickname = "AstrBot日常分析" - nodes = [ - { - "type": "node", - "data": { - "name": nickname, - "content": "⚠️ 图片报告多次生成失败,为您呈现文本版报告:", - }, - }, - { - "type": "node", - "data": {"name": nickname, "content": text_report}, - }, - ] - - # 3. 通过适配器发送结构化消息 - success = await adapter.send_forward_msg(task.group_id, nodes) - - if success: - logger.info(f"[RetryManager] 群 {task.group_id} 文本回退报告发送成功") - else: - # 最终兜底:发送简单文本 - logger.warning("[RetryManager] 结构化发送失败,尝试直接发送文本回退") - await adapter.send_text( - task.group_id, - f"⚠️ 图片报告生成失败,文本报告:\n{text_report}"[:4500], - ) - - except Exception as e: - logger.error(f"[RetryManager] 文本回退流程异常: {e}", exc_info=True)