Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
103 changes: 102 additions & 1 deletion controlmesh/messenger/telegram/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -160,6 +160,33 @@ def _log_text_preview(text: str | None, *, limit: int = 120) -> str:
_WELCOME_IMAGE = Path(__file__).resolve().parent / "controlmesh_images" / "welcome.png"
_CAPTION_LIMIT = 1024
_INBOUND_DRAIN_OWNER = "telegram_frontstage"
# Hard ceiling for keeping a single inbound claim alive. A frontstage run that
# hangs forever (e.g. a provider call that never returns) must not renew the
# lease indefinitely, otherwise recover_stale_claims() can never reclaim the
# lane and that chat stays blocked. When the cap is hit the renewal loop stops,
# the lease expires within claim_ttl, and the backlog recovery reclaims it.
_MAX_CLAIM_LIFETIME_SECONDS = 1800.0


def _poll_restart_backoff_seconds(
*,
consecutive_failures: int,
retry_after_seconds: float,
) -> float:
"""Seconds to wait before re-issuing getUpdates after a poll failure.

``retry_after`` (Telegram 429) wins when Telegram explicitly asked us to
wait; otherwise we back off exponentially with the consecutive failure
count. Returning 0 means "no wait". The cap keeps a long outage from
parking the bot, while still damping the reconnect storm that itself
triggers flood control.
"""
retry_after = max(0.0, float(retry_after_seconds or 0.0))
consecutive = max(0, int(consecutive_failures))
if retry_after <= 0 and consecutive <= 0:
return 0.0
exponential = min(60.0, 0.5 * (2 ** min(consecutive, 7)))
return max(retry_after, exponential)
_CONTROL_COMMAND_PREFIXES = (
"/model",
"/provider",
Expand Down Expand Up @@ -203,6 +230,7 @@ class _TelegramPollDiagnostics:
transport_dirty: bool = False
restart_reason: str | None = None
last_failure_reason: str | None = None
last_retry_after_seconds: float = 0.0
restart_requested: bool = False

def note_poll_started(self, *, offset: int | None) -> None:
Expand All @@ -220,13 +248,22 @@ def note_poll_succeeded(self, *, offset: int | None, update_ids: list[int]) -> N
self.transport_dirty = False
self.restart_reason = None
self.last_failure_reason = None
self.last_retry_after_seconds = 0.0
self.restart_requested = False

def note_poll_failed(self, *, reason: str, offset: int | None, mark_transport_dirty: bool) -> None:
def note_poll_failed(
self,
*,
reason: str,
offset: int | None,
mark_transport_dirty: bool,
retry_after_seconds: float = 0.0,
) -> None:
self.last_poll_finished_at = time.monotonic()
self.last_poll_offset = offset
self.consecutive_failures += 1
self.last_failure_reason = reason
self.last_retry_after_seconds = max(0.0, float(retry_after_seconds or 0.0))
if mark_transport_dirty:
self.transport_dirty = True
self.restart_reason = reason
Expand Down Expand Up @@ -2224,9 +2261,32 @@ async def _watch_restart_marker(self) -> None:

async def _watch_poll_health(self) -> None:
"""Request a fresh Telegram polling transport when getUpdates stalls."""
next_backlog_log = 0.0
try:
while True:
await asyncio.sleep(1.0)
snapshot = self._last_inbound_spool_stats
if (
self._inbound_spool is not None
and snapshot is not None
and (snapshot.pending_count or snapshot.blocked_lane_count)
and time.monotonic() >= next_backlog_log
):
live = self._inbound_spool.stats()
if live.pending_count or live.blocked_lane_count:
oldest_age = (
f"{live.oldest_pending_age_seconds:.0f}"
if live.oldest_pending_age_seconds is not None
else "n/a"
)
logger.info(
"Telegram inbound backlog pending=%s blocked_lanes=%s oldest_age=%ss unhealthy=%s",
live.pending_count,
live.blocked_lane_count,
oldest_age,
live.unhealthy_reason or "no",
)
next_backlog_log = time.monotonic() + 60.0
if self._exit_code == EXIT_RESTART or self._poll_diagnostics.transport_dirty:
continue
inflight_age = self._poll_diagnostics.poll_inflight_age_seconds()
Expand Down Expand Up @@ -2287,9 +2347,33 @@ async def run(self) -> int:
)
if self._exit_code == EXIT_RESTART or not self._poll_diagnostics.transport_dirty:
break
backoff = self._compute_poll_restart_backoff_seconds()
if backoff > 0:
logger.info(
"Telegram backing off %.1fs before poll restart reason=%s failures=%s retry_after=%.1fs",
backoff,
self._poll_diagnostics.restart_reason or "dirty_transport",
self._poll_diagnostics.consecutive_failures,
float(self._poll_diagnostics.last_retry_after_seconds or 0.0),
)
await asyncio.sleep(backoff)
await self._rebuild_poll_transport()
return self._exit_code

def _compute_poll_restart_backoff_seconds(self) -> float:
"""Honor Telegram 429 retry_after and add exponential backoff on failures.

Without this, a burst of transient network errors makes controlmesh
rebuild the transport and immediately re-issue getUpdates, which trips
Telegram's getUpdates flood control (429) and locks inbound delivery
out for an escalating window -- every chat stops receiving replies
while the process stays alive.
"""
return _poll_restart_backoff_seconds(
consecutive_failures=self._poll_diagnostics.consecutive_failures,
retry_after_seconds=self._poll_diagnostics.last_retry_after_seconds,
)

