fix(gateway): close sender-routing gaps found in review
Voice input reused the bound text channel's cached source and replaced only `user_id`, so a second speaker inherited the profile resolved for the first. The ingress gate does not re-resolve a source that already carries a profile. Re-resolve at the voice call site instead of clearing the profile: clearing would drop the receiving bot and could re-home voice arriving on a secondary profile's own bot. `_voice_input_source` reattaches the transport provenance `from_dict` discards, and `_stamp_routed_profile` takes the receiving bot's profile as the fallback when no route matches. Kanban re-subscription could not repair a row created before sender capture: `user_id` was only written by the INSERT, so a legacy row stayed senderless and the notifier's conservative fallback left it undeliverable for good. It now self-heals like `user_id_alt`. An explicit `user_id: null` or empty string is rejected instead of widening the route to every sender, and `to_dict` omits the field when unset so a round-trip cannot reintroduce it. Numeric `0` from an adapter normalizes to "0" rather than being dropped, without changing the shared coercion used by the other fields.
This commit is contained in:
@@ -652,7 +652,9 @@ class GatewayConfig:
|
||||
"streaming": self.streaming.to_dict(),
|
||||
"session_store_max_age_days": self.session_store_max_age_days,
|
||||
"profile_routes": [
|
||||
asdict(r) if is_dataclass(r) and not isinstance(r, type) else r for r in self.profile_routes
|
||||
{k: v for k, v in asdict(r).items() if k != "user_id" or v is not None}
|
||||
if is_dataclass(r) and not isinstance(r, type) else r
|
||||
for r in self.profile_routes
|
||||
],
|
||||
}
|
||||
|
||||
|
||||
@@ -4469,7 +4469,8 @@ class BasePlatformAdapter(ABC):
|
||||
return str(value) if value else None
|
||||
fields = dict(
|
||||
platform=self.platform, chat_id=str(chat_id), chat_name=chat_name, chat_type=chat_type,
|
||||
user_id=_opt(user_id), user_name=user_name, thread_id=_opt(thread_id),
|
||||
user_id=None if user_id is None or user_id == "" else str(user_id),
|
||||
user_name=user_name, thread_id=_opt(thread_id),
|
||||
chat_topic=(chat_topic or "").strip() or None, user_id_alt=user_id_alt,
|
||||
chat_id_alt=chat_id_alt, is_bot=is_bot, scope_id=_opt(scope_id),
|
||||
guild_id=_opt(guild_id), parent_chat_id=_opt(parent_chat_id),
|
||||
|
||||
@@ -87,7 +87,7 @@ class ProfileRoute:
|
||||
return False
|
||||
if _bot_profile_key(self.bot_profile) != _bot_profile_key(adapter_profile):
|
||||
return False
|
||||
if self.user_id is not None and (not str(self.user_id).strip() or self.user_id != user_id):
|
||||
if self.user_id is not None and (not self.user_id.strip() or self.user_id != user_id):
|
||||
return False
|
||||
if self.thread_id and self.thread_id != thread_id:
|
||||
return False
|
||||
@@ -153,12 +153,16 @@ def parse_profile_routes(raw: Optional[List[Dict[str, Any]]]) -> List[ProfileRou
|
||||
except (ValueError, ImportError):
|
||||
logger.warning("Skipping profile route %s: invalid profile name %r", name, profile)
|
||||
continue
|
||||
has_user_id, user_id = "user_id" in entry, entry.get("user_id")
|
||||
if has_user_id and (user_id is None or isinstance(user_id, str) and not user_id.strip()):
|
||||
logger.warning("Skipping profile route %s: user_id cannot be null or empty", name)
|
||||
continue
|
||||
routes.append(ProfileRoute(
|
||||
name=name, platform=platform, profile=profile,
|
||||
guild_id=_coerce_route_id(entry.get("guild_id")),
|
||||
chat_id=_coerce_route_id(entry.get("chat_id")),
|
||||
thread_id=_coerce_route_id(entry.get("thread_id")),
|
||||
user_id=_coerce_route_id(entry.get("user_id")),
|
||||
user_id=_coerce_route_id(user_id),
|
||||
enabled=entry.get("enabled", True),
|
||||
bot_profile=_bot_profile_key(entry.get("bot_profile")),
|
||||
))
|
||||
|
||||
@@ -4310,7 +4310,7 @@ class GatewayRunner(
|
||||
routes, platform=source.platform.value, guild_id=getattr(source, "guild_id", None),
|
||||
chat_id=source.chat_id, thread_id=getattr(source, "thread_id", None),
|
||||
parent_chat_id=getattr(source, "parent_chat_id", None),
|
||||
adapter_profile=adapter_profile, user_id=source.user_id)
|
||||
adapter_profile=adapter_profile, user_id=getattr(source, "user_id", None))
|
||||
except Exception:
|
||||
logger.warning(
|
||||
"Profile route matching failed for %s/%s, falling back to default",
|
||||
@@ -4331,9 +4331,9 @@ class GatewayRunner(
|
||||
raise ProfileRouteRejected(matched.name)
|
||||
return matched.profile
|
||||
logger.debug(
|
||||
"No profile route matched: platform=%s chat_id=%s thread_id=%s parent_chat_id=%s user_id=%s",
|
||||
"No profile route matched: platform=%s chat_id=%s thread_id=%s parent_chat_id=%s",
|
||||
source.platform.value, source.chat_id,
|
||||
getattr(source, "thread_id", None), getattr(source, "parent_chat_id", None), source.user_id)
|
||||
getattr(source, "thread_id", None), getattr(source, "parent_chat_id", None))
|
||||
return None
|
||||
|
||||
def _resolve_profile_home_for_source(self, source: SessionSource) -> "Path":
|
||||
|
||||
@@ -1430,11 +1430,20 @@ class GatewayAdapterLifecycleMixin:
|
||||
except IdentityUnresolved:
|
||||
return None
|
||||
|
||||
def _stamp_routed_profile(self, source) -> bool:
|
||||
"""Stamp ``source.profile`` from ``profile_routes``; False when the route is rejected."""
|
||||
def _stamp_routed_profile(self, source, adapter_profile: Optional[str] = None) -> bool:
|
||||
"""Stamp ``source.profile`` from ``profile_routes``; False when the route is rejected.
|
||||
|
||||
``adapter_profile`` owns the receiving bot: routes are scoped to it and it is the
|
||||
fallback when none matches, so re-stamping a source never crosses a bot boundary.
|
||||
"""
|
||||
from gateway.profile_routing import ProfileRouteRejected
|
||||
try:
|
||||
source.profile = self._profile_name_for_source(source)
|
||||
routed_profile = (
|
||||
self._profile_name_for_source(source, adapter_profile=adapter_profile)
|
||||
if adapter_profile is not None
|
||||
else self._profile_name_for_source(source)
|
||||
)
|
||||
source.profile = routed_profile or adapter_profile
|
||||
except ProfileRouteRejected:
|
||||
return False
|
||||
return True
|
||||
|
||||
@@ -12,6 +12,7 @@ import os
|
||||
import re
|
||||
import sys
|
||||
import time
|
||||
import weakref
|
||||
from contextlib import suppress
|
||||
from difflib import SequenceMatcher
|
||||
from types import SimpleNamespace
|
||||
@@ -228,11 +229,14 @@ class GatewayVoiceMixin:
|
||||
if source_data := getattr(adapter, "_voice_sources", {}).get(guild_id):
|
||||
source = SessionSource.from_dict(source_data)
|
||||
source.user_id = source.user_name = str(user_id)
|
||||
return source
|
||||
return SessionSource(
|
||||
platform=Platform.DISCORD, chat_id=str(text_ch_id), user_id=str(user_id),
|
||||
user_name=str(user_id), chat_type="channel",
|
||||
profile=getattr(adapter, "_owner_profile", None))
|
||||
else:
|
||||
source = SessionSource(
|
||||
platform=Platform.DISCORD, chat_id=str(text_ch_id), user_id=str(user_id),
|
||||
user_name=str(user_id), chat_type="channel",
|
||||
profile=getattr(adapter, "_owner_profile", None))
|
||||
# Serialization drops transport provenance; auth must still follow the receiving bot.
|
||||
source._transport_adapter_ref = weakref.ref(adapter)
|
||||
return source
|
||||
|
||||
async def _handle_voice_channel_input(
|
||||
self, guild_id: int, user_id: int, transcript: str, *, adapter=None
|
||||
@@ -245,6 +249,10 @@ class GatewayVoiceMixin:
|
||||
if not text_ch_id:
|
||||
return
|
||||
source = self._voice_input_source(adapter, guild_id, user_id, text_ch_id)
|
||||
# Cached source still carries the previous speaker's routed profile.
|
||||
if not self._stamp_routed_profile(source, getattr(adapter, "_owner_profile", None)):
|
||||
logger.warning("Dropping voice input: its profile route targets an unserved profile")
|
||||
return
|
||||
# Validate the session owner against the current allowlist before auto-resuming. A session created
|
||||
# before TELEGRAM_ALLOWED_USERS (or equivalent) was configured, or before the owner was removed from
|
||||
# it, must not silently receive a full agent response on gateway restart just because it has a
|
||||
|
||||
@@ -123,9 +123,10 @@ def add_notify_sub(
|
||||
)
|
||||
# chat_type / delivery_mode are last-write-wins; delivery metadata
|
||||
# preserves existing routing fields while supplied fields overwrite them.
|
||||
# user_id_alt and notifier_profile only self-heal legacy rows lacking one.
|
||||
# user_id, user_id_alt and notifier_profile only self-heal legacy rows lacking one.
|
||||
for column, value, fill_only in (
|
||||
("chat_type", chat_type, False),
|
||||
("user_id", user_id, True),
|
||||
("user_id_alt", user_id_alt, True),
|
||||
("notifier_profile", notifier_profile, True),
|
||||
("delivery_mode", valid_mode, False),
|
||||
|
||||
@@ -114,8 +114,8 @@ def test_user_routed_subscription_uses_only_its_authorized_profile(tmp_path, mon
|
||||
assert len(primary.sent) == len(primary.handled) == 1
|
||||
assert primary.handled[0].source.user_id == "creator"
|
||||
|
||||
# Same sender cannot fall through to the primary profile, and a legacy row
|
||||
# without sender identity cannot skip a potentially winning user route.
|
||||
# Same sender can't fall back to the primary profile, and a legacy row with no sender
|
||||
# identity must not skip a user route that could have won.
|
||||
completion(profile="default", chat="shared", thread="")
|
||||
completion(profile="default", chat="shared", thread="", user=None)
|
||||
assert not collect(runner)
|
||||
|
||||
@@ -369,6 +369,20 @@ class TestAdapterToSessionKeyIntegration:
|
||||
# A default-profile key would land in agent:main — must differ.
|
||||
assert key != build_session_key(source, profile=None)
|
||||
|
||||
def test_adapter_preserves_numeric_zero_user_id_for_routing(self, mock_runner):
|
||||
mock_runner.config.profile_routes = [
|
||||
ProfileRoute(name="zero", platform="discord", profile="zero", user_id="0")
|
||||
]
|
||||
adapter = _stub_adapter(Platform.DISCORD, mock_runner)
|
||||
|
||||
with patch(
|
||||
"hermes_cli.profiles.profiles_to_serve",
|
||||
return_value=[("default", Path("/profiles/default")), ("zero", Path("/profiles/zero"))],
|
||||
):
|
||||
source = adapter.build_source(chat_id="channel", user_id=0)
|
||||
|
||||
assert (source.user_id, source.profile) == ("0", "zero")
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_adapter_drops_rejected_route_before_dispatch(self, mock_runner):
|
||||
mock_runner.config.profile_routes = [
|
||||
|
||||
@@ -65,6 +65,15 @@ class TestProfileRouteMatching:
|
||||
"teams", user_id="aad-456", chat_id="conversation-2",
|
||||
)
|
||||
|
||||
def test_sender_route_cannot_cross_the_receiving_bot_boundary(self):
|
||||
route = ProfileRoute(
|
||||
name="owner", platform="telegram", profile="owner",
|
||||
user_id="640466638", bot_profile="team_b",
|
||||
)
|
||||
assert route.matches("telegram", user_id="640466638", adapter_profile="team_b")
|
||||
assert not route.matches("telegram", user_id="640466638", adapter_profile=None)
|
||||
assert not route.matches("telegram", user_id="other", adapter_profile="team_b")
|
||||
|
||||
def test_sender_routes_split_one_chat_and_outrank_its_location_route(self):
|
||||
routes = parse_profile_routes([
|
||||
{"name": "shared", "platform": "teams", "profile": "shared",
|
||||
@@ -109,15 +118,20 @@ class TestParseProfileRoutes:
|
||||
assert (by_name["platform-only"].guild_id, by_name["platform-only"].chat_id,
|
||||
by_name["platform-only"].thread_id) == (None, None, None)
|
||||
|
||||
@pytest.mark.parametrize("blank", ["", " "])
|
||||
def test_blank_user_id_is_inert_never_a_broader_route(self, blank):
|
||||
# Unlike the location fields, a blank sender does not fall back to "unconstrained".
|
||||
routes = parse_profile_routes([
|
||||
{"name": "blank", "platform": "teams", "profile": "owner", "user_id": blank},
|
||||
{"name": "legacy", "platform": "teams", "profile": "shared"},
|
||||
])
|
||||
assert match_profile_route(routes, "teams", user_id=blank).name == "legacy"
|
||||
assert match_profile_route(routes, "teams", user_id="anyone").name == "legacy"
|
||||
@pytest.mark.parametrize("invalid", [None, "", " "])
|
||||
def test_null_or_blank_user_id_rejects_only_that_route(self, invalid, caplog):
|
||||
with caplog.at_level("WARNING", logger="gateway.profile_routing"):
|
||||
routes = parse_profile_routes([
|
||||
{"name": "invalid", "platform": "teams", "profile": "owner", "user_id": invalid},
|
||||
{"name": "missing", "platform": "teams", "profile": "shared"},
|
||||
])
|
||||
assert [route.name for route in routes] == ["missing"]
|
||||
assert match_profile_route(routes, "teams", user_id="anyone").name == "missing"
|
||||
assert "user_id cannot be null or empty" in caplog.text
|
||||
if isinstance(invalid, str):
|
||||
assert not ProfileRoute(
|
||||
name="direct", platform="teams", profile="owner", user_id=invalid,
|
||||
).matches("teams", user_id=invalid)
|
||||
|
||||
def test_non_int_numeric_ids_warn_instead_of_silently_coercing(self, caplog):
|
||||
# #86470 nuance: float/bool stringify to values that can never match
|
||||
@@ -134,6 +148,23 @@ class TestParseProfileRoutes:
|
||||
|
||||
class TestMatchProfileRoute:
|
||||
|
||||
def test_sender_only_route_outranks_the_tightest_location_route(self):
|
||||
routes = parse_profile_routes([
|
||||
{"name": "thread", "platform": "discord", "profile": "thread",
|
||||
"guild_id": "g", "chat_id": "c", "thread_id": "t"},
|
||||
{"name": "sender", "platform": "discord", "profile": "sender", "user_id": "u"},
|
||||
{"name": "sender-chat", "platform": "discord", "profile": "sender-chat",
|
||||
"user_id": "u", "chat_id": "c"},
|
||||
])
|
||||
assert match_profile_route(
|
||||
routes, "discord", guild_id="g", chat_id="c", thread_id="t", user_id="u",
|
||||
).profile == "sender-chat"
|
||||
assert match_profile_route(
|
||||
routes, "discord", guild_id="g", chat_id="other", thread_id="t", user_id="u",
|
||||
).profile == "sender"
|
||||
assert match_profile_route(
|
||||
routes, "discord", guild_id="g", chat_id="c", thread_id="t", user_id="other",
|
||||
).profile == "thread"
|
||||
|
||||
def test_no_match_returns_none(self):
|
||||
routes = [
|
||||
@@ -236,8 +267,6 @@ class TestWhatsAppChatIdIdentityMatching:
|
||||
|
||||
class TestGatewayConfigRoundtrip:
|
||||
def test_routes_survive_to_dict_from_dict_with_user_id_and_enabled(self):
|
||||
# GatewayConfig serializes routes with asdict(); a field added to ProfileRoute but not
|
||||
# re-read by parse_profile_routes would silently widen or disable routes on reload.
|
||||
from gateway.config import GatewayConfig
|
||||
|
||||
config = GatewayConfig(profile_routes=parse_profile_routes([
|
||||
@@ -245,7 +274,9 @@ class TestGatewayConfigRoundtrip:
|
||||
{"name": "off", "platform": "teams", "profile": "owner",
|
||||
"chat_id": "conversation-1", "enabled": False},
|
||||
]))
|
||||
restored = GatewayConfig.from_dict(config.to_dict()).profile_routes
|
||||
raw = config.to_dict()
|
||||
assert "user_id" not in raw["profile_routes"][1]
|
||||
restored = GatewayConfig.from_dict(raw).profile_routes
|
||||
|
||||
assert [(r.name, r.user_id, r.enabled, r.specificity) for r in restored] == [
|
||||
("sender", "aad-456", True, 16), ("off", None, False, 4),
|
||||
|
||||
@@ -632,6 +632,44 @@ class TestVoiceChannelCommands:
|
||||
assert event.source.chat_id == "123"
|
||||
assert event.source.chat_type == "channel"
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_input_reroutes_speaker_without_changing_transport_owner(self, runner, monkeypatch):
|
||||
from gateway.config import Platform
|
||||
from gateway.profile_routing import parse_profile_routes
|
||||
|
||||
runner.config = SimpleNamespace(
|
||||
multiplex_profiles=True,
|
||||
profile_routes=parse_profile_routes([
|
||||
{"name": "second", "platform": "discord", "bot_profile": "team-bot",
|
||||
"user_id": "222", "profile": "second"},
|
||||
]),
|
||||
)
|
||||
monkeypatch.setattr(
|
||||
"gateway.run._multiplex_profile_homes",
|
||||
lambda _config: [("team-bot", None), ("first", None), ("second", None)],
|
||||
)
|
||||
mock_adapter = AsyncMock()
|
||||
mock_adapter._owner_profile = "team-bot"
|
||||
mock_adapter._voice_text_channels = {111: 123}
|
||||
mock_adapter._voice_sources = {111: SessionSource(
|
||||
platform=Platform.DISCORD, chat_id="123", chat_type="channel",
|
||||
user_id="111", profile="first",
|
||||
).to_dict()}
|
||||
mock_adapter._client = MagicMock()
|
||||
mock_adapter._client.get_channel = MagicMock(return_value=AsyncMock())
|
||||
mock_adapter.handle_message = AsyncMock()
|
||||
runner.adapters = {}
|
||||
runner._profile_adapters = {
|
||||
"team-bot": {Platform.DISCORD: mock_adapter},
|
||||
"second": {},
|
||||
}
|
||||
|
||||
await runner._handle_voice_channel_input(111, 222, "Hello from VC", adapter=mock_adapter)
|
||||
|
||||
source = mock_adapter.handle_message.call_args[0][0].source
|
||||
assert (source.user_id, source.profile) == ("222", "second")
|
||||
assert runner._transport_owner(source) == (mock_adapter, "team-bot")
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_input_resolves_channel_prompt(self, runner):
|
||||
"""Voice input must carry the bound text channel's channel_prompt (#50149)."""
|
||||
|
||||
@@ -182,6 +182,27 @@ def test_notify_sub_chat_type_persists_and_last_write_wins(kanban_home):
|
||||
conn.close()
|
||||
|
||||
|
||||
def test_notify_sub_user_id_backfills_legacy_senderless_rows(kanban_home):
|
||||
import hermes_cli.kanban_db as kb
|
||||
from hermes_cli import kanban_db_connect as kbc
|
||||
from hermes_cli import kanban_db_notify as kbn
|
||||
|
||||
conn = kbc.connect()
|
||||
try:
|
||||
tid = kb.create_task(conn, title="legacy sub", assignee="worker1")
|
||||
kbn.add_notify_sub(conn, task_id=tid, platform="telegram", chat_id="chat1")
|
||||
assert kbn.list_notify_subs(conn, tid)[0]["user_id"] is None
|
||||
|
||||
kbn.add_notify_sub(
|
||||
conn, task_id=tid, platform="telegram", chat_id="chat1", user_id="640466638",
|
||||
)
|
||||
subs = kbn.list_notify_subs(conn, tid)
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
assert subs[0]["user_id"] == "640466638"
|
||||
|
||||
|
||||
def test_notify_sub_user_id_alt_persists_and_backfills_legacy_rows(kanban_home):
|
||||
"""user_id_alt is persisted with the notify subscription routing tuple and
|
||||
can backfill a pre-existing row created before the alt id was known."""
|
||||
|
||||
@@ -556,12 +556,32 @@ gateway:
|
||||
profile: owner
|
||||
```
|
||||
|
||||
Routes are matched most-specific-first (`user_id` > `thread_id` > `chat_id` > `guild_id`),
|
||||
all declared fields must hold (AND), and a route keyed on a channel also
|
||||
matches threads/forum posts whose parent is that channel. Messages that match
|
||||
no route stay on the default/active profile. The routed profile gets the full
|
||||
per-profile isolation described above (config, skills, memory, credentials,
|
||||
session namespace). Routing works on every platform adapter, not just Discord.
|
||||
Routes are matched by additive specificity: `user_id` = 16, `thread_id` = 8,
|
||||
`chat_id` = 4, and `guild_id` = 2. Thus `user_id + chat_id` (20) outranks
|
||||
`user_id` alone (16), which outranks every location-only route (at most 14).
|
||||
All declared fields must hold (AND), equal scores keep declaration order, and a
|
||||
route keyed on a channel also matches threads/forum posts whose parent is that
|
||||
channel. Messages that match no route stay on the default/active profile. The
|
||||
routed profile gets the full per-profile isolation described above (config,
|
||||
skills, memory, credentials, session namespace). Routing works on every
|
||||
platform adapter, not just Discord.
|
||||
|
||||
`user_id` is the **sender** of the inbound message, compared for exact equality. It is only
|
||||
as trustworthy as the adapter that reports it, so treat it as an authorization input only on
|
||||
platforms whose ingress authenticates the sender. Sender ids are also namespaced per tenant
|
||||
on some platforms — a Slack user id is workspace-local — so on a gateway serving more than
|
||||
one workspace or server, pair `user_id` with the `guild_id` of that scope (Discord guild,
|
||||
Slack workspace, Matrix server) rather than relying on the id alone.
|
||||
|
||||
Omitting `user_id` keeps the route unconstrained by sender for backward compatibility.
|
||||
Setting it to `null`, an empty string, or whitespace invalidates that route instead of
|
||||
broadening it to every sender on the platform.
|
||||
|
||||
Sender routing selects a profile; it is not deny-by-default authorization. A sender that
|
||||
matches no route falls through to the default/active profile, exactly like an unrouted
|
||||
channel. To give one person a privileged profile and everyone else a restricted one, declare
|
||||
the privileged sender route first, add a platform-wide catch-all route to the restricted
|
||||
profile after it, and keep the platform's own ingress allowlist in place.
|
||||
|
||||
A route applies only to messages received by the **default profile's bot**
|
||||
unless it names another bot with `bot_profile: <profile>`. Telegram DMs use the
|
||||
|
||||
Reference in New Issue
Block a user