fix: TG通过KV回退群列表并增加临时调试日志

This commit is contained in:
clown145
2026-02-12 15:41:07 +08:00
committed by Helian Nuits
parent 926e0df095
commit aa5752d179
2 changed files with 199 additions and 4 deletions
+132
View File
@@ -9,6 +9,7 @@ import asyncio
import os
import re
from collections import Counter
from datetime import datetime, timezone
from astrbot.api import AstrBotConfig, logger
from astrbot.api.event import AstrMessageEvent, filter
@@ -36,6 +37,8 @@ from .src.utils.pdf_utils import PDFInstaller
class QQGroupDailyAnalysis(Star):
"""QQ群日常分析插件主类"""
_TG_GROUP_REGISTRY_KV_KEY = "telegram_seen_groups_v1"
def __init__(self, context: Context, config: AstrBotConfig):
super().__init__(context)
self.config = config
@@ -82,6 +85,7 @@ class QQGroupDailyAnalysis(Star):
self.retry_manager,
self.report_generator,
self.html_render,
plugin_instance=self,
)
self._initialized = False
@@ -271,10 +275,138 @@ class QQGroupDailyAnalysis(Star):
f"parts_count={len(message_parts)}"
)
# Telegram: 记录已见群/话题,用于自动分析拉群回退
if self._is_telegram_event(event, platform_id):
try:
await self._upsert_telegram_group_registry(
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(
"[TEMP][TGRegistry][UpsertFailed] "
f"platform_id={platform_id} group_id={group_id} error={e}"
)
logger.debug(
f"[{platform_id}] 已缓存群 {group_id} 的消息 (发送者: {sender_name})"
)
@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")
async def _upsert_telegram_group_registry(
self,
platform_id: str,
group_id: str,
sender_id: str,
sender_name: str,
event_message_id: str,
) -> None:
"""更新 Telegram 已见群/话题注册表(KV)。"""
registry = await self.get_kv_data(self._TG_GROUP_REGISTRY_KV_KEY, {})
if not isinstance(registry, dict):
registry = {}
platforms = registry.get("platforms")
if not isinstance(platforms, dict):
platforms = {}
registry["platforms"] = platforms
platform_key = str(platform_id).strip()
group_key = str(group_id).strip()
platform_map = platforms.get(platform_key)
if not isinstance(platform_map, dict):
platform_map = {}
platforms[platform_key] = platform_map
now_iso = datetime.now(timezone.utc).isoformat()
existed = group_key in platform_map and isinstance(
platform_map[group_key], dict
)
entry = platform_map.get(group_key)
if not isinstance(entry, dict):
entry = {}
first_seen = entry.get("first_seen")
if not isinstance(first_seen, str) or not first_seen:
first_seen = now_iso
entry.update(
{
"first_seen": first_seen,
"last_seen": now_iso,
"last_sender_id": str(sender_id),
"last_sender_name": str(sender_name),
"last_event_message_id": str(event_message_id),
}
)
platform_map[group_key] = entry
registry["updated_at"] = now_iso
await self.put_kv_data(self._TG_GROUP_REGISTRY_KV_KEY, registry)
platform_targets = len(platform_map)
total_targets = sum(
len(groups) for groups in platforms.values() if isinstance(groups, dict)
)
logger.info(
"[TEMP][TGRegistry][Upsert] "
f"platform_id={platform_key} group_id={group_key} existed={existed} "
f"platform_targets={platform_targets} total_targets={total_targets} "
f"sender_id={sender_id} sender_name={sender_name}"
)
async def get_telegram_seen_group_ids(
self, platform_id: str | None = None
) -> list[str]:
"""读取 Telegram 已见群/话题列表(给调度器回退使用)。"""
registry = await self.get_kv_data(self._TG_GROUP_REGISTRY_KV_KEY, {})
if not isinstance(registry, dict):
logger.info(
"[TEMP][TGRegistry][Read] invalid_registry_type, fallback_empty"
)
return []
platforms = registry.get("platforms")
if not isinstance(platforms, dict):
logger.info("[TEMP][TGRegistry][Read] no_platforms, fallback_empty")
return []
groups: set[str] = set()
if platform_id:
platform_map = platforms.get(str(platform_id).strip(), {})
if isinstance(platform_map, dict):
groups.update(
str(gid).strip() for gid in platform_map.keys() if str(gid).strip()
)
else:
for platform_map in platforms.values():
if not isinstance(platform_map, dict):
continue
groups.update(
str(gid).strip() for gid in platform_map.keys() if str(gid).strip()
)
sorted_groups = sorted(groups)
preview = sorted_groups[:10]
logger.info(
"[TEMP][TGRegistry][Read] "
f"platform_id={platform_id or '*'} count={len(sorted_groups)} "
f"groups_preview={preview}"
)
return sorted_groups
@staticmethod
def _is_placeholder_sender_name(name: str | None, sender_id: str) -> bool:
"""判断 sender_name 是否为空或占位值。"""
+67 -4
View File
@@ -6,6 +6,7 @@
import asyncio
import time as time_mod
import weakref
from typing import Any
from apscheduler.triggers.cron import CronTrigger
@@ -27,6 +28,7 @@ class AutoScheduler:
retry_manager,
report_generator=None,
html_render_func=None,
plugin_instance: Any | None = None,
):
self.config_manager = config_manager
self.analysis_service = analysis_service
@@ -34,6 +36,7 @@ class AutoScheduler:
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)
@@ -831,19 +834,44 @@ class AutoScheduler:
if adapter:
try:
groups = await adapter.get_group_list()
for group_id in groups:
all_groups.add((platform_id, str(group_id)))
groups = [
str(group_id).strip()
for group_id in groups
if str(group_id).strip()
]
# 获取可读的平台名称
p_name = getattr(adapter, "platform_name", None)
# 获取平台名称(用于 Telegram 回退判定)
p_name = None
if hasattr(adapter, "get_platform_name"):
try:
p_name = adapter.get_platform_name()
except Exception:
p_name = None
if not p_name:
p_name = (
self.bot_manager._detect_platform_name(bot_instance)
or "unknown"
)
p_name = str(p_name).strip()
used_tg_kv_fallback = False
if not groups and p_name.lower() == "telegram":
groups = await self._get_telegram_groups_from_plugin_kv(
str(platform_id)
)
used_tg_kv_fallback = bool(groups)
logger.info(
"[TEMP][TGRegistry][SchedulerFallback] "
f"platform_id={platform_id} fallback_used={used_tg_kv_fallback} "
f"fallback_count={len(groups)}"
)
for group_id in groups:
all_groups.add((platform_id, str(group_id)))
logger.info(
f"平台 {platform_id} ({p_name}) 成功获取 {len(groups)} 个群组"
+ (" (KV回退)" if used_tg_kv_fallback else "")
)
continue
@@ -857,3 +885,38 @@ class AutoScheduler:
logger.error(f"平台 {platform_id} 获取群列表异常: {e}")
return list(all_groups)
async def _get_telegram_groups_from_plugin_kv(self, platform_id: str) -> list[str]:
"""从插件 KV 获取 Telegram 已见群/话题,作为 get_group_list 的回退。"""
if not self.plugin_instance:
logger.info(
"[TEMP][TGRegistry][SchedulerFetch] "
f"platform_id={platform_id} skipped=no_plugin_instance"
)
return []
getter = getattr(self.plugin_instance, "get_telegram_seen_group_ids", None)
if not callable(getter):
logger.info(
"[TEMP][TGRegistry][SchedulerFetch] "
f"platform_id={platform_id} skipped=no_getter"
)
return []
try:
groups = await getter(platform_id=platform_id)
normalized = sorted(
{str(group_id).strip() for group_id in groups if str(group_id).strip()}
)
logger.info(
"[TEMP][TGRegistry][SchedulerFetch] "
f"platform_id={platform_id} count={len(normalized)} "
f"groups_preview={normalized[:10]}"
)
return normalized
except Exception as e:
logger.warning(
"[TEMP][TGRegistry][SchedulerFetchFailed] "
f"platform_id={platform_id} error={e}"
)
return []