diff --git a/tests/gateway/test_multiplex_adapter_registry.py b/tests/gateway/test_multiplex_adapter_registry.py index 876697358f..0a16bd2c29 100644 --- a/tests/gateway/test_multiplex_adapter_registry.py +++ b/tests/gateway/test_multiplex_adapter_registry.py @@ -679,36 +679,45 @@ class TestSecondaryProfileConfigHandling: @pytest.mark.asyncio - async def test_secondary_reports_all_port_binding_platforms(self, monkeypatch): - from gateway.run import SecondaryPortBindingConfigError + async def test_secondary_port_binders_run_in_shared_listener_mode(self, monkeypatch): + """A secondary's inbound-port platforms are NOT refused: they are built without a port and + served at /p// by the default's listener; api_server/webhook (already mirrored + there) are skipped, never a second instance.""" from gateway.config import GatewayConfig, Platform, PlatformConfig runner = GatewayRunner.__new__(GatewayRunner) runner.config = GatewayConfig(multiplex_profiles=True) runner._profile_adapters = {} + runner.adapters = {} reviewer_cfg = GatewayConfig(multiplex_profiles=True) reviewer_cfg.platforms = { # connection_mode=webhook: with #52563's conditional check merged, - # default (websocket) Feishu no longer binds a port — only webhook - # mode should be reported here. - Platform.FEISHU: PlatformConfig( - enabled=True, extra={"connection_mode": "webhook"} - ), + # default (websocket) Feishu does not bind a port. + Platform.FEISHU: PlatformConfig(enabled=True, extra={"connection_mode": "webhook"}), Platform.WEBHOOK: PlatformConfig(enabled=True, extra={"port": 8644}), Platform.TELEGRAM: PlatformConfig(enabled=True, token="t"), } - monkeypatch.setattr( - "gateway.config.load_gateway_config", lambda: reviewer_cfg - ) + monkeypatch.setattr("gateway.config.load_gateway_config", lambda: reviewer_cfg) + created = {} - with pytest.raises(SecondaryPortBindingConfigError) as ei: - await runner._start_one_profile_adapters("reviewer", "/tmp/x", {}) - message = str(ei.value) - assert "feishu" in message - assert "webhook" in message - assert "telegram" not in message - assert "reviewer" not in runner._profile_adapters + def fake_create(platform, platform_config): + created[platform] = _FakeAdapter(token=platform_config.token or None) + created[platform].config = platform_config + return created[platform] + + monkeypatch.setattr(runner, "_create_adapter", fake_create) + monkeypatch.setattr(runner, "_wire_adapter_handlers", lambda *a, **k: None) + monkeypatch.setattr(runner, "_bind_voice_input_callback", lambda *a, **k: None) + monkeypatch.setattr(runner, "_sync_voice_mode_state_to_adapter", lambda *a, **k: None) + monkeypatch.setattr(runner, "_connect_initial_adapter_with_timeout", AsyncMock(return_value=True)) + + connected = await runner._start_one_profile_adapters("reviewer", "/tmp/x", {}) + + assert connected == 2 + assert set(created) == {Platform.FEISHU, Platform.TELEGRAM} # webhook is a default-listener mirror + assert created[Platform.FEISHU]._shared_listener_profile == "reviewer" + assert getattr(created[Platform.TELEGRAM], "_shared_listener_profile", None) is None def test_configured_secondary_adapter_namespaces_runtime_status(self): runner = _secondary_recovery_runner() @@ -820,8 +829,7 @@ class TestSecondaryProfileConfigHandling: async def fake_start_one(profile_name, profile_home, claimed): if profile_name == "bad": - from gateway.run import SecondaryPortBindingConfigError - raise SecondaryPortBindingConfigError("bad enables webhook") + raise RuntimeError("bad profile blew up at startup") runner._profile_adapters[profile_name] = {} return 2 @@ -855,7 +863,7 @@ class TestSecondaryProfileConfigHandling: assert status["served_profiles"] == ["default", "bad", "good"] assert "good" in runner._profile_adapters assert "bad" not in runner._profile_adapters - assert "Skipping secondary profile 'bad'" in caplog.text + assert "Failed to start adapters for profile 'bad'" in caplog.text @pytest.mark.asyncio async def test_multiplexer_propagates_security_config_error(self, monkeypatch): @@ -1007,29 +1015,6 @@ class TestSecondaryProfileConfigHandling: assert second == 1 assert runner._profile_adapters["later"][photon] is later - @pytest.mark.asyncio - async def test_secondary_teams_uses_degradable_error(self, monkeypatch): - from gateway.config import GatewayConfig, Platform, PlatformConfig - from gateway.run import SecondaryPortBindingConfigError - - runner = GatewayRunner.__new__(GatewayRunner) - runner.config = GatewayConfig(multiplex_profiles=True) - runner._profile_adapters = {} - - reviewer_cfg = GatewayConfig(multiplex_profiles=True) - reviewer_cfg.platforms = { - Platform("teams"): PlatformConfig(enabled=True, extra={"port": 3978}), - } - monkeypatch.setattr( - "gateway.config.load_gateway_config", lambda: reviewer_cfg - ) - - with pytest.raises(SecondaryPortBindingConfigError) as exc_info: - await runner._start_one_profile_adapters("reviewer", "/tmp/x", {}) - assert "teams" in str(exc_info.value) - assert "reviewer" in str(exc_info.value) - assert "reviewer" not in runner._profile_adapters - @pytest.mark.asyncio async def test_secondary_profile_adapter_start_skips_whatsapp(self, monkeypatch): """WhatsApp is shared process-level ingress like Relay: the bridge is diff --git a/tests/gateway/test_multiplex_shared_ingress.py b/tests/gateway/test_multiplex_shared_ingress.py new file mode 100644 index 0000000000..299dc3a78b --- /dev/null +++ b/tests/gateway/test_multiplex_shared_ingress.py @@ -0,0 +1,137 @@ +"""Inbound-port platforms of a SECONDARY profile are served on the default profile's shared listener +at ``/p//`` (gateway/platforms/shared_ingress.py). + +Invariants: the forwarded request is verified by the NAMED profile's adapter with that profile's +secret and runs under that profile's runtime scope; the un-prefixed path is untouched; a profile +without an adapter for the path is a 404 rather than the default's adapter; a shared-listener +adapter binds no port of its own. +""" +from __future__ import annotations + +import base64 +import hashlib +import hmac +from pathlib import Path +from typing import Any + +import pytest + +pytest.importorskip("aiohttp") +from aiohttp import web # noqa: E402 +from aiohttp.test_utils import TestClient, TestServer # noqa: E402 + +from gateway.config import GatewayConfig, Platform, PlatformConfig # noqa: E402 + + +def _line_sig(body: bytes, secret: str) -> str: + return base64.b64encode(hmac.new(secret.encode(), body, hashlib.sha256).digest()).decode() + + +class _Runner: + """Just enough GatewayRunner for shared_ingress: the served-profile adapter map.""" + + def __init__(self, adapters: dict[str, dict[Platform, Any]]): + self.config = GatewayConfig(multiplex_profiles=True) + self._profile_adapters = adapters + self.adapters: dict = {} + + +def _line_adapter(secret: str, profile: str): + from plugins.platforms.line.adapter import LineAdapter + adapter = LineAdapter(PlatformConfig(enabled=True, extra={ + "channel_access_token": f"tok-{profile}", "channel_secret": secret, "port": 1})) + adapter._shared_listener_profile = profile + adapter.set_owner_profile(profile) + return adapter + + +async def _publish_line(adapter, runner) -> list[tuple[str, Path]]: + """Wire the LINE webhook app the way ``connect()`` does, without the LINE API or a bind, and + record the HERMES_HOME the handler ran under.""" + from gateway.platforms.shared_ingress import bind_listener + from hermes_constants import get_hermes_home + seen: list[tuple[str, Path]] = [] + adapter.gateway_runner = runner + + async def dispatch(event): + seen.append((event.get("type"), Path(get_hermes_home()))) + + adapter._dispatch_event = dispatch + app = web.Application(client_max_size=1024) + app.router.add_post(adapter.webhook_path, adapter._handle_webhook) + bound = await bind_listener(adapter, app, "127.0.0.1", 1, adapter.webhook_path) + assert bound is None # shared-listener mode never binds + return seen + + +@pytest.fixture +def mux_home(tmp_path, monkeypatch): + root = tmp_path / "hermes" + for name in ("coder", "ops"): + (root / "profiles" / name).mkdir(parents=True) + (root / "profiles" / name / ".env").write_text(f"LINE_CHANNEL_SECRET=secret-{name}\n") + monkeypatch.setenv("HERMES_HOME", str(root)) + import hermes_constants + monkeypatch.setattr(hermes_constants, "_default_hermes_root_memo", None) + monkeypatch.setattr(Path, "home", lambda: tmp_path) + from agent import secret_scope + monkeypatch.setattr(secret_scope, "_MULTIPLEX_ACTIVE", True) + monkeypatch.setattr( + "hermes_cli.profiles.profiles_to_serve", + lambda multiplex: [("default", root), ("coder", root / "profiles" / "coder"), ("ops", root / "profiles" / "ops")]) + return root + + +async def _shared_listener(runner) -> TestClient: + """The default profile's listener as the webhook adapter builds it: bare routes + /p/ forwarding.""" + from gateway.platforms.webhook import WebhookAdapter + listener = WebhookAdapter(PlatformConfig(enabled=True, extra={"port": 1, "routes": {}})) + listener.gateway_runner = runner + app = web.Application() + app.router.add_post("/webhooks/{route_name}", listener._handle_webhook) + app.router.add_post("/p/{profile}/webhooks/{route_name}", listener._handle_webhook) + app.router.add_route("*", "/p/{profile}/{tail:.*}", listener._handle_profile_ingress) + return TestClient(TestServer(app)) + + +@pytest.mark.asyncio +async def test_prefixed_line_webhook_is_verified_by_the_named_profiles_secret_under_its_scope(mux_home): + coder, ops = _line_adapter("secret-coder", "coder"), _line_adapter("secret-ops", "ops") + runner = _Runner({"coder": {Platform("line"): coder}, "ops": {Platform("line"): ops}}) + coder_seen, ops_seen = await _publish_line(coder, runner), await _publish_line(ops, runner) + body = b'{"events":[{"type":"probe"}]}' + async with await _shared_listener(runner) as client: + ok = await client.post("/p/coder/line/webhook", data=body, + headers={"X-Line-Signature": _line_sig(body, "secret-coder")}) + assert ok.status == 200 + # ops' secret is rejected at coder's URL — the adapter is per profile, so is the secret. + wrong = await client.post("/p/coder/line/webhook", data=body, + headers={"X-Line-Signature": _line_sig(body, "secret-ops")}) + assert wrong.status == 401 + # The un-prefixed path keeps serving only the default profile's own routes. + bare = await client.post("/line/webhook", data=body, + headers={"X-Line-Signature": _line_sig(body, "secret-coder")}) + assert bare.status == 404 + # An unserved profile, and a served profile without that platform, are 404 — never another adapter. + assert (await client.post("/p/nope/line/webhook", data=body)).status == 404 + runner._profile_adapters["ops"] = {} + assert (await client.post("/p/ops/line/webhook", data=body, + headers={"X-Line-Signature": _line_sig(body, "secret-ops")})).status == 404 + assert [t for t, _ in coder_seen] == ["probe"] and ops_seen == [] + assert coder_seen[0][1] == mux_home / "profiles" / "coder" + + +@pytest.mark.asyncio +async def test_shared_listener_adapter_records_its_public_ingress_url(mux_home, monkeypatch): + """Runtime status carries the /p// URL so `gateway status` / the dashboard can show it.""" + writes: list[dict] = [] + monkeypatch.setattr("gateway.status.write_runtime_status", lambda **kw: writes.append(kw)) + coder = _line_adapter("secret-coder", "coder") + coder._runtime_status_platform_key = "coder:line" + runner = _Runner({"coder": {Platform("line"): coder}}) + runner.adapters = {Platform.API_SERVER: type("L", (), {"_host": "0.0.0.0", "_port": 8642})()} + await _publish_line(coder, runner) + assert coder._shared_ingress_url == "http://127.0.0.1:8642/p/coder/line/webhook" + assert {"platform": "coder:line", "ingress_url": coder._shared_ingress_url} == { + k: v for k, v in writes[-1].items() if k in ("platform", "ingress_url")} + assert coder._media_url("tok", "a.png").startswith("http://127.0.0.1:8642/p/coder/line/media/") diff --git a/tests/hermes_cli/test_web_server_messaging_profiles.py b/tests/hermes_cli/test_web_server_messaging_profiles.py index 193bb4b219..f0770944d8 100644 --- a/tests/hermes_cli/test_web_server_messaging_profiles.py +++ b/tests/hermes_cli/test_web_server_messaging_profiles.py @@ -202,13 +202,9 @@ def _enable_multiplex(default_home): class TestMultiplexPortBindingGuard: - """Enabling a port-binding channel on a secondary multiplexed profile - must be rejected BEFORE anything is persisted. - - The gateway fail-fasts with ``MultiplexConfigError`` when a secondary - profile enables a port-binding platform under - ``gateway.multiplex_profiles`` — but the dashboard used to persist that - exact config, so the next gateway start died for EVERY profile (#62791). + """Enabling api_server/webhook on a secondary multiplexed profile is rejected BEFORE anything + is persisted: the default profile's listener already mirrors them at ``/p//`` (#62791). + Every other inbound-port platform is allowed — the gateway serves it on the shared listener. """ @pytest.fixture(autouse=True) @@ -217,21 +213,25 @@ class TestMultiplexPortBindingGuard: # multiplex flag under test comes from the default profile's config. monkeypatch.delenv("GATEWAY_MULTIPLEX_PROFILES", raising=False) - def test_rejects_every_port_binding_platform_on_secondary( + def test_rejects_only_mirrored_listeners_on_secondary( self, client, isolated_profiles ): - from gateway.config import PORT_BINDING_PLATFORM_VALUES + from gateway.config import PORT_BINDING_PLATFORM_VALUES, SHARED_LISTENER_MIRROR_PLATFORMS _enable_multiplex(isolated_profiles["default"]) - assert PORT_BINDING_PLATFORM_VALUES # guard set must not be empty - for platform_id in sorted(PORT_BINDING_PLATFORM_VALUES): + assert SHARED_LISTENER_MIRROR_PLATFORMS # guard set must not be empty + catalog = {p["id"] for p in client.get("/api/messaging/platforms").json()["platforms"]} + for platform_id in sorted(PORT_BINDING_PLATFORM_VALUES & catalog): resp = client.put( f"/api/messaging/platforms/{platform_id}", params={"profile": "worker_alpha"}, json={"enabled": True}, ) - assert resp.status_code == 409, platform_id - assert "default profile" in resp.json()["detail"] + if platform_id in SHARED_LISTENER_MIRROR_PLATFORMS: + assert resp.status_code == 409, platform_id + assert "default profile" in resp.json()["detail"] + else: # served at /p/worker_alpha/ on the shared listener + assert resp.status_code == 200, (platform_id, resp.text)