mirror of
https://github.com/Nezumi-2711/astrbot_plugin_qq_group_daily_analysis.git
synced 2026-09-22 20:01:04 +00:00
[v4.11.0] - ✨ 新增 QQ 官方机器人群聊分析与专用 Markdown 报告 (@clown145 #206)
* feat: add QQ official bot group analysis support * fix: use text progress replies for QQ official bot * fix: harden QQ official proactive reporting * feat: enhance QQ official markdown reports * refactor: isolate QQ official markdown reporting * fix: restore platform adapter exports * style: format QQ official changes * fix(main): 添加 html_render 类型声明和 Callable 导入 为 self.html_render 属性添加显式类型声明,解决定时任务中可能因缺少类型注解导致的隐式错误。 * fix(analysis): 消除 execute_daily_analysis 中重复的 bot_self_ids 赋值 去除第二次冗余的 config_manager.get_bot_self_ids() 调用,直接复用已获取的变量。 * cleanup(domain): 移除未使用的分析器适配器和死代码服务文件 删除 golden_quote_analyzer.py、topic_analyzer.py、user_title_analyzer.py(均为未被引用的接口+适配器包装)report_generator.py(旧的文本报告生成器)statistics_calculator.py(与 StatisticsService 功能重复)更新 domain/services/__init__.py 移除对上述文件的引用。 * cleanup(domain): 移除冗余的实体和值对象文件,收敛数据模型至 domain/models/data_models.py 删除 analysis_result.py(与 data_models.py 重复的 SummaryTopic/UserTitle/GoldenQuote 等类定义)value_objects/golden_quote.py、statistics.py、topic.py、user_title.py(均为未被引用的 frozen dataclass 迁移残留)。所有活跃代码已统一导入 domain/models/data_models.py。 * refactor(llm): 抽取 _make_session_id 辅助方法消除重复代码 在 5 个方法(analyze_topics、analyze_user_titles、analyze_golden_quotes、analyze_all_concurrent、analyze_incremental_concurrent)中出现完全相同的 datetime.now().strftime(...) + umo 拼接逻辑,已提取为 _make_session_id 静态方法。 * refactor(message): 正则表达式提升为模块级常量避免重复编译 DISCORD_CUSTOM_EMOJI_PATTERN 和 COMMAND_PATTERN 从类属性移至模块级常量,避免每次实例化 MessageCleanerService 时重新编译正则。 * arch(domain): 定义 IActivityVisualizer 接口并通过依赖注入消除领域层反向依赖 创建 IActivityVisualizer 接口于 domain/repositories/visualization_repository.py,StatisticsService 改由依赖注入接收该接口;ActivityVisualizer 继承接口。消除原 StatisticsService 直接 import infrastructure.visualization 的 DDD 违规。 * fix(main): 补全 Callable 导入和 html_render 类型声明 初次提交(877adb6)因 git add -p 交互式分块时 BOM 导致错误的 hunk 被暂存,Callable 导入和 html_render 类型声明丢失。本次补全这两项修改。 * fix(platform): 补充 QQOfficialAdapter、TelegramAdapter、DiscordAdapter 导出 * refactor(main): 提取内嵌的文本报告生成/发送函数为独立类方法 将 _send_analysis_report 方法中的 generate_text_reports() 和 send_text_reports() 内嵌异步函数提取为 _generate_text_reports 和 _send_text_reports 私有方法。减少闭包复杂度,提升可维护性。 * cleanup(config): 移除废弃的 get_qq_official_t2i_activity_histogram_enabled 向后兼容方法 删除 config_manager.py 中的旧名别名方法,简化 qq_official_markdown.py 中对应的 getattr 回退逻辑为直接方法调用,移除测试中专门验证旧名兼容性的 LegacyDisabledConfig 测试用例。该兼容层仅在迁移期间临时存在,现已完成过渡。 * fix(types): 修复 Pylance 类型告警 — platform_key 未绑定、int(object) 和 template.filename 可能为 None platform_group_registry.py: 将 platform_key 定义提前,消除 Pylance reportPossiblyUnboundVariable 告警。 qq_official_markdown.py: 将 int(value or 0) 改为 int(value) if value is not None else 0,消除 reportArgumentType 告警。 templates.py: 在使用 template.filename 前增加 None 检查,消除 str|None 不可分配给 str 的告警。 qq_official_adapter.py: 为 post_group_message 添加类型忽略注解,消除 await 不可等待对象的告警。 * fix(types): 修复残留的 Pylance 类型告警 platform_group_registry.py: 将 (platform_key, group_id) 改为 (str(platform_key), group_id),消除 str|None 不可分配给 str 的 reportArgumentType qq_official_markdown.py: int(value) 添加 # type: ignore[arg-type],消除 object 不可分配给 ConvertibleToInt qq_official_adapter.py: 移除无意义的中间变量,await 行直接添加 # type: ignore[arg-type] * docs(message): 更新 MessageProcessingService 的文档和注释,消除 Telegram 特殊性表述 类 docstring:移除 '维护 Telegram 群组注册表' 等过时描述,补充 QQ 官方去重职责。 group_registry.upsert 注释:改为泛化的跨平台描述,不再限定 Telegram。 _extract_event_timestamp / _reserve_event_id / _commit_event_id / _release_event_id:英文 docstring 统一为中文。 * docs(message): 修正类注释中只提 QQ 官方的问题,明示 Telegram 也由本服务处理 Telegram 和 QQ 官方消息都经过 MessageProcessingService.process_message()。修正前类 docstring 只提了 QQ 官方的事件去重,缺少 Telegram 作为主要调用者的说明。 * fix(types): 修正 _sanitize_analysis_result_for_export 返回类型注解 函数声明 -> dict 但 _sanitize_export_identity_text 可返回 str|dict|list,Pylance 报 reportReturnType。 改为 -> dict[str, Any] 准确表达实际返回类型,补充缺失的 from typing import Any。 * fix(types): 抑制 Pylance reportReturnType 误报 _sanitize_analysis_result_for_export 的 analysis_result 参数运行时始终为 dict,但 _to_plain_export_data 递归返回 Any 导致 Pylance 推导出 str|dict|list 联合类型。 添加 # type: ignore[return-type] 抑制此误报。 * fix(report): 渲染匿名模式下的未知引用 token 不再静默丢弃 当 hide_user_names=True 且 [id] 不在 known_ids 中时,之前返回空 Markup 导致 token 被静默移除,可能扭曲文本语义。改为返回转义后的原始 [id] 字符串,保留文本布局和语义,同时不泄露身份信息。 参考 Sourcery AI code review 建议。 * fix(import): 避免从 astrbot.core 直接导入 File,优先使用公开 API astrbot.core 是内部模块,不在公开 API 契约中,可能随版本变化。 改为 try 优先导入 astrbot.api.message_components.File(公开 API), 失败时回退到 astrbot.core.message.components.File(向后兼容)。 同时修复 main.py 和 qq_official_adapter.py 两处导入。 * fix(import): 回退 try/except 伪装,改为直接导入 + 风险注释 astrbot.api.message_components 不存在,try/except 永远走 except 分支,是无效代码。 改为直接 from astrbot.core.message.components import File,用注释说明这是内部 API 可能变化。 * fix(ruff): 代码质量 * docs(README): 调整文档说明,删除 lark 相关的描述和功能标识 * docs(desc): 更新 desc * fix(message): 将 TG 和 QQ 官方消息缓存成功日志降为 debug * perf(message): 将 _extract_event_timestamp 延迟到 QQ 官方分支内计算 该时间戳仅用于 QQ 官方消息的 history_content 元数据,但对所有平台都执行了深度 getattr 链。 改为只在 QQ 官方分支内延迟计算,消除 Telegram 等平台上每次消息的白算开销。 * chore(CHANGELOG) --------- Co-authored-by: SXP-Simon <sxp20061207@163.com>
This commit is contained in:
@@ -216,7 +216,6 @@ class AnalysisApplicationService:
|
||||
)
|
||||
|
||||
# 4. 用户分析 (Domain Service)
|
||||
bot_self_ids = self.config_manager.get_bot_self_ids()
|
||||
user_activity = await asyncio.to_thread(
|
||||
self.analysis_domain_service.analyze_user_activity,
|
||||
unified_messages,
|
||||
|
||||
@@ -1,10 +1,10 @@
|
||||
import re
|
||||
from collections import Counter
|
||||
from collections import Counter, OrderedDict
|
||||
|
||||
from astrbot.api.event import AstrMessageEvent
|
||||
from astrbot.api.star import Context
|
||||
|
||||
from ...infrastructure.persistence.telegram_group_registry import TelegramGroupRegistry
|
||||
from ...infrastructure.persistence.platform_group_registry import PlatformGroupRegistry
|
||||
from ...utils.logger import logger
|
||||
|
||||
|
||||
@@ -12,27 +12,36 @@ class MessageProcessingService:
|
||||
"""
|
||||
消息处理服务
|
||||
|
||||
负责处理接收到的消息事件:
|
||||
解析收到的群消息事件,提取内容与发送者信息,持久化历史记录,
|
||||
并维护事件驱动平台(Telegram、QQ 官方等)的群组注册表。
|
||||
QQ 官方平台特有的重复消息去重逻辑也在本服务中处理。
|
||||
|
||||
职责:
|
||||
1. 解析消息内容(文本、图片、@提及等)
|
||||
2. 解析发送者信息(跨平台兼容)
|
||||
2. 解析发送者展示名(跨平台兼容)
|
||||
3. 存储消息历史
|
||||
4. 维护 Telegram 群组注册表(回退机制)
|
||||
4. 维护群组注册表,供调度器做群组发现(Telegram、QQ 官方等事件驱动平台)
|
||||
5. QQ 官方事件消息去重(按 message_id 预占 + 确认机制)
|
||||
"""
|
||||
|
||||
def __init__(self, context: Context, telegram_registry: TelegramGroupRegistry):
|
||||
def __init__(self, context: Context, group_registry: PlatformGroupRegistry):
|
||||
self.context = context
|
||||
self.telegram_registry = telegram_registry
|
||||
self.group_registry = group_registry
|
||||
self._seen_event_ids: OrderedDict[str, None] = OrderedDict()
|
||||
self._inflight_event_ids: set[str] = set()
|
||||
self._seen_event_ids_limit = 4096
|
||||
|
||||
async def process_message(self, event: AstrMessageEvent) -> None:
|
||||
"""
|
||||
处理并在历史记录中存储消息。
|
||||
被 main.py 的 Telegram 和 QQ 官方消息拦截器共同调用。
|
||||
|
||||
Args:
|
||||
event: AstrBot 消息事件
|
||||
Args:
|
||||
event: AstrBot 消息事件
|
||||
|
||||
Raises:
|
||||
ValueError: 当必要数据无法获取时
|
||||
RuntimeError: 当消息内容为空时
|
||||
Raises:
|
||||
ValueError: 当必要数据无法获取时
|
||||
RuntimeError: 当消息内容为空时
|
||||
"""
|
||||
# 1. 获取群组 ID(必需)
|
||||
group_id = self._get_group_id_from_event(event)
|
||||
@@ -62,37 +71,63 @@ class MessageProcessingService:
|
||||
f"群 {group_id}: 消息内容为空 (sender={sender_name}),拒绝存储"
|
||||
)
|
||||
|
||||
# 6. 提取事件消息 ID(用于 Telegram 已见群/话题记录)
|
||||
# 6. 提取事件消息 ID 和事件时间
|
||||
msg_obj = getattr(event, "message_obj", None)
|
||||
event_message_id = str(getattr(msg_obj, "message_id", "") or "")
|
||||
|
||||
platform_name = str(event.get_platform_name() or "").strip().lower()
|
||||
reserved_event_id = False
|
||||
if platform_name in {"qq_official", "qq_official_webhook"} and event_message_id:
|
||||
reserved_event_id = self._reserve_event_id(event_message_id)
|
||||
if not reserved_event_id:
|
||||
logger.debug("[QQOfficial] 跳过重复消息事件: %s", event_message_id)
|
||||
return
|
||||
history_content = {
|
||||
"type": "user",
|
||||
"message": message_parts,
|
||||
}
|
||||
if platform_name in {"qq_official", "qq_official_webhook"}:
|
||||
event_timestamp = self._extract_event_timestamp(msg_obj)
|
||||
history_content["_qq_official"] = {
|
||||
"message_id": event_message_id,
|
||||
"timestamp": event_timestamp,
|
||||
}
|
||||
|
||||
# 7. 存储到数据库
|
||||
await self.context.message_history_manager.insert(
|
||||
platform_id=platform_id,
|
||||
user_id=group_id,
|
||||
content={"type": "user", "message": message_parts},
|
||||
sender_id=sender_id,
|
||||
sender_name=sender_name,
|
||||
)
|
||||
try:
|
||||
await self.context.message_history_manager.insert(
|
||||
platform_id=platform_id,
|
||||
user_id=group_id,
|
||||
content=history_content,
|
||||
sender_id=sender_id,
|
||||
sender_name=sender_name,
|
||||
)
|
||||
except BaseException:
|
||||
if reserved_event_id:
|
||||
self._release_event_id(event_message_id)
|
||||
raise
|
||||
else:
|
||||
if reserved_event_id:
|
||||
self._commit_event_id(event_message_id)
|
||||
|
||||
# Telegram: 记录已见群/话题
|
||||
if self._is_telegram_event(event, platform_id):
|
||||
try:
|
||||
await self.telegram_registry.upsert(
|
||||
platform_id=platform_id,
|
||||
group_id=group_id,
|
||||
sender_id=sender_id,
|
||||
sender_name=sender_name,
|
||||
event_message_id=event_message_id,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.warning(
|
||||
"[TGRegistry] Upsert failed: "
|
||||
f"platform_id={platform_id} group_id={group_id} error={e}"
|
||||
)
|
||||
# Register the group so the scheduler can discover platforms that
|
||||
# do not provide a group-list API (Telegram, QQ Official, etc.).
|
||||
try:
|
||||
await self.group_registry.upsert(
|
||||
platform_id=platform_id,
|
||||
group_id=group_id,
|
||||
sender_id=sender_id,
|
||||
sender_name=sender_name,
|
||||
event_message_id=event_message_id,
|
||||
)
|
||||
except Exception as e:
|
||||
logger.warning(
|
||||
"[GroupRegistry] Upsert failed: "
|
||||
f"platform_id={platform_id} group_id={group_id} error={e}"
|
||||
)
|
||||
|
||||
logger.info(
|
||||
f"[Telegram] [{platform_id}] 已缓存群 {group_id} 的消息 (发送者: {sender_name})"
|
||||
logger.debug(
|
||||
f"[{platform_id}] 已缓存群 {group_id} 的消息 (发送者: {sender_name})"
|
||||
)
|
||||
|
||||
def _get_group_id_from_event(self, event: AstrMessageEvent) -> str | None:
|
||||
@@ -208,6 +243,24 @@ class MessageProcessingService:
|
||||
}
|
||||
)
|
||||
|
||||
elif seg_type in ("File", "file"):
|
||||
url = getattr(seg, "url", None) or getattr(seg, "file_", None)
|
||||
message_parts.append(
|
||||
{
|
||||
"type": "file",
|
||||
"url": str(url or ""),
|
||||
"name": str(getattr(seg, "name", "") or ""),
|
||||
}
|
||||
)
|
||||
|
||||
elif seg_type in ("Record", "record", "voice"):
|
||||
url = getattr(seg, "url", None) or getattr(seg, "file", None)
|
||||
message_parts.append({"type": "voice", "url": str(url or "")})
|
||||
|
||||
elif seg_type in ("Video", "video"):
|
||||
url = getattr(seg, "url", None) or getattr(seg, "file", None)
|
||||
message_parts.append({"type": "video", "url": str(url or "")})
|
||||
|
||||
if not message_parts and event.message_str:
|
||||
message_parts.append({"type": "plain", "text": event.message_str})
|
||||
|
||||
@@ -261,9 +314,57 @@ class MessageProcessingService:
|
||||
return normalized == str(sender_id).strip()
|
||||
|
||||
@staticmethod
|
||||
def _is_telegram_event(event: AstrMessageEvent, platform_id: str) -> bool:
|
||||
"""判断当前事件是否为 Telegram 平台"""
|
||||
platform_name = str(event.get_platform_name() or "").strip().lower()
|
||||
if platform_name == "telegram":
|
||||
return True
|
||||
return str(platform_id or "").strip().lower().startswith("telegram")
|
||||
def _extract_event_timestamp(message_obj: object) -> int:
|
||||
"""从消息对象中提取平台事件时间戳。"""
|
||||
raw_message = getattr(message_obj, "raw_message", None)
|
||||
if isinstance(raw_message, dict):
|
||||
candidate = raw_message.get("timestamp")
|
||||
if not candidate:
|
||||
raw_data = raw_message.get("raw_data")
|
||||
if isinstance(raw_data, dict):
|
||||
candidate = raw_data.get("timestamp")
|
||||
else:
|
||||
raw_data = getattr(raw_message, "raw_data", None)
|
||||
candidate = getattr(raw_message, "timestamp", None)
|
||||
if not candidate and isinstance(raw_data, dict):
|
||||
candidate = raw_data.get("timestamp")
|
||||
if isinstance(candidate, (int, float)):
|
||||
return int(candidate)
|
||||
if candidate:
|
||||
try:
|
||||
from datetime import datetime
|
||||
|
||||
return int(
|
||||
datetime.fromisoformat(
|
||||
str(candidate).replace("Z", "+00:00")
|
||||
).timestamp()
|
||||
)
|
||||
except (TypeError, ValueError, OverflowError):
|
||||
pass
|
||||
return 0
|
||||
|
||||
def _reserve_event_id(self, event_message_id: str) -> bool:
|
||||
"""预占事件消息 ID:在历史记录持久化期间防止重复入库。"""
|
||||
if (
|
||||
event_message_id in self._inflight_event_ids
|
||||
or event_message_id in self._seen_event_ids
|
||||
):
|
||||
if event_message_id in self._seen_event_ids:
|
||||
self._seen_event_ids.move_to_end(event_message_id)
|
||||
return False
|
||||
self._inflight_event_ids.add(event_message_id)
|
||||
return True
|
||||
|
||||
def _commit_event_id(self, event_message_id: str) -> None:
|
||||
"""确认事件消息 ID:标记为已持久化,纳入后续去重。"""
|
||||
self._inflight_event_ids.discard(event_message_id)
|
||||
if event_message_id in self._seen_event_ids:
|
||||
self._seen_event_ids.move_to_end(event_message_id)
|
||||
else:
|
||||
self._seen_event_ids[event_message_id] = None
|
||||
if len(self._seen_event_ids) > self._seen_event_ids_limit:
|
||||
self._seen_event_ids.popitem(last=False)
|
||||
|
||||
def _release_event_id(self, event_message_id: str) -> None:
|
||||
"""释放事件消息 ID:持久化失败或取消时清理预占状态。"""
|
||||
self._inflight_event_ids.discard(event_message_id)
|
||||
|
||||
Reference in New Issue
Block a user