fix(relay): authorize send_message targets and surface egress declines

P5 of the relay egress-authorization workstream. The relay path
authenticated the SENDER but never authorized the DESTINATION, and the
gateway compounded it from both ends.

(a) send_message could silently name an arbitrary relay target. Its
`target` parameter is free-form ('platform:chat_id'), so a model could
name ANY chat id and the gateway would emit an outbound frame for it.
gateway/relay/egress.py adds an attestation floor: a relay-routed
destination must have a provenance this gateway can show -- the
operator's home channel, the channel directory, or its own gateway
session origins. Anything else is refused HERE, with a visible tool
error naming the target, before a frame is written. Non-relay platforms
and platforms served by a live native adapter in this process are
untouched (same precedence resolve_delivery_transport applies).

(b) Connector declines were swallowed into apparent successes. The
connector's egress floor answers an unauthorized destination with a
DEFINITE failure whose text is deliberately uniform (F-005). Several
relay lanes degrade a *transport drop* by design and were degrading an
*authorization refusal* the same way:

  - _send_media returned None, sending the caller into
    BasePlatformAdapter's text fallback -- a DIFFERENT op re-addressed at
    the very chat the connector had just refused.
  - _send_prompt returned None, so exec-approval / slash-confirm /
    clarify reported "relay prompt op unavailable" (a wrong reason) and
    ran their numbered-text fallbacks into the refused chat.
  - task_card_stop discarded the error entirely.
  - typing / delete / react / thread ops degraded silently at debug.

is_egress_decline() classifies THAT a decline happened (never why --
the uniform text is not parsed for reasons) and requires a definite,
non-ambiguous failure, so a lost-ack retry is still a transport
outcome. Lanes with an error-carrying contract now report the decline
verbatim; cosmetic bool/None lanes still degrade but log it at WARNING.

Advisory progress drops that legitimately degrade are unchanged: the
task_card send lane, the draft ambiguous/except branches, and every
transport-exception path keep their existing fail-open behaviour.

Tests: 21 mutations of the production source, all KILLED.
This commit is contained in:
Ben Barclay
2026-08-31 16:16:09 +10:00
parent a9c783f219
commit 7cf86188ac
5 changed files with 819 additions and 17 deletions

View File

@@ -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

233
gateway/relay/egress.py Normal file
View File

@@ -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
(``"<platform> 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."
)

View File

@@ -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()

View File

@@ -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."
)
}

View File

@@ -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} <channel_id>"
)
# 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)