Skip to content
Merged
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
2 changes: 2 additions & 0 deletions gateway/platforms/ADDING_A_PLATFORM.md
Original file line number Diff line number Diff line change
Expand Up @@ -158,6 +158,8 @@ def check_<platform>_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
Expand Down
23 changes: 23 additions & 0 deletions gateway/platforms/helpers.py
Original file line number Diff line number Diff line change
Expand Up @@ -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``
Expand Down
16 changes: 12 additions & 4 deletions gateway/run_adapters.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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."""
Expand Down Expand Up @@ -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)
Expand All @@ -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
Expand All @@ -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:
Expand Down Expand Up @@ -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}",
))

Expand Down
19 changes: 19 additions & 0 deletions tests/gateway/test_multiplex_adapter_registry.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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(
Expand Down
25 changes: 25 additions & 0 deletions tests/gateway/test_platform_reconnect.py
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down Expand Up @@ -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 ---


Expand Down
15 changes: 15 additions & 0 deletions website/docs/developer-guide/adding-platform-adapters.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
Loading