From e78c82785780f4a626acdcd58254aa6c81cbce71 Mon Sep 17 00:00:00 2001 From: SXP-Simon Date: Sat, 17 Jan 2026 23:39:17 +0800 Subject: [PATCH] =?UTF-8?q?feat(RetryManager):=20=E6=96=B0=E5=A2=9E?= =?UTF-8?q?=E5=9B=BE=E7=89=87=E6=8A=A5=E5=91=8A=E9=87=8D=E8=AF=95=E6=9C=BA?= =?UTF-8?q?=E5=88=B6=E5=B9=B6=E4=BC=98=E5=8C=96=E8=B0=83=E5=BA=A6=E7=A8=B3?= =?UTF-8?q?=E5=AE=9A=E6=80=A7?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 新增重试管理器 (src/scheduler/retry.py): 实现基于内存的延迟队列与死信队列,支持指数退避 + 随机抖动 (Jitter) 重试策略。 增强容错性:处理渲染服务宕机、网络超时及 Bot API 逻辑错误 (如 retcode=1200)。 兼容性优化:使用 call_action 通用接口适配 OneBot v11。 生命周期管理:增加自动启动保护,防止任务丢失。 --- .gitignore | 1 + main.py | 33 ++++- src/reports/generators.py | 66 ++++++--- src/scheduler/auto_scheduler.py | 53 ++++++- src/scheduler/retry.py | 238 ++++++++++++++++++++++++++++++++ 5 files changed, 359 insertions(+), 32 deletions(-) create mode 100644 src/scheduler/retry.py diff --git a/.gitignore b/.gitignore index 5803258..46f4eed 100644 --- a/.gitignore +++ b/.gitignore @@ -42,3 +42,4 @@ src/utils/__pycache__/helpers.cpython-311.pyc src/utils/__pycache__/pdf_utils.cpython-311.pyc src/visualization/__pycache__/__init__.cpython-311.pyc src/visualization/__pycache__/activity_charts.cpython-311.pyc +src/scheduler/__pycache__/retry.cpython-311.pyc diff --git a/main.py b/main.py index df59d25..8958cf8 100644 --- a/main.py +++ b/main.py @@ -22,6 +22,7 @@ from .src.core.bot_manager import BotManager from .src.core.config import ConfigManager from .src.reports.generators import ReportGenerator 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 @@ -39,12 +40,14 @@ class QQGroupDailyAnalysis(Star): context, self.config_manager, self.bot_manager ) self.report_generator = ReportGenerator(self.config_manager) + self.retry_manager = RetryManager(self.bot_manager, self.html_render) 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.html_render, # 传入html_render函数 ) @@ -77,6 +80,9 @@ class QQGroupDailyAnalysis(Star): status = self.bot_manager.get_status_info() logger.info(f"Bot管理器状态: {status}") + # 始终启动重试管理器,确保手动触发也能使用重试队列 + await self.retry_manager.start() + except Exception as e: logger.debug(f"延迟启动调度器失败,可能由于短时间内多次更新插件配置: {e}") @@ -91,6 +97,9 @@ class QQGroupDailyAnalysis(Star): await self.auto_scheduler.stop_scheduler() logger.info("自动调度器已停止") + if self.retry_manager: + await self.retry_manager.stop() + # 重置实例属性 self.auto_scheduler = None self.bot_manager = None @@ -185,19 +194,35 @@ class QQGroupDailyAnalysis(Star): # 生成报告 output_format = self.config_manager.get_output_format() if output_format == "image": - image_url = await self.report_generator.generate_image_report( + ( + image_url, + html_content, + ) = await self.report_generator.generate_image_report( analysis_result, group_id, self.html_render ) if image_url: yield event.image_result(image_url) + elif html_content: + # 生成失败但有HTML,加入重试队列 + logger.warning("图片报告生成失败,加入重试队列") + yield event.plain_result( + "[AstrBot QQ群日常分析总结插件] ⚠️ 图片报告暂无法生成,已加入重试队列,稍后将自动重试发送。" + ) + # 获取 platform_id + platform_id = await self.auto_scheduler.get_platform_id_for_group( + group_id + ) + await self.retry_manager.add_task( + html_content, group_id, platform_id + ) else: - # 如果图片生成失败,回退到文本报告 - logger.warning("图片报告生成失败,回退到文本报告") + # 如果图片生成失败且无HTML,回退到文本报告 + logger.warning("图片报告生成失败(无HTML),回退到文本报告") text_report = self.report_generator.generate_text_report( analysis_result ) yield event.plain_result( - f"⚠️ 图片报告生成失败,以下是文本版本:\n\n{text_report}" + f"[AstrBot QQ群日常分析总结插件] ⚠️ 图片报告生成失败,以下是文本版本:\n\n{text_report}" ) elif output_format == "pdf": if not self.config_manager.pyppeteer_available: diff --git a/src/reports/generators.py b/src/reports/generators.py index 986ddfb..c801972 100644 --- a/src/reports/generators.py +++ b/src/reports/generators.py @@ -26,8 +26,16 @@ class ReportGenerator: async def generate_image_report( self, analysis_result: dict, group_id: str, html_render_func - ) -> str | None: - """生成图片格式的分析报告""" + ) -> tuple[str | None, str | None]: + """ + 生成图片格式的分析报告 + + Returns: + tuple[str | None, str | None]: (image_url, html_content) + - image_url: 生成的图片URL,如果生成失败则为None + - html_content: 生成的HTML内容,如果渲染失败但HTML生成成功,则返回此内容供重试 + """ + html_content = None try: # 准备渲染数据 render_payload = await self._prepare_render_data( @@ -41,7 +49,7 @@ class ReportGenerator: # 检查HTML内容是否有效 if not html_content: logger.error("图片报告HTML渲染失败:返回空内容") - return None + return None, None logger.info(f"图片报告HTML渲染完成,长度: {len(html_content)} 字符") @@ -59,30 +67,42 @@ class ReportGenerator: image_options, ) - logger.info(f"图片生成成功: {image_url}") - return image_url + if image_url: + logger.info(f"图片生成成功: {image_url}") + return image_url, html_content + else: + # 渲染服务返回None,可能是渲染失败 + logger.warning("渲染服务返回空URL") + return None, html_content except Exception as e: logger.error(f"生成图片报告失败: {e}", exc_info=True) # 尝试使用更简单的选项作为后备方案 - try: - logger.info("尝试使用低质量选项重新生成...") - simple_options = { - "full_page": True, - "type": "jpeg", - "quality": 70, # 降低质量以提高兼容性 - } - image_url = await html_render_func( - html_content, # 使用已渲染的HTML - {}, # 空数据字典 - True, - simple_options, - ) - logger.info(f"使用低质量选项生成成功: {image_url}") - return image_url - except Exception as fallback_e: - logger.error(f"后备低质量方案也失败: {fallback_e}") - return None + if html_content: + try: + logger.info("尝试使用低质量选项重新生成...") + simple_options = { + "full_page": True, + "type": "jpeg", + "quality": 70, # 降低质量以提高兼容性 + } + image_url = await html_render_func( + html_content, # 使用已渲染的HTML + {}, # 空数据字典 + True, + simple_options, + ) + if image_url: + logger.info(f"使用低质量选项生成成功: {image_url}") + return image_url, html_content + else: + logger.warning("低质量作为后备方案也返回空URL") + return None, html_content + except Exception as fallback_e: + logger.error(f"后备低质量方案也失败: {fallback_e}") + return None, html_content + + return None, html_content async def generate_pdf_report( self, analysis_result: dict, group_id: str diff --git a/src/scheduler/auto_scheduler.py b/src/scheduler/auto_scheduler.py index 74f78a8..a305da0 100644 --- a/src/scheduler/auto_scheduler.py +++ b/src/scheduler/auto_scheduler.py @@ -23,6 +23,7 @@ class AutoScheduler: analyzer, report_generator, bot_manager, + retry_manager, # 新增 html_render_func=None, ): self.config_manager = config_manager @@ -30,6 +31,7 @@ class AutoScheduler: self.analyzer = analyzer self.report_generator = report_generator self.bot_manager = bot_manager + self.retry_manager = retry_manager # 保存引用 self.html_render_func = html_render_func self.scheduler_task = None self.last_execution_date = None # 记录上次执行日期,防止重复执行 @@ -410,7 +412,7 @@ class AutoScheduler: return # 生成并发送报告 - await self._send_analysis_report(group_id, analysis_result) + await self._send_analysis_report(group_id, analysis_result, platform_id) # 记录执行时间 end_time = asyncio.get_event_loop().time() @@ -494,10 +496,13 @@ class AutoScheduler: return list(all_groups) - async def _send_analysis_report(self, group_id: str, analysis_result: dict): + async def _send_analysis_report( + self, group_id: str, analysis_result: dict, platform_id: str | None = None + ): logger.info( f"[DEBUG][SEND_REPORT] enter " f"group_id={group_id}, " + f"platform_id={platform_id}, " f"analysis_result_keys={list(analysis_result.keys()) if isinstance(analysis_result, dict) else type(analysis_result)}" ) @@ -510,13 +515,17 @@ class AutoScheduler: # 使用图片格式 logger.info(f"群 {group_id} 自动分析使用图片报告格式") try: - image_url = await self.report_generator.generate_image_report( + ( + image_url, + html_content, + ) = await self.report_generator.generate_image_report( analysis_result, group_id, self.html_render_func ) logger.debug( - f"[DEBUG][SEND_REPORT] 图片生成成功" + f"[DEBUG][SEND_REPORT] 图片生成结果 " f"group_id={group_id}, " - f"image_url={image_url}" + f"image_url={'Success' if image_url else 'Fail'}, " + f"html_content={'Available' if html_content else 'None'}" ) if image_url: @@ -538,6 +547,40 @@ class AutoScheduler: await self._send_text_message( group_id, f"📊 每日群聊分析报告:\n\n{text_report}" ) + elif html_content: + # 生成失败但有HTML,加入重试队列 + logger.warning( + f"群 {group_id} 图片报告生成失败,加入重试队列" + ) + + # 尝试获取 platform_id (如果参数为None) + if not platform_id: + platform_id = await self.get_platform_id_for_group( + group_id + ) + + if platform_id: + # 定时任务静默重试,不发送提示消息,只记录日志 + logger.info( + f"群 {group_id} 图片生成失败,已静默加入重试队列" + ) + await self.retry_manager.add_task( + html_content, group_id, platform_id + ) + else: + logger.error( + f"群 {group_id} 无法获取平台ID,无法加入重试队列" + ) + # Fallback to text + text_report = ( + self.report_generator.generate_text_report( + analysis_result + ) + ) + await self._send_text_message( + group_id, f"📊 每日群聊分析报告:\n\n{text_report}" + ) + else: # 图片生成失败(返回None),回退到文本 logger.warning( diff --git a/src/scheduler/retry.py b/src/scheduler/retry.py new file mode 100644 index 0000000..f09305b --- /dev/null +++ b/src/scheduler/retry.py @@ -0,0 +1,238 @@ +import asyncio +import random +import time +from dataclasses import dataclass +from collections.abc import Callable + +from astrbot.api import logger + + +@dataclass +class RetryTask: + """重试任务数据类""" + + html_content: str + group_id: str + platform_id: str # 需要保存 platform_id 以便找回 Bot + retry_count: int = 0 + max_retries: int = 3 + created_at: float = 0.0 + + 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): + self.bot_manager = bot_manager + self.html_render_func = html_render_func + self.queue = asyncio.Queue() + self.running = False + self.worker_task = None + self._dlq = [] # 死信队列 (Failures) + + 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" + ) + + logger.info("[RetryManager] 图片重试管理器已停止") + + async def add_task(self, html_content: str, group_id: str, platform_id: str): + """添加重试任务""" + if not self.running: + logger.warning( + "[RetryManager] 警告:添加任务时管理器未运行,正在尝试启动..." + ) + await self.start() + + task = RetryTask( + html_content=html_content, + group_id=group_id, + platform_id=platform_id, + 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} 重试成功") + 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._notify_failure(task) + + except asyncio.CancelledError: + break + except Exception as e: + logger.error(f"[RetryManager] Worker 异常: {e}", exc_info=True) + await asyncio.sleep(1) + + async def _requeue_after_delay(self, task: RetryTask, delay: float): + 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} 的图片...") + image_url = await self.html_render_func( + task.html_content, + {}, + True, # 返回 URL + image_options, + ) + + if not image_url: + logger.warning( + f"[RetryManager] 重新渲染失败(返回空 URL){task.group_id}" + ) + return False + + # 2. 获取 Bot 实例 + bot = self.bot_manager.get_bot_instance(task.platform_id) + if not bot: + logger.error( + f"[RetryManager] 平台 {task.platform_id} 的 Bot 实例未找到,无法重试" + ) + return False # 无法重试,因为 Bot 已离线 + + # 3. 发送图片 + logger.info(f"[RetryManager] 正在向群 {task.group_id} 发送重试图片...") + + # 使用 OneBot v11 标准 API + if hasattr(bot, "api") and hasattr(bot.api, "call_action"): + try: + # 构造消息 + # 使用 list 格式兼容性更好 + message = [ + { + "type": "text", + "data": {"text": "📊 每日群聊分析报告(重试发送):\n"}, + }, + {"type": "image", "data": {"file": image_url}}, + ] + + result = await bot.api.call_action( + "send_group_msg", group_id=int(task.group_id), message=message + ) + + # 检查 retcode + if isinstance(result, dict): + retcode = result.get("retcode", 0) + if retcode == 0: + return True + elif retcode == 1200: + logger.warning( + f"[RetryManager] 发送失败 (retcode=1200): 可能是Bot被禁言或不在群内,稍后重试" + ) + return False + else: + logger.warning( + f"[RetryManager] 发送失败 (retcode={retcode}): {result}" + ) + return False + return ( + True # 假设非 dict 类型返回即成功(某些适配器可能返回不同类型) + ) + + except Exception as e: + logger.error(f"[RetryManager] 发送API调用异常: {e}") + return False + + elif hasattr(bot, "send_msg"): # 尝试 AstrBot 抽象接口 + try: + # 尝试直接发送 + await bot.send_msg(image_url, group_id=task.group_id) + return True + except Exception as e: + logger.error(f"[RetryManager] 抽象接口发送失败: {e}") + return False + + else: + logger.warning( + f"[RetryManager] 未知的 Bot 类型 {type(bot)},无法发送消息。" + ) + return False + + except Exception as e: + logger.error(f"[RetryManager] 处理任务时发生意外错误: {e}", exc_info=True) + return False + + async def _notify_failure(self, task: RetryTask): + """通知最终失败""" + try: + bot = self.bot_manager.get_bot_instance(task.platform_id) + if bot and hasattr(bot, "api") and hasattr(bot.api, "call_action"): + await bot.api.call_action( + "send_group_msg", + group_id=int(task.group_id), + message=f"[AstrBot QQ群日常分析总结插件] 报告生成/发送多次失败 (Group: {task.group_id}),请检查服务器日志。", + ) + except Exception: + pass