fix: isolate telegram typing lifecycle

This commit is contained in:
jxxghp
2026-08-24 08:51:49 +08:00
parent 01d29d952f
commit ecbd39a0bb
4 changed files with 647 additions and 498 deletions
+53 -19
View File
@@ -77,18 +77,13 @@ class Telegram:
_bot: TeleBot = None _bot: TeleBot = None
_callback_handlers: Dict[str, Callable] = {} # 存储回调处理器 _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 _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_interval_seconds = 5
_typing_initial_delay_seconds = 1 _typing_initial_delay_seconds = 1
_typing_max_duration_seconds = 10 * 60 _typing_max_duration_seconds = 10 * 60
_typing_command_max_duration_seconds = 30 _typing_command_max_duration_seconds = 30
_typing_callback_max_duration_seconds = 60 _typing_callback_max_duration_seconds = 60
_typing_join_timeout_seconds = 1
def __init__( def __init__(
self, self,
@@ -103,6 +98,13 @@ class Telegram:
self._telegram_token = TELEGRAM_TOKEN self._telegram_token = TELEGRAM_TOKEN
self._telegram_chat_id = TELEGRAM_CHAT_ID self._telegram_chat_id = TELEGRAM_CHAT_ID
self._polling_thread = None 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: if not TELEGRAM_TOKEN or not TELEGRAM_CHAT_ID:
logger.error("Telegram配置不完整!") logger.error("Telegram配置不完整!")
return return
@@ -482,13 +484,24 @@ class Telegram:
chat_id: Union[str, int], chat_id: Union[str, int],
max_duration_seconds: Optional[float] = None, max_duration_seconds: Optional[float] = None,
initial_delay_seconds: Optional[float] = None, initial_delay_seconds: Optional[float] = None,
) -> None: ) -> bool:
""" """
启动持续发送正在输入状态的任务 启动持续发送正在输入状态的任务
:return: 是否取得该会话的唯一 typing owner
""" """
chat_id_str = str(chat_id) chat_id_str = str(chat_id)
# 如果已有任务在运行,先停止 with self._typing_lifecycle_lock:
self._stop_typing_task(chat_id_str) 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 线程互相误改停止标记。 # 使用独立 Event 避免同一 chat 新旧 typing 线程互相误改停止标记。
stop_event = threading.Event() stop_event = threading.Event()
@@ -528,24 +541,44 @@ class Telegram:
self._typing_tasks.pop(chat_id_str, None) self._typing_tasks.pop(chat_id_str, None)
self._typing_stop_flags.pop(chat_id_str, None) self._typing_stop_flags.pop(chat_id_str, None)
thread = threading.Thread(target=typing_worker, daemon=True) thread = threading.Thread(
target=typing_worker,
name=f"MoviePilot-TelegramTyping-{chat_id_str}"[:120],
daemon=True,
)
with self._typing_lock: with self._typing_lock:
self._typing_stop_flags[chat_id_str] = stop_event self._typing_stop_flags[chat_id_str] = stop_event
self._typing_tasks[chat_id_str] = thread self._typing_tasks[chat_id_str] = thread
try:
thread.start() 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]) -> None: def _stop_typing_task(self, chat_id: Union[str, int]) -> bool:
""" """
停止正在输入状态的任务 停止正在输入状态的任务,并保留尚未结束的 owner。
:return: 任务是否已经进入终态
""" """
chat_id_str = str(chat_id) chat_id_str = str(chat_id)
with self._typing_lifecycle_lock:
with self._typing_lock: with self._typing_lock:
stop_event = self._typing_stop_flags.pop(chat_id_str, None) stop_event = self._typing_stop_flags.get(chat_id_str)
task = self._typing_tasks.pop(chat_id_str, None) task = self._typing_tasks.get(chat_id_str)
if stop_event: if stop_event:
stop_event.set() stop_event.set()
if task and task.is_alive() and task is not threading.current_thread(): if task and task.is_alive() and task is not threading.current_thread():
task.join(timeout=1) 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:
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)
return task_finished
def _stop_typing_if_needed( def _stop_typing_if_needed(
self, chat_id: Union[str, int], stop_typing: bool 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) target_chat_id = target_chat_id or (str(userid) if userid else None)
if not target_chat_id: if not target_chat_id:
return False return False
self._start_typing_task(target_chat_id) return self._start_typing_task(target_chat_id)
return True
def stop_typing( def stop_typing(
self, self,
@@ -1714,7 +1746,9 @@ class Telegram:
""" """
停止Telegram消息接收服务 停止Telegram消息接收服务
""" """
# 停止所有typing任务 with self._typing_lifecycle_lock:
self._typing_accepting = False
# 封口与 owner 快照处于同一临界区,停止后不会漏掉并发新增任务。
for chat_id in list(self._typing_tasks.keys()): for chat_id in list(self._typing_tasks.keys()):
self._stop_typing_task(chat_id) self._stop_typing_task(chat_id)
if not self._bot: if not self._bot:
@@ -6,7 +6,7 @@
> 审计范围:宿主后端;排除 `app/plugins/**` 运行时插件副本 > 审计范围:宿主后端;排除 `app/plugins/**` 运行时插件副本
> 规范优先级:`AGENTS.md` 与 `docs/rules/` 高于本文 > 规范优先级:`AGENTS.md` 与 `docs/rules/` 高于本文
> 相关文档:`docs/architecture-overview.md`、`docs/refactor/backend-architecture-governance.md`、`docs/refactor/backend-module-refactor-compatibility.md` > 相关文档:`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 ## 当前复核结论(2026-08-24
@@ -438,6 +438,18 @@
- 该 daemon monitor 故意跨越正常 shutdown drain,以便进程卡死时仍能请求 Docker 重启,因此不纳入 - 该 daemon monitor 故意跨越正常 shutdown drain,以便进程卡死时仍能请求 Docker 重启,因此不纳入
`TaskRegistry` 或普通线程池等待。`SystemHelper` 公开方法、SDK/Compat 映射和 V1/V2/V3 插件 ABI 均未修改。 `TaskRegistry` 或普通线程池等待。`SystemHelper` 公开方法、SDK/Compat 映射和 V1/V2/V3 插件 ABI 均未修改。
### 长期整改阶段 42Telegram 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 映射。
### 总体判断 ### 总体判断
当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**:
+9 -4
View File
@@ -8,13 +8,13 @@ from app.modules.discord import DiscordModule
from app.modules.feishu import FeishuModule from app.modules.feishu import FeishuModule
from app.modules.filter import FilterModule from app.modules.filter import FilterModule
from app.modules.plex import PlexModule 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.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.telegram.telegram import Telegram
from app.modules.themoviedb import TheMovieDbModule from app.modules.themoviedb import TheMovieDbModule
from app.modules.trimemedia import TrimeMediaModule from app.modules.trimemedia.module import TrimeMediaModule
from app.modules.ugreen import UgreenModule from app.modules.ugreen.module import UgreenModule
from app.modules.wechat import WechatModule from app.modules.wechat import WechatModule
from app.modules.wechatclawbot import WechatClawBotModule from app.modules.wechatclawbot import WechatClawBotModule
@@ -177,6 +177,11 @@ def test_telegram_stop_closes_sdk_and_waits_for_polling_thread():
client._bot = bot client._bot = bot
polling_thread = Mock() polling_thread = Mock()
client._polling_thread = polling_thread 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()
client.stop() client.stop()
+175 -77
View File
@@ -1,11 +1,11 @@
import asyncio import asyncio
import threading import threading
import time import time
import unittest
from dataclasses import replace from dataclasses import replace
from types import SimpleNamespace from types import SimpleNamespace
from unittest.mock import AsyncMock, Mock, patch from unittest.mock import AsyncMock, Mock, patch
import pytest
from app.agent import ( # pylint: disable=no-name-in-module from app.agent import ( # pylint: disable=no-name-in-module
AgentManager, AgentManager,
@@ -34,43 +34,51 @@ class _FakeTelegramBot:
"""记录 typing 调用的轻量 bot,避免后台线程与 Mock 内部锁交互。""" """记录 typing 调用的轻量 bot,避免后台线程与 Mock 内部锁交互。"""
def __init__(self): def __init__(self):
"""初始化调用记录和首个动作事件。"""
self.chat_actions = [] self.chat_actions = []
self.action_event = threading.Event() self.action_event = threading.Event()
def send_chat_action(self, chat_id, action): def send_chat_action(self, chat_id, action):
"""记录一次 typing 动作并唤醒等待者。"""
self.chat_actions.append((chat_id, action)) self.chat_actions.append((chat_id, action))
self.action_event.set() self.action_event.set()
class TestTelegramTypingLifecycle(unittest.TestCase): class _BlockingTelegramBot:
def setUp(self): """模拟阻塞中的 Telegram SDK 请求。"""
self._cleanup_typing_tasks()
def tearDown(self): def __init__(self) -> None:
self._cleanup_typing_tasks() """初始化进入和释放请求的事件屏障。"""
self.entered = threading.Event()
self.release = threading.Event()
@staticmethod def send_chat_action(self, _chat_id, _action) -> None:
def _cleanup_typing_tasks(): """阻塞请求,直到测试显式释放。"""
helper = Telegram.__new__(Telegram) self.entered.set()
for chat_id in list(Telegram._typing_tasks.keys()): self.release.wait(timeout=1)
helper._stop_typing_task(chat_id)
Telegram._typing_tasks.clear()
Telegram._typing_stop_flags.clear()
Telegram._user_chat_mapping.clear()
@staticmethod
def _telegram_client() -> Telegram: def _telegram_client(bot=None) -> Telegram:
"""构造不连接外部服务且持有独立运行状态的 Telegram client。"""
telegram = Telegram.__new__(Telegram) telegram = Telegram.__new__(Telegram)
telegram._bot = _FakeTelegramBot() telegram._bot = bot or _FakeTelegramBot()
telegram._telegram_token = "token" telegram._telegram_token = "token"
telegram._telegram_chat_id = "default-chat" 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_interval_seconds = 0.01
telegram._typing_max_duration_seconds = 1 telegram._typing_max_duration_seconds = 1
return telegram return telegram
def test_start_typing_can_stop_by_chat_id(self):
telegram = self._telegram_client() def test_start_typing_can_stop_by_chat_id():
telegram = _telegram_client()
telegram._start_typing_task( telegram._start_typing_task(
"chat-1", "chat-1",
@@ -78,14 +86,15 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
initial_delay_seconds=0, initial_delay_seconds=0,
) )
self.assertIn("chat-1", Telegram._typing_tasks) assert "chat-1" in telegram._typing_tasks
self.assertTrue(telegram._bot.action_event.wait(1.0)) assert telegram._bot.action_event.wait(1.0)
self.assertTrue(telegram.stop_typing(chat_id="chat-1")) assert telegram.stop_typing(chat_id="chat-1")
self.assertNotIn("chat-1", Telegram._typing_tasks) assert "chat-1" not in telegram._typing_tasks
def test_start_typing_can_stop_by_user_mapping(self):
telegram = self._telegram_client() def test_start_typing_can_stop_by_user_mapping():
Telegram._user_chat_mapping["10001"] = "chat-2" telegram = _telegram_client()
telegram._user_chat_mapping["10001"] = "chat-2"
telegram._start_typing_task( telegram._start_typing_task(
"chat-2", "chat-2",
@@ -94,11 +103,12 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
) )
time.sleep(0.03) time.sleep(0.03)
self.assertTrue(telegram.stop_typing(userid="10001")) assert telegram.stop_typing(userid="10001")
self.assertNotIn("chat-2", Telegram._typing_tasks) assert "chat-2" not in telegram._typing_tasks
def test_typing_task_has_max_duration_guard(self):
telegram = self._telegram_client() def test_typing_task_has_max_duration_guard():
telegram = _telegram_client()
telegram._start_typing_task( telegram._start_typing_task(
"chat-3", "chat-3",
@@ -106,14 +116,15 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
initial_delay_seconds=0, initial_delay_seconds=0,
) )
self.assertTrue(_wait_until(lambda: "chat-3" not in Telegram._typing_tasks)) assert _wait_until(lambda: "chat-3" not in telegram._typing_tasks)
self.assertNotIn("chat-3", Telegram._typing_tasks) assert "chat-3" not in telegram._typing_tasks
def test_short_typing_task_can_stop_before_first_chat_action(self):
def test_short_typing_task_can_stop_before_first_chat_action():
""" """
短响应在首次 typing 发出前结束时不应留下客户端自然过期的残留状态 短响应在首次 typing 发出前结束时不应留下客户端自然过期的残留状态
""" """
telegram = self._telegram_client() telegram = _telegram_client()
telegram._start_typing_task( telegram._start_typing_task(
"chat-4", "chat-4",
@@ -123,11 +134,89 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
telegram.stop_typing(chat_id="chat-4") telegram.stop_typing(chat_id="chat-4")
time.sleep(0.08) time.sleep(0.08)
self.assertEqual(telegram._bot.chat_actions, []) assert telegram._bot.chat_actions == []
self.assertNotIn("chat-4", Telegram._typing_tasks) assert "chat-4" not in telegram._typing_tasks
def test_agent_managed_send_msg_keeps_typing_for_worker_cleanup(self):
telegram = self._telegram_client() 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")) sent = SimpleNamespace(message_id=1, chat=SimpleNamespace(id="chat-1"))
with patch.object( with patch.object(
@@ -139,14 +228,15 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
stop_typing=False, stop_typing=False,
) )
self.assertTrue(result["success"]) assert result["success"]
stop_typing.assert_not_called() stop_typing.assert_not_called()
def test_send_msg_does_not_stop_typing_by_default(self):
def test_send_msg_does_not_stop_typing_by_default():
""" """
响应发送不再默认结束 typing由处理状态统一收口 响应发送不再默认结束 typing由处理状态统一收口
""" """
telegram = self._telegram_client() telegram = _telegram_client()
sent = SimpleNamespace(message_id=1, chat=SimpleNamespace(id="chat-1")) sent = SimpleNamespace(message_id=1, chat=SimpleNamespace(id="chat-1"))
with patch.object( with patch.object(
@@ -154,10 +244,11 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
), patch.object(telegram, "_stop_typing_task") as stop_typing: ), patch.object(telegram, "_stop_typing_task") as stop_typing:
result = telegram.send_msg(title="处理中", userid="10001") result = telegram.send_msg(title="处理中", userid="10001")
self.assertTrue(result["success"]) assert result["success"]
stop_typing.assert_not_called() stop_typing.assert_not_called()
def test_telegram_module_processing_status_starts_typing(self):
def test_telegram_module_processing_status_starts_typing():
""" """
Telegram 通过模块处理状态接口启动 typing 保活 Telegram 通过模块处理状态接口启动 typing 保活
""" """
@@ -178,9 +269,10 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
) )
client.start_typing.assert_called_once_with(chat_id="-100", userid="10001") client.start_typing.assert_called_once_with(chat_id="-100", userid="10001")
self.assertEqual(status["metadata"]["kind"], "typing") assert status["metadata"]["kind"] == "typing"
def test_slash_command_defers_processing_status_to_command_handler(self):
def test_slash_command_defers_processing_status_to_command_handler():
chain = MessageChain.__new__(MessageChain) chain = MessageChain.__new__(MessageChain)
chain.eventmanager = Mock() chain.eventmanager = Mock()
status = MessageChain._ProcessingStatus( status = MessageChain._ProcessingStatus(
@@ -207,12 +299,13 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
finish_status.assert_not_called() finish_status.assert_not_called()
chain.eventmanager.send_event.assert_called_once() chain.eventmanager.send_event.assert_called_once()
self.assertEqual( assert (
chain.eventmanager.send_event.call_args.args[1]["processing_status"], chain.eventmanager.send_event.call_args.args[1]["processing_status"]
status.to_dict(), == status.to_dict()
) )
def test_command_handler_finishes_processing_status_after_execute(self):
def test_command_handler_finishes_processing_status_after_execute():
""" """
传统命令响应完成后由命令处理器统一结束 processing status 传统命令响应完成后由命令处理器统一结束 processing status
""" """
@@ -244,7 +337,8 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
user_id="10001", user_id="10001",
) )
def test_finish_command_processing_status_uses_module_interface(self):
def test_finish_command_processing_status_uses_module_interface():
status = { status = {
"channel": NotificationChannel.Telegram.value, "channel": NotificationChannel.Telegram.value,
"source": "telegram-test", "source": "telegram-test",
@@ -261,7 +355,8 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
userid="fallback", userid="fallback",
) )
def test_async_agent_leaves_processing_status_to_worker(self):
def test_async_agent_leaves_processing_status_to_worker():
chain = MessageChain.__new__(MessageChain) chain = MessageChain.__new__(MessageChain)
chain.eventmanager = Mock() chain.eventmanager = Mock()
chain.runtime_config = replace( chain.runtime_config = replace(
@@ -296,15 +391,16 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
start_status.assert_not_called() start_status.assert_not_called()
finish_status.assert_not_called() finish_status.assert_not_called()
process_message.assert_called_once() process_message.assert_called_once()
self.assertNotIn("processing_status", process_message.call_args.kwargs) assert "processing_status" not in process_message.call_args.kwargs
self.assertEqual( assert (
process_message.call_args.kwargs["channel"], process_message.call_args.kwargs["channel"]
NotificationChannel.Telegram.value, == NotificationChannel.Telegram.value
) )
self.assertEqual(process_message.call_args.kwargs["source"], "telegram-test") assert process_message.call_args.kwargs["source"] == "telegram-test"
self.assertEqual(process_message.call_args.kwargs["original_chat_id"], "-100") assert process_message.call_args.kwargs["original_chat_id"] == "-100"
def test_agent_manager_starts_processing_status_when_task_runs(self):
def test_agent_manager_starts_processing_status_when_task_runs():
async def _run(): async def _run():
manager = AgentManager() manager = AgentManager()
task = _MessageTask( task = _MessageTask(
@@ -331,11 +427,12 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
await manager._start_task_processing_status(task) await manager._start_task_processing_status(task)
start_status.assert_awaited_once_with(task) start_status.assert_awaited_once_with(task)
self.assertEqual(task.processing_status, status) assert task.processing_status == status
asyncio.run(_run()) asyncio.run(_run())
def test_agent_start_processing_status_uses_chain_interface(self):
def test_agent_start_processing_status_uses_chain_interface():
async def _run(): async def _run():
task = _MessageTask( task = _MessageTask(
session_id="session-1", session_id="session-1",
@@ -357,26 +454,30 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
calls = [] calls = []
class FakeAgentChain: class FakeAgentChain:
"""记录 processing status 请求的 Agent Chain 替身。"""
def start_message_processing_status(self, **kwargs): def start_message_processing_status(self, **kwargs):
"""记录请求并返回预设状态。"""
calls.append(kwargs) calls.append(kwargs)
return status return status
with patch("app.agent.orchestrator.AgentChain", FakeAgentChain): with patch("app.agent.orchestrator.AgentChain", FakeAgentChain):
result = await _async_start_processing_status(task) result = await _async_start_processing_status(task)
self.assertEqual(calls, [{ assert calls == [{
"channel": NotificationChannel.Telegram, "channel": NotificationChannel.Telegram,
"source": "telegram-test", "source": "telegram-test",
"userid": "10001", "userid": "10001",
"message_id": "10", "message_id": "10",
"chat_id": "-100", "chat_id": "-100",
"text": "第一条", "text": "第一条",
}]) }]
self.assertEqual(result, status) assert result == status
asyncio.run(_run()) asyncio.run(_run())
def test_callback_stops_typing_when_message_handler_returns(self):
def test_callback_stops_typing_when_message_handler_returns():
chain = MessageChain.__new__(MessageChain) chain = MessageChain.__new__(MessageChain)
chain.eventmanager = Mock() chain.eventmanager = Mock()
status = MessageChain._ProcessingStatus( status = MessageChain._ProcessingStatus(
@@ -410,7 +511,8 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
original_chat_id="-100", original_chat_id="-100",
) )
def test_chain_finishes_processing_through_module_interface(self):
def test_chain_finishes_processing_through_module_interface():
chain = MessageChain.__new__(MessageChain) chain = MessageChain.__new__(MessageChain)
status = MessageChain._ProcessingStatus( status = MessageChain._ProcessingStatus(
channel=NotificationChannel.Telegram, channel=NotificationChannel.Telegram,
@@ -438,7 +540,8 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
chat_id="-100", chat_id="-100",
) )
def test_agent_manager_finishes_processing_status_after_each_task(self):
def test_agent_manager_finishes_processing_status_after_each_task():
async def _run(): async def _run():
manager = AgentManager() manager = AgentManager()
status = { status = {
@@ -462,11 +565,12 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
await manager._finish_task_processing_status(task) await manager._finish_task_processing_status(task)
finish_status.assert_awaited_once_with(status, "10001") finish_status.assert_awaited_once_with(status, "10001")
self.assertIsNone(task.processing_status) assert task.processing_status is None
asyncio.run(_run()) asyncio.run(_run())
def test_agent_worker_starts_and_finishes_each_queued_task(self):
def test_agent_worker_starts_and_finishes_each_queued_task():
async def _run(): async def _run():
manager = AgentManager() manager = AgentManager()
manager._session_queues["session-1"] = asyncio.Queue() manager._session_queues["session-1"] = asyncio.Queue()
@@ -520,14 +624,8 @@ class TestTelegramTypingLifecycle(unittest.TestCase):
manager._session_workers["session-1"].cancel() manager._session_workers["session-1"].cancel()
await manager._session_workers["session-1"] await manager._session_workers["session-1"]
self.assertEqual(start_status.await_count, 2) assert start_status.await_count == 2
self.assertEqual( assert finish_status.await_args_list[0].args == (first_status, "10001")
finish_status.await_args_list[0].args, assert finish_status.await_args_list[1].args == (second_status, "10001")
(first_status, "10001"),
)
self.assertEqual(
finish_status.await_args_list[1].args,
(second_status, "10001"),
)
asyncio.run(_run()) asyncio.run(_run())