fix(gateway): schedule secondary-profile reconnect when initial adapter connect fails

When gateway.multiplex_profiles is active, a secondary profile whose platform
adapter fails its initial connect at startup was silently given up on: the
failure branches in _start_one_profile_adapters() logged and disconnected, but
never scheduled recovery. One unlucky connect window during a Telegram API
outage left the profile permanently silent until manual restart (~80 min in
the observed incident), while the mid-run fatal path already recovers via
_handle_profile_adapter_fatal_error() -> _schedule_secondary_profile_reconnect().
The same gap hit both failure shapes: a clean False return from
_connect_initial_adapter_with_timeout() and an exception escaping it.

Fix: call _schedule_secondary_profile_startup_reconnect() from both startup
failure branches after _safe_adapter_disconnect(). Because secondary adapters
are started mid-start(), before self._running flips True, the regular
scheduler's not-self._running guard would silently drop the request — so the
new bridge parks a background task until startup completes (or shutdown
begins) and then hands off to _schedule_secondary_profile_reconnect()
verbatim: backoff, fresh-adapter rebuild under the profile runtime scope,
slot dedupe, and shutdown cancellation all come from the existing path.
Non-retryable failures are dropped at scheduling time exactly as the regular
scheduler drops them, keeping duplicate-credential/auth-failed startups dead
instead of looping.

Fixes #92064
This commit is contained in:
ruangraung
2026-08-22 12:02:07 +07:00
committed by Teknium
parent 682a34a7c6
commit 96489f3c1b
2 changed files with 200 additions and 0 deletions

View File

