Files
hermes-agent/tests/tools/test_mcp_client_cert.py
Teknium a6fdadfcee fix(mcp): cap HTTP/SSE response bodies before SDK parse
Port from openclaw/openclaw#123194: a hostile or misbehaving remote MCP
server could stream an unbounded HTTP catalog/tool-result body that the
MCP SDK buffers and JSON-parses before any of Hermes' post-parse limits
(resource cap, tool-result truncation) run.

New _make_mcp_body_cap_transport wraps the owned httpx AsyncClient's
transport on the Streamable HTTP (mcp >= 1.24) and SSE paths:
- finite HTTP bodies capped at 10 MiB (Content-Length rejected up front,
  streamed bodies capped chunk-by-chunk);
- each SSE event capped at 10 MiB, with accounting reset at completed
  event boundaries so long-lived streams/keepalives are unlimited;
- violations raise httpx.ReadError naming the byte cap, handled by the
  existing transport teardown/reconnect path (#66092).

verify/cert now live on the inner AsyncHTTPTransport (client-level TLS
kwargs are inert once a custom transport is passed); the SSE
httpx_client_factory is always injected so the cap applies with default
TLS too. Legacy mcp < 1.24 path (SDK-internal client, no hook) stays
uncapped — same degradation as strict_redirect_headers.
2026-09-13 20:45:25 -07:00

395 lines
14 KiB
Python

"""Tests for mTLS client certificate config on MCP HTTP/SSE transports.
Covers:
1. ``_resolve_client_cert`` helper — string, tuple, encrypted-key, validation
errors, missing-file errors.
2. HTTP (new SDK ``streamable_http_client``) path forwards ``cert=`` into the
inner ``AsyncHTTPTransport`` wrapped by the wire-body-cap transport that
the user-owned ``httpx.AsyncClient`` is built on.
3. SSE path forwards ``cert`` and ``ssl_verify`` via an ``httpx_client_factory``
(always injected — it also installs the wire-body cap) without breaking the
OAuth/headers/timeout passthrough.
"""
from __future__ import annotations
import asyncio
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
def _patch_sdk_async_client(dummy):
"""Patch ``AsyncClient`` on whichever httpx module the MCP SDK uses.
mcp 2.0 moved the SDK's HTTP stack to ``httpx2``, so patching
``httpx.AsyncClient`` no longer intercepts the client Hermes builds for
the SDK. Resolve the module the same way production does, via
``tools.mcp_tool.sdk_httpx``, so these tests follow the SDK rather than
hardcoding a distribution name.
"""
from tools.mcp_tool import sdk_httpx
return patch.object(sdk_httpx(), "AsyncClient", dummy)
def _patch_sdk_transport(dummy):
"""Patch ``AsyncHTTPTransport`` on the SDK's httpx module.
Since the wire-body cap (_make_mcp_body_cap_transport), verify/cert are
applied to an inner ``AsyncHTTPTransport`` rather than as AsyncClient
kwargs — a custom ``transport=`` makes client-level TLS kwargs inert.
"""
from tools.mcp_tool import sdk_httpx
return patch.object(sdk_httpx(), "AsyncHTTPTransport", dummy)
class _DummyTransport:
"""Capture-only stand-in for AsyncHTTPTransport."""
captured: dict = {}
def __init__(self, **kwargs):
type(self).captured = dict(kwargs)
async def handle_async_request(self, request): # pragma: no cover
raise AssertionError("not dispatched in these tests")
async def aclose(self): # pragma: no cover
pass
# ---------------------------------------------------------------------------
# _resolve_client_cert helper
# ---------------------------------------------------------------------------
class TestResolveClientCert:
def test_returns_none_when_unset(self):
from tools.mcp_tool_errors import _resolve_client_cert
assert _resolve_client_cert("srv", {}) is None
assert _resolve_client_cert("srv", {"url": "https://x"}) is None
def test_string_form_single_pem(self, tmp_path):
from tools.mcp_tool_errors import _resolve_client_cert
pem = tmp_path / "combined.pem"
pem.write_text("dummy")
result = _resolve_client_cert("srv", {"client_cert": str(pem)})
assert result == str(pem)
def test_list_form_two_elements(self, tmp_path):
from tools.mcp_tool_errors import _resolve_client_cert
cert = tmp_path / "client.crt"
key = tmp_path / "client.key"
cert.write_text("cert")
key.write_text("key")
result = _resolve_client_cert("srv", {
"client_cert": [str(cert), str(key)],
})
assert result == (str(cert), str(key))
def test_password_must_be_string(self, tmp_path):
from tools.mcp_tool_errors import _resolve_client_cert
cert = tmp_path / "client.crt"
key = tmp_path / "client.key"
cert.write_text("cert")
key.write_text("key")
with pytest.raises(ValueError, match=r"key passphrase.*must be a string"):
_resolve_client_cert("srv", {
"client_cert": [str(cert), str(key), 42],
})
# ---------------------------------------------------------------------------
# HTTP transport — cert forwarded into httpx.AsyncClient
# ---------------------------------------------------------------------------
class TestHTTPClientCert:
def test_cert_forwarded_to_async_client(self, tmp_path):
"""When client_cert is set, the new-SDK HTTP path passes ``cert=``
into the inner AsyncHTTPTransport under the body-cap wrapper."""
from tools.mcp_tool import MCPServerTask
cert = tmp_path / "client.pem"
cert.write_text("dummy")
server = MCPServerTask("remote")
captured: dict = {}
class DummyAsyncClient:
def __init__(self, **kwargs):
captured.update(kwargs)
async def __aenter__(self):
return self
async def __aexit__(self, *a):
return False
class DummyTransportCtx:
async def __aenter__(self):
return MagicMock(), MagicMock(), (lambda: None)
async def __aexit__(self, *a):
return False
class DummySession:
def __init__(self, *args, **kwargs):
pass
async def __aenter__(self):
return self
async def __aexit__(self, *a):
return False
async def initialize(self):
return None
async def _discover_tools(self):
self._shutdown_event.set()
async def _drive():
with patch("tools.mcp_tool._MCP_HTTP_AVAILABLE", True), \
patch("tools.mcp_tool._MCP_NEW_HTTP", True), \
_patch_sdk_async_client(DummyAsyncClient), \
_patch_sdk_transport(_DummyTransport), \
patch("tools.mcp_tool.streamable_http_client",
return_value=DummyTransportCtx()), \
patch("tools.mcp_tool.ClientSession", DummySession), \
patch.object(MCPServerTask, "_discover_tools", _discover_tools):
await server._run_http({
"url": "https://example.com/mcp",
"client_cert": str(cert),
})
asyncio.run(_drive())
# cert/verify live on the inner transport (client-level TLS kwargs
# are inert once a custom transport is passed); the client itself
# receives the body-cap transport.
assert _DummyTransport.captured.get("cert") == str(cert)
assert _DummyTransport.captured.get("verify") is True
assert "cert" not in captured
assert "transport" in captured
def test_missing_cert_file_surfaces_clear_error(self, tmp_path):
"""A missing cert file fails fast with a server-scoped error message."""
from tools.mcp_tool import MCPServerTask
server = MCPServerTask("remote")
async def _drive():
with patch("tools.mcp_tool._MCP_HTTP_AVAILABLE", True), \
patch("tools.mcp_tool._MCP_NEW_HTTP", True):
await server._run_http({
"url": "https://example.com/mcp",
"client_cert": str(tmp_path / "nope.pem"),
})
with pytest.raises(FileNotFoundError, match=r"remote.*client_cert.*not found"):
asyncio.run(_drive())
# ---------------------------------------------------------------------------
# SSE transport — cert + verify routed via httpx_client_factory
# ---------------------------------------------------------------------------
@pytest.fixture
def patch_sse_client():
"""Replace ``sse_client`` with a MagicMock that records its kwargs.
Returns the captured kwargs dict so tests can assert how ``_run_http``
called it.
"""
captured_kwargs: dict = {}
class _FakeStream:
def __init__(self):
self._read = AsyncMock()
self._write = AsyncMock()
async def __aenter__(self):
return (self._read, self._write)
async def __aexit__(self, *a):
return False
def fake_sse_client(**kwargs):
captured_kwargs.clear()
captured_kwargs.update(kwargs)
return _FakeStream()
class _FakeSession:
def __init__(self, *args, **kwargs):
pass
async def __aenter__(self):
mock_session = MagicMock()
mock_session.initialize = AsyncMock()
return mock_session
async def __aexit__(self, *a):
return False
with patch("tools.mcp_tool.sse_client", new=fake_sse_client), \
patch("tools.mcp_tool.ClientSession", new=_FakeSession):
yield captured_kwargs
class TestSSEClientCert:
def test_factory_always_injected_with_default_tls(self, patch_sse_client):
"""The factory is always injected (it installs the wire-body cap);
with no cert and ssl_verify=True the inner transport keeps default
TLS settings."""
from tools.mcp_tool import MCPServerTask
server = MCPServerTask("sse-test")
server._auth_type = ""
server._sampling = None
async def drive():
with patch.object(MCPServerTask, "_wait_for_lifecycle_event",
new=AsyncMock(return_value="shutdown")), \
patch.object(MCPServerTask, "_discover_tools", new=AsyncMock()):
try:
await asyncio.wait_for(
server._run_http({
"url": "https://example.com/mcp/sse",
"transport": "sse",
}),
timeout=2.0,
)
except (asyncio.TimeoutError, StopAsyncIteration, Exception):
pass
asyncio.run(drive())
factory = patch_sse_client.get("httpx_client_factory")
assert factory is not None, "body-cap factory must always be injected"
captured_client_kwargs: dict = {}
class DummyAsyncClient:
def __init__(self, **kwargs):
captured_client_kwargs.update(kwargs)
with _patch_sdk_async_client(DummyAsyncClient), \
_patch_sdk_transport(_DummyTransport):
factory(headers=None, timeout=None, auth=None)
assert _DummyTransport.captured.get("verify") is True
assert "cert" not in _DummyTransport.captured
assert "transport" in captured_client_kwargs
def test_factory_injected_when_cert_set(self, patch_sse_client, tmp_path):
"""With client_cert set, an httpx_client_factory is injected that
applies the cert (and follow_redirects=True to match the SDK)."""
from tools.mcp_tool import MCPServerTask
cert = tmp_path / "client.pem"
cert.write_text("dummy")
server = MCPServerTask("sse-test")
server._auth_type = ""
server._sampling = None
async def drive():
with patch.object(MCPServerTask, "_wait_for_lifecycle_event",
new=AsyncMock(return_value="shutdown")), \
patch.object(MCPServerTask, "_discover_tools", new=AsyncMock()):
try:
await asyncio.wait_for(
server._run_http({
"url": "https://example.com/mcp/sse",
"transport": "sse",
"client_cert": str(cert),
}),
timeout=2.0,
)
except (asyncio.TimeoutError, StopAsyncIteration, Exception):
pass
asyncio.run(drive())
factory = patch_sse_client.get("httpx_client_factory")
assert factory is not None, "expected httpx_client_factory to be injected"
# Invoke the factory the way the SDK would; capture the resulting
# httpx.AsyncClient kwargs.
captured_client_kwargs: dict = {}
class DummyAsyncClient:
def __init__(self, **kwargs):
captured_client_kwargs.update(kwargs)
from tools.mcp_tool import sdk_httpx
with _patch_sdk_async_client(DummyAsyncClient), \
_patch_sdk_transport(_DummyTransport):
factory(headers={"x": "y"}, timeout=sdk_httpx().Timeout(30.0), auth=None)
assert _DummyTransport.captured["cert"] == str(cert)
assert _DummyTransport.captured["verify"] is True
assert captured_client_kwargs["follow_redirects"] is True
assert captured_client_kwargs["headers"] == {"x": "y"}
assert "transport" in captured_client_kwargs
def test_factory_forwards_custom_ca_bundle(self, patch_sse_client, tmp_path):
"""ssl_verify as a path is forwarded to the factory's httpx client."""
from tools.mcp_tool import MCPServerTask
ca_bundle = tmp_path / "ca.pem"
ca_bundle.write_text("dummy")
server = MCPServerTask("sse-test")
server._auth_type = ""
server._sampling = None
async def drive():
with patch.object(MCPServerTask, "_wait_for_lifecycle_event",
new=AsyncMock(return_value="shutdown")), \
patch.object(MCPServerTask, "_discover_tools", new=AsyncMock()):
try:
await asyncio.wait_for(
server._run_http({
"url": "https://example.com/mcp/sse",
"transport": "sse",
"ssl_verify": str(ca_bundle),
}),
timeout=2.0,
)
except (asyncio.TimeoutError, StopAsyncIteration, Exception):
pass
asyncio.run(drive())
factory = patch_sse_client.get("httpx_client_factory")
assert factory is not None
captured_client_kwargs: dict = {}
class DummyAsyncClient:
def __init__(self, **kwargs):
captured_client_kwargs.update(kwargs)
with _patch_sdk_async_client(DummyAsyncClient), \
_patch_sdk_transport(_DummyTransport):
factory(headers=None, timeout=None, auth=None)
assert _DummyTransport.captured["verify"] == str(ca_bundle)
assert "cert" not in _DummyTransport.captured