diff --git a/gateway/relay/adapter.py b/gateway/relay/adapter.py index 1fdeea8d93..ab68ea7699 100644 --- a/gateway/relay/adapter.py +++ b/gateway/relay/adapter.py @@ -29,6 +29,7 @@ from typing import Any, Callable, Dict, Optional, Tuple, cast from gateway.config import Platform, PlatformConfig from gateway.platforms.base import BasePlatformAdapter, MessageEvent, SendResult from gateway.relay.descriptor import CapabilityDescriptor +from gateway.relay.egress import decline_error, is_egress_decline, log_decline from gateway.relay.media import RelayMediaClient from gateway.relay.transport import RelayTransport from gateway.session import SessionSource @@ -869,7 +870,12 @@ class RelayAdapter(BasePlatformAdapter): return SendResult( success=False, error=f"task_card_stop transport error: {e}" ) - return SendResult(success=bool(result.get("success"))) + if is_egress_decline(result): + log_decline("task_card_stop", chat_id, result) + return SendResult( + success=bool(result.get("success")), + error=result.get("error"), + ) async def abandon_open_draft( self, @@ -2334,6 +2340,10 @@ class RelayAdapter(BasePlatformAdapter): except Exception: logger.debug("relay delete_message failed", exc_info=True) return False + if is_egress_decline(result): + # Cosmetic lane (bool contract) — it degrades by design, but an + # AUTHORIZATION refusal must not vanish into debug-level noise. + log_decline("delete", chat_id, result) return bool(result.get("success")) async def send_typing(self, chat_id: str, metadata=None) -> None: @@ -2396,12 +2406,15 @@ class RelayAdapter(BasePlatformAdapter): if phrase: frame["content"] = str(phrase) try: - await self._transport.send_outbound( + _typing_result = await self._transport.send_outbound( frame, platform=self._platform_by_chat.get(str(chat_id)), ) except Exception: # noqa: BLE001 - typing is cosmetic, never breaks a turn logger.debug("relay send_typing failed for %s", chat_id, exc_info=True) + else: + if is_egress_decline(_typing_result): + log_decline("typing", chat_id, _typing_result) async def stop_typing( self, @@ -2431,7 +2444,7 @@ class RelayAdapter(BasePlatformAdapter): # timeout. Shared helper with send_typing so the two guards cannot drift. md = self._with_status_thread_anchor(chat_id, metadata) try: - await self._transport.send_outbound( + _clear_result = await self._transport.send_outbound( { "op": "typing", "chat_id": chat_id, @@ -2442,6 +2455,9 @@ class RelayAdapter(BasePlatformAdapter): ) except Exception: # noqa: BLE001 - status clear is cosmetic, never breaks a turn logger.debug("relay stop_typing failed for %s", chat_id, exc_info=True) + else: + if is_egress_decline(_clear_result): + log_decline("typing_clear", chat_id, _clear_result) async def get_chat_info(self, chat_id: str) -> Dict[str, Any]: # Proxied to the connector (it owns the platform connection / cache). @@ -2596,6 +2612,19 @@ class RelayAdapter(BasePlatformAdapter): logger.debug("relay send_media transport failure", exc_info=True) return None if not result.get("success"): + if is_egress_decline(result): + # AUTHORIZATION refusal, not a size cap / platform rejection. + # Returning None here would hand the caller back to its + # text-notice fallback — a DIFFERENT op against the SAME + # destination the connector just refused, laundering a + # security decision into "media unsupported". Surface the + # decline verbatim as a failed lane instead. + log_decline("send_media", chat_id, result) + return SendResult( + success=False, + error=decline_error(result), + raw_response=result, + ) # A structured connector decline (size cap, platform rejection). # Surface it as a failed lane so the caller's fallback still # delivers SOMETHING (the caption/notice), mirroring native @@ -2884,6 +2913,17 @@ class RelayAdapter(BasePlatformAdapter): logger.debug("relay prompt transport failure", exc_info=True) return None if not result.get("success"): + if is_egress_decline(result): + # AUTHORIZATION refusal. Returning None here reported "prompt + # op unavailable" to the caller — a WRONG reason — and sent + # the caller's numbered-text fallback into the very chat the + # connector just refused. Surface the decline verbatim. + log_decline("prompt", chat_id, result) + return SendResult( + success=False, + error=decline_error(result), + raw_response=result, + ) logger.warning( "relay prompt declined for %s: %s", chat_id, result.get("error") ) @@ -2949,10 +2989,16 @@ class RelayAdapter(BasePlatformAdapter): options=options, metadata=metadata, ) - if result is not None: + if result is not None and result.success: return result - # Lane unavailable: unregister and let run.py's text fallback run. + # Nothing was delivered — unregister the minted prompt either way. self._pending_prompts.pop(prompt_id, None) + if result is not None: + # A DECLINE (P5b): surface it verbatim. Reporting "op unavailable" + # here would send run.py's text fallback into the very chat the + # connector refused, turning a refusal into an apparent success. + return result + # Lane unavailable: let run.py's text fallback run. return SendResult(success=False, error="relay prompt op unavailable") async def send_slash_confirm( @@ -2992,9 +3038,13 @@ class RelayAdapter(BasePlatformAdapter): options=options, metadata=metadata, ) - if result is not None: + if result is not None and result.success: return result self._pending_prompts.pop(prompt_id, None) + if result is not None: + # DECLINE (P5b): verbatim, not "op unavailable" — the gateway's + # text-intercept fallback must not re-address the refused chat. + return result return SendResult(success=False, error="relay prompt op unavailable") async def send_clarify( @@ -3042,9 +3092,13 @@ class RelayAdapter(BasePlatformAdapter): options=options, metadata=metadata, ) - if result is not None: + if result is not None and result.success: return result self._pending_prompts.pop(prompt_id, None) + if result is not None: + # DECLINE (P5b): the base class's numbered-text fallback would + # send() into the same refused chat. Surface the refusal. + return result return await super().send_clarify( chat_id, question, choices, clarify_id, session_key, metadata=metadata ) @@ -3309,6 +3363,8 @@ class RelayAdapter(BasePlatformAdapter): }, platform=self._platform_by_chat.get(str(chat_id)), ) + if is_egress_decline(result): + log_decline("react", chat_id, result) return bool(result.get("success")) except Exception: # noqa: BLE001 - reactions are cosmetic logger.debug("relay react failed", exc_info=True) @@ -3372,11 +3428,14 @@ class RelayAdapter(BasePlatformAdapter): logger.debug("relay thread_create transport failure", exc_info=True) return None if not result.get("success"): - logger.info( - "relay thread_create declined for %s: %s", - parent_chat_id, - result.get("error"), - ) + if is_egress_decline(result): + log_decline("thread_create", parent_chat_id, result) + else: + logger.info( + "relay thread_create declined for %s: %s", + parent_chat_id, + result.get("error"), + ) return None thread_id = result.get("thread_id") or result.get("message_id") return str(thread_id) if thread_id else None @@ -3434,10 +3493,13 @@ class RelayAdapter(BasePlatformAdapter): logger.debug("relay thread_rename transport failure", exc_info=True) return False if not result.get("success"): - logger.info( - "relay thread_rename declined for %s: %s", - thread_id, - result.get("error"), - ) + if is_egress_decline(result): + log_decline("thread_rename", thread_id, result) + else: + logger.info( + "relay thread_rename declined for %s: %s", + thread_id, + result.get("error"), + ) return False return True diff --git a/gateway/relay/egress.py b/gateway/relay/egress.py new file mode 100644 index 0000000000..4ed8121089 --- /dev/null +++ b/gateway/relay/egress.py @@ -0,0 +1,233 @@ +"""Gateway-side relay EGRESS AUTHORIZATION (P5). + +Two halves of one security surface, kept together because they are the same +concern seen from both ends of the wire: + +**(a) Target hygiene.** The ``send_message`` model tool takes a free-form +``target`` string (``'platform:chat_id'``). Nothing stopped a model from +naming an arbitrary chat id and having the gateway emit an outbound frame for +it. The connector now refuses such a frame (its egress-authorization floor), +but the gateway must not *silently* ask: a destination the gateway has no +record of is refused HERE, with a visible tool error naming the target. +:func:`authorize_relay_target` is that guard; :func:`attested_relay_targets` +is the set of destinations this gateway can show a provenance for (operator +home channel, channel directory, its own gateway session origins). + +**(b) Decline visibility.** The connector answers an unauthorized destination +with a DEFINITE, non-ambiguous failure whose text is deliberately UNIFORM +(``" egress declined: target is not an approved destination for +this connection"``) — it must not leak whether the destination belongs to +another tenant or to nobody. :func:`is_egress_decline` recognises THAT a +decline happened; it never tries to parse WHY. Egress lanes that legitimately +degrade a *transport drop* (advisory progress, cosmetic reactions, media +falling back to a text notice) must NOT degrade a decline — a refusal that +turns into a different op, or into a wrong "op unavailable" reason, is a +security event laundered into an apparent success. +""" + +from __future__ import annotations + +import logging +from typing import Any, Dict, Optional, Set + +logger = logging.getLogger(__name__) + +# The connector stamps this on a structured decline (preferred signal). +EGRESS_DECLINE_CODE = "egress_declined" + +# Fallback signal: the connector's uniform decline sentence. Matched +# case-insensitively on this fragment ONLY — the rest of the sentence is +# deliberately uninformative (finding F-005) and must not be parsed. +EGRESS_DECLINE_MARKER = "egress declined:" + + +def is_egress_decline(result: Any) -> bool: + """True when *result* is the connector REFUSING the destination. + + A decline is distinguished from every other outbound failure by three + properties, all required: + + * it failed (``success`` is falsey), + * it is DEFINITE — an ``ambiguous`` result (lost ack, mid-write drop) says + the frame may well have been applied, so it is a transport outcome, not + an authorization one, + * it carries the connector's decline code, or its uniform decline text. + """ + if not isinstance(result, dict): + return False + if result.get("success"): + return False + if result.get("ambiguous"): + return False + if str(result.get("code") or "") == EGRESS_DECLINE_CODE: + return True + return EGRESS_DECLINE_MARKER in str(result.get("error") or "").lower() + + +def decline_error(result: Any) -> str: + """The connector's decline text, verbatim, for surfacing to the caller. + + Verbatim on purpose: the gateway's job is to report faithfully THAT a + decline happened, not to explain or re-word it. + """ + if isinstance(result, dict): + error = result.get("error") + if error: + return str(error) + return "relay egress declined" + + +def log_decline(op: str, chat_id: Any, result: Any) -> None: + """Record a decline on a lane whose contract cannot carry an error. + + Cosmetic lanes (typing, reactions, delete, thread ops) return ``bool`` / + ``None`` by contract and legitimately degrade. A decline there is still a + security-relevant event, so it is logged at WARNING rather than vanishing + into the lane's debug-level best-effort handling. + """ + logger.warning( + "relay %s DECLINED for %s: %s", op, chat_id, decline_error(result) + ) + + +# --------------------------------------------------------------------------- +# (a) target attestation +# --------------------------------------------------------------------------- + +def _relay_fronted() -> Set[str]: + try: + from gateway.relay import relay_fronted_platforms + + return {str(p) for p in relay_fronted_platforms()} + except Exception: # noqa: BLE001 - env/config absence must never break a send + return set() + + +def _has_live_native_adapter(platform_name: str) -> bool: + """Whether THIS process runs a native (non-relay) adapter for the platform. + + Mirrors ``gateway/delivery.resolve_delivery_transport``'s precedence: a + concrete native adapter always wins, so a platform served natively here is + not a relay egress and this guard does not apply to it. + """ + try: + from gateway.config import Platform + from gateway.run import _gateway_runner_ref + + runner = _gateway_runner_ref() + if runner is None: + return False + adapters = getattr(runner, "adapters", None) or {} + return adapters.get(Platform(platform_name)) is not None + except Exception: # noqa: BLE001 - no runner (cron/CLI) ⇒ no native adapter + return False + + +def relay_routed_platform(platform_name: str) -> bool: + """Whether a send to *platform_name* would egress over the relay connector. + + True for the generic ``relay`` plane itself, and for any logical platform + the connector fronts for this gateway that has no live native adapter in + this process. + """ + name = str(platform_name or "").strip().lower() + if not name: + return False + if name == "relay": + return True + if name not in _relay_fronted(): + return False + return not _has_live_native_adapter(name) + + +def _home_channel_id(platform_name: str) -> Optional[str]: + try: + from gateway.config import Platform, load_gateway_config + + home = load_gateway_config().get_home_channel(Platform(platform_name)) + return str(home.chat_id) if home and home.chat_id else None + except Exception: # noqa: BLE001 - config absence must never break a send + return None + + +def _directory_ids(platform_name: str) -> Set[str]: + try: + from gateway.channel_directory import load_directory + + entries = load_directory().get("platforms", {}).get(platform_name) or [] + except Exception: # noqa: BLE001 + return set() + ids: Set[str] = set() + for entry in entries: + if isinstance(entry, dict) and entry.get("id"): + ids.add(str(entry["id"])) + return ids + + +def _session_ids(platform_name: str) -> Set[str]: + """Chat ids this gateway has actually held a session in for the platform.""" + try: + from gateway.channel_directory import _build_from_sessions + + entries = _build_from_sessions(platform_name) or [] + except Exception: # noqa: BLE001 + return set() + ids: Set[str] = set() + for entry in entries: + if isinstance(entry, dict) and entry.get("id"): + # Session entry ids may be thread-qualified ("chat:thread"); the + # destination the connector authorizes is the CHAT, so attest both. + raw = str(entry["id"]) + ids.add(raw) + ids.add(raw.split(":", 1)[0]) + return ids + + +def attested_relay_targets(platform_name: str) -> Set[str]: + """Chat ids this gateway can show a provenance for on *platform_name*. + + Three provenances, all of them things the gateway already knows rather + than things a model can invent: + + * the operator-configured home channel, + * the channel directory built from live adapters/session origins, + * this gateway's own gateway-session origins. + + For the generic ``relay`` plane the union spans every platform the + connector fronts: a relay session is filed under its LOGICAL platform + (``source = "discord"``), so attesting ``relay`` against only ``relay`` + would refuse chats the agent is demonstrably already talking in. + """ + name = str(platform_name or "").strip().lower() + names = {name} + if name == "relay": + names |= _relay_fronted() + attested: Set[str] = set() + for candidate in names: + home = _home_channel_id(candidate) + if home: + attested.add(home) + attested |= _directory_ids(candidate) + attested |= _session_ids(candidate) + return attested + + +def authorize_relay_target(platform_name: str, chat_id: Any) -> Optional[str]: + """Return an error string when this relay destination may not be named. + + ``None`` means the send may proceed. Non-relay platforms are never + restricted here — their own adapters own their authorization. + """ + if not relay_routed_platform(platform_name): + return None + target = str(chat_id or "").strip() + if not target: + return None + name = str(platform_name).strip().lower() + if target in attested_relay_targets(name): + return None + return ( + f"Refusing to send to unattested relay target '{name}:{target}': " + "this gateway has no record of that destination. Use " + "send_message(action='list') to see the targets it can reach." + ) diff --git a/tests/gateway/relay/test_relay_egress_declines.py b/tests/gateway/relay/test_relay_egress_declines.py new file mode 100644 index 0000000000..a212b99326 --- /dev/null +++ b/tests/gateway/relay/test_relay_egress_declines.py @@ -0,0 +1,271 @@ +"""P5(b): connector egress DECLINES surface as real errors, not swallowed successes. + +The connector's egress-authorization floor answers an unauthorized destination +with a DEFINITE failure whose text is deliberately uniform (finding F-005 — the +caller must not learn *why*). The gateway's job is to report faithfully THAT it +happened. The failure mode these tests pin is specific: several relay lanes +degrade a *transport drop* by design, and that same degradation used to swallow +an *authorization refusal* — either into a wrong reason ("prompt op +unavailable") or, worse, into a DIFFERENT op re-addressed at the very chat the +connector just refused (media falling back to a plain text notice). + +Every lane below drives the REAL `RelayAdapter` built from the REAL descriptor; +the only substitution is the transport, which is what the connector is. +""" + +from __future__ import annotations + +import asyncio +import logging +from typing import Any, Dict, List, Optional + +import pytest + +from gateway.config import PlatformConfig +from gateway.relay.adapter import RelayAdapter +from gateway.relay.descriptor import CONTRACT_VERSION, CapabilityDescriptor +from gateway.relay.egress import ( + EGRESS_DECLINE_CODE, + decline_error, + is_egress_decline, +) + +DECLINE_TEXT = ( + "discord egress declined: target is not an approved destination for this connection" +) +DECLINE: Dict[str, Any] = {"success": False, "error": DECLINE_TEXT} + +ALL_OPS = ( + "send", + "edit", + "typing", + "delete", + "react", + "send_media", + "prompt", + "draft", + "task_card", + "task_card_stop", + "thread_create", + "thread_rename", +) + + +class DecliningConnector: + """A connector that refuses EVERY destination, like the real egress floor.""" + + def __init__(self, descriptor: CapabilityDescriptor) -> None: + self._descriptor = descriptor + self.ops: List[str] = [] + self._identities = [(descriptor.platform, "b1")] + + async def connect(self, *, is_reconnect: bool = False) -> bool: + return True + + async def disconnect(self) -> None: + return None + + async def handshake(self) -> CapabilityDescriptor: + return self._descriptor + + def set_inbound_handler(self, handler) -> None: + return None + + def set_passthrough_handler(self, handler) -> None: + return None + + async def send_outbound( + self, action: Dict[str, Any], *, platform: Optional[str] = None + ) -> Dict[str, Any]: + self.ops.append(str(action.get("op"))) + return dict(DECLINE) + + async def send_follow_up( + self, action: Dict[str, Any], *, platform: Optional[str] = None + ) -> Dict[str, Any]: + return dict(DECLINE) + + async def get_chat_info(self, chat_id: str) -> Dict[str, Any]: + return {"name": chat_id, "type": "dm"} + + async def send_interrupt(self, session_key, reason=None) -> None: + return None + + async def go_idle(self, timeout_s: float = 10.0) -> bool: + return True + + +@pytest.fixture +def relay(): + descriptor = CapabilityDescriptor( + contract_version=CONTRACT_VERSION, + platform="discord", + label="Relay", + max_message_length=4096, + supports_draft_streaming=True, + supports_edit=True, + supports_threads=True, + markdown_dialect="plain", + len_unit="chars", + supported_ops=ALL_OPS, + ) + connector = DecliningConnector(descriptor) + adapter = RelayAdapter( + PlatformConfig(enabled=True, extra={}), descriptor, transport=connector + ) + return adapter, connector + + +# ── the classifier itself ──────────────────────────────────────────────── + +def test_decline_is_recognised_by_code_and_by_uniform_text(): + assert is_egress_decline({"success": False, "code": EGRESS_DECLINE_CODE}) is True + assert is_egress_decline(DECLINE) is True + + +def test_an_ambiguous_failure_is_not_a_decline(): + """A lost ack may well have been APPLIED — it is a transport outcome. + + Classifying it as a decline would convert the relay's deliberate + optimistic-retry behaviour into a hard error on a message that landed. + """ + assert ( + is_egress_decline( + {"success": False, "error": DECLINE_TEXT, "ambiguous": True} + ) + is False + ) + + +def test_an_ordinary_failure_is_not_a_decline(): + assert is_egress_decline({"success": False, "error": "file too large"}) is False + assert is_egress_decline({"success": True}) is False + assert is_egress_decline(None) is False + + +def test_decline_error_is_the_connector_text_verbatim(): + """No re-wording, no reason-parsing: the uniform sentence, unchanged.""" + assert decline_error(DECLINE) == DECLINE_TEXT + + +# ── lanes that must report the decline to their caller ─────────────────── + +def test_send_reports_the_decline(relay): + adapter, connector = relay + result = asyncio.run(adapter.send("C1", "hi")) + + assert result.success is False + assert result.error == DECLINE_TEXT + assert connector.ops == ["send"] + + +def test_media_decline_does_not_fall_back_to_a_text_send(relay): + """The worst swallow: a refused destination re-addressed by a different op. + + `_send_media` returning None hands the caller back to + `BasePlatformAdapter.send_image`, which sends the URL as TEXT — into the + very chat the connector just refused. The lane must report the refusal + instead, and emit exactly ONE op. + """ + adapter, connector = relay + result = asyncio.run(adapter.send_image("C1", "https://x/y.png", caption="cap")) + + assert result.success is False + assert result.error == DECLINE_TEXT + assert connector.ops == ["send_media"] + + +def test_exec_approval_decline_is_not_reported_as_op_unavailable(relay): + """"prompt op unavailable" is a WRONG reason that triggers a text fallback. + + The whole observable: the refusal reaches the caller verbatim, exactly one + op is emitted, and the minted prompt is UNREGISTERED — a prompt left + pending for a card that was never delivered would silently capture the + user's next reply in that chat as an approval press. + """ + adapter, connector = relay + result = asyncio.run(adapter.send_exec_approval("C1", "rm -rf /", "sk1")) + + assert result.success is False + assert result.error == DECLINE_TEXT + assert connector.ops == ["prompt"] + assert adapter._pending_prompts == {} + + +def test_slash_confirm_decline_is_not_reported_as_op_unavailable(relay): + adapter, connector = relay + result = asyncio.run( + adapter.send_slash_confirm("C1", "Title", "Body", "sk1", "cf1") + ) + + assert result.success is False + assert result.error == DECLINE_TEXT + assert connector.ops == ["prompt"] + assert adapter._pending_prompts == {} + + +def test_clarify_decline_does_not_fall_back_to_a_numbered_text_send(relay): + """The base class's numbered-text clarify would `send()` into the refused chat.""" + adapter, connector = relay + result = asyncio.run( + adapter.send_clarify("C1", "Which?", ["a", "b"], "cl1", "sk1") + ) + + assert result.success is False + assert result.error == DECLINE_TEXT + assert connector.ops == ["prompt"] + assert adapter._pending_prompts == {} + + +def test_a_delivered_prompt_stays_registered(relay, monkeypatch): + """Guard the converse: the cleanup must not unregister a LIVE prompt.""" + adapter, connector = relay + + async def _ok(action, *, platform=None): + connector.ops.append(str(action.get("op"))) + return {"success": True, "message_id": "pm1"} + + monkeypatch.setattr(connector, "send_outbound", _ok) + result = asyncio.run(adapter.send_exec_approval("C1", "ls", "sk1")) + + assert result.success is True + assert connector.ops == ["prompt"] + assert len(adapter._pending_prompts) == 1 + + +def test_task_card_stop_decline_carries_the_error(relay): + """The stop lane discarded the error entirely (`success=` only).""" + adapter, connector = relay + result = asyncio.run(adapter.stop_native_task_card_progress("C1")) + + assert result.success is False + assert result.error == DECLINE_TEXT + assert connector.ops == ["task_card_stop"] + + +# ── lanes that legitimately degrade, but must still SAY so ─────────────── + +@pytest.mark.parametrize( + "lane,call,expected", + [ + ("typing", lambda a: a.send_typing("C1"), None), + ("delete", lambda a: a.delete_message("C1", "m1"), False), + ("thread_create", lambda a: a.create_handoff_thread("C1", "n"), None), + ("thread_rename", lambda a: a.rename_thread("T1", "n"), False), + ], +) +def test_cosmetic_lane_degrades_but_logs_the_decline_at_warning( + relay, caplog, lane, call, expected +): + """These return bool/None by contract; a refusal must not vanish silently.""" + adapter, _connector = relay + with caplog.at_level(logging.WARNING, logger="gateway.relay.egress"): + assert asyncio.run(call(adapter)) == expected + + declines = [ + r + for r in caplog.records + if r.name == "gateway.relay.egress" and "DECLINED" in r.getMessage() + ] + assert len(declines) == 1 + assert DECLINE_TEXT in declines[0].getMessage() diff --git a/tests/tools/test_send_message_relay_target_authz.py b/tests/tools/test_send_message_relay_target_authz.py new file mode 100644 index 0000000000..8fb539948d --- /dev/null +++ b/tests/tools/test_send_message_relay_target_authz.py @@ -0,0 +1,200 @@ +"""P5(a): `send_message` cannot silently name an arbitrary relay target. + +The `target` tool parameter is free-form (`'platform:chat_id'`), so before +this guard a model could name ANY chat id and the gateway would emit an +outbound relay frame for it — authenticating the sender while never +authorizing the destination. These tests drive the REAL `send_message_tool` +entrypoint through the REAL production wiring (`gateway.relay.egress`, +`gateway.channel_directory`, `gateway.relay.relay_fronted_platforms`) against +a temp HERMES_HOME; nothing under test is constructed by the test itself. +""" + +from __future__ import annotations + +import json + +import pytest + +from gateway.config import Platform +from tools.send_message_tool import send_message_tool + +ATTESTED_CHAT = "111111111111111111" +ARBITRARY_CHAT = "999999999999999999" +HOME_CHAT = "222222222222222222" + + +@pytest.fixture +def relay_env(tmp_path, monkeypatch): + """A gateway whose ONLY reachable Discord destinations are attested. + + Mirrors the production shape: `GATEWAY_RELAY_PLATFORMS` is the deploy + stamp `gateway.relay.relay_fronted_platforms()` reads, the channel + directory json is the file `channel_directory.load_directory()` reads, and + no live native adapter exists in this process (so the relay owns egress + for `discord`, exactly as `gateway/delivery.resolve_delivery_transport` + decides it). + """ + import gateway.channel_directory as cd + + monkeypatch.setenv("GATEWAY_RELAY_URL", "wss://connector.example/relay") + monkeypatch.setenv("GATEWAY_RELAY_PLATFORMS", "discord") + monkeypatch.setenv("GATEWAY_RELAY_BOT_IDS", json.dumps({"discord": {"botId": "b1"}})) + + directory = tmp_path / "channel_directory.json" + directory.write_text( + json.dumps( + { + "updated_at": None, + "platforms": { + "discord": [ + {"id": ATTESTED_CHAT, "name": "bot-home", "type": "channel"} + ] + }, + } + ), + encoding="utf-8", + ) + monkeypatch.setattr(cd, "DIRECTORY_PATH", directory) + monkeypatch.setattr(cd, "CHANNEL_ALIASES_PATH", tmp_path / "channel_aliases.json") + # No gateway-session origins for discord in this temp home. + monkeypatch.setattr(cd, "_build_from_sessions", lambda _platform: []) + return directory + + +def _send(target: str, sent): + """Invoke the real tool, recording any egress it attempts.""" + from types import SimpleNamespace + from unittest.mock import patch + + import asyncio + + discord_cfg = SimpleNamespace(enabled=True, token="t", extra={}) + config = SimpleNamespace( + platforms={Platform.DISCORD: discord_cfg}, + get_home_channel=lambda _p: SimpleNamespace(chat_id=HOME_CHAT), + ) + + async def _record(platform, pconfig, chat_id, message, **kwargs): + sent.append(chat_id) + return {"success": True, "message_id": "m1"} + + with patch("gateway.config.load_gateway_config", return_value=config), patch( + "tools.interrupt.is_interrupted", return_value=False + ), patch("model_tools._run_async", side_effect=lambda c: asyncio.run(c)), patch( + "tools.send_message_tool._send_to_platform", side_effect=_record + ), patch( + "gateway.mirror.mirror_to_session", return_value=False + ): + return json.loads( + send_message_tool( + {"action": "send", "target": target, "message": "hello"} + ) + ) + + +def test_arbitrary_relay_chat_id_is_refused_and_never_egresses(relay_env): + """The whole observable: refused, naming THAT target, and ZERO egress.""" + sent: list[str] = [] + result = _send(f"discord:{ARBITRARY_CHAT}", sent) + + assert result == { + "error": ( + f"Refusing to send to unattested relay target 'discord:{ARBITRARY_CHAT}': " + "this gateway has no record of that destination. Use " + "send_message(action='list') to see the targets it can reach." + ) + } + assert sent == [] + + +def test_attested_directory_chat_id_still_sends(relay_env): + """The guard must not destroy the feature: an attested chat goes through.""" + sent: list[str] = [] + result = _send(f"discord:{ATTESTED_CHAT}", sent) + + assert result == {"success": True, "message_id": "m1"} + assert sent == [ATTESTED_CHAT] + + +def test_home_channel_is_attested(relay_env): + """The operator-configured home channel is a provenance, not a guess.""" + sent: list[str] = [] + result = _send("discord", sent) + + assert result["success"] is True + assert sent == [HOME_CHAT] + + +def test_session_origin_chat_is_attested(relay_env, monkeypatch): + """A chat this gateway actually holds a session in is reachable.""" + import gateway.channel_directory as cd + + monkeypatch.setattr( + cd, + "_build_from_sessions", + lambda platform: ( + [{"id": ARBITRARY_CHAT, "name": "seen", "type": "channel"}] + if platform == "discord" + else [] + ), + ) + sent: list[str] = [] + result = _send(f"discord:{ARBITRARY_CHAT}", sent) + + assert result["success"] is True + assert sent == [ARBITRARY_CHAT] + + +def test_platform_not_fronted_by_relay_is_untouched(relay_env, monkeypatch): + """Non-relay platforms keep their own adapters' authorization, unchanged.""" + monkeypatch.setenv("GATEWAY_RELAY_PLATFORMS", "telegram") + monkeypatch.setenv( + "GATEWAY_RELAY_BOT_IDS", json.dumps({"telegram": {"botId": "b1"}}) + ) + sent: list[str] = [] + result = _send(f"discord:{ARBITRARY_CHAT}", sent) + + assert result["success"] is True + assert sent == [ARBITRARY_CHAT] + + +def test_live_native_adapter_takes_precedence_over_the_relay_guard( + relay_env, monkeypatch +): + """A platform served by a live NATIVE adapter here is not a relay egress. + + Same precedence `gateway/delivery.resolve_delivery_transport` applies: a + concrete native adapter always wins over the relay, so this guard must not + fire for it. + """ + from types import SimpleNamespace + + import gateway.run + + runner = SimpleNamespace(adapters={Platform.DISCORD: object()}) + monkeypatch.setattr(gateway.run, "_gateway_runner_ref", lambda: runner) + sent: list[str] = [] + result = _send(f"discord:{ARBITRARY_CHAT}", sent) + + assert result["success"] is True + assert sent == [ARBITRARY_CHAT] + + +def test_react_refuses_an_arbitrary_relay_target(relay_env): + """Reactions are outbound acts too — same floor, same refusal.""" + result = json.loads( + send_message_tool( + { + "action": "react", + "target": f"discord:{ARBITRARY_CHAT}", + "emoji": "👍", + } + ) + ) + assert result == { + "error": ( + f"Refusing to send to unattested relay target 'discord:{ARBITRARY_CHAT}': " + "this gateway has no record of that destination. Use " + "send_message(action='list') to see the targets it can reach." + ) + } diff --git a/tools/send_message_tool.py b/tools/send_message_tool.py index 9d64123909..b622dea115 100644 --- a/tools/send_message_tool.py +++ b/tools/send_message_tool.py @@ -155,6 +155,24 @@ def _error(message: str) -> dict: return {"error": _sanitize_error_text(message)} +def _authorize_relay_target(platform_name: str, chat_id) -> str | None: + """Relay egress-authorization guard (P5a); None when the send may proceed. + + Thin, never-raising delegate to ``gateway.relay.egress`` so the tool keeps + working in environments where the gateway package can't be imported. The + guard itself FAILS CLOSED on an unattested target but must not fail closed + on its own import error — a missing gateway module means there is no relay + egress to authorize in the first place. + """ + try: + from gateway.relay.egress import authorize_relay_target + + return authorize_relay_target(platform_name, chat_id) + except Exception: # noqa: BLE001 - no gateway package ⇒ no relay egress + logger.debug("relay target authorization unavailable", exc_info=True) + return None + + def _display_chat_id(platform_name: str, chat_id: str) -> str: """Return a result-safe chat identifier for tool transcripts/log consumers.""" if platform_name == "signal" and str(chat_id).startswith("group:"): @@ -324,6 +342,13 @@ def _handle_react(args, remove=False): ) chat_id = home.chat_id + # P5(a): same egress-authorization floor as the send path — a reaction is + # an outbound act against a named destination, so an unattested relay + # target must be refused here too, not just on `send`. + _relay_denial = _authorize_relay_target(platform_name, chat_id) + if _relay_denial: + return tool_error(_relay_denial) + runner = None try: from gateway.run import _gateway_runner_ref @@ -460,6 +485,17 @@ def _handle_send(args): f"or set a home channel via: hermes config set {home_env} " ) + # P5(a): a relay-routed destination must be one this gateway can show a + # provenance for. The `target` parameter is free-form, so without this a + # model could name ANY chat id and the gateway would dutifully emit an + # outbound frame for it — authenticating the sender while never + # authorizing the destination. Applies to the generic `relay` plane and to + # connector-fronted platforms with no live native adapter; every other + # platform keeps its adapter's own authorization unchanged. + _relay_denial = _authorize_relay_target(platform_name, chat_id) + if _relay_denial: + return tool_error(_relay_denial) + duplicate_skip = _maybe_skip_cron_duplicate_send(platform_name, chat_id, thread_id) if duplicate_skip: return json.dumps(duplicate_skip)