fix(relay): echo the routed profile on every outbound frame and follow_up
The connector stamps `profile` on inbound and passthrough_forward frames but the gateway never sent it back, so the connector had nothing to stamp on the NEXT interaction of a routed chat. `_capture_scope` now remembers the routed profile per chat, `_with_scope` echoes it as `metadata.profile` on chat-addressed frames, and `send_follow_up` derives it from the `agent:<profile>:` key namespace. A single-profile gateway emits no key — frames stay byte-identical. Contract §4 documents the round-trip.
This commit is contained in:
@@ -95,6 +95,17 @@ def _event_ids(event) -> Tuple[Optional[str], Optional[str]]:
|
||||
return message_id, getattr(event.source, "chat_id", None)
|
||||
|
||||
|
||||
def _profile_from_session_key(session_key: str) -> Optional[str]:
|
||||
"""Named profile encoded in an ``agent:<ns>:...`` session key; None for the legacy ``agent:main``
|
||||
namespace (single-profile gateway) so the wire frame stays byte-identical there."""
|
||||
parts = (session_key or "").split(":")
|
||||
if len(parts) < 2 or parts[0] != "agent" or not parts[1]:
|
||||
return None
|
||||
from gateway.session import profile_from_session_key_namespace
|
||||
profile = profile_from_session_key_namespace(parts[1])
|
||||
return None if profile == "default" else profile
|
||||
|
||||
|
||||
class RelayAdapter(BasePlatformAdapter):
|
||||
"""Generic relay adapter advertising a connector-negotiated capability profile."""
|
||||
|
||||
@@ -125,6 +136,10 @@ class RelayAdapter(BasePlatformAdapter):
|
||||
# platforms on one WS and a reply must egress through the platform the
|
||||
# inbound came from. Empty for a single-platform gateway (connector default).
|
||||
self._platform_by_chat: Dict[str, str] = {}
|
||||
# chat_id -> Hermes profile the connector routed the inbound to (multiplex mode). Echoed
|
||||
# on every outbound frame's metadata so the connector can stamp the SAME profile on the
|
||||
# next passthrough_forward for that chat; empty on a single-profile gateway.
|
||||
self._profile_by_chat: Dict[str, str] = {}
|
||||
# Chats the connector has refused (see the terminal-decline latch).
|
||||
# chat_id -> (thread_id, initial_name) of the auto-thread the CONNECTOR
|
||||
# created for our latest send; read by the semantic thread-rename lane.
|
||||
@@ -1042,6 +1057,7 @@ class RelayAdapter(BasePlatformAdapter):
|
||||
for attr, cache in (
|
||||
("user_id", self._dm_user_by_chat), ("scope_id", self._scope_by_chat),
|
||||
("chat_type", self._chat_type_by_chat),
|
||||
("profile", self.__dict__.setdefault("_profile_by_chat", {})),
|
||||
):
|
||||
value = getattr(src, attr, None)
|
||||
if value:
|
||||
@@ -1059,7 +1075,11 @@ class RelayAdapter(BasePlatformAdapter):
|
||||
first and only falls back to user_id on a route miss, so carrying both never
|
||||
overrides routing-table resolution."""
|
||||
meta: Dict[str, Any] = dict(metadata or {})
|
||||
for key, cache in (("scope_id", self._scope_by_chat), ("user_id", self._dm_user_by_chat)):
|
||||
# ``getattr``: relay tests build bare adapters via ``__new__`` without ``__init__``.
|
||||
for key, cache in (
|
||||
("scope_id", self._scope_by_chat), ("user_id", self._dm_user_by_chat),
|
||||
("profile", getattr(self, "_profile_by_chat", {})),
|
||||
):
|
||||
if not meta.get(key):
|
||||
value = cache.get(str(chat_id))
|
||||
if value:
|
||||
@@ -1698,13 +1718,20 @@ class RelayAdapter(BasePlatformAdapter):
|
||||
# default routes it.
|
||||
prefix = kind.split(".", 1)[0] if kind and "." in kind else None
|
||||
follow_up_platform = prefix if prefix and self.fronts_platform(prefix) else None
|
||||
follow_up_metadata = dict(metadata or {})
|
||||
# The session key names the profile namespace the interaction ran under; carry it so the
|
||||
# connector's next passthrough_forward for this interaction routes to the same profile.
|
||||
if not follow_up_metadata.get("profile"):
|
||||
profile = _profile_from_session_key(session_key)
|
||||
if profile:
|
||||
follow_up_metadata["profile"] = profile
|
||||
result = await self._transport.send_follow_up(
|
||||
{
|
||||
"op": "follow_up",
|
||||
"session_key": session_key,
|
||||
"kind": kind,
|
||||
"content": content,
|
||||
"metadata": metadata or {},
|
||||
"metadata": follow_up_metadata,
|
||||
},
|
||||
platform=follow_up_platform,
|
||||
)
|
||||
|
||||
@@ -266,3 +266,42 @@ async def test_dm_interaction_keys_as_discord_dm(adapter, monkeypatch):
|
||||
assert ev.source.delivered_via_upstream_relay is True
|
||||
|
||||
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_routed_profile_round_trips_on_every_egress_frame(adapter, monkeypatch):
|
||||
"""Relay passthrough round-trip keeps ``profile`` (#88715 phase 5): the profile the connector
|
||||
stamped on an inbound interaction is echoed on the chat's outbound frames and on the
|
||||
``follow_up`` addressed by the routed session key, so the connector can stamp it on the NEXT
|
||||
passthrough_forward for that chat; a single-profile gateway emits no ``profile`` key at all."""
|
||||
await adapter.connect()
|
||||
stub = adapter._transport
|
||||
monkeypatch.setattr(adapter, "handle_message", _noop_handle)
|
||||
|
||||
fwd = _interaction_forward(
|
||||
{
|
||||
"id": "interaction-3", "type": 2, "channel_id": "chan-9", "guild_id": "guild-7",
|
||||
"data": {"name": "summarize"}, "member": {"user": {"id": "user-3", "username": "ben"}},
|
||||
},
|
||||
profile="reviewer",
|
||||
)
|
||||
await stub.push_passthrough(fwd, buffer_id=None)
|
||||
await adapter.send("chan-9", "done")
|
||||
await adapter.send_follow_up(
|
||||
session_key="agent:reviewer:discord:group:chan-9", kind="discord.interaction_token", content="x")
|
||||
assert stub.sent[-1]["metadata"]["profile"] == "reviewer"
|
||||
assert stub.follow_ups[-1]["metadata"]["profile"] == "reviewer"
|
||||
|
||||
# Legacy namespace / unrouted chat: byte-identical frames, no profile key.
|
||||
await stub.push_passthrough(_interaction_forward({
|
||||
"id": "interaction-4", "type": 2, "channel_id": "chan-1", "guild_id": "guild-7",
|
||||
"data": {"name": "summarize"}, "member": {"user": {"id": "user-3"}}}), buffer_id=None)
|
||||
await adapter.send("chan-1", "done")
|
||||
await adapter.send_follow_up(
|
||||
session_key="agent:main:discord:group:chan-1", kind="discord.interaction_token", content="x")
|
||||
assert "profile" not in stub.sent[-1]["metadata"]
|
||||
assert "profile" not in stub.follow_ups[-1]["metadata"]
|
||||
|
||||
|
||||
async def _noop_handle(event):
|
||||
return None
|
||||
|
||||
@@ -423,6 +423,16 @@ The gateway calls the transport with action dicts. Source of truth:
|
||||
`get_chat_info(chat_id)` is a separate proxied call returning at least
|
||||
`{name, type}`.
|
||||
|
||||
**`metadata.profile` (multiplex round-trip).** Every chat-addressed outbound
|
||||
frame's `metadata` carries the tenant discriminators the gateway captured from
|
||||
the inbound (`scope_id`, `user_id`) and, on a multiplexed gateway, the Hermes
|
||||
`profile` the connector routed that chat's inbound to; `follow_up` frames carry
|
||||
the profile encoded in their `session_key` namespace. The connector MUST stamp
|
||||
the same `profile` on the next `passthrough_forward` / `inbound` for that chat or
|
||||
interaction, so a Discord button press after a slash command lands in the same
|
||||
profile's session (`gateway/relay/adapter.py::_with_scope`, `send_follow_up`).
|
||||
Single-profile gateways never emit the key — frames stay byte-identical.
|
||||
|
||||
**`send_media` (Phase 2 media egress).** Media crosses the wire BY REFERENCE:
|
||||
`source_url` is either (a) a **connector re-host** the gateway previously
|
||||
uploaded via `POST {connector}/relay/media` (raw bytes body, `Content-Type` +
|
||||
|
||||
Reference in New Issue
Block a user