refactor(gateway): wave-2 config/authz compaction — shared bool-token parser, _extra_choice, authz phase helpers (_chat_scoped_grant, _legacy_telegram_chat_grant), unified adapter flag/policy readers, tristate display normaliser, hand-compacted field docs

This commit is contained in:
Teknium
2026-09-02 21:49:33 -07:00
parent af2ae27695
commit 48a2007092
6 changed files with 387 additions and 599 deletions
+141 -194
View File
@@ -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 -> ``<PLATFORM>_ALLOWED_USERS`` / ``<PLATFORM>_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 -> ``<PLATFORM>_ALLOWED_USERS`` / ``<PLATFORM>_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.<id>.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.<id>.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":
+88 -150
View File
@@ -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)
+63 -76
View File
@@ -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)),
),
),
+33 -50
View File
@@ -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.<key>`` form (what ``hermes config set gateway.<key>`` 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.<key>`` form (what ``hermes config set gateway.<key>`` 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.<key> ...`` 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.<key>`` 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.<key>`` 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.<platform>`` 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.<platform>`` 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 ``<name>:`` 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 ``<name>:`` 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)
+37 -87
View File
@@ -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.<platform>
if setting == "tool_progress":
if setting == "tool_progress": # legacy display.tool_progress_overrides.<platform>
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,
+25 -42
View File
@@ -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,
)