@@ -15618,9 +15618,15 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
else:
logger.warning("✗ %s failed to connect (profile: %s)", platform.value, profile_name)
await self._safe_adapter_disconnect(adapter, platform)
self._schedule_secondary_profile_startup_reconnect(
profile_name, platform, adapter
)
except Exception as e:
logger.error("✗ %s error (profile: %s): %s", platform.value, profile_name, e)
await self._safe_adapter_disconnect(adapter, platform)
self._schedule_secondary_profile_startup_reconnect(
profile_name, platform, adapter
)
return connected
def _configure_profile_adapter(
@@ -15766,6 +15772,53 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
if not profile_pending:
pending.pop(profile_name, None)
def _schedule_secondary_profile_startup_reconnect(
self, profile_name: str, platform: Platform, adapter: BasePlatformAdapter
) -> None:
"""Queue a cold-start reconnect for a secondary adapter.
Startup failure branches run BEFORE ``self._running`` flips True
(``_start_secondary_profile_adapters()`` is called mid-``start()``,
while ``self._running`` is still False), so the regular scheduler's
``not self._running`` guard would silently drop the request and the
runner's ``while self._running`` loop would exit immediately. This
bridge parks a background task across the remainder of startup and
hands off to the regular scheduler once the gateway is live (the
scheduler's own ``_profile_failed_platforms`` slot dedupes at
handoff); if shutdown begins first, the request is released.
Non-retryable failures are dropped here exactly as the regular
scheduler would.
"""
if not getattr(adapter, "fatal_error_retryable", True):
return
async def _await_running_then_schedule() -> None:
if self._running:
self._schedule_secondary_profile_reconnect(
profile_name, platform, adapter
)
return
# Modest poll interval: startup completion has no dedicated event,
# and the reconnect runner's own backoff makes sub-100ms precision
# irrelevant. Bounded so a wedged startup cannot spin the loop.
while not self._running and not self._shutdown_event.is_set():
await asyncio.sleep(0.1)
if self._running and not self._shutdown_event.is_set():
self._schedule_secondary_profile_reconnect(
profile_name, platform, adapter
)
task = asyncio.create_task(
_await_running_then_schedule(),
name=f"secondary-startup-reconnect:{profile_name}:{platform.value}",
)
background_tasks = getattr(self, "_background_tasks", None)
if not isinstance(background_tasks, set):
background_tasks = set()
self._background_tasks = background_tasks
background_tasks.add(task)
task.add_done_callback(background_tasks.discard)
def _schedule_secondary_profile_reconnect(
self, profile_name: str, platform: Platform, adapter: BasePlatformAdapter
) -> None:

View File

@@ -274,6 +274,153 @@ class TestSecondaryProfileFatalRecovery:
assert runner._profile_failed_platforms == {}
class TestSecondaryStartupFailureRecovery:
"""Cold-start connect failures must reach the same reconnect slot as
mid-run fatals — one unlucky connect window must not kill the platform
for the life of the process."""
@pytest.mark.asyncio
async def test_retryable_initial_failure_schedules_reconnect(
self, monkeypatch
):
runner = _secondary_recovery_runner()
failed = _SecondaryRecoveryAdapter()
replacement = _SecondaryRecoveryAdapter()
scoped_homes: list[Path] = []
_install_secondary_reconnect_context(
monkeypatch, runner, replacement, scoped_homes
)
# Startup creates `failed`; the reconnect runner creates `replacement`.
created = [failed, replacement]
monkeypatch.setattr(
runner, "_create_adapter", lambda platform, config: created.pop(0)
)
async def fail_initial_connect(adapter, platform):
return False
monkeypatch.setattr(
runner, "_connect_initial_adapter_with_timeout", fail_initial_connect
)
async def reconnect_ok(adapter, platform, *, is_reconnect=False):
assert is_reconnect is True
assert adapter is replacement
return True
monkeypatch.setattr(runner, "_connect_adapter_with_timeout", reconnect_ok)
connected = await runner._start_one_profile_adapters(
"reviewer", "/tmp/reviewer", {}
)
assert connected == 0
assert failed.disconnected is True
assert Platform.DISCORD not in runner._profile_adapters.get(
"reviewer", {}
)
bridge = list(runner._background_tasks)
assert len(bridge) == 1
# Drive the bridge to completion; it hands off (immediately when the
# gateway is already running) to the regular reconnect task, which
# publishes the replacement and clears its own slot.
await asyncio.wait_for(bridge[0], timeout=0.5)
for _ in range(20):
if (
runner._profile_adapters.get("reviewer", {}).get(Platform.DISCORD)
is replacement
):
break
await asyncio.sleep(0)
assert (
runner._profile_adapters["reviewer"][Platform.DISCORD] is replacement
)
assert Platform.DISCORD not in runner._profile_failed_platforms.get(
"reviewer", {}
)
# Reconnect must have re-entered the profile's own runtime scope.
assert Path("/profiles/reviewer") in scoped_homes
assert all(
path in (Path("/tmp/reviewer"), Path("/profiles/reviewer"))
for path in scoped_homes
)
@pytest.mark.asyncio
async def test_raising_initial_connect_schedules_reconnect(
self, monkeypatch
):
runner = _secondary_recovery_runner()
failed = _SecondaryRecoveryAdapter()
replacement = _SecondaryRecoveryAdapter()
_install_secondary_reconnect_context(monkeypatch, runner, replacement)
created = [failed, replacement]
monkeypatch.setattr(
runner, "_create_adapter", lambda platform, config: created.pop(0)
)
async def explode(adapter, platform):
raise TimeoutError("initial connect budget exhausted")
monkeypatch.setattr(runner, "_connect_initial_adapter_with_timeout", explode)
async def reconnect_ok(adapter, platform, *, is_reconnect=False):
return True
monkeypatch.setattr(runner, "_connect_adapter_with_timeout", reconnect_ok)
connected = await runner._start_one_profile_adapters(
"reviewer", "/tmp/reviewer", {}
)
assert connected == 0
assert failed.disconnected is True
bridge = list(runner._background_tasks)
assert len(bridge) == 1
await asyncio.wait_for(bridge[0], timeout=0.5)
for _ in range(20):
if (
runner._profile_adapters.get("reviewer", {}).get(Platform.DISCORD)
is replacement
):
break
await asyncio.sleep(0)
assert (
runner._profile_adapters["reviewer"][Platform.DISCORD] is replacement
)
assert Platform.DISCORD not in runner._profile_failed_platforms.get(
"reviewer", {}
)
@pytest.mark.asyncio
async def test_non_retryable_initial_failure_does_not_schedule(
self, monkeypatch
):
runner = _secondary_recovery_runner()
failed = _SecondaryRecoveryAdapter(retryable=False)
_install_secondary_reconnect_context(
monkeypatch, runner, _SecondaryRecoveryAdapter()
)
monkeypatch.setattr(runner, "_create_adapter", lambda platform, config: failed)
async def fail_initial_connect(adapter, platform):
return False
monkeypatch.setattr(
runner, "_connect_initial_adapter_with_timeout", fail_initial_connect
)
connected = await runner._start_one_profile_adapters(
"reviewer", "/tmp/reviewer", {}
)
assert connected == 0
assert failed.disconnected is True
assert runner._background_tasks == set()
assert runner._profile_failed_platforms == {}
class TestSecondaryProfileConfigHandling:
"""Secondary config errors degrade only when the profile is safe to skip."""