From 42dc17ec975a178e7ca47392da968538fed7062d Mon Sep 17 00:00:00 2001 From: "Daniel V. Baecker" <295684751+dvbaecker@users.noreply.github.com> Date: Tue, 11 Aug 2026 14:51:19 +0100 Subject: [PATCH] fix(telegram): close hold-lifecycle gaps on permanent fatal and connected drain MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Address review on #83878: - Permanent fatal fences all hold producers and discards pending maps on teardown instead of re-populating a queue that can never drain. - Any hold created while connected schedules a tracked redispatch (cancel- after-pop no longer orphans until a future reconnect). - Redispatch failures re-hold current + remainder without tight-looping. Regression coverage for the three residual paths, plus the interaction with OOF-156's connect-failure classification: the retryable network path (telegram_connect_error) must NOT clear the hold queue — reconnect is precisely what drains it; only non-retryable fatals discard. --- plugins/platforms/telegram/adapter.py | 240 ++++++++++++++----- tests/gateway/test_telegram_text_batching.py | 136 ++++++++++- 2 files changed, 301 insertions(+), 75 deletions(-) diff --git a/plugins/platforms/telegram/adapter.py b/plugins/platforms/telegram/adapter.py index af3578aa71..7f42e38093 100644 --- a/plugins/platforms/telegram/adapter.py +++ b/plugins/platforms/telegram/adapter.py @@ -922,18 +922,7 @@ class TelegramAdapter(BasePlatformAdapter): super()._mark_connected() # Drain anything held while we were down. PTB will not redeliver — # these events exist only in our hold queue now. - if not getattr(self, "_held_inbound_events", None): - return - try: - loop = asyncio.get_running_loop() - except RuntimeError: - return - prior = getattr(self, "_held_inbound_redispatch_task", None) - # Single tracked task so disconnect can cancel+await it (teknium - # lifecycle rule from #72037 review: no untracked dispatch). - self._held_inbound_redispatch_task = loop.create_task( - self._redispatch_held_inbound(prior=prior) - ) + self._schedule_held_inbound_redispatch() def _mark_disconnected(self) -> None: self._drop_delayed_deliveries = True @@ -942,16 +931,26 @@ class TelegramAdapter(BasePlatformAdapter): 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) - # Permanent fatal: no reconnect will drain the queue. Surface the loss - # explicitly instead of letting held messages die silently with the process. - if not retryable and getattr(self, "_held_inbound_events", None): - n = len(self._held_inbound_events) - logger.warning( - "[Telegram] Non-retryable fatal (%s); discarding %d held inbound message(s)", - code, - n, - ) - self._held_inbound_events.clear() + # Permanent fatal: no reconnect will drain. Discard the hold queue now + # and refuse further holds (teardown salvage / late enqueue must not + # re-populate a queue that can never drain — review #83878). + if not retryable: + held = getattr(self, "_held_inbound_events", None) + n = len(held) if held else 0 + if held: + held.clear() + if n: + logger.warning( + "[Telegram] Non-retryable fatal (%s); discarding %d held inbound message(s)", + code, + n, + ) + + def _is_permanent_fatal(self) -> bool: + """True after non-retryable fatal — holds must discard, not queue.""" + if not getattr(self, "_fatal_error_code", None): + return False + return not bool(getattr(self, "_fatal_error_retryable", True)) def _should_drop_delayed_delivery(self) -> bool: """True once teardown/fatal-error started — delayed flushes must not dispatch. @@ -962,11 +961,54 @@ class TelegramAdapter(BasePlatformAdapter): Callers must NOT destroy the event when this returns True: PTB has already advanced the polling offset, so Telegram will never redeliver. - Use ``_hold_inbound_event`` and redispatch on reconnect. + Use ``_hold_inbound_event`` and redispatch on reconnect (unless + permanent fatal, which discards explicitly). """ return bool(getattr(self, "_drop_delayed_deliveries", False)) - def _hold_inbound_event(self, event: "MessageEvent", *, where: str) -> None: + def _schedule_held_inbound_redispatch(self) -> None: + """Ensure a tracked drain runs when held events exist and delivery is live. + + Drain triggers: + - ``_mark_connected`` after reconnect + - any hold created while already connected (e.g. cancel-after-pop) + - end of a drain pass if more events arrived mid-drain + + No-ops while disconnected/tearing down or after permanent fatal. + """ + if self._is_permanent_fatal(): + return + if self._should_drop_delayed_delivery(): + return + held = getattr(self, "_held_inbound_events", None) + if not held: + return + try: + loop = asyncio.get_running_loop() + except RuntimeError: + return + prior = getattr(self, "_held_inbound_redispatch_task", None) + try: + current = asyncio.current_task() + except RuntimeError: + current = None + # Already draining on another task — that pass schedules a follow-up + # if anything remains. Do not stack duplicate tasks. + if prior is not None and not prior.done() and prior is not current: + return + self._held_inbound_redispatch_task = loop.create_task( + self._redispatch_held_inbound( + prior=None if prior is current else prior + ) + ) + + def _hold_inbound_event( + self, + event: "MessageEvent", + *, + where: str, + schedule: bool = True, + ) -> None: """Preserve an inbound event that cannot be dispatched right now. The disconnect drop-guard (#55971) correctly prevents dispatch into a @@ -974,13 +1016,22 @@ class TelegramAdapter(BasePlatformAdapter): enqueue/flush, python-telegram-bot has already acked the update and advanced the offset — silent permanent loss, no log, no error. - Hold the event and redispatch it from ``_mark_connected``. Cap the - queue so a long outage cannot grow without bound. Dedup by object - identity so salvage-after-hold never double-queues the same event. + Hold the event and redispatch from ``_mark_connected`` (or immediately + if already connected). Cap the queue so a long outage cannot grow + without bound. Dedup by object identity so salvage-after-hold never + double-queues the same event. Permanent fatal discards explicitly. + + ``schedule=False`` when the caller is already inside a drain and will + decide follow-up policy (avoids poison-event tight loops). """ - # Dedup: same MessageEvent object already held (flush held, then - # teardown salvage races the same reference — shouldn't happen after - # pop, but identity guard is free and closes the class of bug). + if self._is_permanent_fatal(): + logger.warning( + "[Telegram] Discarding inbound under non-retryable fatal (%s, %d chars)", + where, + len(getattr(event, "text", None) or ""), + ) + return + held = getattr(self, "_held_inbound_events", None) if held is None: self._held_inbound_events = [] @@ -999,17 +1050,23 @@ class TelegramAdapter(BasePlatformAdapter): ) held.append(event) logger.warning( - "[Telegram] Holding inbound during disconnect (%s, %d chars, queue=%d) " - "- will redispatch on reconnect", + "[Telegram] Holding inbound (%s, %d chars, queue=%d)%s", where, len(getattr(event, "text", None) or ""), len(held), + " - will redispatch on reconnect" + if self._should_drop_delayed_delivery() + else (" - scheduling redispatch" if schedule else ""), ) + # Connected cancel-after-pop (and any other live-path hold) must not + # orphan the event waiting for a future reconnect that may never come. + if schedule and not self._should_drop_delayed_delivery(): + self._schedule_held_inbound_redispatch() async def _redispatch_held_inbound( self, prior: Optional[asyncio.Task] = None ) -> None: - """Drain the hold queue after reconnect. + """Drain the hold queue after reconnect or a connected-path hold. ``prior`` is the previous redispatch task, if any — awaited here so ``_mark_connected`` stays synchronous while teardown can still @@ -1023,37 +1080,79 @@ class TelegramAdapter(BasePlatformAdapter): except asyncio.CancelledError: pass + if self._is_permanent_fatal(): + held = getattr(self, "_held_inbound_events", None) + if held: + n = len(held) + held.clear() + logger.warning( + "[Telegram] Redispatch aborted; discarded %d held inbound under non-retryable fatal", + n, + ) + return + held = getattr(self, "_held_inbound_events", None) if not held: return # Take ownership atomically so a concurrent hold during drain appends - # to a fresh list and is picked up by a later connect. + # to a fresh list and is picked up by a follow-up schedule. events = list(held) held.clear() logger.warning( - "[Telegram] Redispatching %d held inbound message(s) after reconnect", + "[Telegram] Redispatching %d held inbound message(s)", len(events), ) - for idx, event in enumerate(events): - if self._should_drop_delayed_delivery(): - # Disconnect returned mid-drain — re-hold current + remainder. - self._hold_inbound_event(event, where="redispatch-interrupted") - for rest in events[idx + 1 :]: - self._hold_inbound_event(rest, where="redispatch-interrupted") - return - try: - await self.handle_message(event) - except asyncio.CancelledError: - # Task cancelled (disconnect / superseding redispatch) — re-hold. - self._hold_inbound_event(event, where="redispatch-cancelled") - for rest in events[idx + 1 :]: - self._hold_inbound_event(rest, where="redispatch-cancelled") - raise - except Exception: - logger.exception( - "[Telegram] Failed to redispatch held inbound (%d chars)", - len(getattr(event, "text", None) or ""), - ) + allow_followup_schedule = True + try: + for idx, event in enumerate(events): + if self._is_permanent_fatal() or self._should_drop_delayed_delivery(): + # Disconnect/fatal mid-drain — re-hold current + remainder + # (hold itself discards under permanent fatal). + self._hold_inbound_event( + event, where="redispatch-interrupted", schedule=False + ) + for rest in events[idx + 1 :]: + self._hold_inbound_event( + rest, where="redispatch-interrupted", schedule=False + ) + return + try: + await self.handle_message(event) + except asyncio.CancelledError: + self._hold_inbound_event( + event, where="redispatch-cancelled", schedule=False + ) + for rest in events[idx + 1 :]: + self._hold_inbound_event( + rest, where="redispatch-cancelled", schedule=False + ) + raise + except Exception: + # Retryable failure: keep current + remainder. Do not + # immediately reschedule — a poison event would tight-loop. + # Next mark_connected or a later connected-path hold drains. + logger.exception( + "[Telegram] Failed to redispatch held inbound (%d chars); re-holding", + len(getattr(event, "text", None) or ""), + ) + self._hold_inbound_event( + event, where="redispatch-failed", schedule=False + ) + for rest in events[idx + 1 :]: + self._hold_inbound_event( + rest, where="redispatch-failed", schedule=False + ) + allow_followup_schedule = False + return + finally: + # Events arrived mid-drain while still connected need another pass. + if ( + allow_followup_schedule + and getattr(self, "_held_inbound_events", None) + and not self._should_drop_delayed_delivery() + and not self._is_permanent_fatal() + ): + self._schedule_held_inbound_redispatch() def _notification_kwargs( self, metadata: Optional[Dict[str, Any]] @@ -4701,16 +4800,27 @@ class TelegramAdapter(BasePlatformAdapter): if awaitable_tasks: await asyncio.gather(*awaitable_tasks, return_exceptions=True) - # Salvage buffered inbound events before clearing maps. Cancel alone - # leaves the event in the map with no flush task — clearing without - # hold would silently destroy user messages (PTB already advanced the - # offset). Hold for redispatch on reconnect. - for event in list(self._pending_text_batches.values()): - self._hold_inbound_event(event, where="text-batch-teardown") - for event in list(self._pending_photo_batches.values()): - self._hold_inbound_event(event, where="photo-batch-teardown") - for event in list(self._media_group_events.values()): - self._hold_inbound_event(event, where="media-group-teardown") + # Salvage buffered inbound events before clearing maps — unless permanent + # fatal, where no reconnect can drain and hold would re-orphan them + # (#83878). Discard pending sources explicitly in that case. + if self._is_permanent_fatal(): + n_pending = ( + len(self._pending_text_batches) + + len(self._pending_photo_batches) + + len(self._media_group_events) + ) + if n_pending: + logger.warning( + "[Telegram] Non-retryable fatal teardown; discarding %d pending inbound batch(es)", + n_pending, + ) + else: + for event in list(self._pending_text_batches.values()): + self._hold_inbound_event(event, where="text-batch-teardown") + for event in list(self._pending_photo_batches.values()): + self._hold_inbound_event(event, where="photo-batch-teardown") + for event in list(self._media_group_events.values()): + self._hold_inbound_event(event, where="media-group-teardown") self._media_group_tasks.clear() self._media_group_events.clear() diff --git a/tests/gateway/test_telegram_text_batching.py b/tests/gateway/test_telegram_text_batching.py index 52c7736ed5..e3e07090a6 100644 --- a/tests/gateway/test_telegram_text_batching.py +++ b/tests/gateway/test_telegram_text_batching.py @@ -238,13 +238,16 @@ class TestHoldInboundAcrossReconnect: """Cancel after pop (before handle_message returns) must hold, not lose. Uses entered/release Events — no sleep timing (teknium #72037 rule). + Connected path then schedules redispatch (#83878). """ adapter = _make_adapter() self._zero_batch_delays(adapter) entered = asyncio.Event() release = asyncio.Event() + seen: list[str] = [] async def _blocking_handle(event): + seen.append(event.text or "") entered.set() await release.wait() @@ -257,15 +260,16 @@ class TestHoldInboundAcrossReconnect: task.cancel() with pytest.raises(asyncio.CancelledError): await task - release.set() # unblock if anything still waiting + release.set() - # handle_message was entered but cancel + CancelledError hold path - # must leave the event recoverable. After cancel during handle, our - # CancelledError handler only holds if event is not None — we set - # event=None only after successful handle. Cancel during handle means - # event is still set → held. + drain = adapter._held_inbound_redispatch_task + assert drain is not None + await asyncio.wait_for(drain, timeout=1.0) + + # Recoverable: held and/or delivered via redispatch (seen may include + # the original in-flight attempt plus the redispatch). held_texts = [e.text for e in adapter._held_inbound_events] - assert "in-flight cancel" in held_texts + assert "in-flight cancel" in seen or "in-flight cancel" in held_texts @pytest.mark.asyncio async def test_cancel_pending_salvages_batches_into_held_queue(self): @@ -386,15 +390,58 @@ class TestHoldInboundAcrossReconnect: async def test_non_retryable_fatal_discards_held_with_warning(self): adapter = _make_adapter() adapter._held_inbound_events = [_make_event("doomed")] - # BasePlatformAdapter._set_fatal_error may need attrs — call ours via - # the override path with a stub super if needed. from gateway.platforms.base import BasePlatformAdapter - with patch.object(BasePlatformAdapter, "_set_fatal_error", lambda *a, **k: None): + def _base_fatal(self, code, message, *, retryable): + self._fatal_error_code = code + self._fatal_error_message = message + self._fatal_error_retryable = retryable + self._running = False + + with patch.object(BasePlatformAdapter, "_set_fatal_error", _base_fatal): adapter._set_fatal_error("auth", "revoked", retryable=False) assert adapter._held_inbound_events == [] assert adapter._drop_delayed_deliveries is True + assert adapter._is_permanent_fatal() is True + + @pytest.mark.asyncio + async def test_retryable_fatal_preserves_held_for_reconnect_drain(self): + """Retryable fatals must NOT clear the hold queue. + + OOF-156's connect-failure classification keeps the common network + path ``retryable=True`` (``telegram_connect_error``) — reconnect is + precisely what must drain a hold queue populated during the outage. + Only non-retryable fatals may discard (covered above). + """ + adapter = _make_adapter() + adapter._held_inbound_events = [_make_event("survives-network-fatal")] + adapter._drop_delayed_deliveries = True # fatal/disconnect already set + + from gateway.platforms.base import BasePlatformAdapter + + def _base_fatal(self, code, message, *, retryable): + self._fatal_error_code = code + self._fatal_error_message = message + self._fatal_error_retryable = retryable + + with patch.object(BasePlatformAdapter, "_set_fatal_error", _base_fatal): + adapter._set_fatal_error( + "telegram_connect_error", "connect timed out", retryable=True + ) + + assert [e.text for e in adapter._held_inbound_events] == [ + "survives-network-fatal" + ] + assert adapter._is_permanent_fatal() is False + + # Reconnect drains what the retryable fatal preserved. + adapter._mark_connected() + await adapter._held_inbound_redispatch_task + adapter.handle_message.assert_called_once() + assert ( + adapter.handle_message.call_args[0][0].text == "survives-network-fatal" + ) @pytest.mark.asyncio async def test_production_text_handler_terminal_step_holds_when_disconnected(self): @@ -414,3 +461,72 @@ class TestHoldInboundAcrossReconnect: await adapter._held_inbound_redispatch_task adapter.handle_message.assert_called_once() assert adapter.handle_message.call_args[0][0].text == "acked-by-ptb-then-held" + + @pytest.mark.asyncio + async def test_permanent_fatal_teardown_discards_pending_not_rehold(self): + """#83878: permanent fatal must not re-populate hold via teardown salvage.""" + adapter = _make_adapter() + adapter._fatal_error_code = "auth" + adapter._fatal_error_retryable = False + adapter._drop_delayed_deliveries = True + adapter._pending_text_batches["t"] = _make_event("pending-text") + adapter._pending_photo_batches["p"] = _make_event("pending-photo") + adapter._media_group_events["m"] = _make_event("pending-media") + + await adapter._cancel_pending_delivery_tasks() + + assert adapter._held_inbound_events == [] + assert adapter._pending_text_batches == {} + assert adapter._pending_photo_batches == {} + assert adapter._media_group_events == {} + + @pytest.mark.asyncio + async def test_permanent_fatal_late_enqueue_discards(self): + """#83878: late enqueue after permanent fatal must discard, not hold.""" + adapter = _make_adapter() + adapter._fatal_error_code = "auth" + adapter._fatal_error_retryable = False + adapter._drop_delayed_deliveries = True + + adapter._enqueue_text_event(_make_event("too-late")) + adapter.handle_message.assert_not_called() + assert adapter._held_inbound_events == [] + + @pytest.mark.asyncio + async def test_connected_hold_schedules_redispatch(self): + """#83878: hold while connected must drain, not orphan until reconnect.""" + adapter = _make_adapter() + adapter._drop_delayed_deliveries = False + adapter.handle_message = AsyncMock() + + adapter._hold_inbound_event( + _make_event("orphan-without-drain"), where="text-flush-cancelled" + ) + + drain = adapter._held_inbound_redispatch_task + assert drain is not None + await asyncio.wait_for(drain, timeout=1.0) + adapter.handle_message.assert_called_once() + assert adapter.handle_message.call_args[0][0].text == "orphan-without-drain" + assert adapter._held_inbound_events == [] + + @pytest.mark.asyncio + async def test_redispatch_exception_reholds_current_and_remainder(self): + """#83878: handle_message failure must not drop current/remainder.""" + adapter = _make_adapter() + adapter._drop_delayed_deliveries = False + adapter._held_inbound_events = [ + _make_event("boom"), + _make_event("after"), + ] + + async def _handle(event): + if event.text == "boom": + raise RuntimeError("dispatch failed") + return None + + adapter.handle_message = _handle + # Direct drain (no auto follow-up on failure) + await adapter._redispatch_held_inbound() + held_texts = [e.text for e in adapter._held_inbound_events] + assert held_texts == ["boom", "after"]