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:
@@ -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
233
gateway/relay/egress.py
Normal 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."
|
||||
)
|
||||
271
tests/gateway/relay/test_relay_egress_declines.py
Normal file
271
tests/gateway/relay/test_relay_egress_declines.py
Normal 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()
|
||||
200
tests/tools/test_send_message_relay_target_authz.py
Normal file
200
tests/tools/test_send_message_relay_target_authz.py
Normal 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."
|
||||
)
|
||||
}
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user