From 55e700abf390fad8da53927ea405fdda2bb9a955 Mon Sep 17 00:00:00 2001 From: EloquentBrush <147827411+EloquentBrush@users.noreply.github.com> Date: Sun, 10 May 2026 22:43:54 +0300 Subject: [PATCH 1/3] fix(yuanbao): clear _processing_msg_ids/_processing_msg_texts after each message MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit _dispatch_inbound_event() writes session_key → msg_id/raw_text into _processing_msg_ids and _processing_msg_texts so RecallGuardMiddleware can find and interrupt the currently-processing message. These entries were never removed after a message finished processing, causing both dicts to grow unboundedly — one persistent entry per unique session key for the lifetime of the bot. Fix: clear both entries in the _process_message_background() finally block, after super() returns. The guard compares the stored msg_id against event.message_id before popping: a concurrent pending message may have already overwritten the entry in _dispatch_inbound_event while we were running, in which case the drain task owns it and we must not clear it. When msg_id is absent (nothing was written at dispatch time) the pop is a safe no-op. Note: _msg_content_cache already bounds itself to 200 entries at the same write site; _processing_msg_ids and _processing_msg_texts had no such bound. --- gateway/platforms/yuanbao.py | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/gateway/platforms/yuanbao.py b/gateway/platforms/yuanbao.py index ad5da066146e..118224fa08ec 100644 --- a/gateway/platforms/yuanbao.py +++ b/gateway/platforms/yuanbao.py @@ -5116,6 +5116,15 @@ async def _process_message_background(self, event, session_key: str) -> None: await super()._process_message_background(event, session_key) finally: self._outbound.cancel_slow_notifier(chat_id) + # Clear the RecallGuard tracking entries for this message only if + # our msg_id is still current. A concurrent pending message may + # have already overwritten the entry in _dispatch_inbound_event + # while we were running; in that case the drain task owns it and + # we must not clear it. + msg_id = event.message_id + if not msg_id or self._processing_msg_ids.get(session_key) == msg_id: + self._processing_msg_ids.pop(session_key, None) + self._processing_msg_texts.pop(session_key, None) # ------------------------------------------------------------------ # Group query (delegate to GroupQueryService) From 4c473862e073d7972d80031fe393c2459b3fbe4f Mon Sep 17 00:00:00 2001 From: EloquentBrush <147827411+EloquentBrush@users.noreply.github.com> Date: Sun, 10 May 2026 22:46:43 +0300 Subject: [PATCH 2/3] fix(yuanbao): evict stale entries from _member_cache on TTL expiry _build_msg_body_with_mentions() checks the TTL of each _member_cache entry and returns an empty member list when the entry is stale, but never removes the entry from the dict. Over time every group_code the bot has ever queried accumulates a permanent entry, retaining the full member list (potentially thousands of records per group) until disconnect(). Fix: delete the stale entry at the point it is detected as expired. The next call to get_group_member_list_raw() for the same group will repopulate the cache with fresh data as before. Symmetric with the existing TTL pattern in MessageDeduplicator, which evicts on access. --- gateway/platforms/yuanbao.py | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/gateway/platforms/yuanbao.py b/gateway/platforms/yuanbao.py index 118224fa08ec..bec2976b1e0c 100644 --- a/gateway/platforms/yuanbao.py +++ b/gateway/platforms/yuanbao.py @@ -4588,7 +4588,11 @@ def _build_msg_body_with_mentions(self, text: str, group_code: str) -> list: cached = self._adapter._member_cache.get(group_code) if cached: ts, member_list = cached - members = member_list if (time.time() - ts < self._adapter.MEMBER_CACHE_TTL_S) else [] + if time.time() - ts < self._adapter.MEMBER_CACHE_TTL_S: + members = member_list + else: + del self._adapter._member_cache[group_code] + members = [] else: members = [] if not members: From 0d8e7cf65d39f3e3fd89ab05e71407391188a8da Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Mon, 3 Aug 2026 17:07:43 +0530 Subject: [PATCH 3/3] fix(yuanbao): pop tracking entries only for truthy matching msg_id + regression tests Follow-up on the salvaged pair: the original guard's `not msg_id` arm let an id-less internal/synthetic event erase a tracking entry a concurrently-queued id-bearing message's drain task still needs for recall matching (id-less events never write entries in _dispatch_inbound_event, so they must never pop). Tests cover: normal cleanup, id-less non-erasure, overwritten-entry ownership handoff, TTL eviction + fresh-entry survival. --- gateway/platforms/yuanbao.py | 8 +- .../platforms/test_yuanbao_state_cleanup.py | 174 ++++++++++++++++++ 2 files changed, 180 insertions(+), 2 deletions(-) create mode 100644 tests/gateway/platforms/test_yuanbao_state_cleanup.py diff --git a/gateway/platforms/yuanbao.py b/gateway/platforms/yuanbao.py index bec2976b1e0c..2e3ed0253b34 100644 --- a/gateway/platforms/yuanbao.py +++ b/gateway/platforms/yuanbao.py @@ -5124,9 +5124,13 @@ async def _process_message_background(self, event, session_key: str) -> None: # our msg_id is still current. A concurrent pending message may # have already overwritten the entry in _dispatch_inbound_event # while we were running; in that case the drain task owns it and - # we must not clear it. + # we must not clear it. Id-less events (internal/synthetic + # messages, pushes without a msg_id) never wrote a tracking entry + # in _dispatch_inbound_event, so they must never pop either — the + # entry they see belongs to a concurrently-queued id-bearing + # message whose drain task still needs it for recall matching. msg_id = event.message_id - if not msg_id or self._processing_msg_ids.get(session_key) == msg_id: + if msg_id and self._processing_msg_ids.get(session_key) == msg_id: self._processing_msg_ids.pop(session_key, None) self._processing_msg_texts.pop(session_key, None) diff --git a/tests/gateway/platforms/test_yuanbao_state_cleanup.py b/tests/gateway/platforms/test_yuanbao_state_cleanup.py new file mode 100644 index 000000000000..d2dc166f2cc3 --- /dev/null +++ b/tests/gateway/platforms/test_yuanbao_state_cleanup.py @@ -0,0 +1,174 @@ +"""Yuanbao per-turn state cleanup: RecallGuard tracking dicts + member cache TTL. + +Covers the salvage of PRs #23383 / #23384: + +* ``_processing_msg_ids`` / ``_processing_msg_texts`` must be cleared when a + turn finishes (they previously leaked forever, letting RecallGuard match a + recall against an already-finished turn). +* The cleanup must pop ONLY when the finishing event's msg_id is truthy AND + still owns the entry. An id-less event (internal/synthetic message, push + without msg_id) never wrote an entry, so it must never erase one either — + the entry it sees belongs to a concurrently-queued id-bearing message whose + drain task still needs it. +* ``_member_cache`` entries past ``MEMBER_CACHE_TTL_S`` must actually be + evicted on read (the dict shrinks), while fresh entries survive. +""" +import asyncio +import time +from types import SimpleNamespace + +from gateway.platforms.base import BasePlatformAdapter +from gateway.platforms.yuanbao import MessageSender, YuanbaoAdapter + + +# --------------------------------------------------------------------------- +# Helpers +# --------------------------------------------------------------------------- + +class _OutboundStub: + async def start_slow_notifier(self, chat_id): # noqa: ANN001 + pass + + def cancel_slow_notifier(self, chat_id): # noqa: ANN001 + pass + + +def _bare_adapter(): + """YuanbaoAdapter instance without running its heavy __init__.""" + adapter = object.__new__(YuanbaoAdapter) + adapter._outbound = _OutboundStub() + adapter._processing_msg_ids = {} + adapter._processing_msg_texts = {} + return adapter + + +def _event(message_id): + return SimpleNamespace( + source=SimpleNamespace(chat_id="chat-1"), + message_id=message_id, + ) + + +def _run_turn(monkeypatch, adapter, event, session_key, during_turn=None): + """Run the yuanbao _process_message_background wrapper with the base + class processing stubbed out (optionally mutating state mid-turn).""" + + async def _base_stub(self, ev, sk): # noqa: ANN001 + if during_turn is not None: + during_turn() + + monkeypatch.setattr( + BasePlatformAdapter, "_process_message_background", _base_stub + ) + asyncio.run( + YuanbaoAdapter._process_message_background(adapter, event, session_key) + ) + + +# --------------------------------------------------------------------------- +# _processing_msg_ids / _processing_msg_texts cleanup (PR #23383) +# --------------------------------------------------------------------------- + +def test_tracking_entries_cleared_after_normal_turn(monkeypatch): + """A turn whose msg_id still owns the tracking entry clears it on exit.""" + adapter = _bare_adapter() + sk = "yuanbao:group:G:user:U" + # _dispatch_inbound_event wrote these before handle_message. + adapter._processing_msg_ids[sk] = "m1" + adapter._processing_msg_texts[sk] = "hello" + + _run_turn(monkeypatch, adapter, _event("m1"), sk) + + assert sk not in adapter._processing_msg_ids + assert sk not in adapter._processing_msg_texts + + +def test_idless_event_must_not_erase_drain_tasks_entry(monkeypatch): + """An id-less outer event finishing must NOT pop the tracking entry a + concurrently-dispatched id-bearing message (queued as pending, to be + handled by a drain task) wrote during the outer turn.""" + adapter = _bare_adapter() + sk = "yuanbao:group:G:user:U" + + def _pending_message_arrives(): + # Simulates _dispatch_inbound_event for msg "m2" arriving while the + # id-less event is still processing: it writes tracking state, then + # handle_message routes it to _pending_messages for the drain task. + adapter._processing_msg_ids[sk] = "m2" + adapter._processing_msg_texts[sk] = "recallable text" + + _run_turn( + monkeypatch, adapter, _event(None), sk, + during_turn=_pending_message_arrives, + ) + + # The drain task for "m2" still needs these for RecallGuard matching. + assert adapter._processing_msg_ids.get(sk) == "m2" + assert adapter._processing_msg_texts.get(sk) == "recallable text" + + +def test_overwritten_entry_not_erased_by_outdated_turn(monkeypatch): + """If a newer message already overwrote the entry, the older finishing + turn must leave it alone (drain task owns it).""" + adapter = _bare_adapter() + sk = "yuanbao:group:G:user:U" + adapter._processing_msg_ids[sk] = "m1" + adapter._processing_msg_texts[sk] = "first" + + def _newer_message_arrives(): + adapter._processing_msg_ids[sk] = "m2" + adapter._processing_msg_texts[sk] = "second" + + _run_turn( + monkeypatch, adapter, _event("m1"), sk, + during_turn=_newer_message_arrives, + ) + + assert adapter._processing_msg_ids.get(sk) == "m2" + assert adapter._processing_msg_texts.get(sk) == "second" + + +# --------------------------------------------------------------------------- +# _member_cache TTL eviction (PR #23384) +# --------------------------------------------------------------------------- + +def _bare_sender(adapter_stub): + sender = object.__new__(MessageSender) + sender._adapter = adapter_stub + return sender + + +def test_member_cache_expired_entry_is_evicted(): + """Reading an expired entry must delete it — the cache dict shrinks.""" + now = time.time() + adapter = SimpleNamespace( + MEMBER_CACHE_TTL_S=300.0, + _member_cache={ + "g-stale": (now - 301.0, [{"nickname": "bob", "user_id": "u1"}]), + }, + ) + sender = _bare_sender(adapter) + + body = sender._build_msg_body_with_mentions("hi @bob", "g-stale") + + # Expired ⇒ no member data ⇒ plain text body, and the key is GONE. + assert body == [{"msg_type": "TIMTextElem", "msg_content": {"text": "hi @bob"}}] + assert "g-stale" not in adapter._member_cache + assert len(adapter._member_cache) == 0 + + +def test_member_cache_fresh_entry_survives_read(): + """A fresh entry is used for mention resolution and stays cached.""" + now = time.time() + members = [{"nickname": "bob", "user_id": "u1"}] + adapter = SimpleNamespace( + MEMBER_CACHE_TTL_S=300.0, + _member_cache={"g-fresh": (now - 10.0, members)}, + ) + sender = _bare_sender(adapter) + + body = sender._build_msg_body_with_mentions("hi @bob", "g-fresh") + + assert "g-fresh" in adapter._member_cache + # Fresh members were actually used: an @mention element is present. + assert any(el.get("msg_type") == "TIMCustomElem" for el in body)