mirror of
https://github.com/Nezumi-2711/astrbot_plugin_qq_group_daily_analysis.git
synced 2026-09-22 13:38:43 +00:00
fix: 移除发送重试模块
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 实现 ====================
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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
|
||||
|
||||
# ================================================================
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
Reference in New Issue
Block a user