mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-10 03:28:53 +00:00
refactor(e2e): drop the unused conformance/prompts/resources foundation
Removes the auth_forwarder key-stamping relay (unused; raw httpx), the stub's /conformance mount and its fixtures, MCP_STUB_CONFORMANCE_URL, the unused prompt/resource client ops, and the unclaimed passes_official_conformance registry cell. None of it was exercised by a test; the suite stays scoped to the core happy paths against the deterministic stub.
This commit is contained in:
parent
644d1ec28f
commit
b92b7cc943
5 changed files with 2 additions and 222 deletions
|
|
@ -151,14 +151,6 @@
|
|||
assertions: [uses_per_user_token]
|
||||
source: "v2 authorization_code egress arm (per-user DB token)"
|
||||
rationale: Calls on an interactive server must present the user's upstream token, never the caller's virtual key or IdP JWT
|
||||
- id: mcp.protocol.api_key.passes_official_conformance
|
||||
module: mcp
|
||||
tier: P0
|
||||
operation: protocol
|
||||
auth_family: api_key
|
||||
assertions: [passes_official_conformance]
|
||||
source: "@modelcontextprotocol/conformance server scenarios via tests/e2e/mcp/test_mcp_conformance_e2e.py"
|
||||
rationale: The gateway re-serves the MCP protocol to hosts, so it must stay spec-conformant as a server; the official suite catches lifecycle/error-shape regressions our behavior tests do not encode
|
||||
- id: mcp.list_tools.oauth.challenges_bearer_key_not_masked
|
||||
module: mcp
|
||||
tier: P0
|
||||
|
|
|
|||
|
|
@ -46,7 +46,6 @@ MCP_STUB_URL = os.environ.get("E2E_MCP_STUB_URL", "http://mcp-stub:8765/mcp")
|
|||
# upstream for the interactive OAuth flow, plus the stub IdP's token endpoint.
|
||||
# Derived from MCP_STUB_URL so one override relocates the whole stub.
|
||||
_MCP_STUB_BASE = MCP_STUB_URL.removesuffix("/mcp")
|
||||
MCP_STUB_CONFORMANCE_URL = f"{_MCP_STUB_BASE}/conformance/mcp"
|
||||
MCP_STUB_OAUTHUSER_URL = f"{_MCP_STUB_BASE}/oauthuser/mcp"
|
||||
MCP_STUB_TOKEN_URL = f"{_MCP_STUB_BASE}/oauth/token"
|
||||
|
||||
|
|
|
|||
|
|
@ -1,88 +0,0 @@
|
|||
"""In-process reverse proxy that stamps the LiteLLM key onto every request.
|
||||
|
||||
The official conformance suite's client offers no way to send custom headers,
|
||||
but the gateway's MCP routes require the virtual key. This forwarder plays the
|
||||
role of an MCP host's HTTP layer configured with the key header: it listens on
|
||||
a loopback port, injects `x-litellm-api-key: Bearer <key>`, and relays
|
||||
everything (including SSE streams, unbuffered) to the proxy under test.
|
||||
|
||||
Runs inside the pytest process on a background uvicorn thread; tests get the
|
||||
local base URL from `serve()` and hand `{base}/{alias}/mcp` to the conformance
|
||||
CLI as the server under test.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import socket
|
||||
import threading
|
||||
import time
|
||||
from collections.abc import Generator
|
||||
from contextlib import contextmanager
|
||||
from typing import cast
|
||||
|
||||
import httpx
|
||||
import uvicorn
|
||||
from starlette.applications import Starlette
|
||||
from starlette.background import BackgroundTask
|
||||
from starlette.requests import Request
|
||||
from starlette.responses import StreamingResponse
|
||||
from starlette.routing import Route
|
||||
|
||||
from e2e_config import PROXY_BASE_URL, REQUEST_TIMEOUT
|
||||
|
||||
_HOP_BY_HOP_REQUEST_HEADERS = frozenset({"host", "content-length"})
|
||||
_HOP_BY_HOP_RESPONSE_HEADERS = frozenset({"content-length", "transfer-encoding", "content-encoding"})
|
||||
|
||||
|
||||
def _build_app(upstream: httpx.AsyncClient, litellm_key: str) -> Starlette:
|
||||
async def forward(request: Request) -> StreamingResponse:
|
||||
url = httpx.URL(path=request.url.path, query=request.url.query.encode())
|
||||
headers = {
|
||||
name: value
|
||||
for name, value in request.headers.items()
|
||||
if name.lower() not in _HOP_BY_HOP_REQUEST_HEADERS
|
||||
}
|
||||
headers["x-litellm-api-key"] = f"Bearer {litellm_key}"
|
||||
proxied = upstream.build_request(request.method, url, headers=headers, content=await request.body())
|
||||
response = await upstream.send(proxied, stream=True)
|
||||
response_headers = {
|
||||
name: value
|
||||
for name, value in response.headers.items()
|
||||
if name.lower() not in _HOP_BY_HOP_RESPONSE_HEADERS
|
||||
}
|
||||
return StreamingResponse(
|
||||
response.aiter_raw(),
|
||||
status_code=response.status_code,
|
||||
headers=response_headers,
|
||||
background=BackgroundTask(response.aclose),
|
||||
)
|
||||
|
||||
methods = ["GET", "POST", "DELETE", "PUT", "PATCH", "HEAD", "OPTIONS"]
|
||||
return Starlette(routes=[Route("/{path:path}", forward, methods=methods)])
|
||||
|
||||
|
||||
def _free_loopback_port() -> int:
|
||||
with socket.socket() as sock:
|
||||
sock.bind(("127.0.0.1", 0))
|
||||
return cast("int", sock.getsockname()[1])
|
||||
|
||||
|
||||
@contextmanager
|
||||
def serve(litellm_key: str) -> Generator[str]:
|
||||
"""Serve the forwarder for the block's duration; yields its base URL."""
|
||||
upstream = httpx.AsyncClient(base_url=PROXY_BASE_URL, timeout=httpx.Timeout(REQUEST_TIMEOUT))
|
||||
port = _free_loopback_port()
|
||||
server = uvicorn.Server(
|
||||
uvicorn.Config(_build_app(upstream, litellm_key), host="127.0.0.1", port=port, log_level="warning")
|
||||
)
|
||||
thread = threading.Thread(target=server.run, daemon=True)
|
||||
thread.start()
|
||||
deadline = time.monotonic() + 15
|
||||
while not server.started:
|
||||
assert time.monotonic() < deadline, "auth forwarder never started"
|
||||
time.sleep(0.05)
|
||||
try:
|
||||
yield f"http://127.0.0.1:{port}"
|
||||
finally:
|
||||
server.should_exit = True
|
||||
thread.join(timeout=10)
|
||||
|
|
@ -121,60 +121,6 @@ async def _call_tool(url: str, headers: dict[str, str], tool: str, arguments: To
|
|||
return McpToolText(text=_first_text(result), is_error=bool(result.isError))
|
||||
|
||||
|
||||
async def _list_prompts(url: str, headers: dict[str, str]) -> tuple[tuple[str, tuple[str, ...]], ...]:
|
||||
"""(name, argument names) per prompt, sorted by name."""
|
||||
async with _http_client(headers) as http_client:
|
||||
async with streamable_http_client(url, http_client=http_client) as (read, write, _):
|
||||
async with ClientSession(read, write) as session:
|
||||
await session.initialize()
|
||||
listed = await session.list_prompts()
|
||||
return tuple(
|
||||
sorted(
|
||||
(prompt.name, tuple(arg.name for arg in (prompt.arguments or [])))
|
||||
for prompt in listed.prompts
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
async def _get_prompt(url: str, headers: dict[str, str], name: str, arguments: dict[str, str]) -> str:
|
||||
"""The text of the rendered prompt's first message."""
|
||||
async with _http_client(headers) as http_client:
|
||||
async with streamable_http_client(url, http_client=http_client) as (read, write, _):
|
||||
async with ClientSession(read, write) as session:
|
||||
await session.initialize()
|
||||
result = await session.get_prompt(name, arguments)
|
||||
first = result.messages[0].content if result.messages else None
|
||||
if isinstance(first, TextContent):
|
||||
return first.text
|
||||
return f"<non-text prompt content: {type(first).__name__}>"
|
||||
|
||||
|
||||
async def _list_resources(url: str, headers: dict[str, str]) -> tuple[tuple[str, str], ...]:
|
||||
"""(uri, name) per resource, sorted by uri."""
|
||||
async with _http_client(headers) as http_client:
|
||||
async with streamable_http_client(url, http_client=http_client) as (read, write, _):
|
||||
async with ClientSession(read, write) as session:
|
||||
await session.initialize()
|
||||
listed = await session.list_resources()
|
||||
return tuple(sorted((str(resource.uri), resource.name) for resource in listed.resources))
|
||||
|
||||
|
||||
async def _read_resource(url: str, headers: dict[str, str], uri: str) -> str:
|
||||
"""The text of the resource's first content block."""
|
||||
from pydantic import AnyUrl
|
||||
|
||||
async with _http_client(headers) as http_client:
|
||||
async with streamable_http_client(url, http_client=http_client) as (read, write, _):
|
||||
async with ClientSession(read, write) as session:
|
||||
await session.initialize()
|
||||
result = await session.read_resource(AnyUrl(uri))
|
||||
first = result.contents[0] if result.contents else None
|
||||
text = getattr(first, "text", None)
|
||||
if isinstance(text, str):
|
||||
return text
|
||||
return f"<non-text resource content: {type(first).__name__}>"
|
||||
|
||||
|
||||
# ---------- interactive (authorization_code) OAuth: the MCP-host side ----------
|
||||
|
||||
# Where the "browser" lands at the end of the authorize dance. Nothing listens
|
||||
|
|
@ -393,22 +339,6 @@ class McpClient:
|
|||
prior poll_oauth_tool_names dance), over its own fresh MCP session."""
|
||||
return asyncio.run(_oauth_call_tool(_mcp_url(alias), headers, storage, tool, arguments))
|
||||
|
||||
def list_prompts(self, alias: str, headers: dict[str, str]) -> tuple[tuple[str, tuple[str, ...]], ...]:
|
||||
"""prompts/list over its own fresh MCP session: (name, argument names) sorted by name."""
|
||||
return asyncio.run(_list_prompts(_mcp_url(alias), headers))
|
||||
|
||||
def get_prompt(self, alias: str, headers: dict[str, str], name: str, arguments: dict[str, str]) -> str:
|
||||
"""prompts/get over its own fresh MCP session: the rendered first message's text."""
|
||||
return asyncio.run(_get_prompt(_mcp_url(alias), headers, name, arguments))
|
||||
|
||||
def list_resources(self, alias: str, headers: dict[str, str]) -> tuple[tuple[str, str], ...]:
|
||||
"""resources/list over its own fresh MCP session: (uri, name) sorted by uri."""
|
||||
return asyncio.run(_list_resources(_mcp_url(alias), headers))
|
||||
|
||||
def read_resource(self, alias: str, headers: dict[str, str], uri: str) -> str:
|
||||
"""resources/read over its own fresh MCP session: the first content block's text."""
|
||||
return asyncio.run(_read_resource(_mcp_url(alias), headers, uri))
|
||||
|
||||
def stub_stats(self, alias: str, headers: dict[str, str], stats_tool: str, marker: str) -> StubToolStats:
|
||||
outcome = self.call_tool(alias, headers, stats_tool, {"marker": marker})
|
||||
return StubToolStats.model_validate_json(outcome.text)
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
"""Deterministic MCP upstreams for the mcp e2e suite.
|
||||
|
||||
One process, one port, three streamable-http MCP mounts plus a deterministic
|
||||
One process, one port, two streamable-http MCP mounts plus a deterministic
|
||||
OAuth2 IdP, so the compose stack keeps a single `mcp-stub` service:
|
||||
|
||||
- `/mcp` — anonymous. `echo` answers immediately so auth tests can assert an
|
||||
|
|
@ -10,13 +10,6 @@ OAuth2 IdP, so the compose stack keeps a single `mcp-stub` service:
|
|||
`max_concurrent_requests` cap must bound); `stats` reads those counters back,
|
||||
so tests observe upstream concurrency through the proxy itself and the stub
|
||||
needs no side-channel port.
|
||||
- `/conformance/mcp` — anonymous, serving the official
|
||||
@modelcontextprotocol/conformance suite's hardcoded fixture contract
|
||||
(test_simple_text / test_error_handling tools, test://static-text,
|
||||
test://static-binary and the test://template/{id}/data resources,
|
||||
test_simple_prompt / test_prompt_with_arguments prompts), so the gateway can
|
||||
be conformance-tested end to end and the suite's own prompt/resource tests
|
||||
have a deterministic upstream that serves all three MCP primitives.
|
||||
- `/oauthuser/mcp` — the interactive (authorization_code) upstream: rejects
|
||||
anything but `Bearer OAUTH_USER_ACCESS_TOKEN`, which only the
|
||||
authorization_code grant hands out, so a served request proves the whole
|
||||
|
|
@ -68,7 +61,6 @@ OAUTH_USER_ACCESS_TOKEN = "e2e-stub-user-access-token"
|
|||
OAUTH_USER_REFRESH_TOKEN = "e2e-stub-user-refresh-token"
|
||||
|
||||
main_mcp = FastMCP("e2e-stub", host="0.0.0.0", port=8765, stateless_http=True)
|
||||
conformance_mcp = FastMCP("e2e-stub-conformance", host="0.0.0.0", port=8765, stateless_http=True)
|
||||
oauthuser_mcp = FastMCP("e2e-stub-oauthuser", host="0.0.0.0", port=8765, stateless_http=True)
|
||||
|
||||
|
||||
|
|
@ -117,50 +109,6 @@ def stats(marker: str) -> str:
|
|||
)
|
||||
|
||||
|
||||
@conformance_mcp.tool()
|
||||
def test_simple_text() -> str:
|
||||
"""Return a plain text result (the conformance suite's simple-text fixture)."""
|
||||
return "Hello from the e2e conformance stub"
|
||||
|
||||
|
||||
@conformance_mcp.tool()
|
||||
def test_error_handling(should_error: bool = True) -> str:
|
||||
"""Raise so the result carries isError (the conformance suite's error fixture)."""
|
||||
if should_error:
|
||||
raise ValueError("Intentional error from the e2e conformance stub")
|
||||
return "no error"
|
||||
|
||||
|
||||
@conformance_mcp.resource("test://static-text")
|
||||
def static_text_resource() -> str:
|
||||
"""The conformance suite's static text resource fixture."""
|
||||
return "Static text resource from the e2e conformance stub"
|
||||
|
||||
|
||||
@conformance_mcp.resource("test://static-binary")
|
||||
def static_binary_resource() -> bytes:
|
||||
"""The conformance suite's static binary resource fixture."""
|
||||
return b"\x89binary-fixture-bytes\x00\x01"
|
||||
|
||||
|
||||
@conformance_mcp.resource("test://template/{id}/data")
|
||||
def template_resource(id: str) -> str:
|
||||
"""The conformance suite's resource template fixture: {id} must substitute."""
|
||||
return f"Template resource data for id={id}"
|
||||
|
||||
|
||||
@conformance_mcp.prompt()
|
||||
def test_simple_prompt() -> str:
|
||||
"""The conformance suite's no-argument prompt fixture."""
|
||||
return "This is a simple prompt from the e2e conformance stub"
|
||||
|
||||
|
||||
@conformance_mcp.prompt()
|
||||
def test_prompt_with_arguments(arg1: str, arg2: str) -> str:
|
||||
"""The conformance suite's argument-substitution prompt fixture."""
|
||||
return f"Prompt rendered with arg1={arg1} and arg2={arg2}"
|
||||
|
||||
|
||||
def _register_guarded_tools(server: FastMCP, mount: str) -> None:
|
||||
def echo(text: str) -> str:
|
||||
"""Return `text` unchanged."""
|
||||
|
|
@ -276,7 +224,7 @@ async def oauth_token(request: Request) -> JSONResponse:
|
|||
|
||||
|
||||
def build_app() -> Starlette:
|
||||
servers = (main_mcp, conformance_mcp, oauthuser_mcp)
|
||||
servers = (main_mcp, oauthuser_mcp)
|
||||
apps = {server.name: server.streamable_http_app() for server in servers}
|
||||
|
||||
@contextlib.asynccontextmanager
|
||||
|
|
@ -299,7 +247,6 @@ def build_app() -> Starlette:
|
|||
expected=f"Bearer {OAUTH_USER_ACCESS_TOKEN}",
|
||||
),
|
||||
),
|
||||
Mount("/conformance", app=apps["e2e-stub-conformance"]),
|
||||
Mount("/", app=apps["e2e-stub"]),
|
||||
],
|
||||
lifespan=lifespan,
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue