diff --git a/astrbot/core/platform/sources/telegram/tg_adapter.py b/astrbot/core/platform/sources/telegram/tg_adapter.py index c9f7a303d9..87096547a2 100644 --- a/astrbot/core/platform/sources/telegram/tg_adapter.py +++ b/astrbot/core/platform/sources/telegram/tg_adapter.py @@ -209,6 +209,7 @@ def __init__( self._polling_restart_delay = delay self._polling_recovery_threshold = 3 self._polling_failure_window = 60.0 + self._drop_pending_updates = True self._application_started = False self._seen_update_ids: OrderedDict[int, None] = OrderedDict() self._callback_bindings: OrderedDict[str, _TelegramCallbackBinding] = ( @@ -837,11 +838,17 @@ async def run(self) -> None: self._application_started = False await asyncio.sleep(self._polling_restart_delay) continue - logger.info("Starting Telegram polling...") + drop_pending_updates = self._drop_pending_updates + logger.info( + "Starting Telegram polling%s...", + " (dropping pending updates)" if drop_pending_updates else "", + ) await updater.start_polling( allowed_updates=TELEGRAM_ALLOWED_UPDATES, + drop_pending_updates=drop_pending_updates, error_callback=self._on_polling_error, ) + self._drop_pending_updates = False logger.info("Telegram Platform Adapter is running.") while updater.running and not self._terminating: # noqa: ASYNC110 if self._polling_recovery_requested.is_set(): diff --git a/docs/en/platform/telegram.md b/docs/en/platform/telegram.md index 7587ddbd73..b6c412a478 100644 --- a/docs/en/platform/telegram.md +++ b/docs/en/platform/telegram.md @@ -63,6 +63,8 @@ Whenever polling starts or its client is rebuilt, AstrBot explicitly subscribes The handler accepts only the three supported message variants even if unacknowledged updates from an older subscription still arrive. Each adapter retains the latest 4096 admitted Update IDs, so repeated delivery does not execute a command or agent again while its ID remains cached. The cache survives polling-client rebuilds but is not persisted across process restarts. Edited updates are always ignored. +On an adapter instance's first polling start, AstrBot drops unacknowledged pending updates on Telegram's servers so offline history is not flushed into the pipeline and the per-session rate limiter. After an in-process client rebuild, later polling keeps updates from that gap. Pending updates after a process restart are still dropped and are not replayed. + ### Business Sessions and Replies A Business chat is independent of an ordinary Bot chat with the same chat ID. AstrBot uses `business::` as its Business route, appending `#` for topics. Save the complete session target for proactive sends: the connection, chat, and topic are restored, and albums are isolated by that route too. diff --git a/docs/zh/platform/telegram.md b/docs/zh/platform/telegram.md index c64b5106d3..b7660d60f3 100644 --- a/docs/zh/platform/telegram.md +++ b/docs/zh/platform/telegram.md @@ -61,6 +61,8 @@ async def ask(self, event: AstrMessageEvent): 即使旧订阅中尚未确认的更新在切换后到达,处理器也只接收上述三类消息。每个适配器保留最近 4096 个已接纳的 Update ID,重复投递不会在缓存有效期内再次执行指令或 agent;缓存随轮询客户端重建保留,但不跨进程重启持久化。编辑更新始终忽略。 +适配器实例首次启动轮询时,会丢弃 Telegram 服务器上尚未确认的积压更新,避免离线期间的历史消息在接入后瞬间灌入 pipeline 并触发会话限流。同一进程内因网络错误重建客户端后再轮询时,会保留这段缺口中的更新。进程重启后的积压仍会被丢弃,不会补处理。 + ### Business 会话与回复 Business 聊天与相同 chat ID 的普通 Bot 聊天相互独立。AstrBot 使用 `business:<经过百分号编码的连接 ID>:` 作为 Business 路由,有主题时追加 `#`。请保存完整会话目标用于主动发送;连接、聊天和主题信息都会恢复,相册也按该路由隔离。 diff --git a/tests/unit/platform/test_telegram_adapter.py b/tests/unit/platform/test_telegram_adapter.py index 9006d226d3..b9b3ccbda1 100644 --- a/tests/unit/platform/test_telegram_adapter.py +++ b/tests/unit/platform/test_telegram_adapter.py @@ -2873,9 +2873,12 @@ async def test_telegram_run_rebuilds_application_after_repeated_polling_errors() builder.build.side_effect = created_apps adapter = None + first_poll_kwargs: dict[str, object] = {} + second_poll_kwargs: dict[str, object] = {} def start_polling_side_effect(*args, **kwargs): nonlocal adapter + first_poll_kwargs.update(kwargs) error_callback = kwargs["error_callback"] assert adapter is not None @@ -2890,6 +2893,7 @@ async def _emit_errors(): app_one.updater.start_polling.side_effect = start_polling_side_effect async def second_start_polling(*args, **kwargs): + second_poll_kwargs.update(kwargs) assert adapter is not None adapter._terminating = True @@ -2913,6 +2917,8 @@ async def second_start_polling(*args, **kwargs): await adapter.run() assert builder.build.call_count == 2 + assert first_poll_kwargs["drop_pending_updates"] is True + assert second_poll_kwargs["drop_pending_updates"] is False app_one.updater.stop.assert_awaited() app_one.bot.delete_my_commands.assert_awaited_once() app_one.stop.assert_awaited() @@ -2955,9 +2961,12 @@ async def test_telegram_run_rebuilds_fresh_application_after_recreate_init_failu builder.build.side_effect = created_apps adapter = None + first_poll_kwargs: dict[str, object] = {} + final_poll_kwargs: dict[str, object] = {} def first_start_polling(*args, **kwargs): nonlocal adapter + first_poll_kwargs.update(kwargs) error_callback = kwargs["error_callback"] assert adapter is not None @@ -2973,6 +2982,7 @@ async def _emit_errors(): app_two.initialize.side_effect = TimeoutError("init timeout") async def final_start_polling(*args, **kwargs): + final_poll_kwargs.update(kwargs) assert adapter is not None adapter._terminating = True @@ -2999,6 +3009,8 @@ async def final_start_polling(*args, **kwargs): await adapter.run() assert builder.build.call_count == 3 + assert first_poll_kwargs["drop_pending_updates"] is True + assert final_poll_kwargs["drop_pending_updates"] is False app_two.stop.assert_awaited() app_two.shutdown.assert_awaited() app_three.initialize.assert_awaited() @@ -3902,6 +3914,7 @@ async def stop_after_start(**kwargs): "business_message", "callback_query", ), + drop_pending_updates=True, error_callback=adapter._on_polling_error, )