From 9bdc016409145a62412b9205a5570d9fccd27020 Mon Sep 17 00:00:00 2001 From: qbit-mirror-bot Date: Wed, 1 Jul 2026 02:29:31 +0000 Subject: [PATCH] fix(telegram): recover when polling updater stops while process stays alive (#55769) --- plugins/platforms/telegram/adapter.py | 116 ++++---------------------- 1 file changed, 16 insertions(+), 100 deletions(-) diff --git a/plugins/platforms/telegram/adapter.py b/plugins/platforms/telegram/adapter.py index e816bbcbf2cc..e4d0aba2d299 100644 --- a/plugins/platforms/telegram/adapter.py +++ b/plugins/platforms/telegram/adapter.py @@ -412,7 +412,6 @@ def __init__(self, config: PlatformConfig): ) self._pending_text_batches: Dict[str, MessageEvent] = {} self._pending_text_batch_tasks: Dict[str, asyncio.Task] = {} - self._drop_delayed_deliveries = False self._polling_error_task: Optional[asyncio.Task] = None self._polling_conflict_count: int = 0 self._polling_network_error_count: int = 0 @@ -500,27 +499,6 @@ def __init__(self, config: PlatformConfig): # same key edit the same message instead of appending new ones (#30045). self._status_message_ids: Dict[tuple, str] = {} - def _mark_connected(self) -> None: - self._drop_delayed_deliveries = False - super()._mark_connected() - - def _mark_disconnected(self) -> None: - self._drop_delayed_deliveries = True - super()._mark_disconnected() - - def _set_fatal_error(self, code: str, message: str, *, retryable: bool) -> None: - self._drop_delayed_deliveries = True - super()._set_fatal_error(code, message, retryable=retryable) - - def _should_drop_delayed_delivery(self) -> bool: - """True once teardown/fatal-error started — delayed flushes must drop. - - Buffered text/photo/media-group flushes sit behind an asyncio.sleep(). - If disconnect wins the race, dispatching them spawns an agent on a - torn-down session, producing stale/duplicate deliveries. - """ - return bool(getattr(self, "_drop_delayed_deliveries", False)) - def _notification_kwargs( self, metadata: Optional[Dict[str, Any]] ) -> Dict[str, Any]: @@ -2937,60 +2915,8 @@ async def _set_status_indicator(self, online: bool) -> None: self.name, text, e, ) - async def _cancel_pending_delivery_tasks(self) -> None: - """Cancel every delayed-delivery task family before disconnect completes. - - Covers media-group, photo-batch and text-batch flush tasks plus the - polling-error recovery task. Each sits behind an ``asyncio.sleep()``; - if teardown leaves them running they dispatch ``handle_message`` into a - torn-down session. Skips the current task so the coroutine driving - teardown does not cancel itself. - """ - current_task = asyncio.current_task() - pending_tasks: list[asyncio.Task] = [] - awaitable_tasks: list[asyncio.Task] = [] - seen: set[int] = set() - - def collect(task: Optional[asyncio.Task]) -> None: - if not task or task.done() or task is current_task: - return - marker = id(task) - if marker in seen: - return - seen.add(marker) - pending_tasks.append(task) - if asyncio.isfuture(task) or asyncio.iscoroutine(task): - awaitable_tasks.append(task) - - for task in list(self._media_group_tasks.values()): - collect(task) - for task in list(self._pending_photo_batch_tasks.values()): - collect(task) - for task in list(self._pending_text_batch_tasks.values()): - collect(task) - collect(self._polling_error_task) - - for task in pending_tasks: - task.cancel() - if awaitable_tasks: - await asyncio.gather(*awaitable_tasks, return_exceptions=True) - - self._media_group_tasks.clear() - self._media_group_events.clear() - self._pending_photo_batch_tasks.clear() - self._pending_photo_batches.clear() - self._pending_text_batch_tasks.clear() - self._pending_text_batches.clear() - if self._polling_error_task is not current_task: - self._polling_error_task = None - async def disconnect(self) -> None: - """Stop polling/webhook, cancel pending delayed deliveries, and disconnect.""" - # Mark disconnected first so the drop guard short-circuits any flush - # that wins the race against teardown and prevents new delayed tasks - # from being scheduled by late update handlers. - self._mark_disconnected() - + """Stop polling/webhook, cancel pending album flushes, and disconnect.""" # Cancel the heartbeat before tearing down the app so the probe task # cannot fire get_me() into a half-shutdown bot client. if self._polling_heartbeat_task and not self._polling_heartbeat_task.done(): @@ -3011,7 +2937,13 @@ async def disconnect(self) -> None: except Exception: pass - await self._cancel_pending_delivery_tasks() + pending_media_group_tasks = list(self._media_group_tasks.values()) + for task in pending_media_group_tasks: + task.cancel() + if pending_media_group_tasks: + await asyncio.gather(*pending_media_group_tasks, return_exceptions=True) + self._media_group_tasks.clear() + self._media_group_events.clear() if self._app: try: @@ -3025,6 +2957,13 @@ async def disconnect(self) -> None: logger.warning("[%s] Error during Telegram disconnect: %s", self.name, e, exc_info=True) self._release_platform_lock() + for task in self._pending_photo_batch_tasks.values(): + if task and not task.done(): + task.cancel() + self._pending_photo_batch_tasks.clear() + self._pending_photo_batches.clear() + + self._mark_disconnected() self._app = None self._bot = None logger.info("[%s] Disconnected from Telegram", self.name) @@ -6935,10 +6874,6 @@ def _enqueue_text_event(self, event: MessageEvent) -> None: concatenates them and waits for a short quiet period before dispatching the combined message. """ - if self._should_drop_delayed_delivery(): - logger.debug("[Telegram] Dropping text batch enqueue after disconnect started") - return - key = self._text_batch_key(event) existing = self._pending_text_batches.get(key) chunk_len = len(event.text or "") @@ -6998,9 +6933,6 @@ async def _flush_text_batch(self, key: str) -> None: event = self._pending_text_batches.pop(key, None) if not event: return - if self._should_drop_delayed_delivery(): - logger.debug("[Telegram] Dropping text batch flush after disconnect started") - return logger.info( "[Telegram] Flushing text batch %s (%d chars)", key, len(event.text or ""), @@ -7035,9 +6967,6 @@ async def _flush_photo_batch(self, batch_key: str) -> None: event = self._pending_photo_batches.pop(batch_key, None) if not event: return - if self._should_drop_delayed_delivery(): - logger.debug("[Telegram] Dropping photo batch flush after disconnect started") - return logger.info("[Telegram] Flushing photo batch %s with %d image(s)", batch_key, len(event.media_urls)) await self.handle_message(event) finally: @@ -7046,10 +6975,6 @@ async def _flush_photo_batch(self, batch_key: str) -> None: def _enqueue_photo_event(self, batch_key: str, event: MessageEvent) -> None: """Merge photo events into a pending batch and schedule flush.""" - if self._should_drop_delayed_delivery(): - logger.debug("[Telegram] Dropping photo batch enqueue after disconnect started") - return - existing = self._pending_photo_batches.get(batch_key) if existing is None: self._pending_photo_batches[batch_key] = event @@ -7354,10 +7279,6 @@ async def _queue_media_group_event(self, media_group_id: str, event: MessageEven new user message and interrupts the first. We debounce briefly and merge the attachments into a single MessageEvent. """ - if self._should_drop_delayed_delivery(): - logger.debug("[Telegram] Dropping media group enqueue after disconnect started") - return - existing = self._media_group_events.get(media_group_id) if existing is None: self._media_group_events[media_group_id] = event @@ -7376,20 +7297,15 @@ async def _queue_media_group_event(self, media_group_id: str, event: MessageEven ) async def _flush_media_group_event(self, media_group_id: str) -> None: - current_task = asyncio.current_task() try: await asyncio.sleep(self.MEDIA_GROUP_WAIT_SECONDS) event = self._media_group_events.pop(media_group_id, None) if event is not None: - if self._should_drop_delayed_delivery(): - logger.debug("[Telegram] Dropping media group flush after disconnect started") - return await self.handle_message(event) except asyncio.CancelledError: return finally: - if self._media_group_tasks.get(media_group_id) is current_task: - self._media_group_tasks.pop(media_group_id, None) + self._media_group_tasks.pop(media_group_id, None) async def _handle_sticker(self, msg: Message, event: "MessageEvent") -> None: """