diff --git a/gateway/config.py b/gateway/config.py index ef78ed0e15..d316cf5509 100644 --- a/gateway/config.py +++ b/gateway/config.py @@ -237,15 +237,18 @@ 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 binds: a SECONDARY +# profile's port-binder is built in shared-listener mode and served at /p// on the +# default's listener (gateway/platforms/shared_ingress.py); api_server/webhook are mirrored there. 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-binders whose /p// surface is a MIRROR served by the default's own adapter; a secondary +# never gets an instance of these (api_server: /p//v1/..., webhook: profile-bound routes). +SHARED_LISTENER_MIRROR_PLATFORMS = frozenset({"api_server", "webhook"}) def platform_binds_port(platform_value: str, extra: Optional[dict] = None) -> bool: diff --git a/gateway/config_env.py b/gateway/config_env.py index e3d384bbf3..a03843d774 100644 --- a/gateway/config_env.py +++ b/gateway/config_env.py @@ -21,7 +21,7 @@ from gateway.config import ( PlatformConfig, _getenv_str, _has_usable_api_server_key, - platform_binds_port, + SHARED_LISTENER_MIRROR_PLATFORMS, ) from utils import is_truthy_value @@ -183,11 +183,10 @@ def _enable_from_env( ) -> PlatformConfig: """Enable *platform* on env credentials 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. By default the ``_enabled_explicit`` marker is READ (the plugin-enable - and relay passes still need it) and the disable is warned once; port-binding platforms POP it - (terminal branch) and stay silent. + A multiplex secondary profile may pin ``enabled: false`` yet inherit the process env; without + this guard env presence would force-enable it. By default the ``_enabled_explicit`` marker is + READ (the plugin-enable and relay passes still need it) and the disable is warned once; + api_server/webhook POP it (terminal branch) and stay silent. """ platform_config = config.platforms.setdefault(platform, PlatformConfig()) extra = platform_config.extra @@ -195,12 +194,12 @@ def _enable_from_env( if platform_config.enabled: return platform_config if not explicit and not ( - platform_binds_port(platform.value, extra) and _loading_secondary_under_multiplexer() + platform.value in SHARED_LISTENER_MIRROR_PLATFORMS and _loading_secondary_under_multiplexer() ): - # A secondary's port-binding credential (the docs require API_SERVER_KEY in its .env for - # /p// auth) must not turn into listener intent: the default profile owns the one - # shared listener and ``_load_secondary_profile_config`` skips the WHOLE profile for it (#100397). - # The credential itself still lands in ``extra`` for the shared adapter to authenticate with. + # A secondary's API_SERVER_KEY / WEBHOOK_ENABLED (the docs require the key in its .env for + # /p// auth) must not turn into listener intent: the default profile's listener already + # mirrors those two at /p// (#100397). The credential still lands in ``extra`` for it. + # Every other inbound-port platform IS enabled for a secondary: it runs in shared-listener mode. platform_config.enabled = True elif warn: _warn_explicit_disable_beats_env(platform) diff --git a/gateway/platforms/api_server.py b/gateway/platforms/api_server.py index b0995de452..3e4aaecc89 100644 --- a/gateway/platforms/api_server.py +++ b/gateway/platforms/api_server.py @@ -1487,6 +1487,14 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter): from hermes_cli.profiles import get_profile_dir return _profile_runtime_scope(get_profile_dir(profile)) + async def _handle_profile_ingress(self, request: "web.Request") -> "web.StreamResponse": + """``/p//`` → the served profile's shared-listener adapter (already scoped by the + prefix middleware); a profile with no adapter for the path is a 404, never the default's.""" + from gateway.platforms.shared_ingress import dispatch_profile_ingress + return await dispatch_profile_ingress( + self.gateway_runner, _api_request_profile.get(), request.match_info.get("tail", ""), request, + scoped=True) + def _make_profile_prefix_middleware(self): """Reject unknown /p// prefixes and scope the request home.""" @@ -3896,6 +3904,9 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter): for method, path, handler in self._http_route_table(): self._app.router.add_route(method, path, handler) self._app.router.add_route(method, f"/p/{{profile}}{path}", handler) + # Registered LAST so every native mirror above wins: anything else under /p// is a + # secondary profile's inbound-port platform (Twilio, LINE, Teams, ...) served on this listener. + self._app.router.add_route("*", "/p/{profile}/{tail:.*}", self._handle_profile_ingress) # After native routes: Relay bootstrap shims feature-detect on this key and must # no-op rather than shadow the native session-control handlers. self._app["api_server_adapter"] = self diff --git a/gateway/platforms/base.py b/gateway/platforms/base.py index 36172b831d..1a06ceae1b 100644 --- a/gateway/platforms/base.py +++ b/gateway/platforms/base.py @@ -1884,6 +1884,9 @@ class BasePlatformAdapter(ABC): self._busy_session_handler: Optional[Callable[[MessageEvent, str], Awaitable[bool]]] = None # Owning multiplex profile (None on primary); see _session_key_profile. self._owner_profile: Optional[str] = None + # Set by the runner on a secondary's port-binding adapter: serve via the default profile's + # shared listener (/p//...) instead of binding a port (gateway/platforms/shared_ingress.py). + self._shared_listener_profile: Optional[str] = None # Registered by GatewayRunner (see set_authorization_check). self._authorization_check: Optional[Callable[[str, Optional[str], Optional[str]], bool]] = None # Auto-TTS on voice input: ``voice.auto_tts`` default plus per-chat /voice on|tts / off. diff --git a/gateway/platforms/shared_ingress.py b/gateway/platforms/shared_ingress.py new file mode 100644 index 0000000000..89b70925b5 --- /dev/null +++ b/gateway/platforms/shared_ingress.py @@ -0,0 +1,141 @@ +"""Shared-listener ingress for inbound-port platforms under ``gateway.multiplex_profiles``. + +The default profile owns the ONE HTTP listener (api_server, and the webhook adapter's port). A +secondary profile's port-binding adapter (Twilio SMS, LINE, Teams, BlueBubbles, Microsoft Graph, +WhatsApp Cloud, WeCom callback, Feishu webhook mode) therefore cannot bind its own port; instead the +runner constructs it in *shared-listener mode*: ``bind_listener`` publishes the adapter's fully wired +``web.Application`` instead of starting a ``TCPSite``, and the default listener forwards +``/p//`` to it through ``dispatch_profile_ingress``. + +Invariants: the forwarded request runs under the NAMED profile's runtime scope and is verified by +that profile's adapter with that profile's secret; the un-prefixed path keeps serving the default +profile untouched; a profile with no adapter for the path gets a 404, never the default's adapter. +""" +from __future__ import annotations + +import logging +from typing import TYPE_CHECKING, Any, Optional + +if TYPE_CHECKING: # aiohttp is an optional dependency of every adapter using this module + from aiohttp import web + +logger = logging.getLogger(__name__) + +_WILDCARD_HOSTS = frozenset({"", "0.0.0.0", "::", "*"}) + + +def shared_ingress_profile(adapter: Any) -> Optional[str]: + """Profile name when *adapter* was constructed in shared-listener mode, else None.""" + return getattr(adapter, "_shared_listener_profile", None) or None + + +def shared_listener_base(runner: Any) -> Optional[str]: + """``http://host:port`` of the default profile's live listener (api_server first, then webhook).""" + from gateway.config import Platform + adapters = getattr(runner, "adapters", None) or {} + for platform in (Platform.API_SERVER, Platform.WEBHOOK): + adapter = adapters.get(platform) + if adapter is None: + continue + host = getattr(adapter, "_host", None) + host = "127.0.0.1" if host is None or str(host).strip() in _WILDCARD_HOSTS else str(host) + if ":" in host and not host.startswith("["): + host = f"[{host}]" + return f"http://{host}:{getattr(adapter, '_port', 0)}" + return None + + +async def bind_listener( + adapter: Any, app: "web.Application", host: Optional[str], port: int, ingress_path: str, *, + reuse_address: Optional[bool] = None, access_log: Any = ..., +) -> Optional["web.AppRunner"]: + """Start *app* on ``host:port`` and return its ``AppRunner`` — or, in shared-listener mode, + publish *app* for ``/p//`` forwarding and return None (nothing bound). ``ingress_path`` + is the adapter's primary callback path, used for the log line and runtime status.""" + from aiohttp import web + profile = shared_ingress_profile(adapter) + if profile: + publish_shared_ingress(adapter, app, ingress_path) + return None + runner = web.AppRunner(app) if access_log is ... else web.AppRunner(app, access_log=access_log) + await runner.setup() + site = web.TCPSite(runner, host, port, reuse_address=reuse_address) + try: + await site.start() + except BaseException: + await runner.cleanup() + raise + return runner + + +def publish_shared_ingress(adapter: Any, app: "web.Application", ingress_path: str) -> None: + """Freeze *app* and expose it to the default listener; records the ``/p//`` URL.""" + profile = shared_ingress_profile(adapter) + app.freeze() + adapter._shared_ingress_app = app + base = shared_listener_base(getattr(adapter, "gateway_runner", None)) + prefix = f"/p/{profile}" + adapter._shared_ingress_base = f"{base}{prefix}" if base else prefix + adapter._shared_ingress_url = f"{adapter._shared_ingress_base}{ingress_path}" + platform = getattr(getattr(adapter, "platform", None), "value", "adapter") + if base: + logger.info( + "[%s] profile '%s' is served on the default profile's shared listener: %s " + "(point the vendor's callback URL at this path behind your public host)", + platform, profile, adapter._shared_ingress_url, + ) + else: + logger.warning( + "[%s] profile '%s' is in shared-listener mode but the default profile has no live api_server " + "or webhook listener yet; it will be reachable at %s%s once one is up", + platform, profile, prefix, ingress_path, + ) + write = getattr(adapter, "_write_runtime_status_safe", None) + if callable(write): + write("shared_ingress", ingress_url=adapter._shared_ingress_url) + + +def shared_ingress_apps(runner: Any, profile: Optional[str]) -> list[tuple[Any, "web.Application"]]: + """``(adapter, app)`` for every shared-listener adapter of a NAMED served profile. ``default`` and + unknown profiles yield nothing: the default's port-binders own their own ports, and a profile + without a live adapter must never fall back to another profile's.""" + if not profile or profile == "default": + return [] + adapters = (getattr(runner, "_profile_adapters", None) or {}).get(profile) or {} + return [ + (adapter, app) for adapter in adapters.values() + if (app := getattr(adapter, "_shared_ingress_app", None)) is not None + ] + + +async def dispatch_profile_ingress( + runner: Any, profile: Optional[str], tail: str, request: "web.Request", *, scoped: bool = False, +) -> "web.StreamResponse": + """Forward ``/p//`` to the served profile's adapter app that routes ``/``, + under that profile's runtime scope (``scoped=True`` when the caller already entered it). 404 when + no adapter of *profile* serves the path.""" + from aiohttp import web + from aiohttp.web_urldispatcher import MatchInfoError + candidates = shared_ingress_apps(runner, profile) + if not profile or not candidates: + raise web.HTTPNotFound(text="Unknown or unconfigured profile") + rel_url = request.rel_url.with_path("/" + tail.lstrip("/"), keep_query=True, keep_fragment=True) + forwarded = request.clone(rel_url=rel_url) + chosen = None + for _adapter, app in candidates: + match = await app.router.resolve(forwarded) + # A 404 means "not my path": try the profile's next adapter; 405 and real matches belong here. + if isinstance(match, MatchInfoError) and match.http_exception.status == 404: + continue + chosen = app + break + if chosen is None: + raise web.HTTPNotFound(text="No adapter serves this path for the profile") + # Each adapter chose its own body cap when it built its Application; honour it for the forwarded read. + forwarded = request.clone(rel_url=rel_url, client_max_size=chosen._client_max_size) + if scoped: + return await chosen._handle(forwarded) + from gateway.run import _profile_runtime_scope + from hermes_cli.profiles import get_profile_dir + with _profile_runtime_scope(get_profile_dir(profile)): + return await chosen._handle(forwarded) diff --git a/gateway/platforms/webhook.py b/gateway/platforms/webhook.py index 955c76867a..d70cde0e6b 100644 --- a/gateway/platforms/webhook.py +++ b/gateway/platforms/webhook.py @@ -217,6 +217,9 @@ class WebhookAdapter(BasePlatformAdapter): app.router.add_post("/webhooks/{route_name}", self._handle_webhook) # /p// routes the event to that profile (honored only under gateway.multiplex_profiles). app.router.add_post("/p/{profile}/webhooks/{route_name}", self._handle_webhook) + # Without an api_server listener this port is the shared listener: forward a secondary's + # inbound-port platforms (Twilio, LINE, Teams, ...) registered in shared-listener mode. + app.router.add_route("*", "/p/{profile}/{tail:.*}", self._handle_profile_ingress) self._runner = web.AppRunner(app) await self._runner.setup() # SO_REUSEADDR: on macOS (BSD) two wildcard/specific sockets can silently split traffic while @@ -374,6 +377,14 @@ class WebhookAdapter(BasePlatformAdapter): except Exception as e: logger.error("[webhook] Failed to reload dynamic routes: %s", e) + async def _handle_profile_ingress(self, request: "web.Request") -> "web.StreamResponse": + profile = self._resolve_request_profile(request) + if profile is _PROFILE_REJECTED or profile is None: + return _json_error("Unknown or unconfigured profile", 404) + from gateway.platforms.shared_ingress import dispatch_profile_ingress + return await dispatch_profile_ingress( + self.gateway_runner, profile, request.match_info.get("tail", ""), request) + def _resolve_request_profile(self, request: "web.Request"): """Resolve + validate the /p// URL prefix: None (no prefix, or multiplexing off and the prefix names this gateway's own profile), the profile name (served under multiplexing), or diff --git a/gateway/run.py b/gateway/run.py index 61d795d64f..52d6a41bad 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -1650,11 +1650,6 @@ class MultiplexConfigError(RuntimeError): startup guard instead of being treated as retryable adapter-connect noise.""" -class SecondaryPortBindingConfigError(MultiplexConfigError): - """A secondary profile enabled a port-binding platform: the default profile owns the single shared - listener (/p//), so this is always a misconfiguration and is skipped, not fatal.""" - - class HygieneTurnHoldExceeded(Exception): """Hygiene-compression turn-hold budget elapsed mid-stream. Availability boundary, not a failure: must NOT take the idle-timeout path (AGENT_COMPRESSION_TIMEOUT, "no output", failure cooldown).""" diff --git a/gateway/run_adapters.py b/gateway/run_adapters.py index 2897c97e1a..1ce511b478 100644 --- a/gateway/run_adapters.py +++ b/gateway/run_adapters.py @@ -19,7 +19,7 @@ import weakref as _weakref from agent.async_utils import consume_detached_task_result from contextvars import Context from datetime import datetime, timedelta, timezone -from gateway.config import Platform, platform_binds_port as _platform_binds_port +from gateway.config import SHARED_LISTENER_MIRROR_PLATFORMS, Platform, platform_binds_port as _platform_binds_port from gateway.platforms.base import BasePlatformAdapter from gateway.restart import is_global_startup_conflict from gateway.run_shutdown import _log_suppressed @@ -826,9 +826,7 @@ class GatewayAdapterLifecycleMixin: """Bring up adapters for every non-active profile (multiplex only); returns connected count. Each profile connects under its own HERMES_HOME + secret scope; credential/listener collisions are refused here — the only point seeing every profile's credentials together.""" - from gateway.run import ( - MultiplexConfigError, SecondaryPortBindingConfigError, _multiplex_profile_homes - ) + from gateway.run import MultiplexConfigError, _multiplex_profile_homes if not self._multiplex_on(): return 0 try: @@ -844,10 +842,6 @@ class GatewayAdapterLifecycleMixin: continue # handled by the primary startup loop try: connected += await self._start_one_profile_adapters(profile_name, profile_home, claimed) - except SecondaryPortBindingConfigError as e: - logger.warning( - "Skipping secondary profile '%s' due to port-binding config error: %s", profile_name, e, - ) except MultiplexConfigError: raise except Exception as e: @@ -888,10 +882,11 @@ class GatewayAdapterLifecycleMixin: async def _load_secondary_profile_config(self, profile_name: str, profile_home: "Path"): """Hydrate + enter ``profile_home``'s scope once; return its gateway config. Raises - ``MultiplexConfigError`` (open dm/group policy) or ``SecondaryPortBindingConfigError`` (the - default profile owns the single shared HTTP listener).""" + ``MultiplexConfigError`` (open dm/group policy). Port-binding platforms are NOT refused: the + default profile owns the single shared listener and a secondary's port-binders are built in + shared-listener mode (``/p//...``) by ``_start_one_profile_adapters``.""" from gateway.run import ( - MultiplexConfigError, SecondaryPortBindingConfigError, _load_gateway_runtime_config, + MultiplexConfigError, _load_gateway_runtime_config, _own_policy_open_startup_violation, _profile_runtime_scope, ) from gateway.config import load_gateway_config @@ -915,20 +910,6 @@ class GatewayAdapterLifecycleMixin: "Enable GATEWAY_ALLOW_ALL_USERS or the platform allow-all flag " "for that profile, or change dm_policy/group_policy away from 'open'." ) - port_binding_platforms = sorted( - platform.value - for platform, platform_config in profile_cfg.platforms.items() - if platform_config.enabled and _platform_binds_port(platform.value, platform_config.extra) - ) - if port_binding_platforms: - raise SecondaryPortBindingConfigError( - f"Profile '{profile_name}' enables port-binding platform(s) " - f"{', '.join(port_binding_platforms)}, but gateway.multiplex_profiles is on. The default " - f"profile owns the single shared HTTP listener and serves every " - f"profile through the /p/{profile_name}/ URL prefix. Remove " - f"these platform entries from profile '{profile_name}'s config.yaml " - f"or configure them only on the default profile." - ) return profile_cfg def _refuse_duplicate_claim( @@ -985,6 +966,14 @@ class GatewayAdapterLifecycleMixin: # Relay/WhatsApp are shared process-level ingress under multiplex; a secondary would retry-loop. if multiplex and platform in (Platform.RELAY, Platform.WHATSAPP): continue + # api_server / webhook: the default's listener already mirrors them at /p//; a second + # instance here would fight the default for the port (#100397). + if multiplex and platform.value in SHARED_LISTENER_MIRROR_PLATFORMS: + logger.info( + "[MULTIPLEX] Profile '%s': %s is served by the default profile's listener at /p/%s/ — " + "not starting a second listener", profile_name, platform.value, profile_name, + ) + continue adapter = None with _log_suppressed( logging.ERROR, "[MULTIPLEX] Profile '%s': _create_adapter('%s') raised %s", profile_name, @@ -1083,6 +1072,11 @@ class GatewayAdapterLifecycleMixin: # Secondary adapters carry their profile so prune paths namespace topic bindings correctly. # See #76423. adapter._hermes_profile_name = profile_name + # A secondary's port-binding adapter never binds: the default profile owns the one shared + # listener, which forwards /p// to this adapter's app (shared_ingress.py). + if self._multiplex_on() and platform.value not in SHARED_LISTENER_MIRROR_PLATFORMS \ + and _platform_binds_port(platform.value, getattr(getattr(adapter, "config", None), "extra", None)): + adapter._shared_listener_profile = profile_name async def _secondary_reconnect_attempt(self, profile_name: str, platform: Platform): """One scoped attempt to rebuild+connect a secondary adapter → ``(adapter, success)``; diff --git a/gateway/status.py b/gateway/status.py index b928b99407..32ba861376 100644 --- a/gateway/status.py +++ b/gateway/status.py @@ -804,7 +804,7 @@ def write_runtime_status( active_agents: Any = _UNSET, platform: Any = _UNSET, platform_state: Any = _UNSET, error_code: Any = _UNSET, error_message: Any = _UNSET, needs_attention: Any = _UNSET, retrying_since: Any = _UNSET, served_profiles: Any = _UNSET, session_store: Any = _UNSET, - clear_profile_platforms: bool = False, + ingress_url: Any = _UNSET, clear_profile_platforms: bool = False, ) -> None: """Persist gateway runtime health information for diagnostics/status.""" path = _get_runtime_status_path() @@ -842,6 +842,8 @@ def write_runtime_status( ("needs_attention", needs_attention, bool), # ISO start of the current retry episode; None clears it. ("retrying_since", retrying_since, None), + # Shared-listener secondaries: the /p// callback URL the vendor console must target. + ("ingress_url", ingress_url, None), )) # Per-entry writer provenance: top-level pid/start_time only identify the most recent # writer; /api/status tells "live" from "preserved" by exact (pid, start_time) equality.