From a189286aea8abbf33dc93e2871cf67d8d87b0d2b Mon Sep 17 00:00:00 2001 From: Ben Barclay Date: Wed, 12 Aug 2026 10:04:26 +1000 Subject: [PATCH] fix(relay): stop sibling gateways answering another instance's button press (#83677) * fix(relay): stop sibling gateways answering another instance's button press A Discord button press arrives on the passthrough plane, and the connector fans a passthrough forward out to EVERY live gateway session of the tenant (relayServer.routeBusMessage delivers `passthrough` via sessionsByTenant), unlike a message, which it narrows to the admitted instance set. The prompt went out from exactly one instance and _pending_prompts is process-local, so every sibling gateway saw an answer for a prompt it never minted, could not tell that from its own prompt expiring, and fell through to chat dispatch -- where the option-shaped text ("/c1") is not a real command and run.py replied "Unknown command `/c1`". One copy per sibling, under the single real ack. Prompt ids are now minted as `.<8 hex>`, so an answer can be attributed to the process that minted it. A prompt answer is always consumed, never re-dispatched as chat: a sibling's prompt and a repeat answer are both dropped silently, and an expired prompt of our own gets a short "no longer waiting" notice from the owning gateway only. Ids stay inside the connector codec's contract ([A-Za-z0-9_.-], <=32 chars, 64-byte callback budget -- verified against promptCodec.ts: 52 bytes worst case with a full-length option id). An id with no nonce segment (a prompt in flight across an in-place upgrade) is still treated as ours. Tests: 4 added, each verified to fail without the fix. Full relay suite green (160 tests). * style(tests): ruff-format the added relay prompt tests --- docs/relay-connector-contract.md | 19 ++- gateway/relay/adapter.py | 140 ++++++++++++++++-- tests/gateway/relay/test_relay_interactive.py | 97 ++++++++++++ 3 files changed, 236 insertions(+), 20 deletions(-) diff --git a/docs/relay-connector-contract.md b/docs/relay-connector-contract.md index 184ce905e6..1d6242cfd4 100644 --- a/docs/relay-connector-contract.md +++ b/docs/relay-connector-contract.md @@ -440,21 +440,30 @@ clarify pickers) with NATIVE controls: Discord button components, Telegram inline keyboards, Slack Block Kit actions, WhatsApp button messages (≤3 options) / list messages (4–10; >10 degrades to the numbered-text fallback). `prompt_kind` (`approval`/`clarify`/`choice`) is a styling hint only. -`prompt_id` is gateway-minted (8 hex) and opaque to the connector; each +`prompt_id` is gateway-minted and opaque to the connector; each option's callback payload carries the token `hp1::` (≤64 bytes — Telegram's `callback_data` cap binds every lane; option ids are -`[A-Za-z0-9_.-]`, ≤32 chars). `style` maps per-platform +`[A-Za-z0-9_.-]`, ≤32 chars). The gateway mints the prompt id as +`.<8 hex>` within that same alphabet and budget: the +connector fans a passthrough forward (a Discord press) out to EVERY live +gateway session of the tenant, unlike a message, which it narrows to the +admitted instance set, so the nonce is how a gateway tells its OWN prompt +from a sibling's. `style` maps per-platform (primary/success/danger/secondary). `timeout_s` is advisory on the wire — expiry is enforced GATEWAY-side (the pending-prompt registry drops expired -entries; a stale press falls through as typed text, mirroring the native -adapters' "approval expired" edit). +entries; the owning gateway then replies with a short "no longer waiting" +notice). **`prompt_response` (Phase 3 inbound).** The user's press crosses back as a normal inbound MessageEvent carrying `prompt_response: {prompt_id, option_id, label?, prompt_message_id?}` — never a bare platform `custom_id`. The event's `text` mirrors `/{option_id}` with `message_type: "command"` so a gateway predating the field routes the press -as a typed reply instead of dropping it. The SOURCE is the authentic +as a typed reply instead of dropping it. A gateway that DOES understand the +field always consumes the press instead: a prompt id it did not mint belongs +to a sibling gateway that the same fan-out also reached, and letting the +`/{option_id}` text reach the chat lane made every sibling answer +"Unknown command" under the owner's single ack. The SOURCE is the authentic CLICKING user (connector-observed: Telegram `callback_query.from`, Slack `block_actions.user`, WhatsApp `messages[].from`, Discord interaction member/user), so gateway-side authorization gates apply to a button press diff --git a/gateway/relay/adapter.py b/gateway/relay/adapter.py index 5c9507031f..703cea0470 100644 --- a/gateway/relay/adapter.py +++ b/gateway/relay/adapter.py @@ -20,7 +20,9 @@ from __future__ import annotations import asyncio import logging +import secrets import time +from collections import OrderedDict from typing import Any, Callable, Dict, Optional, Tuple, cast from gateway.config import Platform, PlatformConfig @@ -40,6 +42,13 @@ logger = logging.getLogger(__name__) _RELAY_GO_IDLE_ON_DISCONNECT_TIMEOUT_S = 2.0 _RELAY_REVOCATION_MONITOR_TEARDOWN_TIMEOUT_S = 1.0 +# How many already-answered prompt ids to remember, so a duplicate answer for +# one of them (a double tap, or a connector redelivery of the same forward) is +# recognized as a repeat rather than mistaken for a stale prompt. One short +# string per entry; sized well above the number of prompts a single gateway can +# have in flight, and far below anything that matters for memory. +_RESOLVED_PROMPT_MEMORY = 256 + def _utf16_len(text: str) -> int: """Count UTF-16 code units (Telegram's length unit).""" @@ -133,8 +142,29 @@ class RelayAdapter(BasePlatformAdapter): # waiting primitive (approval / slash-confirm / clarify) so the click # resolves EXACTLY like the native adapters' button callbacks. # Entries expire lazily (see _pop_prompt) so an unanswered prompt - # never leaks. Keyed by our own minted 8-hex ids. + # never leaks. Keyed by our own minted ids (see _prompt_owner_nonce). self._pending_prompts: Dict[str, Dict[str, Any]] = {} + # Per-process marker prefixed onto every prompt id this adapter mints, + # so an answer can be recognized as OURS before it is acted on. + # + # WHY: a button press arrives on the passthrough plane, which the + # connector fans out to EVERY live gateway session of the tenant + # (relayServer.routeBusMessage delivers `passthrough` via + # sessionsByTenant, unlike `message`, which narrows to the admitted + # instance set). The prompt itself went out from exactly one instance, + # and _pending_prompts is process-local — so every OTHER instance sees + # an answer for a prompt it never minted. Without this marker those + # instances cannot tell "a sibling owns this" from "my own prompt + # expired", and the id-shaped text ("/c1") falls through to chat + # dispatch, where run.py answers "Unknown command `/c1`" — one copy per + # sibling gateway. Siblings are the common case in a DM: the connector + # resolves the invoker's bindings to DISTINCT TENANTS, so several + # instances of one tenant all receive the forward. + self._prompt_owner_nonce: str = secrets.token_hex(3) + # Prompt ids this process has already resolved, newest last. A repeat + # answer for the same id (double tap, or a connector redelivery) is + # then consumed silently instead of being treated as a stale prompt. + self._resolved_prompts: "OrderedDict[str, float]" = OrderedDict() # ── capability surface (from descriptor) ───────────────────────────── @property @@ -1681,17 +1711,22 @@ class RelayAdapter(BasePlatformAdapter): def _mint_prompt( self, kind: str, state: Dict[str, Any], timeout_s: float = 3600.0 ) -> str: - """Register a pending prompt and return its 8-hex id. + """Register a pending prompt and return its id. ``state`` carries what the resolver needs when the answer comes back (session_key, resolver kind, per-kind extras). Expiry is enforced gateway-side on consumption (_pop_prompt) — the wire's timeout_s is advisory only. + + The id is ``.<8 hex>``: the nonce marks the minting + process so a sibling gateway that receives the same fanned-out answer + can tell it is not the owner and stay quiet (see + ``_prompt_owner_nonce``). Both segments use the callback alphabet the + connector's prompt codec accepts ([A-Za-z0-9_.-], <=32 chars). """ - import secrets import time - prompt_id = secrets.token_hex(4) + prompt_id = f"{self._prompt_owner_nonce}.{secrets.token_hex(4)}" self._pending_prompts[prompt_id] = { **state, "kind": kind, @@ -1706,6 +1741,16 @@ class RelayAdapter(BasePlatformAdapter): self._pending_prompts.pop(stale, None) return prompt_id + def _minted_here(self, prompt_id: str) -> bool: + """True when this process minted ``prompt_id``. + + Ids minted before the owner nonce existed (a prompt still pending + across an in-place upgrade) carry no ``.`` segment; treat them as ours + so an in-flight prompt from the previous build still resolves. + """ + head, sep, _ = str(prompt_id).partition(".") + return head == self._prompt_owner_nonce if sep else True + def _pop_prompt(self, prompt_id: str) -> Optional[Dict[str, Any]]: """Consume a pending prompt: one answer wins, expired entries miss.""" import time @@ -1717,6 +1762,17 @@ class RelayAdapter(BasePlatformAdapter): return None return state + def _note_prompt_resolved(self, prompt_id: str) -> None: + """Remember that this process already answered ``prompt_id``. + + Bounded FIFO: a repeat answer is only interesting for as long as a + redelivery or a double tap can plausibly arrive, so old ids are + dropped rather than retained for the process lifetime. + """ + self._resolved_prompts[str(prompt_id)] = time.time() + while len(self._resolved_prompts) > _RESOLVED_PROMPT_MEMORY: + self._resolved_prompts.popitem(last=False) + async def _send_prompt( self, chat_id: str, @@ -1944,11 +2000,25 @@ class RelayAdapter(BasePlatformAdapter): """Route an inbound prompt_response to its waiting primitive. Returns True when the event was a prompt answer (consumed — do NOT - dispatch it as a chat message), False otherwise. Unknown/expired - prompt ids fall through to normal dispatch: the command-shaped text - ("/once", "/deny", …) then behaves like a typed reply, which the - approval/confirm text lanes already understand — the same stale-tap - degradation the native adapters implement with an "expired" edit. + dispatch it as a chat message), False otherwise. + + A prompt answer is ALWAYS consumed, whoever it belongs to. The three + non-resolving cases each stay silent rather than falling through to + chat dispatch: + + * **A sibling's prompt** (``_minted_here`` false). The connector fans a + passthrough forward out to every live session of the tenant, so one + button press reaches every gateway of that tenant while only the + minting one can resolve it. Falling through here is what produced a + wall of ``Unknown command `/c1``` — one per sibling — under the + single ``✅`` from the real owner. + * **A repeat answer** for an id this process already resolved (a double + tap, or a redelivered forward). The first answer won; a second must + not re-run the resolver or re-ack. + * **Our own expired/unknown prompt.** Answer with a short expiry notice + instead of a command-not-found error: the option ids a prompt renders + ("c1", "once", "other") are not commands, so the text lane cannot do + anything useful with them. """ pr = getattr(event, "prompt_response", None) if not isinstance(pr, dict): @@ -1957,15 +2027,35 @@ class RelayAdapter(BasePlatformAdapter): option_id = str(pr.get("option_id") or "") if not prompt_id or not option_id: return False - state = self._pop_prompt(prompt_id) - if state is None: - logger.info( - "relay prompt_response for unknown/expired prompt %s (option=%s) — " - "falling through to text dispatch", + if not self._minted_here(prompt_id): + # A sibling gateway of this tenant owns this prompt and is + # resolving it right now. Consume silently — two gateways must + # never both answer one press. + logger.debug( + "relay prompt_response %s (option=%s) belongs to another " + "gateway instance — ignoring", prompt_id, option_id, ) - return False + return True + if prompt_id in self._resolved_prompts: + logger.debug( + "relay prompt_response %s (option=%s) already resolved — ignoring " + "repeat", + prompt_id, + option_id, + ) + return True + state = self._pop_prompt(prompt_id) + if state is None: + logger.info( + "relay prompt_response for unknown/expired prompt %s (option=%s)", + prompt_id, + option_id, + ) + await self._notify_prompt_expired(event) + return True + self._note_prompt_resolved(prompt_id) kind = state.get("kind") chat_id = str(state.get("chat_id") or getattr(event.source, "chat_id", "")) @@ -2056,6 +2146,26 @@ class RelayAdapter(BasePlatformAdapter): logger.warning("relay prompt_response resolution failed", exc_info=True) return True + async def _notify_prompt_expired(self, event) -> None: + """Tell the presser their prompt is no longer waiting. + + Only the OWNING gateway reaches this (siblings return earlier), so the + user sees exactly one notice. Best-effort: a send failure here must + not break the reader. + """ + chat_id = str(getattr(event.source, "chat_id", "") or "") + if not chat_id: + return + try: + await self.send( + chat_id, + "⌛ That prompt is no longer waiting for an answer. " + "Send your reply as a normal message.", + metadata=self._prompt_reply_metadata(event), + ) + except Exception: # noqa: BLE001 - notification is best-effort + logger.debug("relay expired-prompt notice failed", exc_info=True) + def _prompt_reply_metadata(self, event) -> Dict[str, Any]: """Thread/topic metadata so prompt acks land where the prompt lives.""" meta: Dict[str, Any] = {} diff --git a/tests/gateway/relay/test_relay_interactive.py b/tests/gateway/relay/test_relay_interactive.py index 7cabcb6eb3..80151d32a5 100644 --- a/tests/gateway/relay/test_relay_interactive.py +++ b/tests/gateway/relay/test_relay_interactive.py @@ -257,3 +257,100 @@ async def test_processing_lifecycle_reacts_eyes_then_check(): assert all(r["message_id"] == "m42" and r["chat_id"] == "ch1" for r in reacts) +# ── fanned-out prompt answers (one press, many gateways) ───────────────── +# +# The connector delivers a passthrough forward (a Discord button press) to +# EVERY live gateway session of the tenant, unlike a message, which it narrows +# to the admitted instance set. So one press reaches every sibling gateway +# while only the minting one can resolve it. These pin that a non-owner stays +# silent and that the owner still answers exactly once. + + +@pytest.mark.asyncio +async def test_sibling_gateway_ignores_another_instances_prompt_answer(monkeypatch): + """A press for a prompt this process didn't mint is consumed silently.""" + owner, owner_stub = _adapter() + sibling, sibling_stub = _adapter() + await owner.send_clarify("c1", "Which?", ["alpha", "beta"], "cl-1", "s") + prompt_id = owner_stub.sent[-1]["prompt_id"] + + resolved: list[tuple] = [] + monkeypatch.setattr( + "tools.clarify_gateway.resolve_gateway_clarify", + lambda cid, resp: resolved.append((cid, resp)) or True, + ) + monkeypatch.setattr("tools.clarify_gateway.mark_awaiting_text", lambda cid: None) + + event = _event({"prompt_id": prompt_id, "option_id": "c1"}) + # Consumed (True) so the "/c1"-shaped text is never dispatched as chat -- + # that fall-through is what produced one "Unknown command `/c1`" per + # sibling gateway. And the sibling neither resolves nor says anything. + assert await sibling._consume_prompt_response(event) is True + assert resolved == [] + assert sibling_stub.sent == [] + + # The owner still resolves the same press normally. + assert await owner._consume_prompt_response(event) is True + assert resolved == [("cl-1", "beta")] + + +@pytest.mark.asyncio +async def test_repeat_answer_for_resolved_prompt_is_ignored(monkeypatch): + """A double tap / redelivered forward must not resolve twice.""" + adapter, stub = _adapter() + await adapter.send_clarify("c1", "Which?", ["alpha", "beta"], "cl-2", "s") + prompt_id = stub.sent[-1]["prompt_id"] + + resolved: list[tuple] = [] + monkeypatch.setattr( + "tools.clarify_gateway.resolve_gateway_clarify", + lambda cid, resp: resolved.append((cid, resp)) or True, + ) + event = _event({"prompt_id": prompt_id, "option_id": "c0"}) + assert await adapter._consume_prompt_response(event) is True + assert resolved == [("cl-2", "alpha")] + + sent_after_first = len(stub.sent) + assert await adapter._consume_prompt_response(event) is True + assert resolved == [("cl-2", "alpha")] # not resolved a second time + assert len(stub.sent) == sent_after_first # and no second ack / notice + + +@pytest.mark.asyncio +async def test_expired_own_prompt_notifies_instead_of_unknown_command(): + """An expired prompt of OURS gets an expiry notice, not chat dispatch. + + Falling through would hand run.py a command-shaped "/c1", which is not a + real command, so the user got "Unknown command `/c1`". + """ + adapter, stub = _adapter() + prompt_id = adapter._mint_prompt("clarify", {"chat_id": "c1"}, timeout_s=-1.0) + + event = _event({"prompt_id": prompt_id, "option_id": "c1"}) + assert await adapter._consume_prompt_response(event) is True + notices = [a for a in stub.sent if a["op"] == "send"] + assert len(notices) == 1 + assert "no longer waiting" in notices[0]["content"] + + +def test_minted_prompt_ids_are_instance_scoped_and_callback_safe(): + """Ids carry the minting process's nonce and stay codec-legal. + + The connector's promptCodec validates each id as [A-Za-z0-9_.-]{1,32} and + caps "hp1::" at Telegram's 64-byte callback budget. + """ + import re + + a, _ = _adapter() + b, _ = _adapter() + id_a = a._mint_prompt("clarify", {"chat_id": "c1"}) + id_b = b._mint_prompt("clarify", {"chat_id": "c1"}) + + assert re.fullmatch(r"[A-Za-z0-9_.\-]{1,32}", id_a) + assert len(f"hp1:{id_a}:option_id_up_to_32_chars_here") <= 64 + assert a._minted_here(id_a) is True + assert b._minted_here(id_a) is False + assert a._minted_here(id_b) is False + # A legacy id minted before the nonce existed (no "." segment) is still + # treated as ours, so a prompt in flight across an upgrade resolves. + assert a._minted_here("a1b2c3d4") is True