fix(telegram): close hold-lifecycle gaps on permanent fatal and connected drain

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.
This commit is contained in:
Daniel V. Baecker
2026-08-11 14:51:19 +01:00
committed by kshitij
parent b4da6b15e1
commit 42dc17ec97
2 changed files with 301 additions and 75 deletions
+175 -65
View File
@@ -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()
+126 -10
View File
@@ -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"]