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:
@@ -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:
|
||||
|
||||
@@ -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."""
|
||||
|
||||
|
||||
Reference in New Issue
Block a user