diff --git a/gateway/platforms/ADDING_A_PLATFORM.md b/gateway/platforms/ADDING_A_PLATFORM.md index 9887ad5dd9..3fb182a3c7 100644 --- a/gateway/platforms/ADDING_A_PLATFORM.md +++ b/gateway/platforms/ADDING_A_PLATFORM.md @@ -141,7 +141,12 @@ def check__requirements() -> bool: ### Key patterns to follow -- Use `self.build_source(...)` to construct `SessionSource` objects +- Use `self.build_source(...)` to construct `SessionSource` objects (never `SessionSource(...)` + directly — the transport provenance and profile route are stamped there) +- Derive every adapter-side session key (batching, per-chat queues, busy detection) through + `self._event_session_key(event)` / `self._source_session_key(source)`, never the free + `build_session_key()` — the seam keys in the owning profile's namespace under a multiplexed + gateway; the advisory lint (`scripts/check_profile_scope_patterns.py`, pattern P32) flags both - Call `self.handle_message(event)` to dispatch inbound messages to the gateway - 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 diff --git a/gateway/platforms/weixin.py b/gateway/platforms/weixin.py index d8c3e835d3..c338b14d0d 100644 --- a/gateway/platforms/weixin.py +++ b/gateway/platforms/weixin.py @@ -913,12 +913,6 @@ class WeixinAdapter(OwnAccessPolicyMixin, BasePlatformAdapter): else: await self.handle_message(event) - def _text_batch_key(self, event: MessageEvent) -> str: - from gateway.session import build_session_key - return build_session_key( - event.source, group_sessions_per_user=self.config.extra.get("group_sessions_per_user", True), - thread_sessions_per_user=self.config.extra.get("thread_sessions_per_user", False), profile=event.source.profile) - async def _collect_media(self, item: Dict[str, Any], media_paths: List[str], media_types: List[str]) -> None: spec = _INBOUND_MEDIA.get(item.get("type")) path, mime = await self._download_media(item, spec) if spec else (None, "") diff --git a/gateway/platforms/yuanbao.py b/gateway/platforms/yuanbao.py index 7c4be70d82..ad1d60284a 100644 --- a/gateway/platforms/yuanbao.py +++ b/gateway/platforms/yuanbao.py @@ -64,7 +64,6 @@ from gateway.platforms.yuanbao_proto import ( encode_send_private_heartbeat, encode_send_group_heartbeat, encode_query_group_info, encode_get_group_member_list, next_seq_no, ) -from gateway.session import build_session_key from gateway.session_transcript import TranscriptReadError logger = logging.getLogger(__name__) @@ -1666,11 +1665,8 @@ class DispatchMiddleware(InboundMiddleware): async def handle(self, ctx: InboundContext, next_fn) -> None: adapter = ctx.adapter - _sk = build_session_key( - ctx.source, - group_sessions_per_user=adapter.config.extra.get("group_sessions_per_user", True), - thread_sessions_per_user=adapter.config.extra.get("thread_sessions_per_user", False), - ) + # The adapter seam: keyed in the owner profile's namespace, same as ``handle_message``. + _sk = adapter._source_session_key(ctx.source) async def _dispatch_inbound_event() -> None: if any(mt.startswith(("application/", "text/")) for mt in ctx.media_types): diff --git a/plugins/platforms/raft/adapter.py b/plugins/platforms/raft/adapter.py index d3d45a5c18..2f4d26ee83 100644 --- a/plugins/platforms/raft/adapter.py +++ b/plugins/platforms/raft/adapter.py @@ -39,7 +39,6 @@ sys.path.insert(0, str(_Path(__file__).resolve().parents[3])) from gateway.config import Platform, PlatformConfig from gateway.platforms.base import BasePlatformAdapter, SendResult, merge_pending_message_event from gateway.platforms.event import MessageEvent, MessageType -from gateway.session import build_session_key from gateway.platforms._shared import coerce_port, profile_scoped as _profile_scoped logger = logging.getLogger(__name__) @@ -491,10 +490,7 @@ class RaftAdapter(BasePlatformAdapter): return if not self._message_handler: return - session_key = build_session_key( - event.source, group_sessions_per_user=self.config.extra.get("group_sessions_per_user", True), - thread_sessions_per_user=self.config.extra.get("thread_sessions_per_user", False), - profile=self._session_key_profile(event.source)) + session_key = self._event_session_key(event) if session_key in self._active_sessions: logger.debug("[raft] Wake queued for busy session %s", session_key) merge_pending_message_event(self._pending_messages, session_key, event) diff --git a/plugins/platforms/slack/adapter.py b/plugins/platforms/slack/adapter.py index ea8374041f..7d39d80b0a 100644 --- a/plugins/platforms/slack/adapter.py +++ b/plugins/platforms/slack/adapter.py @@ -5961,20 +5961,14 @@ class SlackAdapter(BasePlatformAdapter): def _build_thread_session_key( self, channel_id: str, thread_ts: str, user_id: str, team_id: str = "", *, chat_type: str = "group") -> Optional[str]: - """Thread session key via ``build_session_key()`` (honours per-user isolation). - ``chat_type`` must come from the event's ``channel_type``, not the ID prefix (MPIM ids - start with ``G``).""" - session_store = getattr(self, "_session_store", None) - if not session_store: + """Thread session key through the adapter seam (``_source_session_key``: per-user isolation + from the adapter config the runner seeded, owner-profile namespace). ``chat_type`` must come + from the event's ``channel_type``, not the ID prefix (MPIM ids start with ``G``).""" + if not getattr(self, "_session_store", None): return None try: - from gateway.session import build_session_key source = self._thread_session_source(channel_id, thread_ts, user_id, team_id, chat_type) - store_cfg = getattr(session_store, "config", None) - return build_session_key( - source, group_sessions_per_user=getattr(store_cfg, "group_sessions_per_user", True), - thread_sessions_per_user=getattr(store_cfg, "thread_sessions_per_user", False), - profile=self._session_key_profile(source)) + return self._source_session_key(source) except Exception: return None diff --git a/plugins/platforms/telegram/adapter.py b/plugins/platforms/telegram/adapter.py index 4d1c7f0d9e..1ba65f1925 100644 --- a/plugins/platforms/telegram/adapter.py +++ b/plugins/platforms/telegram/adapter.py @@ -6271,11 +6271,7 @@ class TelegramAdapter(BasePlatformAdapter): def _photo_batch_key(self, event: MessageEvent, msg: Message) -> str: """Return a batching key for Telegram photos/albums.""" - from gateway.session import build_session_key - session_key = build_session_key( - event.source, group_sessions_per_user=self.config.extra.get("group_sessions_per_user", True), - thread_sessions_per_user=self.config.extra.get("thread_sessions_per_user", False), - profile=self._session_key_profile(event.source)) + session_key = self._event_session_key(event) media_group_id = getattr(msg, "media_group_id", None) return f"{session_key}:album:{media_group_id}" if media_group_id else f"{session_key}:photo-burst" diff --git a/scripts/check_profile_scope_patterns.py b/scripts/check_profile_scope_patterns.py index dc7e7f3821..79dbccfcad 100644 --- a/scripts/check_profile_scope_patterns.py +++ b/scripts/check_profile_scope_patterns.py @@ -49,7 +49,9 @@ def load_patterns(path: Path = PATTERNS) -> list[dict]: data = json.loads(path.read_text(encoding="utf-8")) out = [] for p in data["patterns"]: - out.append({**p, "_rx": re.compile(p["pattern_regex"], re.M)}) + # ``path_regex`` (optional) restricts a pattern to files whose repo-relative path matches. + path_rx = re.compile(p["path_regex"]) if p.get("path_regex") else None + out.append({**p, "_rx": re.compile(p["pattern_regex"], re.M), "_path_rx": path_rx}) return out @@ -63,6 +65,9 @@ def scan_text(rel: str, text: str, patterns: list[dict], lines: set[int] | None findings: list[Finding] = [] src_lines = text.split("\n") for p in patterns: + path_rx = p.get("_path_rx") + if path_rx is not None and not path_rx.search(rel): + continue for m in p["_rx"].finditer(text): line_no = text.count("\n", 0, m.start()) + 1 if lines is not None and line_no not in lines: diff --git a/scripts/ci/profile_scope_patterns.json b/scripts/ci/profile_scope_patterns.json index 50baf25932..94e786cc51 100644 --- a/scripts/ci/profile_scope_patterns.json +++ b/scripts/ci/profile_scope_patterns.json @@ -32,7 +32,7 @@ "P24: 62 hits on main", "P26: 252 hits on main" ], - "usage": "scripts/check_profile_scope_patterns.py --base origin/main [--head HEAD] | --files " + "usage": "scripts/check_profile_scope_patterns.py --base origin/main [--head HEAD] | --files ; optional path_regex restricts a pattern to matching repo-relative paths" }, "patterns": [ { @@ -158,8 +158,16 @@ "id": "P31", "class": "C6", "pattern_regex": "f\"agent:\\{|[\"']agent:[\"']\\s*\\+|session_key\\s*=\\s*f\"[a-z]+:", - "scope_hint": "Every adapter-built session key carries the agent:: namespace (profile 'main' is 'agent:main~'); yuanbao still builds keys with no profile component per the MindDragon probe.", + "scope_hint": "Every adapter-built session key carries the agent:: namespace (profile 'main' is 'agent:main~'); adapters derive keys through _source_session_key / _event_session_key (P32), never a hand-built prefix.", "why": "Rows for a served profile land in the root store; browser/computer_use caches never saw the namespace because turns pass the bare session id." + }, + { + "id": "P32", + "class": "C4", + "path_regex": "^(gateway|plugins)/platforms/(?!base\\.py$).+\\.py$", + "pattern_regex": "\\bbuild_session_key\\(|\\bSessionSource\\(", + "scope_hint": "Inside an adapter derive every key through self._source_session_key(source) / self._event_session_key(event) (owner-profile namespace, runner-seeded isolation flags) and build sources with self.build_source(...) so the transport provenance is kept; only platforms/base.py owns the free calls.", + "why": "Yuanbao keyed its per-group queue and RecallGuard with the free build_session_key() (no profile) while handle_message keyed under agent:: - two derivations of one identity, one lane shared across bots (#88715)." } ] } diff --git a/tests/gateway/test_adapter_session_key_seam.py b/tests/gateway/test_adapter_session_key_seam.py new file mode 100644 index 0000000000..1e807b6fa4 --- /dev/null +++ b/tests/gateway/test_adapter_session_key_seam.py @@ -0,0 +1,69 @@ +"""Every adapter-side session key goes through one seam (#88715, invariant 4). + +A key an adapter derives for batching / queueing / busy detection must be the key the runner +derives for the same event, for a primary bot and for a secondary-owned bot alike. Yuanbao's +``DispatchMiddleware`` used the free ``build_session_key()`` (no profile) while ``handle_message`` +keyed under ``agent::``, so a secondary bot's per-group queue and RecallGuard entries lived +in a lane the runner never popped. +""" + +import asyncio +from types import SimpleNamespace + +import pytest + +from gateway.config import GatewayConfig, Platform, PlatformConfig +from gateway.platforms.event import MessageEvent +from gateway.platforms.yuanbao import DispatchMiddleware, InboundContext, YuanbaoAdapter +from gateway.profile_routing import parse_profile_routes + + +def _yuanbao(owner): + adapter = YuanbaoAdapter(PlatformConfig(extra={ + "app_id": "k", "app_secret": "s", "ws_url": "wss://x", "api_domain": "https://x", + "group_sessions_per_user": True, "thread_sessions_per_user": False})) + adapter.set_owner_profile(owner) + return adapter + + +def _runner(adapter, owner): + from gateway.run import GatewayRunner + + runner = object.__new__(GatewayRunner) + runner.config = GatewayConfig(multiplex_profiles=True) + runner.config.profile_routes = parse_profile_routes([]) + runner._primary_profile_name = "default" + runner.adapters = {} if owner else {Platform.YUANBAO: adapter} + runner._profile_adapters = {owner: {Platform.YUANBAO: adapter}} if owner else {} + adapter.gateway_runner = runner + return runner + + +@pytest.mark.parametrize("owner", [None, "acme"]) +def test_adapter_batch_key_equals_runner_session_key(owner): + adapter = _yuanbao(owner) + runner = _runner(adapter, owner) + source = adapter.build_source(chat_id="grp-1", chat_type="group", user_id="u1") + ctx = InboundContext(adapter=adapter, chat_type="group", chat_id="grp-1", raw_text="hi", msg_id="m1", source=source) + seen = [] + + async def go(): + async def _next(): + pass + + async def _capture(event): + seen.append(adapter._event_session_key(event)) + + adapter.handle_message = _capture + await DispatchMiddleware().handle(ctx, _next) + queued = list(adapter._group_queues) + await asyncio.sleep(0.05) # let the group consumer dispatch once + return queued + + queued = asyncio.run(go()) + runner_key = runner._session_key_for_source(source) + expected_ns = f"agent:{owner}" if owner else "agent:main" + assert runner_key.startswith(expected_ns + ":"), runner_key + assert queued == [runner_key], (queued, runner_key) + assert seen == [runner_key] + assert adapter._processing_msg_ids == {runner_key: "m1"} diff --git a/tests/gateway/test_slack.py b/tests/gateway/test_slack.py index 851acf5005..790014273f 100644 --- a/tests/gateway/test_slack.py +++ b/tests/gateway/test_slack.py @@ -3070,7 +3070,10 @@ class TestThreadReplyHandling: from gateway.session import SessionEntry # Deserialize a legacy routing entry so lifecycle flags have real defaults. + # The thread key with a per-user suffix comes from the adapter's isolation flags (the runner + # seeds them into PlatformConfig.extra); this store has no bearing on the key any more. session_key = "agent:main:slack:group:T_TEAM:C123:123.000:U_USER" + adapter_with_session_store.config.extra["thread_sessions_per_user"] = True mock_session_store._entries = {session_key: SessionEntry.from_dict({ "session_key": session_key, "session_id": "slack-thread-session", diff --git a/tests/scripts/test_check_profile_scope_patterns.py b/tests/scripts/test_check_profile_scope_patterns.py index 7e5bb66322..1e13b8bb03 100644 --- a/tests/scripts/test_check_profile_scope_patterns.py +++ b/tests/scripts/test_check_profile_scope_patterns.py @@ -59,6 +59,30 @@ def test_child_env_from_environ_is_flagged_and_the_scoped_builder_is_not(): assert [f.line for f in mod.scan_text("tools/x.py", hazard, patterns, lines={5})] == [5] +def test_adapter_key_outside_the_seam_is_flagged_only_under_platforms(): + """An adapter that derives a session key with the free ``build_session_key()`` bypasses the + owner-profile seam (``_source_session_key``); the same call in ``platforms/base.py`` (the seam + itself) or in the runner is legitimate and must stay silent.""" + mod = _load() + patterns = mod.load_patterns() + free_key = textwrap.dedent(''' + from gateway.session import build_session_key + + def _batch_key(self, event): + return build_session_key(event.source, profile=event.source.profile) + ''') + seam = textwrap.dedent(''' + def _batch_key(self, event): + return self._event_session_key(event) + ''') + flagged = mod.scan_text("plugins/platforms/acme/adapter.py", free_key, patterns) + assert [(f.line, f.pattern_id, f.pattern_class) for f in flagged] == [(5, "P32", "C4")] + assert [f.pattern_id for f in mod.scan_text("gateway/platforms/acme.py", free_key, patterns)] == ["P32"] + assert mod.scan_text("plugins/platforms/acme/adapter.py", seam, patterns) == [] + for owner in ("gateway/platforms/base.py", "gateway/run_startup.py", "gateway/session_recovery.py"): + assert not [f for f in mod.scan_text(owner, free_key, patterns) if f.pattern_id == "P32"], owner + + def test_lint_is_advisory_and_exits_zero_with_findings(tmp_path, capsys): mod = _load() bad = tmp_path / "bad.py"