diff --git a/gateway/platforms/ADDING_A_PLATFORM.md b/gateway/platforms/ADDING_A_PLATFORM.md index 63ab2e3103025..917482cf92cb4 100644 --- a/gateway/platforms/ADDING_A_PLATFORM.md +++ b/gateway/platforms/ADDING_A_PLATFORM.md @@ -158,6 +158,8 @@ def check__requirements() -> bool: - Use `MessageEvent`, `MessageType` from `gateway.platforms.event` and `SendResult` from base - Use `cache_image_from_bytes`, `cache_audio_from_bytes`, `cache_document_from_bytes` for attachments - Filter self-messages (prevent reply loops) +- Drop redelivered inbound IDs with `MessageDeduplicator` (`gateway/platforms/helpers.py`) held as an adapter + attribute; the runner's reconnect copies its live IDs into the rebuilt adapter, a hand-rolled cache starts empty - Filter sync/echo messages if the platform has them - Redact sensitive identifiers (phone numbers, tokens) in all log output - Implement reconnection with exponential backoff + jitter for streaming connections diff --git a/gateway/platforms/helpers.py b/gateway/platforms/helpers.py index ee5699a5ee75d..0069b10cb50fd 100644 --- a/gateway/platforms/helpers.py +++ b/gateway/platforms/helpers.py @@ -60,6 +60,29 @@ def discard(self, msg_id: str) -> None: def clear(self): self._seen.clear() + def absorb(self, other: "MessageDeduplicator") -> None: + """Adopt *other*'s still-live IDs (at their original seen times) into this cache.""" + cutoff = time.time() - self._ttl + self._seen.update({k: v for k, v in other._seen.items() if v > cutoff and k not in self._seen}) + + +def inbound_dedup_caches(adapter: Any) -> dict[str, MessageDeduplicator]: + """The adapter's ``MessageDeduplicator`` attributes, by name (held by reference, so IDs the old + adapter admits after this call still reach its replacement).""" + return {name: v for name, v in vars(adapter).items() if isinstance(v, MessageDeduplicator)} + + +def carry_inbound_dedup(caches: Optional[dict], adapter: Any) -> None: + """Seed a rebuilt adapter's dedup caches from the instance it replaces. + + The runner's reconnect path builds a NEW adapter; without this a platform replaying a recent + inbound ID after the reconnect (websocket resume, webhook retry, unacked poll batch) is + admitted and answered a second time.""" + for name, previous in (caches or {}).items(): + current = getattr(adapter, name, None) + if isinstance(current, MessageDeduplicator) and current is not previous: + current.absorb(previous) + # Worker-thread handoff used by the off-loop persist paths. A module attribute # so tests can replace THIS seam instead of patching ``asyncio.to_thread`` diff --git a/gateway/run_adapters.py b/gateway/run_adapters.py index 418c8f122b1c9..36033a374c0e1 100644 --- a/gateway/run_adapters.py +++ b/gateway/run_adapters.py @@ -21,6 +21,7 @@ from datetime import datetime, timedelta, timezone from gateway.config import SHARED_LISTENER_MIRROR_PLATFORMS, Platform, platform_binds_port as _platform_binds_port from gateway.platforms.base import BasePlatformAdapter +from gateway.platforms.helpers import carry_inbound_dedup, inbound_dedup_caches from gateway.restart import is_global_startup_conflict from gateway.run_shutdown import _log_suppressed from gateway.session import SessionSource @@ -219,6 +220,7 @@ def _reconnect_queue_entry( **({"queued_at": now} if queued else {}), "credential_claim": self._adapter_credential_claim(platform, adapter), "listener_claim": self._adapter_listener_claim(platform, adapter), + "inbound_dedup": inbound_dedup_caches(adapter), } def _queue_retryable_fatal_platform(self, adapter: BasePlatformAdapter) -> bool: @@ -739,6 +741,7 @@ async def _reconnect_failed_platform(self, platform, now: float) -> None: if not adapter: self._drop_from_reconnect_queue(platform, "adapter creation returned None") return + carry_inbound_dedup(info.get("inbound_dedup"), adapter) self._wire_adapter_handlers(adapter) # is_reconnect keeps the server-side update queue so offline-period messages are delivered. success = await self._connect_adapter_with_timeout(adapter, platform, is_reconnect=True) @@ -1213,7 +1216,7 @@ def _configure_profile_adapter( and _platform_binds_port(platform.value, getattr(getattr(adapter, "config", None), "extra", None)): adapter._shared_listener_profile = profile_name - async def _secondary_reconnect_attempt(self, profile_name: str, platform: Platform): + async def _secondary_reconnect_attempt(self, profile_name: str, platform: Platform, inbound_dedup=None): """One scoped attempt to rebuild+connect a secondary adapter → ``(adapter, success)``; ``(None, None)`` = give up for good (disabled, credential removed, adapter unavailable). Caller tears down a RETURNED adapter; one whose configure/connect raised is torn down here.""" @@ -1245,6 +1248,7 @@ async def _secondary_reconnect_attempt(self, profile_name: str, platform: Platfo platform.value, profile_name, ) return None, None + carry_inbound_dedup(inbound_dedup, adapter) try: self._configure_profile_adapter(adapter, profile_name, platform) success = await self._connect_adapter_with_timeout(adapter, platform, is_reconnect=True) @@ -1254,7 +1258,9 @@ async def _secondary_reconnect_attempt(self, profile_name: str, platform: Platfo raise return adapter, success - async def _run_secondary_profile_reconnect(self, profile_name: str, platform: Platform) -> None: + async def _run_secondary_profile_reconnect( + self, profile_name: str, platform: Platform, inbound_dedup=None + ) -> None: """Reconnect a retryable secondary adapter under its own profile scope.""" from gateway.run import _profile_runtime_scope, _reconnect_backoff attempts = 0 @@ -1265,7 +1271,9 @@ async def _run_secondary_profile_reconnect(self, profile_name: str, platform: Pl while self._running: adapter = None try: - adapter, success = await self._secondary_reconnect_attempt(profile_name, platform) + adapter, success = await self._secondary_reconnect_attempt( + profile_name, platform, inbound_dedup + ) if adapter is None: return if success and self._running: @@ -1386,7 +1394,7 @@ def _schedule_secondary_profile_reconnect( if platform in profile_pending: return profile_pending[platform] = self._retain_background_task(asyncio.create_task( - self._run_secondary_profile_reconnect(profile_name, platform), + self._run_secondary_profile_reconnect(profile_name, platform, inbound_dedup_caches(adapter)), name=f"secondary-reconnect:{profile_name}:{platform.value}", )) diff --git a/tests/gateway/test_multiplex_adapter_registry.py b/tests/gateway/test_multiplex_adapter_registry.py index 7d4453e4ae592..f21b12a946af4 100644 --- a/tests/gateway/test_multiplex_adapter_registry.py +++ b/tests/gateway/test_multiplex_adapter_registry.py @@ -12,6 +12,7 @@ import gateway.run as gateway_run from gateway.config import GatewayConfig, Platform, PlatformConfig +from gateway.platforms.helpers import MessageDeduplicator from gateway.run import GatewayRunner from gateway.status import flush_runtime_status @@ -186,6 +187,7 @@ def __init__(self, *, retryable=True): self.fatal_error_message = "Gateway transport stale" self.connected = False self.disconnected = False + self._dedup = MessageDeduplicator() async def disconnect(self): self.disconnected = True @@ -396,6 +398,23 @@ async def redeliver(platform, *, profile=None): assert all(path != Path("/profiles/reviewer") for path in redelivery_homes) + @pytest.mark.asyncio + async def test_secondary_reconnect_keeps_inbound_dedup(self, monkeypatch): + """A secondary profile's rebuilt adapter still drops an inbound ID the stale one admitted.""" + runner = _secondary_recovery_runner() + stale, replacement = _SecondaryRecoveryAdapter(), _SecondaryRecoveryAdapter() + runner._profile_adapters["reviewer"] = {Platform.DISCORD: stale} + _install_secondary_reconnect_context(monkeypatch, runner, replacement) + monkeypatch.setattr(runner, "_connect_adapter_with_timeout", AsyncMock(return_value=True)) + assert stale._dedup.is_duplicate("m1") is False + + await runner._handle_profile_adapter_fatal_error("reviewer", Platform.DISCORD, stale) + await asyncio.gather(*runner._background_tasks) + + assert runner._profile_adapters["reviewer"][Platform.DISCORD] is replacement + assert replacement._dedup.is_duplicate("m1") is True + assert replacement._dedup.is_duplicate("m2") is False + @pytest.mark.asyncio @pytest.mark.parametrize("connect_result", [True, False], ids=["success", "failure"]) async def test_secondary_reconnect_does_not_publish_after_shutdown( diff --git a/tests/gateway/test_platform_reconnect.py b/tests/gateway/test_platform_reconnect.py index 7ee9a433a7ee4..b194dd55b36db 100644 --- a/tests/gateway/test_platform_reconnect.py +++ b/tests/gateway/test_platform_reconnect.py @@ -8,6 +8,7 @@ from gateway.config import GatewayConfig, Platform, PlatformConfig from gateway.platforms.base import BasePlatformAdapter, SendResult +from gateway.platforms.helpers import MessageDeduplicator from gateway.run import GatewayRunner @@ -388,6 +389,30 @@ async def test_retryable_error_keeps_gateway_alive_when_all_down(self): assert Platform.TELEGRAM in runner._failed_platforms +class TestReconnectKeepsInboundDedup: + @pytest.mark.asyncio + async def test_replayed_inbound_id_after_runner_reconnect_is_dropped(self): + """The watcher builds a NEW adapter; an inbound ID the old one already admitted must still + read as a duplicate there, or a platform replay after the reconnect is answered twice.""" + runner = _make_runner() + runner.stop = AsyncMock() + runner._sync_voice_mode_state_to_adapter = MagicMock() + old, new = StubAdapter(), StubAdapter() + for a in (old, new): + a._dedup = MessageDeduplicator() + runner.adapters[Platform.TELEGRAM] = old + assert old._dedup.is_duplicate("m1") is False # handled before the drop + + old._set_fatal_error("network_error", "socket closed", retryable=True) + await runner._handle_adapter_fatal_error(old) + with patch.object(runner, "_create_adapter", return_value=new): + await runner._reconnect_failed_platform(Platform.TELEGRAM, time.monotonic() + 1) + + assert runner.adapters[Platform.TELEGRAM] is new + assert new._dedup.is_duplicate("m1") is True + assert new._dedup.is_duplicate("m2") is False + + # --- Pause / resume circuit breaker --- diff --git a/website/docs/developer-guide/adding-platform-adapters.md b/website/docs/developer-guide/adding-platform-adapters.md index fab10a5811fdd..536a7f6646895 100644 --- a/website/docs/developer-guide/adding-platform-adapters.md +++ b/website/docs/developer-guide/adding-platform-adapters.md @@ -760,6 +760,21 @@ async def _handle_callback(self, request): For platforms with tight response deadlines (e.g., WeCom's 5-second limit), always acknowledge immediately and deliver the agent's reply proactively via API later. Agent sessions run 3–30 minutes — inline replies within a callback response window are not feasible. +### Inbound Deduplication + +Platforms redeliver: websocket resumes replay recent events, webhooks retry, and an unacknowledged poll batch comes back. Drop repeats with the shared helper, keyed on the platform's message ID: + +```python +from gateway.platforms.helpers import MessageDeduplicator + +self._dedup = MessageDeduplicator(ttl_seconds=600) # in __init__ + +if self._dedup.is_duplicate(msg_id): # in the inbound handler + return +``` + +When the gateway's reconnect watcher replaces a failed adapter with a new instance, it copies every `MessageDeduplicator` attribute's live IDs from the old instance to the new one, so a replay right after the reconnect is still dropped. A cache kept in another structure (a plain dict or set) starts empty on the new instance. + ### Token Locks If the adapter holds a persistent connection with a unique credential, add a scoped lock to prevent two profiles from using the same credential: