diff --git a/app/modules/telegram/telegram.py b/app/modules/telegram/telegram.py index 5d592398a..a5bcc7e98 100644 --- a/app/modules/telegram/telegram.py +++ b/app/modules/telegram/telegram.py @@ -77,18 +77,13 @@ class Telegram: _bot: TeleBot = None _callback_handlers: Dict[str, Callable] = {} # 存储回调处理器 - _user_chat_mapping: Dict[ - str, str - ] = {} # userid -> chat_id mapping for reply targeting _bot_username: Optional[str] = None # Bot username for mention detection - _typing_tasks: Dict[str, threading.Thread] = {} # chat_id -> typing任务 - _typing_stop_flags: Dict[str, threading.Event] = {} # chat_id -> 停止信号 - _typing_lock = threading.RLock() _typing_interval_seconds = 5 _typing_initial_delay_seconds = 1 _typing_max_duration_seconds = 10 * 60 _typing_command_max_duration_seconds = 30 _typing_callback_max_duration_seconds = 60 + _typing_join_timeout_seconds = 1 def __init__( self, @@ -103,6 +98,13 @@ class Telegram: self._telegram_token = TELEGRAM_TOKEN self._telegram_chat_id = TELEGRAM_CHAT_ID self._polling_thread = None + # 一个 Telegram 配置对应一个 SDK client,运行状态不能被其他配置共享。 + self._user_chat_mapping: Dict[str, str] = {} + self._typing_tasks: Dict[str, threading.Thread] = {} + self._typing_stop_flags: Dict[str, threading.Event] = {} + self._typing_lock = threading.RLock() + self._typing_lifecycle_lock = threading.RLock() + self._typing_accepting = True if not TELEGRAM_TOKEN or not TELEGRAM_CHAT_ID: logger.error("Telegram配置不完整!") return @@ -482,70 +484,101 @@ class Telegram: chat_id: Union[str, int], max_duration_seconds: Optional[float] = None, initial_delay_seconds: Optional[float] = None, - ) -> None: + ) -> bool: """ - 启动持续发送正在输入状态的任务 + 启动持续发送正在输入状态的任务。 + + :return: 是否取得该会话的唯一 typing owner """ chat_id_str = str(chat_id) - # 如果已有任务在运行,先停止 - self._stop_typing_task(chat_id_str) + with self._typing_lifecycle_lock: + if not self._typing_accepting: + logger.debug("Telegram client已停止,拒绝启动typing任务") + return False + # 如果已有任务在运行,先停止;阻塞的旧 SDK 请求不能被新 owner 覆盖。 + if not self._stop_typing_task(chat_id_str): + logger.warning( + "Telegram typing旧任务尚未结束,拒绝并行启动: chat_id=%s", + chat_id_str, + ) + return False - # 使用独立 Event 避免同一 chat 新旧 typing 线程互相误改停止标记。 - stop_event = threading.Event() - max_duration = max_duration_seconds or self._typing_max_duration_seconds - initial_delay = ( - self._typing_initial_delay_seconds - if initial_delay_seconds is None - else max(initial_delay_seconds, 0) - ) + # 使用独立 Event 避免同一 chat 新旧 typing 线程互相误改停止标记。 + stop_event = threading.Event() + max_duration = max_duration_seconds or self._typing_max_duration_seconds + initial_delay = ( + self._typing_initial_delay_seconds + if initial_delay_seconds is None + else max(initial_delay_seconds, 0) + ) - def typing_worker(): - """延迟首发并定期发送 typing 状态的后台线程。""" - started_at = time.monotonic() - try: - # Telegram 没有撤销 typing 的接口,短响应先等待一小段时间, - # 避免回复已经发出后客户端仍残留几秒“正在输入”。 - if initial_delay and stop_event.wait(initial_delay): - return - while not stop_event.is_set(): - if time.monotonic() - started_at >= max_duration: - logger.warning( - "Telegram typing状态超过最大续期,自动停止: chat_id=%s", - chat_id_str, - ) - break - try: - if self._bot: - self._bot.send_chat_action(chat_id, "typing") - except Exception as e: - logger.debug(f"发送typing状态失败: {e}") - # Telegram 客户端约 5-6 秒后会隐藏 typing,需要周期性续发。 - stop_event.wait(self._typing_interval_seconds) - finally: + def typing_worker(): + """延迟首发并定期发送 typing 状态的后台线程。""" + started_at = time.monotonic() + try: + # Telegram 没有撤销 typing 的接口,短响应先等待一小段时间, + # 避免回复已经发出后客户端仍残留几秒“正在输入”。 + if initial_delay and stop_event.wait(initial_delay): + return + while not stop_event.is_set(): + if time.monotonic() - started_at >= max_duration: + logger.warning( + "Telegram typing状态超过最大续期,自动停止: chat_id=%s", + chat_id_str, + ) + break + try: + if self._bot: + self._bot.send_chat_action(chat_id, "typing") + except Exception as e: + logger.debug(f"发送typing状态失败: {e}") + # Telegram 客户端约 5-6 秒后会隐藏 typing,需要周期性续发。 + stop_event.wait(self._typing_interval_seconds) + finally: + with self._typing_lock: + current = self._typing_tasks.get(chat_id_str) + if current is threading.current_thread(): + self._typing_tasks.pop(chat_id_str, None) + self._typing_stop_flags.pop(chat_id_str, None) + + thread = threading.Thread( + target=typing_worker, + name=f"MoviePilot-TelegramTyping-{chat_id_str}"[:120], + daemon=True, + ) + with self._typing_lock: + self._typing_stop_flags[chat_id_str] = stop_event + self._typing_tasks[chat_id_str] = thread + try: + thread.start() + except BaseException: + self._typing_stop_flags.pop(chat_id_str, None) + self._typing_tasks.pop(chat_id_str, None) + raise + return True + + def _stop_typing_task(self, chat_id: Union[str, int]) -> bool: + """ + 停止正在输入状态的任务,并保留尚未结束的 owner。 + + :return: 任务是否已经进入终态 + """ + chat_id_str = str(chat_id) + with self._typing_lifecycle_lock: + with self._typing_lock: + stop_event = self._typing_stop_flags.get(chat_id_str) + task = self._typing_tasks.get(chat_id_str) + if stop_event: + stop_event.set() + if task and task.is_alive() and task is not threading.current_thread(): + task.join(timeout=self._typing_join_timeout_seconds) + task_finished = task is None or not task.is_alive() + if task_finished: with self._typing_lock: - current = self._typing_tasks.get(chat_id_str) - if current is threading.current_thread(): + if self._typing_tasks.get(chat_id_str) is task: self._typing_tasks.pop(chat_id_str, None) self._typing_stop_flags.pop(chat_id_str, None) - - thread = threading.Thread(target=typing_worker, daemon=True) - with self._typing_lock: - self._typing_stop_flags[chat_id_str] = stop_event - self._typing_tasks[chat_id_str] = thread - thread.start() - - def _stop_typing_task(self, chat_id: Union[str, int]) -> None: - """ - 停止正在输入状态的任务 - """ - chat_id_str = str(chat_id) - with self._typing_lock: - stop_event = self._typing_stop_flags.pop(chat_id_str, None) - task = self._typing_tasks.pop(chat_id_str, None) - if stop_event: - stop_event.set() - if task and task.is_alive() and task is not threading.current_thread(): - task.join(timeout=1) + return task_finished def _stop_typing_if_needed( self, chat_id: Union[str, int], stop_typing: bool @@ -574,8 +607,7 @@ class Telegram: target_chat_id = target_chat_id or (str(userid) if userid else None) if not target_chat_id: return False - self._start_typing_task(target_chat_id) - return True + return self._start_typing_task(target_chat_id) def stop_typing( self, @@ -1714,9 +1746,11 @@ class Telegram: """ 停止Telegram消息接收服务 """ - # 停止所有typing任务 - for chat_id in list(self._typing_tasks.keys()): - self._stop_typing_task(chat_id) + with self._typing_lifecycle_lock: + self._typing_accepting = False + # 封口与 owner 快照处于同一临界区,停止后不会漏掉并发新增任务。 + for chat_id in list(self._typing_tasks.keys()): + self._stop_typing_task(chat_id) if not self._bot: return diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index ace6dcc7f..4f4d4ba11 100644 --- a/docs/refactor/backend-architecture-next-stage.md +++ b/docs/refactor/backend-architecture-next-stage.md @@ -6,7 +6,7 @@ > 审计范围:宿主后端;排除 `app/plugins/**` 运行时插件副本 > 规范优先级:`AGENTS.md` 与 `docs/rules/` 高于本文 > 相关文档:`docs/architecture-overview.md`、`docs/refactor/backend-architecture-governance.md`、`docs/refactor/backend-module-refactor-compatibility.md` -> 实施进度:阶段 0~6 的宿主架构能力已完成收口;API/Application 公共复杂度基线已清零,启动组合根的 SystemConfigOper 构造点已由 14 降至 1;API 进程内后台任务已完成首批统一登记,插件仓适配和 Outbox 外围扩展仍按风险切片推进。Model/Base 查询与写装饰器、legacy 隐式会话外壳均已清零,插件 SDK 也不再导出宿主 Model。2026-08-23 的长期整改阶段 0 已恢复宿主、启动性能、官方插件和 SDK 契约门禁的可信基线;阶段 1a 已补齐 TaskRegistry owner 零债务门禁和诚实的关停超时语义;阶段 1b1 已收口整理 worker、pending 回放、失败通知、进程内 AI 重试、插件监控与事件投递的生命周期所有权;2026-08-24 的阶段 2 已将 212 个已观察宿主模块方法的 legacy aggregation 清零,并补齐可执行 fanout 与下载器文件 DTO 边界;阶段 3 已将消息交互和远程命令的订阅删除统一到 Application/UoW/outbox,宿主不再调用裸线程统计入口;阶段 4 已统一七种消息渠道的宿主回环与后台执行边界;阶段 5 已补齐事件窗口聚合任务的生命周期所有权;阶段 6 已统一插件文件操作的取消完成语义;阶段 7 已统一插件协程补偿的终态等待;阶段 8 已统一宿主同步函数的异步线程池入口;阶段 9 已统一工作流运行时的宿主获取路径;阶段 10 已统一模块、插件与调度运行时的显式 getter 调用;阶段 11 已清除系统配置 getter 的 Oper 形别名;阶段 12 已完成工作流域的显式 Chain 数据端口迁移;阶段 13 已收口用户、交互与消息链的数据端口;阶段 14 已收口音乐订阅数据端口;阶段 15 已收口站点数据端口;阶段 16 已收口媒体服务器数据端口;阶段 17 已收口下载数据端口;阶段 18 已收口主订阅数据端口;阶段 19 已收口整理数据端口;阶段 20 已收口 Agent 数据端口;阶段 21 已收口监控历史端口;阶段 22 已统一服务配置应用边界;阶段 23 已补齐媒体服务器 API 遗留的类形配置读取路径;阶段 24 已清除 Scheduler 内部无 owner 的协程提交双轨;阶段 25 已补齐 TaskRegistry 跨线程 owner 并迁移整理 AI 接管;阶段 26 已统一 Agent 会话清理提交;阶段 27 已统一历史 AI 进度 owner;阶段 28 已托管旧插件订阅统计线程;阶段 29 已统一 Emby 系条目转换并清零重复代码白名单;阶段 30 已收口插件市场请求级子任务;阶段 31 已托管搜索 AI 推荐任务;阶段 32 已清除事件调度器绕过生命周期 owner 的投递回退;阶段 33 已统一宿主 Agent 运行时的获取路径;阶段 34 已统一 durable-required 事件与 Outbox topic 事实源;阶段 35 已统一 LLM provider 管理 API 的运行时解析路径;阶段 36 已统一 WebAgent 音频能力访问边界;阶段 37 已统一插件输入事件发布路径;阶段 38 已统一 WebAgent 通知事件监听与队列边界;阶段 39 已补齐搜索 SSE 断线时的上游任务清理;阶段 40 已补齐异步防抖取消的终态所有权;阶段 41 已统一优雅重启兜底线程的唯一所有权。 +> 实施进度:阶段 0~6 的宿主架构能力已完成收口;API/Application 公共复杂度基线已清零,启动组合根的 SystemConfigOper 构造点已由 14 降至 1;API 进程内后台任务已完成首批统一登记,插件仓适配和 Outbox 外围扩展仍按风险切片推进。Model/Base 查询与写装饰器、legacy 隐式会话外壳均已清零,插件 SDK 也不再导出宿主 Model。2026-08-23 的长期整改阶段 0 已恢复宿主、启动性能、官方插件和 SDK 契约门禁的可信基线;阶段 1a 已补齐 TaskRegistry owner 零债务门禁和诚实的关停超时语义;阶段 1b1 已收口整理 worker、pending 回放、失败通知、进程内 AI 重试、插件监控与事件投递的生命周期所有权;2026-08-24 的阶段 2 已将 212 个已观察宿主模块方法的 legacy aggregation 清零,并补齐可执行 fanout 与下载器文件 DTO 边界;阶段 3 已将消息交互和远程命令的订阅删除统一到 Application/UoW/outbox,宿主不再调用裸线程统计入口;阶段 4 已统一七种消息渠道的宿主回环与后台执行边界;阶段 5 已补齐事件窗口聚合任务的生命周期所有权;阶段 6 已统一插件文件操作的取消完成语义;阶段 7 已统一插件协程补偿的终态等待;阶段 8 已统一宿主同步函数的异步线程池入口;阶段 9 已统一工作流运行时的宿主获取路径;阶段 10 已统一模块、插件与调度运行时的显式 getter 调用;阶段 11 已清除系统配置 getter 的 Oper 形别名;阶段 12 已完成工作流域的显式 Chain 数据端口迁移;阶段 13 已收口用户、交互与消息链的数据端口;阶段 14 已收口音乐订阅数据端口;阶段 15 已收口站点数据端口;阶段 16 已收口媒体服务器数据端口;阶段 17 已收口下载数据端口;阶段 18 已收口主订阅数据端口;阶段 19 已收口整理数据端口;阶段 20 已收口 Agent 数据端口;阶段 21 已收口监控历史端口;阶段 22 已统一服务配置应用边界;阶段 23 已补齐媒体服务器 API 遗留的类形配置读取路径;阶段 24 已清除 Scheduler 内部无 owner 的协程提交双轨;阶段 25 已补齐 TaskRegistry 跨线程 owner 并迁移整理 AI 接管;阶段 26 已统一 Agent 会话清理提交;阶段 27 已统一历史 AI 进度 owner;阶段 28 已托管旧插件订阅统计线程;阶段 29 已统一 Emby 系条目转换并清零重复代码白名单;阶段 30 已收口插件市场请求级子任务;阶段 31 已托管搜索 AI 推荐任务;阶段 32 已清除事件调度器绕过生命周期 owner 的投递回退;阶段 33 已统一宿主 Agent 运行时的获取路径;阶段 34 已统一 durable-required 事件与 Outbox topic 事实源;阶段 35 已统一 LLM provider 管理 API 的运行时解析路径;阶段 36 已统一 WebAgent 音频能力访问边界;阶段 37 已统一插件输入事件发布路径;阶段 38 已统一 WebAgent 通知事件监听与队列边界;阶段 39 已补齐搜索 SSE 断线时的上游任务清理;阶段 40 已补齐异步防抖取消的终态所有权;阶段 41 已统一优雅重启兜底线程的唯一所有权;阶段 42 已补齐 Telegram typing 的多实例隔离和终态 owner。 ## 当前复核结论(2026-08-24) @@ -438,6 +438,18 @@ - 该 daemon monitor 故意跨越正常 shutdown drain,以便进程卡死时仍能请求 Docker 重启,因此不纳入 `TaskRegistry` 或普通线程池等待。`SystemHelper` 公开方法、SDK/Compat 映射和 V1/V2/V3 插件 ABI 均未修改。 +### 长期整改阶段 42:Telegram typing 多实例与终态所有权收口(2026-08-24) + +- `TelegramModule` 支持多个通知配置实例,但用户到 chat 的映射、typing 线程、停止信号和锁原先都在 + `Telegram` 类级共享;两个配置遇到相同 chat ID 时会相互停止或覆盖 owner。生产 client 现在各自持有 + 完整运行状态,同一配置内则由 lifecycle 锁串行替换,保留一个 chat 一个 typing owner。 +- 停止等待超过预算时不再提前删除仍存活的线程句柄,新请求会拒绝覆盖阻塞中的旧 owner;线程启动失败会 + 回滚 owner 和停止信号,client 停止时先封住新增任务再取得完整快照。SDK 请求恢复后,线程仍由自己的 + `finally` 在真实终态释放登记。 +- `Telegram` 类路径、构造参数、`start_typing()`/`stop_typing()` 布尔合同、消息格式、模块方法及配置字段 + 均保持不变;私有类级可变状态已清除,正常 V1/V2/V3 模块实例不再共享运行状态。本阶段未修改插件仓、 + SDK 或 Compat 映射。 + ### 总体判断 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: diff --git a/tests/test_module_lifecycle.py b/tests/test_module_lifecycle.py index 597bbc04f..7c78658c0 100644 --- a/tests/test_module_lifecycle.py +++ b/tests/test_module_lifecycle.py @@ -8,13 +8,13 @@ from app.modules.discord import DiscordModule from app.modules.feishu import FeishuModule from app.modules.filter import FilterModule from app.modules.plex import PlexModule -from app.modules.qqbot import QQBotModule +from app.modules.qqbot.module import QQBotModule from app.modules.slack import SlackModule -from app.modules.telegram import TelegramModule +from app.modules.telegram.module import TelegramModule from app.modules.telegram.telegram import Telegram from app.modules.themoviedb import TheMovieDbModule -from app.modules.trimemedia import TrimeMediaModule -from app.modules.ugreen import UgreenModule +from app.modules.trimemedia.module import TrimeMediaModule +from app.modules.ugreen.module import UgreenModule from app.modules.wechat import WechatModule from app.modules.wechatclawbot import WechatClawBotModule @@ -177,6 +177,11 @@ def test_telegram_stop_closes_sdk_and_waits_for_polling_thread(): client._bot = bot polling_thread = Mock() client._polling_thread = polling_thread + client._typing_tasks = {} + client._typing_stop_flags = {} + client._typing_lock = threading.RLock() + client._typing_lifecycle_lock = threading.RLock() + client._typing_accepting = True client.stop() client.stop() diff --git a/tests/test_telegram_typing_lifecycle.py b/tests/test_telegram_typing_lifecycle.py index 16c66949d..07e718eb4 100644 --- a/tests/test_telegram_typing_lifecycle.py +++ b/tests/test_telegram_typing_lifecycle.py @@ -1,11 +1,11 @@ import asyncio import threading import time -import unittest from dataclasses import replace from types import SimpleNamespace from unittest.mock import AsyncMock, Mock, patch +import pytest from app.agent import ( # pylint: disable=no-name-in-module AgentManager, @@ -34,217 +34,383 @@ class _FakeTelegramBot: """记录 typing 调用的轻量 bot,避免后台线程与 Mock 内部锁交互。""" def __init__(self): + """初始化调用记录和首个动作事件。""" self.chat_actions = [] self.action_event = threading.Event() def send_chat_action(self, chat_id, action): + """记录一次 typing 动作并唤醒等待者。""" self.chat_actions.append((chat_id, action)) self.action_event.set() -class TestTelegramTypingLifecycle(unittest.TestCase): - def setUp(self): - self._cleanup_typing_tasks() +class _BlockingTelegramBot: + """模拟阻塞中的 Telegram SDK 请求。""" - def tearDown(self): - self._cleanup_typing_tasks() + def __init__(self) -> None: + """初始化进入和释放请求的事件屏障。""" + self.entered = threading.Event() + self.release = threading.Event() - @staticmethod - def _cleanup_typing_tasks(): - helper = Telegram.__new__(Telegram) - for chat_id in list(Telegram._typing_tasks.keys()): - helper._stop_typing_task(chat_id) - Telegram._typing_tasks.clear() - Telegram._typing_stop_flags.clear() - Telegram._user_chat_mapping.clear() + def send_chat_action(self, _chat_id, _action) -> None: + """阻塞请求,直到测试显式释放。""" + self.entered.set() + self.release.wait(timeout=1) - @staticmethod - def _telegram_client() -> Telegram: - telegram = Telegram.__new__(Telegram) - telegram._bot = _FakeTelegramBot() - telegram._telegram_token = "token" - telegram._telegram_chat_id = "default-chat" - # 缩短测试中的等待时间,不改变生产默认续发间隔。 - telegram._typing_interval_seconds = 0.01 - telegram._typing_max_duration_seconds = 1 - return telegram - def test_start_typing_can_stop_by_chat_id(self): - telegram = self._telegram_client() +def _telegram_client(bot=None) -> Telegram: + """构造不连接外部服务且持有独立运行状态的 Telegram client。""" + telegram = Telegram.__new__(Telegram) + telegram._bot = bot or _FakeTelegramBot() + telegram._telegram_token = "token" + telegram._telegram_chat_id = "default-chat" + telegram._user_chat_mapping = {} + telegram._typing_tasks = {} + telegram._typing_stop_flags = {} + telegram._typing_lock = threading.RLock() + telegram._typing_lifecycle_lock = threading.RLock() + telegram._typing_accepting = True + telegram._typing_join_timeout_seconds = 0.01 + # 缩短测试中的等待时间,不改变生产默认续发间隔。 + telegram._typing_interval_seconds = 0.01 + telegram._typing_max_duration_seconds = 1 + return telegram - telegram._start_typing_task( - "chat-1", - max_duration_seconds=1, - initial_delay_seconds=0, + +def test_start_typing_can_stop_by_chat_id(): + telegram = _telegram_client() + + telegram._start_typing_task( + "chat-1", + max_duration_seconds=1, + initial_delay_seconds=0, + ) + + assert "chat-1" in telegram._typing_tasks + assert telegram._bot.action_event.wait(1.0) + assert telegram.stop_typing(chat_id="chat-1") + assert "chat-1" not in telegram._typing_tasks + + +def test_start_typing_can_stop_by_user_mapping(): + telegram = _telegram_client() + telegram._user_chat_mapping["10001"] = "chat-2" + + telegram._start_typing_task( + "chat-2", + max_duration_seconds=1, + initial_delay_seconds=0, + ) + time.sleep(0.03) + + assert telegram.stop_typing(userid="10001") + assert "chat-2" not in telegram._typing_tasks + + +def test_typing_task_has_max_duration_guard(): + telegram = _telegram_client() + + telegram._start_typing_task( + "chat-3", + max_duration_seconds=0.02, + initial_delay_seconds=0, + ) + + assert _wait_until(lambda: "chat-3" not in telegram._typing_tasks) + assert "chat-3" not in telegram._typing_tasks + + +def test_short_typing_task_can_stop_before_first_chat_action(): + """ + 短响应在首次 typing 发出前结束时,不应留下客户端自然过期的残留状态。 + """ + telegram = _telegram_client() + + telegram._start_typing_task( + "chat-4", + max_duration_seconds=1, + initial_delay_seconds=0.05, + ) + telegram.stop_typing(chat_id="chat-4") + time.sleep(0.08) + + assert telegram._bot.chat_actions == [] + assert "chat-4" not in telegram._typing_tasks + + +def test_typing_owner_is_isolated_between_config_instances(): + """不同 Telegram 配置即使 chat_id 相同也不得互相停止 typing。""" + first = _telegram_client() + second = _telegram_client() + try: + assert first._start_typing_task("shared-chat", initial_delay_seconds=0) + assert second._start_typing_task("shared-chat", initial_delay_seconds=0) + assert first._bot.action_event.wait(timeout=1) + assert second._bot.action_event.wait(timeout=1) + + assert first.stop_typing(chat_id="shared-chat") + + assert "shared-chat" not in first._typing_tasks + assert "shared-chat" in second._typing_tasks + finally: + first.stop_typing(chat_id="shared-chat") + second.stop_typing(chat_id="shared-chat") + + +def test_typing_stop_keeps_blocked_thread_owner_until_terminal(): + """SDK 请求阻塞超过等待预算时,不得提前删除线程 owner。""" + bot = _BlockingTelegramBot() + telegram = _telegram_client(bot) + try: + assert telegram._start_typing_task("blocked-chat", initial_delay_seconds=0) + assert bot.entered.wait(timeout=1) + + assert telegram._stop_typing_task("blocked-chat") is False + owner = telegram._typing_tasks["blocked-chat"] + assert owner.is_alive() + assert telegram._start_typing_task( + "blocked-chat", initial_delay_seconds=0 + ) is False + + bot.release.set() + owner.join(timeout=1) + assert not owner.is_alive() + assert "blocked-chat" not in telegram._typing_tasks + finally: + bot.release.set() + telegram.stop_typing(chat_id="blocked-chat") + + +def test_typing_start_failure_releases_registered_owner(monkeypatch): + """线程启动失败时应清理已登记的 owner 和停止信号。""" + telegram = _telegram_client() + + class FailingThread: + """模拟登记成功后无法启动的线程对象。""" + + def __init__(self, **_kwargs) -> None: + """接收真实 Thread 构造参数。""" + + def start(self) -> None: + """模拟系统拒绝创建新线程。""" + raise RuntimeError("thread start failed") + + monkeypatch.setattr("app.modules.telegram.telegram.threading.Thread", FailingThread) + + with pytest.raises(RuntimeError, match="thread start failed"): + telegram._start_typing_task("failed-chat", initial_delay_seconds=0) + + assert telegram._typing_tasks == {} + assert telegram._typing_stop_flags == {} + + +def test_typing_start_is_rejected_after_client_stop(): + """client 停止封口后不得再接受新的 typing 线程。""" + telegram = _telegram_client() + telegram._bot = None + + telegram.stop() + + assert telegram._start_typing_task("closed-chat", initial_delay_seconds=0) is False + assert telegram._typing_tasks == {} + + +def test_agent_managed_send_msg_keeps_typing_for_worker_cleanup(): + telegram = _telegram_client() + sent = SimpleNamespace(message_id=1, chat=SimpleNamespace(id="chat-1")) + + with patch.object( + telegram, "_Telegram__send_request", return_value=sent + ), patch.object(telegram, "_stop_typing_task") as stop_typing: + result = telegram.send_msg( + title="处理中", + userid="10001", + stop_typing=False, ) - self.assertIn("chat-1", Telegram._typing_tasks) - self.assertTrue(telegram._bot.action_event.wait(1.0)) - self.assertTrue(telegram.stop_typing(chat_id="chat-1")) - self.assertNotIn("chat-1", Telegram._typing_tasks) + assert result["success"] + stop_typing.assert_not_called() - def test_start_typing_can_stop_by_user_mapping(self): - telegram = self._telegram_client() - Telegram._user_chat_mapping["10001"] = "chat-2" - telegram._start_typing_task( - "chat-2", - max_duration_seconds=1, - initial_delay_seconds=0, - ) - time.sleep(0.03) +def test_send_msg_does_not_stop_typing_by_default(): + """ + 响应发送不再默认结束 typing,由处理状态统一收口。 + """ + telegram = _telegram_client() + sent = SimpleNamespace(message_id=1, chat=SimpleNamespace(id="chat-1")) - self.assertTrue(telegram.stop_typing(userid="10001")) - self.assertNotIn("chat-2", Telegram._typing_tasks) + with patch.object( + telegram, "_Telegram__send_request", return_value=sent + ), patch.object(telegram, "_stop_typing_task") as stop_typing: + result = telegram.send_msg(title="处理中", userid="10001") - def test_typing_task_has_max_duration_guard(self): - telegram = self._telegram_client() + assert result["success"] + stop_typing.assert_not_called() - telegram._start_typing_task( - "chat-3", - max_duration_seconds=0.02, - initial_delay_seconds=0, - ) - self.assertTrue(_wait_until(lambda: "chat-3" not in Telegram._typing_tasks)) - self.assertNotIn("chat-3", Telegram._typing_tasks) +def test_telegram_module_processing_status_starts_typing(): + """ + Telegram 通过模块处理状态接口启动 typing 保活。 + """ + module = TelegramModule() + module._channel = NotificationChannel.Telegram + client = Mock() + client.start_typing.return_value = True - def test_short_typing_task_can_stop_before_first_chat_action(self): - """ - 短响应在首次 typing 发出前结束时,不应留下客户端自然过期的残留状态。 - """ - telegram = self._telegram_client() - - telegram._start_typing_task( - "chat-4", - max_duration_seconds=1, - initial_delay_seconds=0.05, - ) - telegram.stop_typing(chat_id="chat-4") - time.sleep(0.08) - - self.assertEqual(telegram._bot.chat_actions, []) - self.assertNotIn("chat-4", Telegram._typing_tasks) - - def test_agent_managed_send_msg_keeps_typing_for_worker_cleanup(self): - telegram = self._telegram_client() - sent = SimpleNamespace(message_id=1, chat=SimpleNamespace(id="chat-1")) - - with patch.object( - telegram, "_Telegram__send_request", return_value=sent - ), patch.object(telegram, "_stop_typing_task") as stop_typing: - result = telegram.send_msg( - title="处理中", - userid="10001", - stop_typing=False, - ) - - self.assertTrue(result["success"]) - stop_typing.assert_not_called() - - def test_send_msg_does_not_stop_typing_by_default(self): - """ - 响应发送不再默认结束 typing,由处理状态统一收口。 - """ - telegram = self._telegram_client() - sent = SimpleNamespace(message_id=1, chat=SimpleNamespace(id="chat-1")) - - with patch.object( - telegram, "_Telegram__send_request", return_value=sent - ), patch.object(telegram, "_stop_typing_task") as stop_typing: - result = telegram.send_msg(title="处理中", userid="10001") - - self.assertTrue(result["success"]) - stop_typing.assert_not_called() - - def test_telegram_module_processing_status_starts_typing(self): - """ - Telegram 通过模块处理状态接口启动 typing 保活。 - """ - module = TelegramModule() - module._channel = NotificationChannel.Telegram - client = Mock() - client.start_typing.return_value = True - - with patch.object( - module, "get_config", return_value=SimpleNamespace(name="telegram-test") - ), patch.object(module, "get_instance", return_value=client): - status = module.mark_message_processing_started( - channel=NotificationChannel.Telegram, - source="telegram-test", - userid="10001", - chat_id="-100", - text="hello", - ) - - client.start_typing.assert_called_once_with(chat_id="-100", userid="10001") - self.assertEqual(status["metadata"]["kind"], "typing") - - def test_slash_command_defers_processing_status_to_command_handler(self): - chain = MessageChain.__new__(MessageChain) - chain.eventmanager = Mock() - status = MessageChain._ProcessingStatus( + with patch.object( + module, "get_config", return_value=SimpleNamespace(name="telegram-test") + ), patch.object(module, "get_instance", return_value=client): + status = module.mark_message_processing_started( channel=NotificationChannel.Telegram, source="telegram-test", userid="10001", chat_id="-100", - metadata={"kind": "typing"}, + text="hello", ) - with patch.object(chain, "_record_user_message"), patch.object( - chain, "_mark_message_processing_started", return_value=status - ), patch.object( - chain, "_mark_message_processing_finished" - ) as finish_status: - chain.handle_message( - channel=NotificationChannel.Telegram, - source="telegram-test", - userid="10001", - username="tester", - text="/sites", - original_chat_id="-100", - ) + client.start_typing.assert_called_once_with(chat_id="-100", userid="10001") + assert status["metadata"]["kind"] == "typing" - finish_status.assert_not_called() - chain.eventmanager.send_event.assert_called_once() - self.assertEqual( - chain.eventmanager.send_event.call_args.args[1]["processing_status"], - status.to_dict(), + +def test_slash_command_defers_processing_status_to_command_handler(): + chain = MessageChain.__new__(MessageChain) + chain.eventmanager = Mock() + status = MessageChain._ProcessingStatus( + channel=NotificationChannel.Telegram, + source="telegram-test", + userid="10001", + chat_id="-100", + metadata={"kind": "typing"}, + ) + + with patch.object(chain, "_record_user_message"), patch.object( + chain, "_mark_message_processing_started", return_value=status + ), patch.object( + chain, "_mark_message_processing_finished" + ) as finish_status: + chain.handle_message( + channel=NotificationChannel.Telegram, + source="telegram-test", + userid="10001", + username="tester", + text="/sites", + original_chat_id="-100", ) - def test_command_handler_finishes_processing_status_after_execute(self): - """ - 传统命令响应完成后由命令处理器统一结束 processing status。 - """ - command = Command.__new__(Command) - command.get = Mock(return_value={"func": Mock()}) - command.execute = Mock() - event = SimpleNamespace( - event_data={ - "cmd": "/sites", - "user": "10001", - "channel": NotificationChannel.Telegram, + finish_status.assert_not_called() + chain.eventmanager.send_event.assert_called_once() + assert ( + chain.eventmanager.send_event.call_args.args[1]["processing_status"] + == status.to_dict() + ) + + +def test_command_handler_finishes_processing_status_after_execute(): + """ + 传统命令响应完成后由命令处理器统一结束 processing status。 + """ + command = Command.__new__(Command) + command.get = Mock(return_value={"func": Mock()}) + command.execute = Mock() + event = SimpleNamespace( + event_data={ + "cmd": "/sites", + "user": "10001", + "channel": NotificationChannel.Telegram, + "source": "telegram-test", + "processing_status": { + "channel": NotificationChannel.Telegram.value, "source": "telegram-test", - "processing_status": { - "channel": NotificationChannel.Telegram.value, - "source": "telegram-test", - "userid": "10001", - "chat_id": "-100", - "metadata": {"kind": "typing"}, - }, - } + "userid": "10001", + "chat_id": "-100", + "metadata": {"kind": "typing"}, + }, + } + ) + + with patch("app.command._finish_command_processing_status") as finish_status: + command.command_event(event) + + command.execute.assert_called_once() + finish_status.assert_called_once_with( + event.event_data["processing_status"], + user_id="10001", + ) + + +def test_finish_command_processing_status_uses_module_interface(): + status = { + "channel": NotificationChannel.Telegram.value, + "source": "telegram-test", + "userid": "10001", + "chat_id": "-100", + "metadata": {"kind": "typing"}, + } + + with patch("app.command.CommandChain") as chain_cls: + _finish_command_processing_status(status, user_id="fallback") + + chain_cls.return_value.finish_message_processing_status.assert_called_once_with( + status=status, + userid="fallback", + ) + + +def test_async_agent_leaves_processing_status_to_worker(): + chain = MessageChain.__new__(MessageChain) + chain.eventmanager = Mock() + chain.runtime_config = replace( + chain.runtime_config, + ai_agent_enable=True, + ) + + loop = Mock(**{"is_running.return_value": True, "is_closed.return_value": False}) + with patch.object(global_vars, "CURRENT_EVENT_LOOP", loop), patch.object( + chain, "_record_user_message" + ), patch.object( + chain, "_mark_message_processing_started" + ) as start_status, patch( + "app.chain.message.get_running_agent_manager", + ) as get_running_manager, patch( + "app.chain.message.asyncio.run_coroutine_threadsafe", + side_effect=lambda coro, _loop: (coro.close(), Mock())[1], + ), patch.object( + chain, "_mark_message_processing_finished" + ) as finish_status: + process_message = AsyncMock() + get_running_manager.return_value.process_message = process_message + chain.handle_message( + channel=NotificationChannel.Telegram, + source="telegram-test", + userid="10001", + username="tester", + text="/ai 搜索电影", + original_chat_id="-100", ) - with patch("app.command._finish_command_processing_status") as finish_status: - command.command_event(event) + start_status.assert_not_called() + finish_status.assert_not_called() + process_message.assert_called_once() + assert "processing_status" not in process_message.call_args.kwargs + assert ( + process_message.call_args.kwargs["channel"] + == NotificationChannel.Telegram.value + ) + assert process_message.call_args.kwargs["source"] == "telegram-test" + assert process_message.call_args.kwargs["original_chat_id"] == "-100" - command.execute.assert_called_once() - finish_status.assert_called_once_with( - event.event_data["processing_status"], + +def test_agent_manager_starts_processing_status_when_task_runs(): + async def _run(): + manager = AgentManager() + task = _MessageTask( + session_id="session-1", user_id="10001", + message="第一条", + channel=NotificationChannel.Telegram.value, + source="telegram-test", + original_chat_id="-100", ) - - def test_finish_command_processing_status_uses_module_interface(self): status = { "channel": NotificationChannel.Telegram.value, "source": "telegram-test", @@ -253,281 +419,213 @@ class TestTelegramTypingLifecycle(unittest.TestCase): "metadata": {"kind": "typing"}, } - with patch("app.command.CommandChain") as chain_cls: - _finish_command_processing_status(status, user_id="fallback") + with patch( + "app.agent.orchestrator._async_start_processing_status", + new_callable=AsyncMock, + return_value=status, + ) as start_status: + await manager._start_task_processing_status(task) - chain_cls.return_value.finish_message_processing_status.assert_called_once_with( - status=status, - userid="fallback", + start_status.assert_awaited_once_with(task) + assert task.processing_status == status + + asyncio.run(_run()) + + +def test_agent_start_processing_status_uses_chain_interface(): + async def _run(): + task = _MessageTask( + session_id="session-1", + user_id="10001", + message="第一条", + channel=NotificationChannel.Telegram.value, + source="telegram-test", + original_message_id="10", + original_chat_id="-100", ) + status = { + "channel": NotificationChannel.Telegram.value, + "source": "telegram-test", + "userid": "10001", + "message_id": "10", + "chat_id": "-100", + "metadata": {"kind": "typing"}, + } + calls = [] - def test_async_agent_leaves_processing_status_to_worker(self): - chain = MessageChain.__new__(MessageChain) - chain.eventmanager = Mock() - chain.runtime_config = replace( - chain.runtime_config, - ai_agent_enable=True, - ) + class FakeAgentChain: + """记录 processing status 请求的 Agent Chain 替身。""" - loop = Mock(**{"is_running.return_value": True, "is_closed.return_value": False}) - with patch.object(global_vars, "CURRENT_EVENT_LOOP", loop), patch.object( - chain, "_record_user_message" - ), patch.object( - chain, "_mark_message_processing_started" - ) as start_status, patch( - "app.chain.message.get_running_agent_manager", - ) as get_running_manager, patch( - "app.chain.message.asyncio.run_coroutine_threadsafe", - side_effect=lambda coro, _loop: (coro.close(), Mock())[1], - ), patch.object( - chain, "_mark_message_processing_finished" - ) as finish_status: - process_message = AsyncMock() - get_running_manager.return_value.process_message = process_message - chain.handle_message( - channel=NotificationChannel.Telegram, - source="telegram-test", - userid="10001", - username="tester", - text="/ai 搜索电影", - original_chat_id="-100", - ) + def start_message_processing_status(self, **kwargs): + """记录请求并返回预设状态。""" + calls.append(kwargs) + return status - start_status.assert_not_called() - finish_status.assert_not_called() - process_message.assert_called_once() - self.assertNotIn("processing_status", process_message.call_args.kwargs) - self.assertEqual( - process_message.call_args.kwargs["channel"], - NotificationChannel.Telegram.value, - ) - self.assertEqual(process_message.call_args.kwargs["source"], "telegram-test") - self.assertEqual(process_message.call_args.kwargs["original_chat_id"], "-100") + with patch("app.agent.orchestrator.AgentChain", FakeAgentChain): + result = await _async_start_processing_status(task) - def test_agent_manager_starts_processing_status_when_task_runs(self): - async def _run(): - manager = AgentManager() - task = _MessageTask( - session_id="session-1", - user_id="10001", - message="第一条", - channel=NotificationChannel.Telegram.value, - source="telegram-test", - original_chat_id="-100", - ) - status = { - "channel": NotificationChannel.Telegram.value, - "source": "telegram-test", - "userid": "10001", - "chat_id": "-100", - "metadata": {"kind": "typing"}, - } + assert calls == [{ + "channel": NotificationChannel.Telegram, + "source": "telegram-test", + "userid": "10001", + "message_id": "10", + "chat_id": "-100", + "text": "第一条", + }] + assert result == status - with patch( - "app.agent.orchestrator._async_start_processing_status", - new_callable=AsyncMock, - return_value=status, - ) as start_status: - await manager._start_task_processing_status(task) + asyncio.run(_run()) - start_status.assert_awaited_once_with(task) - self.assertEqual(task.processing_status, status) - asyncio.run(_run()) +def test_callback_stops_typing_when_message_handler_returns(): + chain = MessageChain.__new__(MessageChain) + chain.eventmanager = Mock() + status = MessageChain._ProcessingStatus( + channel=NotificationChannel.Telegram, + source="telegram-test", + userid="10001", + chat_id="-100", + metadata={"kind": "typing"}, + ) - def test_agent_start_processing_status_uses_chain_interface(self): - async def _run(): - task = _MessageTask( - session_id="session-1", - user_id="10001", - message="第一条", - channel=NotificationChannel.Telegram.value, - source="telegram-test", - original_message_id="10", - original_chat_id="-100", - ) - status = { - "channel": NotificationChannel.Telegram.value, - "source": "telegram-test", - "userid": "10001", - "message_id": "10", - "chat_id": "-100", - "metadata": {"kind": "typing"}, - } - calls = [] - - class FakeAgentChain: - def start_message_processing_status(self, **kwargs): - calls.append(kwargs) - return status - - with patch("app.agent.orchestrator.AgentChain", FakeAgentChain): - result = await _async_start_processing_status(task) - - self.assertEqual(calls, [{ - "channel": NotificationChannel.Telegram, - "source": "telegram-test", - "userid": "10001", - "message_id": "10", - "chat_id": "-100", - "text": "第一条", - }]) - self.assertEqual(result, status) - - asyncio.run(_run()) - - def test_callback_stops_typing_when_message_handler_returns(self): - chain = MessageChain.__new__(MessageChain) - chain.eventmanager = Mock() - status = MessageChain._ProcessingStatus( + with patch.object(chain, "_record_user_message"), patch.object( + chain, "_mark_message_processing_started", return_value=status + ), patch.object(chain, "_handle_message_core"), patch.object( + chain, "_mark_message_processing_finished" + ) as finish_status: + chain.handle_message( channel=NotificationChannel.Telegram, source="telegram-test", userid="10001", - chat_id="-100", - metadata={"kind": "typing"}, - ) - - with patch.object(chain, "_record_user_message"), patch.object( - chain, "_mark_message_processing_started", return_value=status - ), patch.object(chain, "_handle_message_core"), patch.object( - chain, "_mark_message_processing_finished" - ) as finish_status: - chain.handle_message( - channel=NotificationChannel.Telegram, - source="telegram-test", - userid="10001", - username="tester", - text="CALLBACK:sites:req-1:refresh", - original_chat_id="-100", - ) - - finish_status.assert_called_once_with( - channel=NotificationChannel.Telegram, - source="telegram-test", - userid="10001", - status=status, - original_message_id=None, + username="tester", + text="CALLBACK:sites:req-1:refresh", original_chat_id="-100", ) - def test_chain_finishes_processing_through_module_interface(self): - chain = MessageChain.__new__(MessageChain) - status = MessageChain._ProcessingStatus( + finish_status.assert_called_once_with( + channel=NotificationChannel.Telegram, + source="telegram-test", + userid="10001", + status=status, + original_message_id=None, + original_chat_id="-100", + ) + + +def test_chain_finishes_processing_through_module_interface(): + chain = MessageChain.__new__(MessageChain) + status = MessageChain._ProcessingStatus( + channel=NotificationChannel.Telegram, + source="telegram-test", + userid="10001", + chat_id="-100", + metadata={"kind": "typing"}, + ) + + with patch.object(chain, "finish_message_processing_status") as finish_status: + chain._mark_message_processing_finished( channel=NotificationChannel.Telegram, source="telegram-test", userid="10001", - chat_id="-100", - metadata={"kind": "typing"}, + status=status, + original_chat_id="-100", ) - with patch.object(chain, "finish_message_processing_status") as finish_status: - chain._mark_message_processing_finished( - channel=NotificationChannel.Telegram, - source="telegram-test", - userid="10001", - status=status, - original_chat_id="-100", - ) + finish_status.assert_called_once_with( + status=status.to_dict(), + channel=NotificationChannel.Telegram, + source="telegram-test", + userid="10001", + message_id=None, + chat_id="-100", + ) - finish_status.assert_called_once_with( - status=status.to_dict(), - channel=NotificationChannel.Telegram, + +def test_agent_manager_finishes_processing_status_after_each_task(): + async def _run(): + manager = AgentManager() + status = { + "channel": NotificationChannel.Telegram.value, + "source": "telegram-test", + "userid": "10001", + "chat_id": "-100", + "metadata": {"kind": "typing"}, + } + task = _MessageTask( + session_id="session-1", + user_id="10001", + message="第一条", + processing_status=status, + ) + + with patch( + "app.agent.orchestrator._async_finish_processing_status", + new_callable=AsyncMock, + ) as finish_status: + await manager._finish_task_processing_status(task) + + finish_status.assert_awaited_once_with(status, "10001") + assert task.processing_status is None + + asyncio.run(_run()) + + +def test_agent_worker_starts_and_finishes_each_queued_task(): + async def _run(): + manager = AgentManager() + manager._session_queues["session-1"] = asyncio.Queue() + first_status = { + "channel": NotificationChannel.Telegram.value, + "source": "telegram-test", + "userid": "10001", + "chat_id": "-100", + "metadata": {"kind": "typing", "seq": 1}, + } + second_status = { + "channel": NotificationChannel.Telegram.value, + "source": "telegram-test", + "userid": "10001", + "chat_id": "-100", + "metadata": {"kind": "typing", "seq": 2}, + } + await manager._session_queues["session-1"].put(_MessageTask( + session_id="session-1", + user_id="10001", + message="第一条", + channel=NotificationChannel.Telegram.value, source="telegram-test", - userid="10001", - message_id=None, - chat_id="-100", - ) + original_chat_id="-100", + )) + await manager._session_queues["session-1"].put(_MessageTask( + session_id="session-1", + user_id="10001", + message="第二条", + channel=NotificationChannel.Telegram.value, + source="telegram-test", + original_chat_id="-100", + )) - def test_agent_manager_finishes_processing_status_after_each_task(self): - async def _run(): - manager = AgentManager() - status = { - "channel": NotificationChannel.Telegram.value, - "source": "telegram-test", - "userid": "10001", - "chat_id": "-100", - "metadata": {"kind": "typing"}, - } - task = _MessageTask( - session_id="session-1", - user_id="10001", - message="第一条", - processing_status=status, + with patch( + "app.agent.orchestrator._async_start_processing_status", + new_callable=AsyncMock, + side_effect=[first_status, second_status], + ) as start_status, patch.object( + manager, + "_process_message_internal", + new_callable=AsyncMock, + ), patch( + "app.agent.orchestrator._async_finish_processing_status", + new_callable=AsyncMock, + ) as finish_status: + manager._session_workers["session-1"] = asyncio.create_task( + manager._session_worker("session-1") ) + await manager._session_queues["session-1"].join() + manager._session_workers["session-1"].cancel() + await manager._session_workers["session-1"] - with patch( - "app.agent.orchestrator._async_finish_processing_status", - new_callable=AsyncMock, - ) as finish_status: - await manager._finish_task_processing_status(task) + assert start_status.await_count == 2 + assert finish_status.await_args_list[0].args == (first_status, "10001") + assert finish_status.await_args_list[1].args == (second_status, "10001") - finish_status.assert_awaited_once_with(status, "10001") - self.assertIsNone(task.processing_status) - - asyncio.run(_run()) - - def test_agent_worker_starts_and_finishes_each_queued_task(self): - async def _run(): - manager = AgentManager() - manager._session_queues["session-1"] = asyncio.Queue() - first_status = { - "channel": NotificationChannel.Telegram.value, - "source": "telegram-test", - "userid": "10001", - "chat_id": "-100", - "metadata": {"kind": "typing", "seq": 1}, - } - second_status = { - "channel": NotificationChannel.Telegram.value, - "source": "telegram-test", - "userid": "10001", - "chat_id": "-100", - "metadata": {"kind": "typing", "seq": 2}, - } - await manager._session_queues["session-1"].put(_MessageTask( - session_id="session-1", - user_id="10001", - message="第一条", - channel=NotificationChannel.Telegram.value, - source="telegram-test", - original_chat_id="-100", - )) - await manager._session_queues["session-1"].put(_MessageTask( - session_id="session-1", - user_id="10001", - message="第二条", - channel=NotificationChannel.Telegram.value, - source="telegram-test", - original_chat_id="-100", - )) - - with patch( - "app.agent.orchestrator._async_start_processing_status", - new_callable=AsyncMock, - side_effect=[first_status, second_status], - ) as start_status, patch.object( - manager, - "_process_message_internal", - new_callable=AsyncMock, - ), patch( - "app.agent.orchestrator._async_finish_processing_status", - new_callable=AsyncMock, - ) as finish_status: - manager._session_workers["session-1"] = asyncio.create_task( - manager._session_worker("session-1") - ) - await manager._session_queues["session-1"].join() - manager._session_workers["session-1"].cancel() - await manager._session_workers["session-1"] - - self.assertEqual(start_status.await_count, 2) - self.assertEqual( - finish_status.await_args_list[0].args, - (first_status, "10001"), - ) - self.assertEqual( - finish_status.await_args_list[1].args, - (second_status, "10001"), - ) - - asyncio.run(_run()) + asyncio.run(_run())