fix(retry): 增强重试管理器,尝试避免重试地狱和重复发送

This commit is contained in:
SXP-Simon
2026-02-23 00:34:55 +08:00
parent a6e539b9f6
commit 544458aadd
5 changed files with 197 additions and 85 deletions
+10 -9
View File
@@ -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(
@@ -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),
+5 -1
View File
@@ -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:
+26 -40
View File
@@ -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"""
# 为每个群聊使用独立的锁
+86 -35
View File
@@ -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: