diff --git a/gateway/run.py b/gateway/run.py index 4107f53b03..b8ffe3eeff 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -85,6 +85,13 @@ _PLATFORM_CONNECT_TIMEOUT_SECS_DEFAULT = 30.0 # returns. Leave enough outer budget for initialize/deleteWebhook/start_polling # wall deadlines plus readiness; other platforms retain the 30s isolation bound. _TELEGRAM_CONNECT_TIMEOUT_SECS_DEFAULT = 180.0 +# Cold-start cap for Telegram (#85993): the initial connect awaited before the +# gateway reaches `running` must not spend the full 180s budget — an +# unreachable Telegram would hold EVERY platform's serving state hostage for +# the whole window. The initial attempt gets one bounded try; on timeout the +# platform is queued for the reconnect watcher, which retries with the full +# 180s budget (is_reconnect=True preserves the offline update queue, #46621). +_TELEGRAM_INITIAL_CONNECT_TIMEOUT_SECS_DEFAULT = 45.0 _ADAPTER_DISCONNECT_TIMEOUT_SECS_DEFAULT = 5.0 # End reasons that mean the USER deliberately closed this thread of work # (/new -> session_reset / new_session, an explicit exit, or a /switch). @@ -7194,8 +7201,18 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew return max(0.0, timeout) return _ADAPTER_DISCONNECT_TIMEOUT_SECS_DEFAULT - def _platform_connect_timeout_secs(self, platform=None) -> float: - """Return the per-platform connect timeout used during startup/retry.""" + def _platform_connect_timeout_secs(self, platform=None, *, initial: bool = False) -> float: + """Return the per-platform connect timeout used during startup/retry. + + ``initial=True`` marks the cold-start connect awaited before the + gateway reaches ``running``. Telegram's full connect budget (180s, + raised for #67498 so cold polling can prove getUpdates readiness) is + deliberately NOT spent there: an unreachable Telegram would hold the + whole gateway out of the ``running`` state for the full budget + (#85993). The cold-start wait is capped and the platform is handed to + the reconnect watcher, which retries with the full budget (and + ``is_reconnect=True``, preserving the offline update queue — #46621). + """ raw = os.getenv("HERMES_GATEWAY_PLATFORM_CONNECT_TIMEOUT", "").strip() if raw: try: @@ -7208,11 +7225,13 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew else: return max(0.0, timeout) if platform == Platform.TELEGRAM: + if initial: + return _TELEGRAM_INITIAL_CONNECT_TIMEOUT_SECS_DEFAULT return _TELEGRAM_CONNECT_TIMEOUT_SECS_DEFAULT return _PLATFORM_CONNECT_TIMEOUT_SECS_DEFAULT async def _connect_adapter_with_timeout( - self, adapter, platform, *, is_reconnect: bool = False + self, adapter, platform, *, is_reconnect: bool = False, initial: bool = False ) -> bool: """Connect an adapter without allowing one platform to block others. @@ -7221,8 +7240,12 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew server-side queue) from a watcher reconnect after a prolonged outage (preserve the queue so messages sent during the outage are delivered rather than silently dropped — #46621). + + ``initial`` selects the capped cold-start budget for platforms whose + full connect budget is too long to spend before the gateway reaches + ``running`` (#85993 — Telegram's 180s). """ - timeout = self._platform_connect_timeout_secs(platform) + timeout = self._platform_connect_timeout_secs(platform, initial=initial) if timeout <= 0: return await adapter.connect(is_reconnect=is_reconnect) # Use the detach-on-timeout pattern instead of plain asyncio.wait_for: @@ -7261,7 +7284,9 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew self._platform_lock_takeover_on_start ) try: - return await self._connect_adapter_with_timeout(adapter, platform) + return await self._connect_adapter_with_timeout( + adapter, platform, initial=True + ) finally: adapter._platform_lock_takeover_allowed = False diff --git a/tests/gateway/test_startup_connect_parallel.py b/tests/gateway/test_startup_connect_parallel.py index ccc53681d5..1c81112cd7 100644 --- a/tests/gateway/test_startup_connect_parallel.py +++ b/tests/gateway/test_startup_connect_parallel.py @@ -181,3 +181,108 @@ async def test_startup_one_failing_platform_does_not_block_others(monkeypatch, t assert Platform.DISCORD in runner.adapters # The failed platform is queued for retry, not silently dropped. assert Platform.TELEGRAM in runner._failed_platforms + + +class TestTelegramColdStartCap: + """The initial (pre-`running`) Telegram connect uses a capped budget (#85993). + + The full 180s Telegram connect budget (#67498) still applies to reconnect + watcher retries; only the cold-start attempt awaited before the gateway + reaches `running` is capped, so an unreachable Telegram can't hold every + other platform's serving state hostage for 3 minutes. + """ + + def _runner(self, tmp_path): + config = GatewayConfig( + platforms={}, sessions_dir=tmp_path / "sessions" + ) + return GatewayRunner(config) + + def test_initial_telegram_budget_is_capped(self, tmp_path, monkeypatch): + monkeypatch.delenv("HERMES_GATEWAY_PLATFORM_CONNECT_TIMEOUT", raising=False) + runner = self._runner(tmp_path) + initial = runner._platform_connect_timeout_secs( + Platform.TELEGRAM, initial=True + ) + full = runner._platform_connect_timeout_secs(Platform.TELEGRAM) + assert initial < full, ( + "cold-start Telegram budget must be shorter than the reconnect " + f"budget (initial={initial}, full={full})" + ) + assert full == 180.0 # #67498 reconnect budget unchanged + assert initial <= 60.0 # gateway reaches `running` within a minute + + def test_other_platforms_unchanged(self, tmp_path, monkeypatch): + monkeypatch.delenv("HERMES_GATEWAY_PLATFORM_CONNECT_TIMEOUT", raising=False) + runner = self._runner(tmp_path) + assert runner._platform_connect_timeout_secs( + Platform.DISCORD, initial=True + ) == runner._platform_connect_timeout_secs(Platform.DISCORD) + + def test_env_override_applies_to_initial(self, tmp_path, monkeypatch): + monkeypatch.setenv("HERMES_GATEWAY_PLATFORM_CONNECT_TIMEOUT", "12") + runner = self._runner(tmp_path) + assert runner._platform_connect_timeout_secs( + Platform.TELEGRAM, initial=True + ) == 12.0 + + @pytest.mark.asyncio + async def test_initial_connect_times_out_at_cap_and_queues_retry( + self, tmp_path, monkeypatch + ): + """A wedged Telegram connect is abandoned at the capped budget and the + platform lands in the reconnect queue instead of blocking startup.""" + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + monkeypatch.delenv("HERMES_GATEWAY_PLATFORM_CONNECT_TIMEOUT", raising=False) + + class _WedgedAdapter(BasePlatformAdapter): + def __init__(self): + super().__init__( + PlatformConfig(enabled=True, token="***"), Platform.TELEGRAM + ) + + async def connect(self, *, is_reconnect: bool = False) -> bool: + await asyncio.sleep(3600) + return True + + async def disconnect(self) -> None: + self._mark_disconnected() + + async def send(self, chat_id, content, reply_to=None, metadata=None): + raise NotImplementedError + + async def get_chat_info(self, chat_id): + return {"id": chat_id} + + config = GatewayConfig( + platforms={ + Platform.TELEGRAM: PlatformConfig(enabled=True, token="***"), + Platform.DISCORD: PlatformConfig(enabled=True, token="***"), + }, + sessions_dir=tmp_path / "sessions", + ) + runner = GatewayRunner(config) + + # Shrink the capped budget so the test is fast; the assertion is that + # the INITIAL path (initial=True) is the one that fires, not the 180s + # reconnect budget. + import gateway.run as gateway_run + + monkeypatch.setattr( + gateway_run, "_TELEGRAM_INITIAL_CONNECT_TIMEOUT_SECS_DEFAULT", 0.2 + ) + + def _make_adapter(platform, platform_config): + if platform is Platform.TELEGRAM: + return _WedgedAdapter() + return _TimingAdapter(platform, 0.0) + + monkeypatch.setattr(runner, "_create_adapter", _make_adapter) + monkeypatch.setattr(runner, "_start_secondary_profile_adapters", lambda: 0) + + await asyncio.wait_for(runner.start(), timeout=30) + + # Discord served; Telegram queued for the watcher's full-budget retry. + assert Platform.DISCORD in runner.adapters + assert Platform.TELEGRAM not in runner.adapters + assert Platform.TELEGRAM in runner._failed_platforms