async def shutdown(self) -> None:
await _cancel_task(self._restart_watcher)
await _cancel_task(self._poll_watchdog)
Expand Down Expand Up @@ -2358,10 +2442,17 @@ async def _note_poll_failed(self, method: GetUpdates, exc: Exception) -> None:
reason = self._recoverable_poll_reason(exc)
if reason is None:
return
retry_after = 0.0
if isinstance(exc, TelegramRetryAfter):
try:
retry_after = float(getattr(exc, "retry_after", 0) or 0)
except (TypeError, ValueError):
retry_after = 0.0
self._poll_diagnostics.note_poll_failed(
reason=reason,
offset=method.offset,
mark_transport_dirty=True,
retry_after_seconds=retry_after,
)
logger.warning(
"Telegram poll marked transport dirty reason=%s offset=%s last_success_age=%s failures=%s",
Expand Down Expand Up @@ -2619,7 +2710,17 @@ async def _keep_inbound_claim_alive(self, claim: TelegramInboundClaim) -> None:
return
current_claim = claim
interval = max(1.0, self._inbound_spool.claim_ttl_seconds / 3.0)
deadline = time.monotonic() + _MAX_CLAIM_LIFETIME_SECONDS
while True:
if time.monotonic() >= deadline:
logger.warning(
"Telegram inbound claim exceeded max lifetime; letting lease expire "
"so the backlog can be reclaimed lane=%s chat_id=%s spool_id=%s",
claim.lane_key,
claim.entry.chat_id,
claim.spool_id,
)
return
await asyncio.sleep(interval)
renewed = self._inbound_spool.renew(current_claim)
if renewed is None:
Expand Down
89 changes: 89 additions & 0 deletions tests/messenger/telegram/test_poll_backoff.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
"""Tests for Telegram poll-restart backoff and inbound claim diagnostics.

Regression coverage for the fleet-wide "active but unresponsive" incident:
a burst of transient getUpdates errors caused an immediate transport rebuild
+ re-poll, which tripped Telegram's getUpdates flood control (429) and locked
inbound delivery out. The fix honors ``retry_after`` and backs off.
"""

from __future__ import annotations

import pytest

from controlmesh.messenger.telegram.app import (
_TelegramPollDiagnostics,
_poll_restart_backoff_seconds,
)


class TestPollRestartBackoff:
def test_no_failures_means_no_wait(self) -> None:
assert _poll_restart_backoff_seconds(consecutive_failures=0, retry_after_seconds=0.0) == 0.0

def test_exponential_growth_then_cap(self) -> None:
assert _poll_restart_backoff_seconds(consecutive_failures=1, retry_after_seconds=0.0) == pytest.approx(1.0)
assert _poll_restart_backoff_seconds(consecutive_failures=3, retry_after_seconds=0.0) == pytest.approx(4.0)
assert _poll_restart_backoff_seconds(consecutive_failures=5, retry_after_seconds=0.0) == pytest.approx(16.0)
# 0.5 * 2**7 = 64 -> clamped to the 60s cap
assert _poll_restart_backoff_seconds(consecutive_failures=7, retry_after_seconds=0.0) == 60.0
# cap holds for larger counts
assert _poll_restart_backoff_seconds(consecutive_failures=20, retry_after_seconds=0.0) == 60.0

def test_retry_after_wins_when_larger(self) -> None:
assert _poll_restart_backoff_seconds(consecutive_failures=1, retry_after_seconds=12.0) == 12.0

def test_exponential_wins_when_larger(self) -> None:
assert _poll_restart_backoff_seconds(consecutive_failures=5, retry_after_seconds=1.0) == pytest.approx(16.0)

def test_negative_inputs_clamped_to_no_wait(self) -> None:
assert _poll_restart_backoff_seconds(consecutive_failures=-3, retry_after_seconds=-2.0) == 0.0

def test_retry_after_zero_with_failures_still_backs_off(self) -> None:
assert _poll_restart_backoff_seconds(consecutive_failures=2, retry_after_seconds=0.0) == pytest.approx(2.0)


class TestPollDiagnosticsRetryAfter:
def test_failed_records_retry_after(self) -> None:
d = _TelegramPollDiagnostics()
d.note_poll_failed(
reason="recoverable_http_429",
offset=42,
mark_transport_dirty=True,
retry_after_seconds=7.0,
)
assert d.last_retry_after_seconds == 7.0
assert d.consecutive_failures == 1
assert d.transport_dirty is True
assert d.restart_reason == "recoverable_http_429"

def test_success_resets_retry_after(self) -> None:
d = _TelegramPollDiagnostics()
d.note_poll_failed(
reason="recoverable_http_429",
offset=1,
mark_transport_dirty=True,
retry_after_seconds=9.0,
)
d.note_poll_succeeded(offset=2, update_ids=[])
assert d.last_retry_after_seconds == 0.0
assert d.consecutive_failures == 0
assert d.transport_dirty is False

def test_failed_defaults_retry_after_to_zero(self) -> None:
d = _TelegramPollDiagnostics()
d.note_poll_failed(
reason="recoverable_network_error",
offset=1,
mark_transport_dirty=True,
)
assert d.last_retry_after_seconds == 0.0

def test_non_positive_retry_after_is_clamped(self) -> None:
d = _TelegramPollDiagnostics()
d.note_poll_failed(
reason="recoverable_http_429",
offset=1,
mark_transport_dirty=True,
retry_after_seconds=-5.0,
)
assert d.last_retry_after_seconds == 0.0
Loading