diff --git a/gateway/authz_mixin.py b/gateway/authz_mixin.py index 61e15e257b..990d306303 100644 --- a/gateway/authz_mixin.py +++ b/gateway/authz_mixin.py @@ -1,8 +1,8 @@ -"""User-authorization mixin for ``GatewayRunner``: may this user/chat talk to the -agent, the per-adapter DM policy, and the unauthorized-DM behavior. +"""User-authorization mixin for ``GatewayRunner``: may this user/chat talk to the agent, +the per-adapter DM policy, and the unauthorized-DM behavior. -``gateway.run`` is never imported at module import time (cycle); the one method -that logs imports its ``logger`` lazily so records keep the ``"gateway.run"`` name. +``gateway.run`` is never imported at module import time (cycle); the one method that logs +imports its ``logger`` lazily so records keep the ``"gateway.run"`` name. """ from __future__ import annotations @@ -23,18 +23,12 @@ _GROUP_CHAT_TYPES = frozenset({"group", "forum", "channel"}) _GROUP_FORUM_TYPES = frozenset({"group", "forum"}) _TRUTHY = frozenset({"true", "1", "yes"}) -# Platform -> ``_ALLOWED_USERS`` / ``_ALLOW_ALL_USERS``. -# Shared with the pairing store's allowlist mirror (single source of truth); -# plugin platforms are added per-call from the platform registry. +# Platform -> ``_ALLOWED_USERS`` / ``_ALLOW_ALL_USERS``. Shared with the pairing +# store's allowlist mirror (single source of truth); plugin platforms are added per-call from the registry. _ALLOWED_USERS_ENV = {Platform(k): v for k, v in _PLATFORM_ALLOWLIST_ENV.items()} -_ALLOW_ALL_ENV = { - p: v.replace("_ALLOWED_USERS", "_ALLOW_ALL_USERS") for p, v in _ALLOWED_USERS_ENV.items() -} +_ALLOW_ALL_ENV = {p: v.replace("_ALLOWED_USERS", "_ALLOW_ALL_USERS") for p, v in _ALLOWED_USERS_ENV.items()} _GROUP_USER_ENV = {Platform.TELEGRAM: "TELEGRAM_GROUP_ALLOWED_USERS"} -_GROUP_CHAT_ENV = { - Platform.TELEGRAM: "TELEGRAM_GROUP_ALLOWED_CHATS", - Platform.QQBOT: "QQ_GROUP_ALLOWED_USERS", -} +_GROUP_CHAT_ENV = {Platform.TELEGRAM: "TELEGRAM_GROUP_ALLOWED_CHATS", Platform.QQBOT: "QQ_GROUP_ALLOWED_USERS"} _ALLOW_BOTS_ENV = { Platform.DISCORD: "DISCORD_ALLOW_BOTS", Platform.FEISHU: "FEISHU_ALLOW_BOTS", @@ -46,10 +40,9 @@ _ALLOW_BOTS_ENV = { def _platform_gate_env(name: str, default: str = "") -> str: """Read an allow/deny gate env var with per-profile isolation. - With a profile secret scope installed AND multiplexing active, a scoped miss - returns ``default`` instead of falling through to ``os.environ``, which may - hold ANOTHER profile's first-writer bridged value (allowlist leak). - Single-profile deployments behave exactly like ``os.getenv``. + With a profile secret scope installed AND multiplexing active, a scoped miss returns ``default`` + instead of falling through to ``os.environ``, which may hold ANOTHER profile's first-writer + bridged value (allowlist leak). Single-profile deployments behave exactly like ``os.getenv``. """ if not name: return default @@ -98,20 +91,19 @@ def _adapter_config_extra(adapter) -> dict: return getattr(getattr(adapter, "config", None), "extra", None) or {} -# Nostr npub -> hex (Buzz): ``BUZZ_ALLOWED_USERS`` accepts hex or ``npub1…`` but -# inbound pubkeys are always hex. Pure stdlib; mirrors plugins/platforms/buzz/adapter.py. - +# Nostr npub -> hex (Buzz): ``BUZZ_ALLOWED_USERS`` accepts hex or ``npub1…`` but inbound pubkeys +# are always hex. Pure stdlib; mirrors plugins/platforms/buzz/adapter.py. _BECH32_CHARSET = "qpzry9x8gf2tvdw0s3jn54khce6mua7l" +_BECH32_GENERATOR = (0x3B6A57B2, 0x26508E6D, 0x1EA119FA, 0x3D4233DD, 0x2A1462B3) def _bech32_polymod(values): chk = 1 - generator = [0x3B6A57B2, 0x26508E6D, 0x1EA119FA, 0x3D4233DD, 0x2A1462B3] for value in values: top = chk >> 25 chk = (chk & 0x1FFFFFF) << 5 ^ value - for i in range(5): - chk ^= generator[i] if ((top >> i) & 1) else 0 + for i, gen in enumerate(_BECH32_GENERATOR): + chk ^= gen if ((top >> i) & 1) else 0 return chk @@ -132,10 +124,9 @@ def _convertbits(data, frombits: int, tobits: int, pad: bool = True): while bits >= tobits: bits -= tobits ret.append((acc >> bits) & maxv) - if pad: - if bits: - ret.append((acc << (tobits - bits)) & maxv) - elif bits >= frombits or ((acc << (tobits - bits)) & maxv): + if pad and bits: + ret.append((acc << (tobits - bits)) & maxv) + elif not pad and (bits >= frombits or ((acc << (tobits - bits)) & maxv)): return None return ret @@ -159,13 +150,7 @@ def _npub_to_hex(npub: str) -> Optional[str]: def _normalize_nostr_allow_entries(entries: set) -> set: """Add the hex form of every valid ``npub1…`` entry; invalid entries are kept as-is.""" - expanded = set(entries) - for entry in entries: - if entry.lower().startswith("npub1"): - hex_key = _npub_to_hex(entry) - if hex_key: - expanded.add(hex_key) - return expanded + return set(entries) | {h for e in entries if e.lower().startswith("npub1") and (h := _npub_to_hex(e))} def _principal_matches_allowlist(source, user_id: str, allowed_ids: set) -> bool: @@ -198,26 +183,28 @@ def _principal_matches_allowlist(source, user_id: str, allowed_ids: set) -> bool class GatewayAuthorizationMixin: """User/chat authorization methods for ``GatewayRunner``.""" - # ``getattr(self, ...)`` throughout: test helpers build bare runners via - # ``object.__new__`` without ``adapters`` / ``config``. + # ``getattr(self, ...)`` throughout: test helpers build bare runners via ``object.__new__`` + # without ``adapters`` / ``config``. def _primary_adapters(self) -> dict: return getattr(self, "adapters", None) or {} + def _profile_adapters_map(self) -> dict: + return getattr(self, "_profile_adapters", None) or {} + def _authorization_adapter(self, platform: Optional[Platform], profile: Optional[str] = None): """Live adapter whose intake policy gates authorization. - Secondary-profile adapters live in ``_profile_adapters[profile]``; the primary - profile owns ``self.adapters``. ``_profile_adapters`` is consulted BEFORE the - active profile name: multiplex turns override ``HERMES_HOME`` so - ``_active_profile_name()`` reports the secondary profile mid-turn, and treating - it as primary would hand it the default bot. + Secondary-profile adapters live in ``_profile_adapters[profile]``; the primary profile owns + ``self.adapters``. ``_profile_adapters`` is consulted BEFORE the active profile name: multiplex + turns override ``HERMES_HOME`` so ``_active_profile_name()`` reports the secondary profile + mid-turn, and treating it as primary would hand it the default bot. """ if not platform: return None profile_name = (profile or "").strip() or None if profile_name and profile_name != "default": - profile_adapters = getattr(self, "_profile_adapters", None) or {} + profile_adapters = self._profile_adapters_map() if profile_name in profile_adapters: return profile_adapters[profile_name].get(platform) # Identity captured at construction, not the per-turn HERMES_HOME-derived name. @@ -239,9 +226,9 @@ class GatewayAuthorizationMixin: transport_adapter = self._registered_transport_adapter(source) if transport_adapter is not None: return transport_adapter - # Relay ingress keeps the underlying platform on the source, but delivery must - # use the one process-level RelayAdapter owning the connector socket; a - # profile-aware lookup would silently disable streaming/typing/tool progress. + # Relay ingress keeps the underlying platform on the source, but delivery must use the one + # process-level RelayAdapter owning the connector socket; a profile-aware lookup would + # silently disable streaming/typing/tool progress. if getattr(source, "delivered_via_upstream_relay", False) is True: return self._primary_adapters().get(Platform.RELAY) # ``getattr``: test fixtures build bare SimpleNamespace sources without ``profile``. @@ -251,7 +238,7 @@ class GatewayAuthorizationMixin: """Return (registered, profile) for a live adapter: profile is None for primary.""" if adapter is self._primary_adapters().get(platform): return True, None - for profile, profile_adapters in (getattr(self, "_profile_adapters", None) or {}).items(): + for profile, profile_adapters in self._profile_adapters_map().items(): if adapter is profile_adapters.get(platform): return True, profile return False, None @@ -259,10 +246,9 @@ class GatewayAuthorizationMixin: def _registered_transport_adapter(self, source: SessionSource): """The registered adapter that created *source*, if retained. - ``source.profile`` may differ from the adapter profile when one shared credential - serves several routed runtimes; ``build_source`` keeps the receiving adapter as - provenance so replies stay on that transport. Restored/hand-built sources fall - back (fail-closed) to profile lookup. + ``source.profile`` may differ from the adapter profile when one shared credential serves + several routed runtimes; ``build_source`` keeps the receiving adapter as provenance so replies + stay on that transport. Restored/hand-built sources fall back (fail-closed) to profile lookup. """ adapter_ref = getattr(source, "_transport_adapter_ref", None) adapter = adapter_ref() if callable(adapter_ref) else None @@ -282,27 +268,15 @@ class GatewayAuthorizationMixin: return getattr(source, "profile", None) def _adapter_flag(self, platform, name: str, profile) -> bool: + """Adapter-declared boolean, False when unknown. ``authorization_is_upstream`` (relay: a trusted + authenticated upstream decides) is honored directly; ``enforces_own_access_policy`` (WeCom, Weixin, + Yuanbao, QQBot, WhatsApp gate at intake) is NOT "already authorized" — those adapters default to + ``open``, so ``_is_user_authorized`` only trusts them under an actual ``allowlist`` policy.""" if not platform: return False adapter = self._authorization_adapter(platform, profile) return adapter is not None and bool(getattr(adapter, name, False)) - def _adapter_authorization_is_upstream(self, platform: Optional[Platform], *, profile: Optional[str] = None) -> bool: - """Whether the adapter delegates authz to a trusted authenticated upstream (relay). - - Unlike ``_adapter_enforces_own_access_policy`` (a LOCAL policy only trusted when - it is an allowlist) this UPSTREAM decision is honored directly. False when unknown. - """ - return self._adapter_flag(platform, "authorization_is_upstream", profile) - - def _adapter_enforces_own_access_policy(self, platform: Optional[Platform], *, profile: Optional[str] = None) -> bool: - """Whether the adapter gates access at intake itself (WeCom, Weixin, Yuanbao, QQBot, WhatsApp). - - The flag alone is NOT "already authorized": these adapters default to ``open``, - so ``_is_user_authorized`` only trusts them under an actual ``allowlist`` policy. - """ - return self._adapter_flag(platform, "enforces_own_access_policy", profile) - def _config_extra(self, platform) -> dict: """``config.platforms[platform].extra`` as a dict ({} when absent).""" platforms = getattr(getattr(self, "config", None), "platforms", None) @@ -318,33 +292,24 @@ class GatewayAuthorizationMixin: value = self._config_extra(platform).get(extra_key) return value - def _adapter_policy(self, platform, attr: str, extra_key: str, profile) -> str: + def _adapter_policy(self, platform, kind: str, profile) -> str: + """Lowercased effective ``dm_policy`` (open/allowlist/disabled/pairing) or ``group_policy`` + (open/allowlist/disabled) for *kind* in {"dm", "group"}; ``""`` if unknown.""" if not platform: return "" - return str(self._adapter_setting(platform, attr, extra_key, profile) or "").strip().lower() - - def _adapter_dm_policy(self, platform: Optional[Platform], *, profile: Optional[str] = None) -> str: - """Lowercased effective ``dm_policy`` (open/allowlist/disabled/pairing), ``""`` if unknown.""" - return self._adapter_policy(platform, "_dm_policy", "dm_policy", profile) - - def _adapter_group_policy(self, platform: Optional[Platform], *, profile: Optional[str] = None) -> str: - """Lowercased effective ``group_policy`` (open/allowlist/disabled), ``""`` if unknown.""" - return self._adapter_policy(platform, "_group_policy", "group_policy", profile) + return str(self._adapter_setting(platform, f"_{kind}_policy", f"{kind}_policy", profile) or "").strip().lower() def _adapter_group_has_sender_allowlist( self, platform: Optional[Platform], chat_id: Optional[str], *, profile: Optional[str] = None ) -> bool: - """Whether a per-group sender allowlist (WeCom ``groups..allow_from``) gated this message. - - A group may be open at the chat level while restricting senders; reaching the - gateway then means the adapter already checked that list. - """ + """Whether a per-group sender allowlist (WeCom ``groups..allow_from``) gated this message: + a group may be open at the chat level while restricting senders, so reaching the gateway + means the adapter already checked that list.""" if not platform or not chat_id: return False groups = self._adapter_setting(platform, "_groups", "groups", profile) if not isinstance(groups, dict): return False - chat_id_str = str(chat_id) group_cfg = groups.get(chat_id_str) if not isinstance(group_cfg, dict): @@ -355,13 +320,10 @@ class GatewayAuthorizationMixin: ) if not isinstance(group_cfg, dict): return False - sender_allow = group_cfg.get("allow_from") or group_cfg.get("allowFrom") if isinstance(sender_allow, str): return bool(sender_allow.strip()) - if isinstance(sender_allow, (list, tuple, set)): - return any(str(item).strip() for item in sender_allow) - return False + return isinstance(sender_allow, (list, tuple, set)) and any(str(item).strip() for item in sender_allow) def _pairing_store_for(self, source: "SessionSource"): """Per-profile PairingStore for a source, else the global ``self.pairing_store``.""" @@ -377,22 +339,18 @@ class GatewayAuthorizationMixin: def _own_policy_authorizes(self, source, user_id, is_group, adapter_profile) -> Optional[bool]: """Own-policy adapter verdict when no env allowlist exists; None = no verdict. - Trusted only when the effective policy for THIS chat type is ``allowlist``: - ``open`` forwards EVERY sender (the fail-open SECURITY.md §2.6 forbids), - ``disabled`` never forwards, ``pairing`` forwards unpaired DMs for the handshake - (already denied by the pairing-store check). Anything else → default-deny. + Trusted only when the effective policy for THIS chat type is ``allowlist``: ``open`` forwards + EVERY sender (the fail-open SECURITY.md §2.6 forbids), ``disabled`` never forwards, ``pairing`` + forwards unpaired DMs for the handshake (already denied by the pairing-store check). + Anything else → default-deny. """ - if is_group: - if self._adapter_group_has_sender_allowlist(source.platform, source.chat_id, profile=adapter_profile): - return True - effective_policy = self._adapter_group_policy(source.platform, profile=adapter_profile) - else: - effective_policy = self._adapter_dm_policy(source.platform, profile=adapter_profile) - if effective_policy != "allowlist": + if is_group and self._adapter_group_has_sender_allowlist(source.platform, source.chat_id, profile=adapter_profile): + return True + if self._adapter_policy(source.platform, "group" if is_group else "dm", adapter_profile) != "allowlist": return None - # Re-check DMs via the live adapter's ``_is_dm_allowed`` when present: pairing - # revoke can clear WHATSAPP_ALLOWED_USERS while a construction-time snapshot - # would keep authorizing until restart. Others keep the historical rubber-stamp. + # Re-check DMs via the live adapter's ``_is_dm_allowed`` when present: pairing revoke can clear + # WHATSAPP_ALLOWED_USERS while a construction-time snapshot would keep authorizing until + # restart. Others keep the historical rubber-stamp. if not is_group: adapter = self._authorization_adapter(source.platform, profile=adapter_profile) dm_check = getattr(adapter, "_is_dm_allowed", None) if adapter is not None else None @@ -409,9 +367,9 @@ class GatewayAuthorizationMixin: extra = _adapter_config_extra(adapter) adapter_allow = extra.get("group_allow_from" if is_group else "allow_from") if not adapter_allow: - # Plugin platforms (Buzz, DingTalk) spell their env allowlist as - # ``extra.allowed_users``; under multiplex only the default profile's list - # reaches the env (first-writer-wins bridge), so read the live adapter's. + # Plugin platforms (Buzz, DingTalk) spell their env allowlist as ``extra.allowed_users``; + # under multiplex only the default profile's list reaches the env (first-writer-wins + # bridge), so read the live adapter's. entry = _registry_entry(source.platform) if entry and entry.allowed_users_env: adapter_allow = extra.get("allowed_users") @@ -427,11 +385,11 @@ class GatewayAuthorizationMixin: def _adapter_resolved_allowlist_ids(self, source) -> set[str]: """IDs an adapter resolved from username-shaped allowlist entries at connect time (Discord). - The per-turn .env hot-reload restores RAW usernames, so from the second turn on the - env allowlist holds usernames while user_id is numeric. Never a widening: the - empty-allowlist branch already returned and adapters only resolve operator-written - entries. Only called with a non-empty platform allowlist so group/global-only - configs never consult adapter memory; type-checked so mocks cannot auto-truthy in. + The per-turn .env hot-reload restores RAW usernames, so from the second turn on the env + allowlist holds usernames while user_id is numeric. Never a widening: the empty-allowlist + branch already returned and adapters only resolve operator-written entries. Only called with + a non-empty platform allowlist so group/global-only configs never consult adapter memory; + type-checked so mocks cannot auto-truthy in. """ adapter = resolved_ids = None with contextlib.suppress(Exception): @@ -444,14 +402,64 @@ class GatewayAuthorizationMixin: return set() return {str(entry).strip() for entry in resolved_ids if isinstance(entry, (str, int)) and str(entry).strip()} + def _chat_scoped_grant(self, source, adapter_profile, is_group: bool, allow_adapter_delegation: bool) -> bool: + """Grants that need no ``user_id`` (checked before the no-user-id guard).""" + # Trusted-upstream delegation (relay): the connector authenticates this gateway's WS and + # resolves owner bindings BEFORE delivering, so there is no local RELAY_ALLOWED_USERS. Not a + # fail-open: fires only for events actually delivered over the relay WS + # (``delivered_via_upstream_relay``) or whose adapter declares ``authorization_is_upstream``. + # The delivery marker is PRIMARY because a relayed message carries the UNDERLYING platform, + # not ``Platform.RELAY``. ``is True``: a MagicMock stand-in must not auto-truthy into authz. + if allow_adapter_delegation and ( + source.delivered_via_upstream_relay is True + or self._adapter_flag(source.platform, "authorization_is_upstream", adapter_profile) + ): + return True + # Chat-scoped group allowlists must work with ``user_id is None`` (anonymous admins, + # sender_chat posts, channel broadcasts). + if is_group and source.chat_id: + chat_allowlist_env = _GROUP_CHAT_ENV.get(source.platform, "") + if chat_allowlist_env and _allows(_coerce_allow_set(_platform_gate_env(chat_allowlist_env)), source.chat_id): + return True + # config.yaml fallback (``extra.group_allowed_chats``): Telegram observe-unmentioned mode + # strips user_id, so the env-only check above misses it. + with contextlib.suppress(Exception): + adapter_group_allowed = self._adapter_extra_for_source(source).get("group_allowed_chats") + if adapter_group_allowed and _allows(_coerce_allow_set(adapter_group_allowed), source.chat_id): + return True + # Bots admitted by {PLATFORM}_ALLOW_BOTS bypass the human allowlist (Slack Workflow Builder + # posts arrive with user=None). + if getattr(source, "is_bot", False): + allow_bots_var = _ALLOW_BOTS_ENV.get(source.platform) + if allow_bots_var and _platform_gate_env(allow_bots_var, "none").lower().strip() in {"mentions", "all"}: + return True + return False + + def _legacy_telegram_chat_grant(self, source, group_user_allowlist: str) -> bool: + """TELEGRAM_GROUP_ALLOWED_USERS was once (mis)used as a chat-ID allowlist; "-"-prefixed + values are chat IDs, honor them and warn once.""" + from gateway.run import logger + legacy_chat_ids = {v.strip() for v in group_user_allowlist.split(",") if v.strip().startswith("-")} + if not legacy_chat_ids: + return False + if not getattr(self, "_warned_telegram_group_users_legacy", False): + logger.warning( + "TELEGRAM_GROUP_ALLOWED_USERS contains chat-ID-shaped values " + "(%s). Treating them as chat IDs for backward compatibility. " + "Move chat IDs to TELEGRAM_GROUP_ALLOWED_CHATS — the _USERS var " + "is now for sender user IDs.", + ",".join(sorted(legacy_chat_ids)), + ) + self._warned_telegram_group_users_legacy = True + return source.chat_id in legacy_chat_ids + def _is_user_authorized(self, source: SessionSource, *, allow_adapter_delegation: bool = True) -> bool: """Whether a user may use the bot. - Order: trusted-upstream delegation, chat-scoped group allowlists, - ``{PLATFORM}_ALLOW_BOTS``, per-platform allow-all, adapter role auth, pairing - store, env/config allowlists, ``GATEWAY_ALLOW_ALL_USERS``, default deny. + Order: trusted-upstream delegation, chat-scoped group allowlists, ``{PLATFORM}_ALLOW_BOTS``, + per-platform allow-all, adapter role auth, pairing store, env/config allowlists, + ``GATEWAY_ALLOW_ALL_USERS``, default deny. """ - from gateway.run import logger # HA events are system-generated (HASS_TOKEN); webhook events are HMAC-verified. if source.platform in {Platform.HOMEASSISTANT, Platform.WEBHOOK}: return True @@ -459,42 +467,9 @@ class GatewayAuthorizationMixin: adapter_profile = self._adapter_profile_for_source(source) is_group = source.chat_type in _GROUP_CHAT_TYPES is_group_or_forum = source.chat_type in _GROUP_FORUM_TYPES - - # Trusted-upstream delegation (relay): the connector authenticates this - # gateway's WS and resolves owner bindings BEFORE delivering, so there is no - # local RELAY_ALLOWED_USERS. Not a fail-open: fires only for events actually - # delivered over the relay WS (``delivered_via_upstream_relay``) or whose adapter - # declares ``authorization_is_upstream``. The delivery marker is PRIMARY because - # a relayed message carries the UNDERLYING platform, not ``Platform.RELAY``. - # ``is True``: a MagicMock stand-in must not auto-truthy into authz. - if allow_adapter_delegation and ( - source.delivered_via_upstream_relay is True - or self._adapter_authorization_is_upstream(source.platform, profile=adapter_profile) - ): + if self._chat_scoped_grant(source, adapter_profile, is_group, allow_adapter_delegation): return True - user_id = source.user_id - - # Chat-scoped group allowlists must work with ``user_id is None`` (anonymous - # admins, sender_chat posts, channel broadcasts): run before the no-user-id guard. - if is_group and source.chat_id: - chat_allowlist_env = _GROUP_CHAT_ENV.get(source.platform, "") - if chat_allowlist_env and _allows(_coerce_allow_set(_platform_gate_env(chat_allowlist_env)), source.chat_id): - return True - # config.yaml fallback (``extra.group_allowed_chats``): Telegram observe- - # unmentioned mode strips user_id, so the env-only check above misses it. - with contextlib.suppress(Exception): - adapter_group_allowed = self._adapter_extra_for_source(source).get("group_allowed_chats") - if adapter_group_allowed and _allows(_coerce_allow_set(adapter_group_allowed), source.chat_id): - return True - - # Bots admitted by {PLATFORM}_ALLOW_BOTS bypass the human allowlist. Also before - # the no-user-id guard: Slack Workflow Builder posts arrive with user=None. - if getattr(source, "is_bot", False): - allow_bots_var = _ALLOW_BOTS_ENV.get(source.platform) - if allow_bots_var and _platform_gate_env(allow_bots_var, "none").lower().strip() in {"mentions", "all"}: - return True - if not user_id: return False @@ -505,34 +480,25 @@ class GatewayAuthorizationMixin: with contextlib.suppress(Exception): platform_allow_env = getattr(entry, "allowed_users_env", "") or platform_allow_env platform_allow_all_var = getattr(entry, "allow_all_env", "") or platform_allow_all_var - if platform_allow_all_var and _env_truthy(platform_allow_all_var): return True - # Adapter-verified role auth (Discord DISCORD_ALLOWED_ROLES). ``is True``: no MagicMock pass. if allow_adapter_delegation and getattr(source, "role_authorized", False) is True: return True - - # Pairing store: a first-class grant created only by an operator approving a - # code. Honored as a UNION with the allowlist (approval also mirrors into it). - platform_name = source.platform.value if source.platform else "" + # Pairing store: a first-class grant created only by an operator approving a code. Honored as + # a UNION with the allowlist (approval also mirrors into it). pairing_store = self._pairing_store_for(source) - if pairing_store is not None and pairing_store.is_approved(platform_name, user_id): + if pairing_store is not None and pairing_store.is_approved(source.platform.value if source.platform else "", user_id): return True platform_allowlist = _auth_env(platform_allow_env) - group_user_allowlist = "" - group_chat_allowlist = "" - if is_group_or_forum: - group_user_allowlist = _auth_env(_GROUP_USER_ENV.get(source.platform, "")) - group_chat_allowlist = _auth_env(_GROUP_CHAT_ENV.get(source.platform, "")) + group_user_allowlist = _auth_env(_GROUP_USER_ENV.get(source.platform, "")) if is_group_or_forum else "" + group_chat_allowlist = _auth_env(_GROUP_CHAT_ENV.get(source.platform, "")) if is_group_or_forum else "" global_allowlist = _auth_env("GATEWAY_ALLOWED_USERS") if not (platform_allowlist or group_user_allowlist or group_chat_allowlist or global_allowlist): # No env allowlist: own-policy adapters gate at intake (see _own_policy_authorizes). - if allow_adapter_delegation and self._adapter_enforces_own_access_policy( - source.platform, profile=adapter_profile - ): + if allow_adapter_delegation and self._adapter_flag(source.platform, "enforces_own_access_policy", adapter_profile): verdict = self._own_policy_authorizes(source, user_id, is_group, adapter_profile) if verdict is not None: return verdict @@ -540,30 +506,15 @@ class GatewayAuthorizationMixin: return True return _env_truthy("GATEWAY_ALLOW_ALL_USERS") - # Telegram group traffic authorized by chat ID (separate from - # TELEGRAM_GROUP_ALLOWED_USERS, which gates the sender). - if ( - group_chat_allowlist and is_group_or_forum and source.chat_id - and _allows(_coerce_allow_set(group_chat_allowlist), source.chat_id) - ): - return True - - # Backward-compat: TELEGRAM_GROUP_ALLOWED_USERS was once (mis)used as a chat-ID - # allowlist; "-"-prefixed values are chat IDs, honor them and warn once. - if source.platform == Platform.TELEGRAM and group_user_allowlist and is_group_or_forum and source.chat_id: - legacy_chat_ids = {v.strip() for v in group_user_allowlist.split(",") if v.strip().startswith("-")} - if legacy_chat_ids: - if not getattr(self, "_warned_telegram_group_users_legacy", False): - logger.warning( - "TELEGRAM_GROUP_ALLOWED_USERS contains chat-ID-shaped values " - "(%s). Treating them as chat IDs for backward compatibility. " - "Move chat IDs to TELEGRAM_GROUP_ALLOWED_CHATS — the _USERS var " - "is now for sender user IDs.", - ",".join(sorted(legacy_chat_ids)), - ) - self._warned_telegram_group_users_legacy = True - if source.chat_id in legacy_chat_ids: - return True + if is_group_or_forum and source.chat_id: + # Telegram group traffic authorized by chat ID (TELEGRAM_GROUP_ALLOWED_USERS gates the sender). + if group_chat_allowlist and _allows(_coerce_allow_set(group_chat_allowlist), source.chat_id): + return True + if ( + source.platform == Platform.TELEGRAM and group_user_allowlist + and self._legacy_telegram_chat_grant(source, group_user_allowlist) + ): + return True # TELEGRAM_GROUP_ALLOWED_USERS is group-scoped (no DM access); TELEGRAM_ALLOWED_USERS is platform-wide. allowed_ids = ( @@ -578,28 +529,24 @@ class GatewayAuthorizationMixin: def _get_unauthorized_dm_behavior(self, platform: Optional[Platform], *, profile: Optional[str] = None) -> str: """How unauthorized DMs are handled ("pair" / "ignore") for a platform. - Order: explicit per-platform config; Email → "ignore" (inboxes hold arbitrary - mail); explicit non-default global; adapter dm_policy (pairing → "pair", - allowlist/disabled → "ignore"); any configured allowlist → "ignore" (spamming - unknown contacts with codes is noisy and leaks); else "pair". + Order: explicit per-platform config; Email → "ignore" (inboxes hold arbitrary mail); explicit + non-default global; adapter dm_policy (pairing → "pair", allowlist/disabled → "ignore"); any + configured allowlist → "ignore" (spamming unknown contacts with codes is noisy and leaks); else "pair". """ config = getattr(self, "config", None) - if ( config and hasattr(config, "get_unauthorized_dm_behavior") and platform and "unauthorized_dm_behavior" in self._config_extra(platform) ): return config.get_unauthorized_dm_behavior(platform) - if platform == Platform.EMAIL: return "ignore" - if config and hasattr(config, "unauthorized_dm_behavior") and config.unauthorized_dm_behavior != "pair": return config.unauthorized_dm_behavior allowlist_keys = ["GATEWAY_ALLOWED_USERS"] if platform: - dm_policy = self._adapter_dm_policy(platform, profile=profile) + dm_policy = self._adapter_policy(platform, "dm", profile) if not dm_policy: dm_policy = str(self._config_extra(platform).get("dm_policy") or "").strip().lower() if dm_policy == "pairing": diff --git a/gateway/config.py b/gateway/config.py index c428da7be5..f330b3c65f 100644 --- a/gateway/config.py +++ b/gateway/config.py @@ -26,17 +26,19 @@ _TRUTHY_STRINGS = frozenset({"1", "true", "yes", "on"}) _FALSY_STRINGS = frozenset({"0", "false", "no", "off"}) +def _bool_token(value: Any) -> Optional[bool]: + """True/False for a recognized truthy/falsy token, else None.""" + token = str(value).strip().lower() + return True if token in _TRUTHY_STRINGS else False if token in _FALSY_STRINGS else None + + def _coerce_bool(value: Any, default: bool = True) -> bool: """Coerce bool-ish config values, preserving a caller-provided default.""" if value is None: return default if isinstance(value, str): - lowered = value.strip().lower() - if lowered in _TRUTHY_STRINGS: - return True - if lowered in _FALSY_STRINGS: - return False - return default + parsed = _bool_token(value) + return default if parsed is None else parsed return is_truthy_value(value, default=default) @@ -80,19 +82,16 @@ def _env_multiplex_profiles_override() -> "bool | None": a config.yaml opt-in. """ raw = os.getenv("GATEWAY_MULTIPLEX_PROFILES") - token = (raw or "").strip().lower() - if not token: + if not (raw or "").strip(): return None - if token in _TRUTHY_STRINGS: - return True - if token in _FALSY_STRINGS: - return False - logger.warning( - "Ignoring unrecognized GATEWAY_MULTIPLEX_PROFILES=%r " - "(expected one of %s or %s); falling back to config.yaml.", - raw, sorted(_TRUTHY_STRINGS), sorted(_FALSY_STRINGS), - ) - return None + parsed = _bool_token(raw) + if parsed is None: + logger.warning( + "Ignoring unrecognized GATEWAY_MULTIPLEX_PROFILES=%r " + "(expected one of %s or %s); falling back to config.yaml.", + raw, sorted(_TRUTHY_STRINGS), sorted(_FALSY_STRINGS), + ) + return parsed def _normalize_transport_token(value: Any) -> str: @@ -148,13 +147,9 @@ def coerce_systemd_watchdog_seconds( parsed: Optional[int] = None if isinstance(value, int) and not isinstance(value, bool): parsed = value - elif isinstance(value, str): - raw = value.strip() - if raw and raw.isascii() and raw.isdecimal(): - try: - parsed = int(raw, 10) - except (TypeError, ValueError, OverflowError): - parsed = None + elif isinstance(value, str) and value.strip().isascii() and value.strip().isdecimal(): + with contextlib.suppress(TypeError, ValueError, OverflowError): # int() digit limit + parsed = int(value.strip(), 10) if parsed is None: logger.warning("Ignoring invalid %s (expected a positive integer)", key) return 0 @@ -200,8 +195,7 @@ def _getenv_str(name: str, default: str = "") -> str: return val if val is not None else default -# Bundled platform plugin names, cached outside the enum so it never becomes a member. -_Platform__bundled_plugin_names: Optional[set] = None +_Platform__bundled_plugin_names: Optional[set] = None # cached outside the enum: never a member class Platform(Enum): @@ -279,19 +273,15 @@ class Platform(Enum): # Built-in values snapshotted before any dynamic _missing_ lookup. _BUILTIN_PLATFORM_VALUES = frozenset(m.value for m in Platform.__members__.values()) - -# Platforms that bind a host TCP port. In a multiplexer only the default profile owns -# the shared listener, so a SECONDARY profile enabling one is a misconfiguration. -# Single source of truth for gateway/run.py and hermes_cli/web_server.py validation. +# Platforms that bind a host TCP port. In a multiplexer only the default profile owns the +# shared listener, so a SECONDARY profile enabling one is a misconfiguration (single source +# of truth for gateway/run.py and hermes_cli/web_server.py validation). PORT_BINDING_PLATFORM_VALUES = frozenset({ "webhook", "api_server", "msgraph_webhook", "feishu", "wecom_callback", "bluebubbles", "sms", "whatsapp_cloud", "line", "teams", }) - # Platforms that only bind in one connection mode (Feishu's default websocket mode is outbound). -PORT_BINDING_CONDITIONAL_MODES: dict[str, str] = { - "feishu": "webhook", -} +PORT_BINDING_CONDITIONAL_MODES: dict[str, str] = {"feishu": "webhook"} def platform_binds_port(platform_value: str, extra: Optional[dict] = None) -> bool: @@ -299,10 +289,7 @@ def platform_binds_port(platform_value: str, extra: Optional[dict] = None) -> bo if platform_value not in PORT_BINDING_PLATFORM_VALUES: return False expected_mode = PORT_BINDING_CONDITIONAL_MODES.get(platform_value) - if expected_mode is not None: - actual = str((extra or {}).get("connection_mode", "websocket")).strip().lower() - return actual == expected_mode - return True + return expected_mode is None or str((extra or {}).get("connection_mode", "websocket")).strip().lower() == expected_mode @dataclass @@ -313,8 +300,7 @@ class HomeChannel: chat_id: str name: str thread_id: Optional[str] = None - # Authenticated logical-target provenance; relay egress re-attaches these but the - # connector remains the authorization boundary. + # Authenticated logical-target provenance (relay egress re-attaches; connector stays the authz boundary). user_id: Optional[str] = None scope_id: Optional[str] = None @@ -350,8 +336,7 @@ class SessionResetPolicy: idle_minutes: int = 1440 notify: bool = True # Notify the user when auto-reset occurs notify_exclude_platforms: tuple = ("api_server", "webhook") - # A background process this old no longer blocks reset (a forgotten preview server - # must not pin a session forever); it is NOT killed, only ignored by the guard. + # A background process this old no longer blocks reset (not killed, only ignored by the guard). bg_process_max_age_hours: int = 24 def to_dict(self) -> Dict[str, Any]: @@ -385,9 +370,7 @@ class ChannelOverride: @classmethod def from_dict(cls, data: Dict[str, Any]) -> "ChannelOverride": - if not data: - return cls() - return cls(model=data.get("model"), provider=data.get("provider"), system_prompt=data.get("system_prompt")) + return cls(**{f.name: data.get(f.name) for f in fields(cls)}) if data else cls() # Platforms whose primary credential is ``PlatformConfig.token`` → its env var (empty-token @@ -410,15 +393,10 @@ class PlatformConfig: token: Optional[str] = None api_key: Optional[str] = None # API key if different from token home_channel: Optional[HomeChannel] = None - # Reply threading: "off" never threads, "first" threads only the first chunk, "all" every chunk. - reply_to_mode: str = "first" - # "♻️ Gateway online/restarted" pings; False on end-user platforms where they are noise. - gateway_restart_notification: bool = True - # "typing…" indicator (drives _keep_typing in platforms/base.py); False where unwanted - # (Slack's setStatus disables the compose box). - typing_indicator: bool = True - # Working-state text for text-rendering indicators (Slack assistant status, Google - # Chat marker); None = platform default; textless indicators ignore it. + reply_to_mode: str = "first" # "off" never threads, "first" only the first chunk, "all" every chunk + gateway_restart_notification: bool = True # "♻️ Gateway online/restarted" pings; noise on end-user platforms + typing_indicator: bool = True # drives _keep_typing; False where unwanted (Slack setStatus blocks compose) + # Working-state text for text-rendering indicators (Slack status, Google Chat marker); None = platform default. typing_status_text: Optional[str] = None channel_overrides: Dict[str, ChannelOverride] = field(default_factory=dict) extra: Dict[str, Any] = field(default_factory=dict) # Platform-specific settings @@ -428,10 +406,9 @@ class PlatformConfig: "enabled": self.enabled, "extra": self.extra, "reply_to_mode": self.reply_to_mode, "gateway_restart_notification": self.gateway_restart_notification, "typing_indicator": self.typing_indicator, + **({"typing_status_text": self.typing_status_text} if self.typing_status_text is not None else {}), + **{k: v for k in ("token", "api_key") if (v := getattr(self, k))}, } - if self.typing_status_text is not None: - result["typing_status_text"] = self.typing_status_text - result.update({k: v for k in ("token", "api_key") if (v := getattr(self, k))}) if self.home_channel: result["home_channel"] = self.home_channel.to_dict() if self.channel_overrides: @@ -441,10 +418,7 @@ class PlatformConfig: @classmethod def from_dict(cls, data: Dict[str, Any]) -> "PlatformConfig": data = _coerce_dict(data) - home_channel = None - if isinstance(data.get("home_channel"), dict): - home_channel = HomeChannel.from_dict(data["home_channel"]) - + home = data.get("home_channel") # The typing/restart-notification keys may be top-level or bridged into ``extra``; top-level wins. extra = _coerce_dict(data.get("extra", {})) @@ -463,7 +437,7 @@ class PlatformConfig: enabled=_coerce_bool(data.get("enabled"), False), token=data.get("token"), api_key=data.get("api_key"), - home_channel=home_channel, + home_channel=HomeChannel.from_dict(home) if isinstance(home, dict) else None, reply_to_mode=data.get("reply_to_mode", "first"), gateway_restart_notification=_coerce_bool(toplevel_or_extra("gateway_restart_notification"), True), typing_indicator=_coerce_bool(toplevel_or_extra("typing_indicator"), True), @@ -484,15 +458,13 @@ DEFAULT_STREAMING_CURSOR: str = " ▉" class StreamingConfig: """Real-time token streaming to messaging platforms.""" enabled: bool = False - # "auto" prefers native drafts (Telegram sendMessageDraft) with edit fallback — safe - # globally since adapters without draft support use the edit path unchanged; - # "draft" / "edit" force one; "off" disables. + # "auto" prefers native drafts (Telegram sendMessageDraft) with edit fallback (adapters without + # draft support use the edit path unchanged); "draft" / "edit" force one; "off" disables. transport: str = "auto" edit_interval: float = DEFAULT_STREAMING_EDIT_INTERVAL buffer_threshold: int = DEFAULT_STREAMING_BUFFER_THRESHOLD cursor: str = DEFAULT_STREAMING_CURSOR - # >0: deliver the final edit as a fresh message when the preview has been visible - # this long, so the timestamp reflects completion. Telegram only; 0 disables. + # >0: final edit becomes a fresh message once the preview was visible this long (Telegram only; 0 = off). fresh_final_after_seconds: float = 0.0 def to_dict(self) -> Dict[str, Any]: @@ -508,18 +480,13 @@ class StreamingConfig: # ``streaming.enabled`` is the documented master switch. raw_transport = data.get("transport") raw_mode = data.get("mode") - transport = _normalize_transport_token(raw_transport if raw_transport is not None else raw_mode) - if "enabled" in data: enabled = _coerce_bool(data.get("enabled"), False) - elif raw_mode is not None: - enabled = _normalize_transport_token(raw_mode) != "off" else: - enabled = False - + enabled = raw_mode is not None and _normalize_transport_token(raw_mode) != "off" return cls( enabled=enabled, - transport=transport, + transport=_normalize_transport_token(raw_transport if raw_transport is not None else raw_mode), edit_interval=_coerce_float(data.get("edit_interval"), DEFAULT_STREAMING_EDIT_INTERVAL), buffer_threshold=_coerce_int(data.get("buffer_threshold"), DEFAULT_STREAMING_BUFFER_THRESHOLD), cursor=data.get("cursor", DEFAULT_STREAMING_CURSOR), @@ -575,47 +542,41 @@ class GatewayConfig: reset_by_type: Dict[str, SessionResetPolicy] = field(default_factory=dict) reset_by_platform: Dict[Platform, SessionResetPolicy] = field(default_factory=dict) reset_triggers: List[str] = field(default_factory=lambda: ["/new", "/reset"]) - # Slash commands that bypass the agent loop. - quick_commands: Dict[str, Any] = field(default_factory=dict) + quick_commands: Dict[str, Any] = field(default_factory=dict) # slash commands that bypass the agent loop sessions_dir: Path = field(default_factory=lambda: get_hermes_home() / "sessions") - # Keep the legacy sessions.json mirror of the routing index (primary: state.db) - # for external tooling and downgrade safety. + # Legacy sessions.json mirror of the routing index (primary: state.db) for external tooling / downgrades. write_sessions_json: bool = True always_log_local: bool = True # Always save cron outputs to local files - # Drop outbound "silence narration" (*(silent)*, 🔇, a bare ".") that ping-pongs in - # bot-to-bot channels; a substrate guard that survives prompt drift. + # Drop outbound "silence narration" (*(silent)*, 🔇, a bare ".") that ping-pongs in bot-to-bot + # channels; a substrate guard that survives prompt drift. filter_silence_narration: bool = True stt_enabled: bool = True # Auto-transcribe inbound voice messages stt_echo_transcripts: bool = True # Echo raw STT transcripts back to the user group_sessions_per_user: bool = True # Isolate group sessions per participant when user IDs exist thread_sessions_per_user: bool = False # False = threads shared across participants max_concurrent_sessions: Optional[int] = None # Positive int caps simultaneous active sessions - # Opt-in: the default profile's gateway serves every profile on the host - # (profiles stamped into session keys, per-profile adapters/credentials). + # Opt-in: the default profile's gateway serves every profile on the host (profiles stamped into + # session keys, per-profile adapters/credentials). Allowlist None = serve all; [] = default only. multiplex_profiles: bool = False - # None = historical serve-all; [] = default profile only. multiplex_profile_allowlist: Optional[List[str]] = None - # Public HTTPS endpoint for scoped RoomLink calls; an API key alone must never - # advertise a route. HERMES_ROOM_LINK_URL overrides. + # Public HTTPS endpoint for scoped RoomLink calls (an API key alone must never advertise a + # route); HERMES_ROOM_LINK_URL overrides. room_link_url: Optional[str] = None - # Opt-in systemd event-loop watchdog; zero keeps Type=simple and disables sd_notify. - systemd_watchdog_seconds: int = 0 - # In-process loop liveness watchdog: after consecutive missed probes it dumps - # all-thread stacks and hard-exits with the service-restart code. The knobs - # tolerate transient self-recovering stalls (adapter reconnect doing sync socket - # I/O) so a short block does not cause restart churn; a genuine wedge still escalates. + systemd_watchdog_seconds: int = 0 # opt-in; zero keeps Type=simple and disables sd_notify + # In-process loop liveness watchdog: after consecutive missed probes it dumps all-thread stacks + # and hard-exits with the service-restart code. The knobs tolerate transient self-recovering + # stalls (adapter reconnect doing sync socket I/O) so a short block does not cause restart churn. + # max_strikes ~= 90-120s sustained block; the heartbeat-fsync false positive is fixed at the root + # (off-loop write + two-witness probe), so raising it would only delay recovery. loop_watchdog: bool = True loop_watchdog_probe_interval_s: float = DEFAULT_LOOP_WATCHDOG_INTERVAL_S loop_watchdog_probe_timeout_s: float = DEFAULT_LOOP_WATCHDOG_TIMEOUT_S - # ~90-120s sustained block. The heartbeat-fsync false positive is fixed at the root - # (off-loop write + two-witness probe), so raising this would only delay recovery. loop_watchdog_max_strikes: int = DEFAULT_LOOP_WATCHDOG_MAX_STRIKES unauthorized_dm_behavior: str = "pair" # "pair" or "ignore" streaming: StreamingConfig = field(default_factory=StreamingConfig) # Prune SessionEntry records older than this (a resumed chat gets a fresh session). 0 = off. session_store_max_age_days: int = 90 - # Route guilds/channels/threads to profiles (gateway/profile_routing.py). - profile_routes: list = field(default_factory=list) + profile_routes: list = field(default_factory=list) # gateway/profile_routing.py # Scalar fields serialized verbatim by ``to_dict`` (in output order). _SCALAR_DICT_FIELDS = ( @@ -634,11 +595,7 @@ class GatewayConfig: def get_connected_platforms(self) -> List[Platform]: """Enabled + configured platforms, sorted by value so the rendered "Connected Platforms" prompt block is byte-stable (a reorder busts the prompt cache).""" - connected = [ - platform - for platform, config in self.platforms.items() - if config.enabled and self._is_platform_connected(platform, config) - ] + connected = [p for p, c in self.platforms.items() if c.enabled and self._is_platform_connected(p, c)] return sorted(connected, key=lambda p: str(p.value)) def _is_platform_connected(self, platform: Platform, config: PlatformConfig) -> bool: @@ -659,11 +616,8 @@ class GatewayConfig: discover_plugins() entry = platform_registry.get(platform.value) if entry: - if entry.is_connected is not None: - return entry.is_connected(config) - if entry.validate_config is not None: - return entry.validate_config(config) - return True + check = entry.is_connected if entry.is_connected is not None else entry.validate_config + return True if check is None else check(config) except Exception: pass # Registry not yet initialised during early import return False @@ -681,7 +635,7 @@ class GatewayConfig: return self.default_reset_policy def to_dict(self) -> Dict[str, Any]: - result: Dict[str, Any] = { + return { "platforms": {p.value: c.to_dict() for p, c in self.platforms.items()}, "default_reset_policy": self.default_reset_policy.to_dict(), "reset_by_type": {k: v.to_dict() for k, v in self.reset_by_type.items()}, @@ -689,16 +643,13 @@ class GatewayConfig: "reset_triggers": self.reset_triggers, "quick_commands": self.quick_commands, "sessions_dir": str(self.sessions_dir), + **{name: getattr(self, name) for name in self._SCALAR_DICT_FIELDS}, + "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 + ], } - for name in self._SCALAR_DICT_FIELDS: - result[name] = getattr(self, name) - result["streaming"] = self.streaming.to_dict() - result["session_store_max_age_days"] = self.session_store_max_age_days - result["profile_routes"] = [ - asdict(r) if is_dataclass(r) and not isinstance(r, type) else r - for r in self.profile_routes - ] - return result @classmethod def from_dict(cls, data: Dict[str, Any]) -> "GatewayConfig": @@ -742,17 +693,15 @@ class GatewayConfig: systemd_watchdog_seconds = coerce_systemd_watchdog_seconds( pick("systemd_watchdog_seconds"), key_label("systemd_watchdog_seconds") ) - - # env > config.yaml > False: a recognized GATEWAY_MULTIPLEX_PROFILES wins (hosted - # deployments stamp it on the container); blank/unrecognized falls through to the - # top-level VALUE when not None, else ``gateway.multiplex_profiles``. + # env > config.yaml > False: a recognized GATEWAY_MULTIPLEX_PROFILES wins (hosted deployments + # stamp it on the container); blank/unrecognized falls through to the top-level VALUE when + # not None, else ``gateway.multiplex_profiles``. multiplex_profiles = data.get("multiplex_profiles") if multiplex_profiles is None: multiplex_profiles = nested_gateway.get("multiplex_profiles") env_multiplex = _env_multiplex_profiles_override() if env_multiplex is not None: multiplex_profiles = env_multiplex - max_concurrent_sessions = _coerce_optional_positive_int( pick("max_concurrent_sessions"), key_label("max_concurrent_sessions") ) @@ -789,39 +738,31 @@ class GatewayConfig: loop_watchdog_probe_timeout_s=bounded_float("loop_watchdog_probe_timeout_s", DEFAULT_LOOP_WATCHDOG_TIMEOUT_S, 1.0, 600.0), loop_watchdog_max_strikes=max_strikes, max_concurrent_sessions=max_concurrent_sessions, - unauthorized_dm_behavior=_normalize_choice( - data.get("unauthorized_dm_behavior"), {"pair", "ignore"}, "pair" - ), + unauthorized_dm_behavior=_normalize_choice(data.get("unauthorized_dm_behavior"), {"pair", "ignore"}, "pair"), streaming=StreamingConfig.from_dict(data.get("streaming", {})), session_store_max_age_days=session_store_max_age_days, profile_routes=parse_profile_routes(data.get("profile_routes") or []), ) + def _extra_choice(self, platform: Optional[Platform], key: str, choices: set, default: str) -> Optional[str]: + """Normalized ``platforms[platform].extra[key]`` when the key is present, else None.""" + platform_cfg = self.platforms.get(platform) if platform else None + if platform_cfg and key in platform_cfg.extra: + return _normalize_choice(platform_cfg.extra.get(key), choices, default) + return None + def get_unauthorized_dm_behavior(self, platform: Optional[Platform] = None) -> str: - """Effective unauthorized-DM behavior. Email is inbox-shaped so it defaults to - ``"ignore"`` unless its own ``unauthorized_dm_behavior`` opts in (a global - default does not).""" - if platform: - platform_cfg = self.platforms.get(platform) - if platform_cfg and "unauthorized_dm_behavior" in platform_cfg.extra: - return _normalize_choice( - platform_cfg.extra.get("unauthorized_dm_behavior"), - {"pair", "ignore"}, - self.unauthorized_dm_behavior, - ) - if platform == Platform.EMAIL: - return "ignore" - return self.unauthorized_dm_behavior + """Effective unauthorized-DM behavior. Email is inbox-shaped so it defaults to ``"ignore"`` + unless its own ``unauthorized_dm_behavior`` opts in (a global default does not).""" + choice = self._extra_choice(platform, "unauthorized_dm_behavior", {"pair", "ignore"}, self.unauthorized_dm_behavior) + if choice is not None: + return choice + return "ignore" if platform == Platform.EMAIL else self.unauthorized_dm_behavior def get_notice_delivery(self, platform: Optional[Platform] = None) -> str: """Effective notice-delivery mode ("public"/"private") for a platform.""" - if platform: - platform_cfg = self.platforms.get(platform) - if platform_cfg and "notice_delivery" in platform_cfg.extra: - return _normalize_choice( - platform_cfg.extra.get("notice_delivery"), {"public", "private"}, "public" - ) - return "public" + choice = self._extra_choice(platform, "notice_delivery", {"public", "private"}, "public") + return "public" if choice is None else choice def load_gateway_config() -> GatewayConfig: @@ -866,14 +807,12 @@ def _validate_gateway_config(config: "GatewayConfig") -> None: (p, c, PLATFORM_TOKEN_ENV_NAMES[p]) for p, c in config.platforms.items() if c.enabled and p in PLATFORM_TOKEN_ENV_NAMES and c.token is not None ] - # An empty token won't connect; say so. - for platform, pconfig, env_name in token_platforms: + for platform, pconfig, env_name in token_platforms: # an empty token won't connect; say so if not pconfig.token.strip(): logger.warning("%s is enabled but %s is empty. The adapter will likely fail to connect.", platform.value, env_name) if has_usable_secret is None: return - # Reject placeholder tokens (copied .env.example) with a clear startup error. - for platform, pconfig, env_name in token_platforms: + for platform, pconfig, env_name in token_platforms: # reject placeholder tokens (copied .env.example) token = pconfig.token if token.strip() and not has_usable_secret(token, min_length=4): logger.error( @@ -888,5 +827,4 @@ def _validate_gateway_config(config: "GatewayConfig") -> None: def _apply_env_overrides(config: GatewayConfig) -> None: """Apply environment variable overrides to config (see ``gateway.config_env``).""" from gateway.config_env import _apply_env_overrides as _impl - _impl(config) diff --git a/gateway/config_env.py b/gateway/config_env.py index f448a0d7e0..eb67de3967 100644 --- a/gateway/config_env.py +++ b/gateway/config_env.py @@ -1,10 +1,9 @@ """Environment-variable overrides for the gateway config (``_apply_env_overrides``). -Runs after ``GatewayConfig.from_dict`` so env always wins over config.yaml / -gateway.json. Most platforms follow one shape — "credentials present in env -⇒ enable the platform and copy the values into ``extra``" — declared as -``_Cred`` rows in ``_ENV_STEPS`` (source order = application order). -Platforms with unique gating are small functions in the same table. +Runs after ``GatewayConfig.from_dict`` so env always wins over config.yaml / gateway.json. +Most platforms follow one shape — "credentials present in env ⇒ enable the platform and +copy the values into ``extra``" — declared as ``_Cred`` rows in ``_ENV_STEPS`` (source +order = application order). Platforms with unique gating are small functions in the same table. """ import contextlib @@ -30,13 +29,11 @@ logger = logging.getLogger("gateway.config") getenv = _getenv_str -# Platforms already warned about "explicitly disabled in config.yaml, but -# credentials present in env". Config reloads every turn, so the notice is -# one-time per platform per process. +# Platforms already warned "explicitly disabled in config.yaml, but credentials present in +# env": config reloads every turn, so the notice is one-time per platform per process. _EXPLICIT_DISABLE_WARNED: set = set() -# Env var(s) whose presence drives each platform's env-enable branch, named in -# the explicit-disable WARNING. +# Env var(s) whose presence drives each platform's env-enable branch, named in the explicit-disable WARNING. _ENV_ENABLE_CREDENTIALS: dict = { Platform.TELEGRAM: ("TELEGRAM_BOT_TOKEN",), Platform.DISCORD: ("DISCORD_BOT_TOKEN",), @@ -107,6 +104,14 @@ def _truthy_token(value: str) -> bool: return value.lower() in {"true", "1", "yes", "on"} +def _csv_extras(extra: Dict[str, Any], spec) -> None: + """``extra[key] = _csv_list(raw)`` for each ``(key, raw)`` whose list is non-empty.""" + for key, raw in spec: + items = _csv_list(raw) + if items: + extra[key] = items + + def _mention_patterns(value: str) -> Any: try: return json.loads(value) @@ -181,9 +186,9 @@ def _enable_from_env(config: GatewayConfig, platform: Platform) -> PlatformConfi def _enable_port_bound_from_env(config: GatewayConfig, platform: Platform) -> PlatformConfig: """Enable a port-binding platform unless config.yaml explicitly disabled it. - A multiplex secondary profile pins ``enabled: false`` to share the default profile's - listener yet inherits the process env; without this guard env presence would - force-enable it and trip MultiplexConfigError. POPs the marker (terminal branch). + A multiplex secondary profile pins ``enabled: false`` to share the default profile's listener + yet inherits the process env; without this guard env presence would force-enable it and trip + MultiplexConfigError. POPs the marker (terminal branch). """ platform_config = config.platforms.setdefault(platform, PlatformConfig()) explicit = platform_config.extra.pop("_enabled_explicit", False) @@ -196,12 +201,11 @@ def _enable_port_bound_from_env(config: GatewayConfig, platform: Platform) -> Pl class _Cred: """Credential-gated platform enable, applied as ``step(config)``. - ``creds``: env names that must ALL be truthy (inner tuple = ANY of). ``token``: env - stored as ``PlatformConfig.token`` even when yaml disables the adapter (sending skills - use it). ``fixed``: ``(extra_key, env[, default[, fn]])`` always written once enabled. - ``optional*``: ``_env_extras`` specs. ``warn_missing``: ``(env, msg)`` logged BEFORE - enabling when blank. ``then``: tail ``fn(config, platform_config)``. ``home``: - ``_env_home_channel`` env base applied only when the gate passed. + ``creds``: env names that must ALL be truthy (inner tuple = ANY of). ``token``: env stored as + ``PlatformConfig.token`` even when yaml disables the adapter (sending skills use it). ``fixed``: + ``(extra_key, env[, default[, fn]])`` always written once enabled. ``optional*``: ``_env_extras`` + specs. ``warn_missing``: ``(env, msg)`` logged BEFORE enabling when blank. ``then``: tail + ``fn(config, platform_config)``. ``home``: ``_env_home_channel`` env base applied only when the gate passed. """ platform: Platform creds: tuple @@ -220,10 +224,8 @@ class _Cred: if self.warn_missing and not getenv(self.warn_missing[0]): logger.warning(self.warn_missing[1]) platform_config = _enable_from_env(config, self.platform) - if self.token: - token = getenv(self.token) - if token: - platform_config.token = token + if self.token and (token := getenv(self.token)): + platform_config.token = token extra = platform_config.extra for key, env, *rest in self.fixed: default = rest[0] if rest else "" @@ -258,18 +260,17 @@ def _whatsapp(config: GatewayConfig) -> None: raw = getenv("WHATSAPP_ENABLED") enabled = is_truthy_value(raw) wa_cfg = config.platforms.get(Platform.WHATSAPP) - if wa_cfg is not None: - if raw.lower() in {"false", "0", "no"}: - wa_cfg.enabled = False - elif enabled: - wa_cfg.enabled = True + if wa_cfg is None: + if enabled: + config.platforms[Platform.WHATSAPP] = PlatformConfig(enabled=True) + elif raw.lower() in {"false", "0", "no"}: + wa_cfg.enabled = False elif enabled: - config.platforms[Platform.WHATSAPP] = PlatformConfig(enabled=True) + wa_cfg.enabled = True def _slack_home(config: GatewayConfig) -> None: - """SLACK_HOME_CHANNEL creates a disabled Slack entry if needed and keeps - user_id/scope_id provenance when the chat_id is unchanged.""" + """SLACK_HOME_CHANNEL creates a disabled Slack entry if needed; user_id/scope_id provenance survives an unchanged chat_id.""" slack_home = getenv("SLACK_HOME_CHANNEL") if not slack_home: return @@ -306,9 +307,7 @@ def _api_server(config: GatewayConfig) -> None: return extra = _enable_port_bound_from_env(config, Platform.API_SERVER).extra extra["key"] = key - origins = _csv_list(getenv("API_SERVER_CORS_ORIGINS")) - if origins: - extra["cors_origins"] = origins + _csv_extras(extra, (("cors_origins", getenv("API_SERVER_CORS_ORIGINS")),)) _env_extras(extra, (("port", "API_SERVER_PORT", _INT), ("host", "API_SERVER_HOST"), ("model_name", "API_SERVER_MODEL_NAME"))) @@ -333,24 +332,19 @@ def _msgraph_webhook(config: GatewayConfig) -> None: _env_extras(msgraph_cfg.extra, (("port", "MSGRAPH_WEBHOOK_PORT", _INT),)) if client_state: msgraph_cfg.extra["client_state"] = client_state - for key, raw in (("accepted_resources", resources), ("allowed_source_cidrs", allowed_cidrs)): - items = _csv_list(raw) - if items: - msgraph_cfg.extra[key] = items + _csv_extras(msgraph_cfg.extra, (("accepted_resources", resources), ("allowed_source_cidrs", allowed_cidrs))) def _qq_home(config: GatewayConfig, qq_config: PlatformConfig) -> None: qq_home = getenv("QQBOT_HOME_CHANNEL").strip() name_env = "QQBOT_HOME_CHANNEL_NAME" - if not qq_home: - # Back-compat: accept the pre-rename name and log a one-time warning. - qq_home = getenv("QQ_HOME_CHANNEL").strip() - if qq_home: - name_env = "QQ_HOME_CHANNEL_NAME" - logger.warning( - "QQ_HOME_CHANNEL is deprecated; rename to QQBOT_HOME_CHANNEL " - "in your .env for consistency with the platform key." - ) + if not qq_home and (qq_home := getenv("QQ_HOME_CHANNEL").strip()): + # Back-compat: accept the pre-rename name and warn. + name_env = "QQ_HOME_CHANNEL_NAME" + logger.warning( + "QQ_HOME_CHANNEL is deprecated; rename to QQBOT_HOME_CHANNEL " + "in your .env for consistency with the platform key." + ) if qq_home: qq_config.home_channel = HomeChannel( platform=Platform.QQBOT, chat_id=qq_home, @@ -380,10 +374,8 @@ def _plugin_probe_seed(entry) -> Optional[dict]: def _plugin_is_configured(entry, existing_extra: dict, seed: Optional[dict]) -> bool: - """``entry.is_connected`` on a transient ``enabled=True`` view seeded with env extras.""" + """``entry.is_connected`` on a transient ``enabled=True`` view seeded with env extras (never the real config).""" try: - # Seed extras so ``is_connected`` implementations reading ``config.extra`` see - # post-enable state; never mutate the real config for the probe. for k, v in (seed or {}).items(): if k != "home_channel": existing_extra.setdefault(k, v) @@ -413,10 +405,10 @@ def _enable_plugin_platform(config: GatewayConfig, entry) -> None: if existing_cfg is not None and not already_enabled and existing_extra.get("_enabled_explicit", False): return seed = _plugin_probe_seed(entry) - # Only consult is_connected for platforms not already enabled by YAML/env. + # is_connected gates only platforms not already enabled by YAML/env. if not already_enabled and entry.is_connected is not None and not _plugin_is_configured(entry, existing_extra, seed): return - # Verify dependencies LAST — only for platforms already enabled or past the credential gate. + # Dependencies LAST — only for platforms already enabled or past the credential gate. try: deps_ok = bool(entry.check_fn()) except Exception as e: @@ -426,8 +418,7 @@ def _enable_plugin_platform(config: GatewayConfig, entry) -> None: return platform_config = config.platforms.setdefault(platform, PlatformConfig()) platform_config.enabled = True - if seed: - # Commit the env-seeded extras (reuse the probe result; don't call env_enablement_fn twice). + if seed: # commit the probe's env-seeded extras (env_enablement_fn is never called twice) seed = dict(seed) home = seed.pop("home_channel", None) platform_config.extra.update(seed) @@ -441,10 +432,9 @@ def _enable_plugin_platform(config: GatewayConfig, entry) -> None: def _enable_plugin_platforms_from_env(config: GatewayConfig) -> None: """Registry-driven enable for plugin platforms (built-ins have rows in ``_ENV_STEPS``). - Enabled when credentials are configured (``is_connected`` MUST gate: ``check_fn`` - alone would enable unconfigured platforms that retry-connect forever) and deps are - present (``check_fn``) or installable later by ``create_adapter()`` — never here: - installing in this sweep pip-installed SDKs on every load and boot-looped the app. + Enabled when credentials are configured (``is_connected`` MUST gate: ``check_fn`` alone would + enable unconfigured platforms that retry-connect forever) and deps are present (``check_fn``) or + installable later by ``create_adapter()`` — never here: installing in this sweep boot-looped the app. """ try: from hermes_cli.plugins import discover_plugins @@ -457,16 +447,15 @@ def _enable_plugin_platforms_from_env(config: GatewayConfig) -> None: def _relay(config: GatewayConfig) -> None: - """Relay (connector-fronted, EXPERIMENTAL): enabled by GATEWAY_RELAY_URL or - gateway.relay_url; the URL is mirrored into extra["relay_url"] for the connected-checker. + """Relay (connector-fronted, EXPERIMENTAL): enabled by GATEWAY_RELAY_URL or gateway.relay_url; + the URL is mirrored into extra["relay_url"] for the connected-checker. - Relay-exclusive: the GATEWAY_RELAY_URL env stamp means the connector owns every - platform connection; a directly-connected adapter would be a second unmanaged ingress - (duplicate deliveries, split sessions, a socket disarming scale-to-zero). So the env - stamp disables all other messaging platforms — even ones explicitly enabled — except - non-messaging surfaces (local, api_server, webhook; same set as the scale-to-zero arm - gate). relay_url from YAML only keeps the additive behavior. Opt out with - GATEWAY_RELAY_ALLOW_DIRECT_PLATFORMS=true. + Relay-exclusive: the GATEWAY_RELAY_URL env stamp means the connector owns every platform + connection; a directly-connected adapter would be a second unmanaged ingress (duplicate + deliveries, split sessions, a socket disarming scale-to-zero). So the env stamp disables all + other messaging platforms — even explicitly enabled ones — except non-messaging surfaces (local, + api_server, webhook; same set as the scale-to-zero arm gate). relay_url from YAML only keeps the + additive behavior. Opt out with GATEWAY_RELAY_ALLOW_DIRECT_PLATFORMS=true. """ relay_url_env = getenv("GATEWAY_RELAY_URL").strip() existing_relay = config.platforms.get(Platform.RELAY) @@ -477,9 +466,9 @@ def _relay(config: GatewayConfig) -> None: if not relay_url_env or is_truthy_value(getenv("GATEWAY_RELAY_ALLOW_DIRECT_PLATFORMS")): return - non_messaging = {Platform.LOCAL, Platform.API_SERVER, Platform.WEBHOOK} + skip = {Platform.RELAY, Platform.LOCAL, Platform.API_SERVER, Platform.WEBHOOK} for platform, platform_config in config.platforms.items(): - if platform is Platform.RELAY or platform in non_messaging or not platform_config.enabled: + if platform in skip or not platform_config.enabled: continue if platform_config.extra.get("_enabled_explicit"): logger.warning( @@ -505,9 +494,9 @@ def _scrub_explicit_markers(config: GatewayConfig) -> None: platform_config.extra.pop("_enabled_explicit", None) -# Order is significant: a home channel only attaches to a platform that already exists -# (Telegram's reply mode may create the entry first; Discord reads home first). Relay -# disabling runs after the plugin pass; the marker scrub must be last. +# Order is significant: a home channel only attaches to a platform that already exists (Telegram's +# reply mode may create the entry first; Discord reads home first). Relay disabling runs after the +# plugin pass; the marker scrub must be last. _ENV_STEPS: tuple = ( _Cred(Platform.TELEGRAM, ("TELEGRAM_BOT_TOKEN",), token="TELEGRAM_BOT_TOKEN"), _ReplyMode(Platform.TELEGRAM, "TELEGRAM_REPLY_TO_MODE"), @@ -518,8 +507,7 @@ _ENV_STEPS: tuple = ( _ReplyMode(Platform.DISCORD, "DISCORD_REPLY_TO_MODE"), _whatsapp, _Home(Platform.WHATSAPP, "WHATSAPP_HOME_CHANNEL"), - # WhatsApp Cloud API (Meta Business Platform). Distinct from the Baileys bridge; - # both adapters can run in parallel against different phone numbers. + # WhatsApp Cloud API (Meta). Distinct from the Baileys bridge; both may run against different numbers. _Cred( Platform.WHATSAPP_CLOUD, ("WHATSAPP_CLOUD_PHONE_NUMBER_ID", "WHATSAPP_CLOUD_ACCESS_TOKEN"), fixed=(("phone_number_id", "WHATSAPP_CLOUD_PHONE_NUMBER_ID"), ("access_token", "WHATSAPP_CLOUD_ACCESS_TOKEN")), @@ -592,8 +580,7 @@ _ENV_STEPS: tuple = ( ("corp_id", "WECOM_CALLBACK_CORP_ID"), ("corp_secret", "WECOM_CALLBACK_CORP_SECRET"), ("agent_id", "WECOM_CALLBACK_AGENT_ID"), ("token", "WECOM_CALLBACK_TOKEN"), ("encoding_aes_key", "WECOM_CALLBACK_ENCODING_AES_KEY"), - # No host default: a falsy extra.host lets the adapter's dual-stack DEFAULT_HOST=None - # apply (binds IPv4 + IPv6; "0.0.0.0" was IPv4-only). + # No host default: falsy extra.host lets the adapter's dual-stack DEFAULT_HOST=None bind v4+v6. ("host", "WECOM_CALLBACK_HOST"), ("port", "WECOM_CALLBACK_PORT", "", _int_or(8645)), ), ), diff --git a/gateway/config_loader.py b/gateway/config_loader.py index 6ef62d927b..7f35552d5b 100644 --- a/gateway/config_loader.py +++ b/gateway/config_loader.py @@ -1,11 +1,12 @@ """config.yaml / gateway.json → ``GatewayConfig.from_dict`` schema (the ``load_gateway_config`` phases). -Precedence contract for top-level keys: key-presence at the TOP LEVEL of config.yaml wins; the -nested ``gateway.`` form (what ``hermes config set gateway.`` produces) is consulted only -when the top-level key is absent — not merely falsy/mistyped — so a present-but-empty top-level -value is never silently replaced by the nested one. Both overwrite whatever legacy gateway.json set. +Precedence for top-level keys: key-presence at the TOP LEVEL of config.yaml wins; the nested +``gateway.`` form (what ``hermes config set gateway.`` produces) is consulted only when the +top-level key is absent — not merely falsy/mistyped — so a present-but-empty top-level value is never +silently replaced by the nested one. Both overwrite whatever legacy gateway.json set. """ +import contextlib import json import logging import os @@ -35,11 +36,10 @@ def load_legacy_gateway_json(home: Path) -> Any: # --- top-level key bridging ---------------------------------------------------- # -# Settings meant to be top-level keys are also accepted nested under ``gateway:`` (what -# ``hermes config set gateway. ...`` naturally produces). This loader builds gw_data FLAT -# and never forwards the yaml ``gateway:`` section, so even keys GatewayConfig.from_dict can -# fall back on itself (loop_watchdog*, multiplex_profiles, ...) must be bridged here or they -# are silently ignored on the real gateway startup path. +# Top-level settings are also accepted nested under ``gateway:`` (what ``hermes config set +# gateway.`` produces). This loader builds gw_data FLAT and never forwards the yaml ``gateway:`` +# section, so even keys GatewayConfig.from_dict can fall back on itself (loop_watchdog*, +# multiplex_profiles, ...) must be bridged here or they are silently ignored on real startup. # # Fallback modes (how the nested ``gateway.`` form is consulted): # "presence": top-level key present → its value; else nested key present. @@ -101,9 +101,7 @@ def _bridge_lookup(yaml_cfg: dict, gateway_section: Any, gw_data: dict, key: str return True, gateway_section[key] return False, None if mode == "nested": - if nested and key in gateway_section: - return True, gateway_section[key] - return False, None + return (True, gateway_section[key]) if nested and key in gateway_section else (False, None) value = yaml_cfg.get(key) if mode == "none": if value is None and nested: @@ -129,9 +127,9 @@ def merge_platform_sections(yaml_cfg: dict, gateway_cfg: Any, gw_data: dict) -> """Merge every place a platform block may live into ``gw_data["platforms"]`` and return it. Order (later wins on shared keys, ``extra`` deep-merged so gateway.json defaults survive): - ``gateway.platforms.*`` → top-level ``platforms.*`` → ``gateway.`` subsections - (nested first so top-level config keeps precedence, matching the gateway.streaming fallback). - An ``enabled`` key in any block sets the ``_enabled_explicit`` marker consumed by the env pass. + ``gateway.platforms.*`` → top-level ``platforms.*`` → ``gateway.`` subsections (nested + first so top-level config keeps precedence, matching the gateway.streaming fallback). An + ``enabled`` key in any block sets the ``_enabled_explicit`` marker consumed by the env pass. Finally api_server's port/key/host/cors_origins/model_name are bridged into ``extra`` so ``gateway.api_server.port: 8642`` reaches the adapter (mirrors the env path). """ @@ -144,8 +142,7 @@ def merge_platform_sections(yaml_cfg: dict, gateway_cfg: Any, gw_data: dict) -> if not isinstance(plat_block, dict): continue existing = platforms_data.get(plat_name, {}) - if not isinstance(existing, dict): - existing = {} + existing = existing if isinstance(existing, dict) else {} merged_extra = {**existing.get("extra", {}), **plat_block.get("extra", {})} if "enabled" in plat_block: merged_extra["_enabled_explicit"] = True @@ -177,12 +174,9 @@ def _is_platform_name(key: Any) -> bool: def platform_section(yaml_cfg: dict, name: str, gateway_platforms: Any) -> tuple: - """Return ``(section, is_toplevel)`` for platform *name*. - - A top-level ``:`` block wins; otherwise fall back to the block under ``gateway.platforms`` - / ``platforms`` so shared-key bridging and adapter hooks still run when the user configured the - platform only under those nested paths. - """ + """``(section, is_toplevel)`` for platform *name*: a top-level ``:`` block wins; otherwise + the block under ``gateway.platforms`` / ``platforms`` so shared-key bridging and adapter hooks + still run for nested-only configs.""" section = yaml_cfg.get(name) toplevel = isinstance(section, dict) if not toplevel: @@ -225,9 +219,8 @@ _SHARED_KEYS: tuple = ( *_plain("gateway_restart_notification", "typing_indicator", "typing_status_text"), ) -# Top-level port/host/secret bridged into ``extra`` for adapters that read them from config.extra. -# Without this ``platforms.webhook.port: 8649`` silently falls back to the hardcoded DEFAULT_PORT, -# because PlatformConfig.from_dict only extracts ``extra`` from the ``extra:`` sub-key. +# Top-level port/host/secret bridged into ``extra`` for adapters that read them from config.extra +# (PlatformConfig.from_dict only reads the ``extra:`` sub-key, so ``platforms.webhook.port`` would be lost). _PORT_BRIDGE_KEYS: dict = { Platform.WEBHOOK: ("port", "host", "secret"), Platform.MSGRAPH_WEBHOOK: ("port", "host", "secret"), @@ -253,13 +246,9 @@ def _bridged_keys(plat: Platform, platform_cfg: dict, gw_data: dict) -> dict: def shared_loop_targets(registry) -> list: """Built-in platforms plus registered plugin platforms (so plugin authors get shared-key bridging).""" targets: list = list(Platform) - if registry is not None: - for entry in registry.plugin_entries(): - try: - plat = Platform(entry.name) - except (ValueError, KeyError): - continue - if plat not in targets: + for entry in registry.plugin_entries() if registry is not None else (): + with contextlib.suppress(ValueError, KeyError): + if (plat := Platform(entry.name)) not in targets: targets.append(plat) return targets @@ -269,11 +258,10 @@ def bridge_platform_shared_keys( ) -> None: """Copy shared keys (allow_from, require_mention, …) from each platform's YAML section into ``extra``. - ``enabled`` is only written from a TOP-LEVEL block: for nested-only configs - ``merge_platform_sections`` already merged it with the correct precedence, and re-applying it here - would overwrite that. An explicit top-level enable/disable sets ``_enabled_explicit`` so the env - pass honors ``enabled: false`` for migrated plugin platforms instead of re-enabling them on - token/SDK presence. + ``enabled`` is only written from a TOP-LEVEL block: for nested-only configs ``merge_platform_sections`` + already merged it with the correct precedence. An explicit top-level enable/disable sets + ``_enabled_explicit`` so the env pass honors ``enabled: false`` for migrated plugin platforms + instead of re-enabling them on token/SDK presence. """ for plat in targets: if plat == Platform.LOCAL: @@ -303,11 +291,8 @@ def bridge_platform_shared_keys( def apply_plugin_yaml_hooks(yaml_cfg: dict, gateway_platforms: Any, platforms_data: dict, registry) -> None: - """Plugin-owned YAML→env config bridges (``PlatformEntry.apply_yaml_config_fn``). - - Order: shared-key loop → this dispatch → core-only bridges (require_mention/signal) → - ``_apply_env_overrides()`` after ``GatewayConfig.from_dict``. - """ + """Plugin-owned YAML→env config bridges (``PlatformEntry.apply_yaml_config_fn``). Order: shared-key + loop → this dispatch → core-only bridges (require_mention/signal) → ``_apply_env_overrides()``.""" if registry is None: return for entry in registry.all_entries(): @@ -326,12 +311,11 @@ def apply_plugin_yaml_hooks(yaml_cfg: dict, gateway_platforms: Any, platforms_da def bridge_core_env_settings(yaml_cfg: dict, platforms_data: dict) -> None: - """The two YAML→env bridges that stay in core (the per-platform ones live in plugin hooks). + """The two YAML→env bridges that stay in core (per-platform ones live in plugin hooks). Top-level ``require_mention`` → Telegram when the ``telegram:`` section has none: users write it - alongside ``group_sessions_per_user`` expecting it to work. It keys off the TOP-LEVEL key, so - the telegram plugin's hook (which only runs when a telegram block exists) can't cover it. - Signal ``require_mention`` → ``SIGNAL_REQUIRE_MENTION`` (env wins when already set). + alongside ``group_sessions_per_user`` expecting it to work, and the telegram plugin's hook only + runs when a telegram block exists. Signal ``require_mention`` → ``SIGNAL_REQUIRE_MENTION`` (env wins). """ tl_require_mention = yaml_cfg.get("require_mention") if tl_require_mention is not None and "require_mention" not in (yaml_cfg.get("telegram") or {}): @@ -355,9 +339,8 @@ def load_yaml_layer(home: Path, gw_data: dict) -> None: with open(config_yaml_path, encoding="utf-8") as f: yaml_cfg = yaml.safe_load(f) or {} - # Managed scope: overlay administrator-pinned values so the gateway honors them too. This - # loader builds its own dict instead of going through hermes_cli.config.load_config, so - # without this a managed session_reset / quick_commands / stt would be ignored. Fail-open. + # Managed scope: overlay administrator-pinned values (this loader bypasses + # hermes_cli.config.load_config, so a managed session_reset / quick_commands / stt would otherwise be ignored). from hermes_cli import managed_scope yaml_cfg = managed_scope.apply_managed_overlay(yaml_cfg) diff --git a/gateway/display_config.py b/gateway/display_config.py index 85ad1fb3e5..0003ce7ea3 100644 --- a/gateway/display_config.py +++ b/gateway/display_config.py @@ -16,80 +16,57 @@ _GLOBAL_DEFAULTS: dict[str, Any] = { "tool_progress": "all", "tool_progress_grouping": "accumulate", # "accumulate" = edit one bubble; "separate" = one msg per tool "show_reasoning": False, - # "code" (💭 **Reasoning:** + fenced block), "blockquote" ("> "), "subtext" ("-# " Discord). - "reasoning_style": "code", + "reasoning_style": "code", # "code" (💭 **Reasoning:** + fence), "blockquote" ("> "), "subtext" ("-# " Discord) "tool_preview_length": 0, "streaming": None, # None = follow top-level streaming config # Gateway-only assistant/status chatter; mobile platforms opt down to final-answer-first. "interim_assistant_messages": True, "long_running_notifications": True, "busy_ack_detail": True, - # busy_input_mode=steer echo ("Steered into current run"); the text still lands in the run. - "busy_steer_ack_enabled": True, - # Delete tool-progress / "⏳ Working" bubbles after a SUCCESSFUL final response on - # platforms that support deletion (Telegram); failed runs keep them as breadcrumbs. + "busy_steer_ack_enabled": True, # busy_input_mode=steer echo; the text still lands in the run + # Delete tool-progress / "⏳ Working" bubbles after a SUCCESSFUL final response where deletion is + # supported (Telegram); failed runs keep them as breadcrumbs. "cleanup_progress": False, - # Working-state text on text-rendering indicators (Slack assistant status): - # "full"/true = verb + argument preview, "verb" = verb only (keeps paths out of - # shared channels), "off"/false = static text. + # Working-state text on text-rendering indicators (Slack assistant status): "full"/true = verb + + # argument preview, "verb" = verb only (keeps paths out of shared channels), "off"/false = static. "live_status": "full", } # Tiers: HIGH = editing, personal/team use; MEDIUM = editing but customer-facing; # LOW = no edit support (progress messages are permanent); MINIMAL = batch delivery. _TIER_HIGH = { - "tool_progress": "all", - "show_reasoning": False, - "tool_preview_length": 40, + "tool_progress": "all", "show_reasoning": False, "tool_preview_length": 40, "streaming": None, # follow global - "interim_assistant_messages": True, - "long_running_notifications": True, - "busy_ack_detail": True, + "interim_assistant_messages": True, "long_running_notifications": True, "busy_ack_detail": True, } _TIER_MEDIUM = {**_TIER_HIGH, "tool_progress": "new"} _TIER_LOW = { - **_TIER_HIGH, - "tool_progress": "off", - "streaming": False, - "interim_assistant_messages": False, - "long_running_notifications": False, - "busy_ack_detail": False, + **_TIER_HIGH, "tool_progress": "off", "streaming": False, + "interim_assistant_messages": False, "long_running_notifications": False, "busy_ack_detail": False, } _TIER_MINIMAL = {**_TIER_LOW, "tool_preview_length": 0} _PLATFORM_DEFAULTS: dict[str, dict[str, Any]] = { - # Mobile inbox: quiet tool_progress / busy-ack, but keep interim commentary and - # heartbeats so it doesn't look like "typing..." for 30 minutes. + # Mobile inbox: quiet tool_progress / busy-ack, but keep interim commentary and heartbeats so it + # doesn't look like "typing..." for 30 minutes. "telegram": {**_TIER_HIGH, "tool_progress": "off", "busy_ack_detail": False}, - # Discord's "-# " subtext reads as metadata, so reasoning defaults to it. - "discord": {**_TIER_HIGH, "reasoning_style": "subtext"}, - + "discord": {**_TIER_HIGH, "reasoning_style": "subtext"}, # "-# " subtext reads as metadata # Slack: Bolt posts cannot be edited like CLI; "new"/"all" spam permanent lines. - "slack": { - **_TIER_MEDIUM, - "tool_progress": "off", - "long_running_notifications": False, - "busy_ack_detail": False, - }, + "slack": {**_TIER_MEDIUM, "tool_progress": "off", "long_running_notifications": False, "busy_ack_detail": False}, "mattermost": _TIER_MEDIUM, "matrix": _TIER_MEDIUM, "feishu": _TIER_MEDIUM, - # Buzz (Nostr) can edit in place but channels are shared community spaces. - "buzz": _TIER_MEDIUM, - + "buzz": _TIER_MEDIUM, # Nostr: edits in place but channels are shared community spaces "signal": _TIER_LOW, "whatsapp": _TIER_MEDIUM, # Baileys bridge supports /edit "whatsapp_cloud": _TIER_LOW, # adapter lacks edit_message; promote once it lands - # Permanent-message iMessage inboxes (no edit). - "photon": _TIER_LOW, + "photon": _TIER_LOW, # permanent-message iMessage inboxes (no edit) "bluebubbles": _TIER_LOW, "weixin": _TIER_LOW, - # Non-editable, but its native "stream" msgtype gives a typing animation + - # cumulative updates instead of a one-shot markdown drop. + # Non-editable, but its native "stream" msgtype gives a typing animation + cumulative updates. "wecom": {**_TIER_LOW, "streaming": True}, "wecom_callback": _TIER_LOW, "dingtalk": _TIER_LOW, - "email": _TIER_MINIMAL, "sms": _TIER_MINIMAL, "webhook": _TIER_MINIMAL, @@ -101,36 +78,22 @@ _PLATFORM_DEFAULTS: dict[str, dict[str, Any]] = { OVERRIDEABLE_KEYS = frozenset(_GLOBAL_DEFAULTS.keys()) -def resolve_display_setting( - user_config: dict, - platform_key: str, - setting: str, - fallback: Any = None, -) -> Any: - """Resolve a display setting with per-platform override support. +def resolve_display_setting(user_config: dict, platform_key: str, setting: str, fallback: Any = None) -> Any: + """Resolve a display setting with per-platform override support (see module docstring for order). - ``platform_key`` is the platform config key (``"telegram"``, ``"slack"``; - see ``_platform_config_key`` in gateway/run.py). Returns *fallback* when - nothing is configured. + ``platform_key`` is the platform config key (``"telegram"``; see ``_platform_config_key`` in + gateway/run.py). Returns *fallback* when nothing is configured. """ display_cfg = user_config.get("display") or {} - - # 1. Explicit per-platform override plat_overrides = (display_cfg.get("platforms") or {}).get(platform_key) if isinstance(plat_overrides, dict) and plat_overrides.get(setting) is not None: return _normalise(setting, plat_overrides[setting]) - - # 1b. Backward compat: display.tool_progress_overrides. - if setting == "tool_progress": + if setting == "tool_progress": # legacy display.tool_progress_overrides. legacy = display_cfg.get("tool_progress_overrides") if isinstance(legacy, dict) and legacy.get(platform_key) is not None: return _normalise(setting, legacy[platform_key]) - - # 2. Global user setting (display.streaming is CLI-only, see module docstring). - if setting != "streaming" and display_cfg.get(setting) is not None: + if setting != "streaming" and display_cfg.get(setting) is not None: # display.streaming is CLI-only return _normalise(setting, display_cfg[setting]) - - # 3. Built-in platform default, 4. built-in global default val = _PLATFORM_DEFAULTS.get(platform_key, {}).get(setting) if val is None: val = _GLOBAL_DEFAULTS.get(setting) @@ -143,17 +106,18 @@ _TRUTHY = {"true", "1", "yes", "on"} _FALSY = {"false", "0", "no"} -def _norm_tool_progress(value: Any) -> str: - if value is False: - return "off" - if value is True: - return "all" - val = str(value).strip().lower() - if val in _FALSY: - return "off" - if val in _TRUTHY: - return "all" - return val if val in {"off", "new", "all", "verbose", "log"} else "all" +def _norm_tristate(on: str, off: str, choices: set, extra_truthy: set = frozenset()): + """Normaliser for bool-or-keyword settings: bools/truthy tokens → *on*, falsy → *off*, else a known choice or *on*.""" + def norm(value: Any) -> str: + if isinstance(value, bool): + return on if value else off + val = str(value).strip().lower() + if val in _FALSY: + return off + if val in _TRUTHY | extra_truthy: + return on + return val if val in choices else on + return norm def _norm_bool(value: Any) -> bool: @@ -174,20 +138,6 @@ def _norm_cleanup_progress(value: Any) -> bool: return bool(value) -def _norm_live_status(value: Any) -> str: - """Tri-state: "full" (verb + preview), "verb" (verb only), "off".""" - if value is True: - return "full" - if value is False: - return "off" - val = str(value).strip().lower() - if val in _TRUTHY | {"all"}: - return "full" - if val in _FALSY: - return "off" - return val if val in {"full", "verb", "off"} else "full" - - def _norm_choice(choices: tuple[str, ...]) -> Any: def norm(value: Any) -> str: val = str(value).lower() @@ -204,7 +154,7 @@ def _norm_int(value: Any) -> int: _NORMALISERS: dict[str, Any] = { - "tool_progress": _norm_tool_progress, + "tool_progress": _norm_tristate("all", "off", {"off", "new", "all", "verbose", "log"}), "show_reasoning": _norm_bool, "streaming": _norm_bool, "interim_assistant_messages": _norm_bool, @@ -213,7 +163,7 @@ _NORMALISERS: dict[str, Any] = { "busy_steer_ack_enabled": _norm_bool, "thinking_progress": _norm_bool, "cleanup_progress": _norm_cleanup_progress, - "live_status": _norm_live_status, + "live_status": _norm_tristate("full", "off", {"full", "verb", "off"}, extra_truthy={"all"}), "tool_progress_grouping": _norm_choice(("accumulate", "separate")), "reasoning_style": _norm_choice(("code", "blockquote", "subtext")), "tool_preview_length": _norm_int, diff --git a/gateway/profile_routing.py b/gateway/profile_routing.py index b5d248af64..da403e2dad 100644 --- a/gateway/profile_routing.py +++ b/gateway/profile_routing.py @@ -8,11 +8,10 @@ direct parent, so a channel route also matches any thread/post under it. from __future__ import annotations +import logging from dataclasses import dataclass from typing import Any, Dict, List, Optional -import logging - logger = logging.getLogger(__name__) # Baileys and Cloud share phone/JID/LID identity rules; other platforms compare exactly. @@ -22,18 +21,15 @@ _WHATSAPP_NON_USER_SUFFIXES = ("@g.us", "@broadcast", "@newsletter") def _is_whatsapp_non_user_chat(chat_id: Optional[str]) -> bool: """True for group / broadcast / newsletter JIDs — not a sender identity.""" - if not chat_id: - return False - cid = str(chat_id).strip().lower() - return any(cid.endswith(suffix) for suffix in _WHATSAPP_NON_USER_SUFFIXES) + return bool(chat_id) and str(chat_id).strip().lower().endswith(_WHATSAPP_NON_USER_SUFFIXES) def _whatsapp_user_chat_ids_match(platform: str, left: Optional[str], right: Optional[str]) -> bool: """True when two WhatsApp *user* chat_ids refer to the same person. - Uses ``expand_whatsapp_aliases`` (same helper as session keys and adapter - allowlists) so a bare number, JID and LID collapse to one identity. - Group/broadcast JIDs are chats, not senders; non-WhatsApp platforms → False. + Uses ``expand_whatsapp_aliases`` (same helper as session keys and adapter allowlists) so a bare + number, JID and LID collapse to one identity. Group/broadcast JIDs are chats, not senders; + non-WhatsApp platforms → False. """ if ( (platform or "").strip().lower() not in _WHATSAPP_IDENTITY_PLATFORMS @@ -69,18 +65,13 @@ class ProfileRoute: return 2 * bool(self.guild_id) + 4 * bool(self.chat_id) + 8 * bool(self.thread_id) def matches( - self, - platform: str, - guild_id: Optional[str] = None, - chat_id: Optional[str] = None, - thread_id: Optional[str] = None, - parent_chat_id: Optional[str] = None, + self, platform: str, guild_id: Optional[str] = None, chat_id: Optional[str] = None, + thread_id: Optional[str] = None, parent_chat_id: Optional[str] = None, ) -> bool: """True if every discriminator the route declares holds (AND). - ``chat_id`` matches the channel directly or as the parent of a thread/forum - post; WhatsApp ``chat_id`` also matches across number/JID/LID after the - exact check (groups/broadcasts stay exact-only). + ``chat_id`` matches the channel directly or as the parent of a thread/forum post; WhatsApp + ``chat_id`` also matches across number/JID/LID after the exact check (groups/broadcasts stay exact-only). """ if not self.enabled or self.platform != platform: return False @@ -99,9 +90,9 @@ class ProfileRoute: def _coerce_route_id(value: Any) -> Optional[str]: """Normalize a route discriminator to str for strict equality matching. - PyYAML loads unquoted numeric IDs as ``int`` while ``SessionSource`` fields are - ``str``. Only ``int`` (not ``bool``) is coerced; floats stringify to something - (``"123.0"``) that can never match, so they get a load-time warning instead. + PyYAML loads unquoted numeric IDs as ``int`` while ``SessionSource`` fields are ``str``. Only + ``int`` (not ``bool``) is coerced; floats stringify to something (``"123.0"``) that can never + match, so they get a load-time warning instead. """ if value is None or isinstance(value, str): return value @@ -138,32 +129,24 @@ 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 - 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")), - enabled=entry.get("enabled", True), - ) - ) + 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")), + enabled=entry.get("enabled", True), + )) routes.sort(key=lambda r: r.specificity, reverse=True) logger.debug("Loaded %d profile routes (most-specific-first)", len(routes)) return routes def match_profile_route( - routes: List[ProfileRoute], - platform: str, - guild_id: Optional[str] = None, - chat_id: Optional[str] = None, - thread_id: Optional[str] = None, - parent_chat_id: Optional[str] = None, + routes: List[ProfileRoute], platform: str, guild_id: Optional[str] = None, chat_id: Optional[str] = None, + thread_id: Optional[str] = None, parent_chat_id: Optional[str] = None, ) -> Optional[ProfileRoute]: """Return the first (most specific) matching route, or None.""" - for route in routes: - if route.matches(platform, guild_id=guild_id, chat_id=chat_id, thread_id=thread_id, parent_chat_id=parent_chat_id): - return route - return None + return next( + (r for r in routes if r.matches(platform, guild_id=guild_id, chat_id=chat_id, thread_id=thread_id, parent_chat_id=parent_chat_id)), + None, + )