From aa5752d179fafa00897767ea2bd603ddb91f4104 Mon Sep 17 00:00:00 2001 From: clown145 Date: Thu, 12 Feb 2026 02:13:47 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20TG=E9=80=9A=E8=BF=87KV=E5=9B=9E=E9=80=80?= =?UTF-8?q?=E7=BE=A4=E5=88=97=E8=A1=A8=E5=B9=B6=E5=A2=9E=E5=8A=A0=E4=B8=B4?= =?UTF-8?q?=E6=97=B6=E8=B0=83=E8=AF=95=E6=97=A5=E5=BF=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- main.py | 132 ++++++++++++++++++ .../scheduler/auto_scheduler.py | 71 +++++++++- 2 files changed, 199 insertions(+), 4 deletions(-) diff --git a/main.py b/main.py index c81abdf..2544735 100644 --- a/main.py +++ b/main.py @@ -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 是否为空或占位值。""" diff --git a/src/infrastructure/scheduler/auto_scheduler.py b/src/infrastructure/scheduler/auto_scheduler.py index a6f0455..3ae9103 100644 --- a/src/infrastructure/scheduler/auto_scheduler.py +++ b/src/infrastructure/scheduler/auto_scheduler.py @@ -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 []