test(multiplex): shared-listener ingress invariants
- /p/<profile>/line/webhook is verified with the NAMED profile's channel secret under its runtime scope; another profile's secret is 401 at that URL; the bare path is untouched; unknown profile / profile without the adapter is 404. - A shared-listener adapter binds no port and records its /p/<profile>/ ingress_url in runtime status; LINE media URLs use the shared prefix. - Runner: a secondary's port-binders are constructed in shared-listener mode instead of refusing the whole profile; api_server/webhook are skipped as mirrors. - Dashboard: only the mirrored pair is refused (409) on a secondary.
This commit is contained in:
@@ -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/<profile>/ 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
|
||||
|
||||
137
tests/gateway/test_multiplex_shared_ingress.py
Normal file
137
tests/gateway/test_multiplex_shared_ingress.py
Normal file
@@ -0,0 +1,137 @@
|
||||
"""Inbound-port platforms of a SECONDARY profile are served on the default profile's shared listener
|
||||
at ``/p/<profile>/<path>`` (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/<profile>/ 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/")
|
||||
@@ -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/<profile>/`` (#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/<path> on the shared listener
|
||||
assert resp.status_code == 200, (platform_id, resp.text)
|
||||
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user