mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-07 02:59:05 +00:00
* test(e2e): wait for MCP tool discovery instead of racing it
/v1/mcp/server returns as soon as the DB row is written, but the gateway runs
the initialize + tools/list handshake against the upstream lazily, on the first
request that needs it. Every MCP test read tools/list immediately after
registering, so it raced that handshake.
The gateway reports a server it has not discovered yet exactly like a dead one:
it catches the per-server handshake exception and returns an empty tool list.
The tests asserted on a single read, so the race surfaced as "granted key never
saw search_datadog_logs; tools=frozenset()" while a sibling test against the
same upstream in the same run passed.
Add McpClient.await_tool, which polls tools/list to the suite's existing
poll_timeout and returns the qualified tool name, and route the four discovery
sites through it. An unreachable upstream or an unapplied grant still fails, and
the failure now names the last tools/list result.
Refs LIT-4821
* test(e2e): wait for a2a agents to reach the data plane after registration
POST /v1/agents is a control-plane write; the /a2a/{agent_id} routes that serve
the card and run message/send are data plane and only see the agent after the
next DB reload. Every test registered an agent and immediately read its card or
sent it a message, so the first data-plane touch could 404 on the agent it had
just created.
register_agent now waits for the card to become servable before returning, the
same way ProxyClient.create_model waits for a new model, so callers do not each
have to poll. Registration failures skip the wait, leaving the two rejection
tests unchanged. A genuine propagation failure now fails naming the agent id and
the last card read rather than as a bare 404 on whichever /a2a call ran first.
Refs LIT-4821
* test(e2e): wait for presidio guardrails to sync before asserting masking
Registering a guardrail is a control-plane write; the data-plane worker that
serves /chat/completions only picks it up on its next periodic DB sync (~30s), so
the first call after the create ran against a worker with no guardrail and passed
the raw email straight through. The tests asserted on that first call, so they
read in-flight propagation as a PII leak.
Confirmed directly against a live proxy: the same call is unmasked at t=0s and
masked at t=8s, and the presidio analyzer itself correctly returns EMAIL_ADDRESS
with score 1.0 the whole time. The MCP guardrail suite already documents and
waits out this exact sync delay; presidio never got the same treatment.
Poll the call until the placeholder replaces the PII, so the assertions judge the
synced state. A guardrail that never masks still fails, on the last unmasked
content. pre_call and post_call now pass repeatably.
Refs LIT-4821
* test(e2e): drop the presidio logging_only check pending LIT-4841
pre_call and post_call masking both pass once the guardrail-sync wait is in place,
but logging_only left the raw email in the OTEL span's gen_ai.input.messages on
every attempt across a full poll deadline. Keeping an assertion against
known-failing behavior just turns every run red, so the cell is tracked in
LIT-4841 instead.
The registry row stays, so guardrail.presidio.logging_only.masks now reports as an
uncovered gap rather than silently disappearing.
Refs LIT-4821, LIT-4841
* test(e2e): wait for guardrail sync in bedrock, moderation and block-code checks
All three asserted on the first call after registering a guardrail, so they were
served by a data-plane worker that had not synced it yet (~30s DB poll) and read
in-flight propagation as a guardrail that failed to block. Verified directly: the
openai_moderation guardrail lets a flagged prompt through at t=0s and returns
"Violated OpenAI moderation policy" at t=8s.
The reasoning-only responses noted in triage (content=None with reasoning_tokens
set) were a symptom of the same thing, not the cause; these are pre_call
guardrails, so a synced guardrail rejects the request before the model runs.
Add poll_until_blocked to guardrails_client for the two that surface a non-success
status, and poll on the block marker in the block_code_execution check, which
replaces the reply rather than erroring. All eight guardrail tests now pass.
Refs LIT-4821
* test(e2e): drop the openai prompt-cache check pending LIT-4841
Prompt caching never engages through the proxy: cached_tokens is 0 on every
repeat, while the identical payload sent straight to OpenAI reports 3615 cached
tokens on the second call. Pinning prompt_cache_key on the proxy request restores
caching (3328 tokens), so something varying per request is defeating OpenAI's
automatic prefix cache.
That is a product bug with a direct billing cost, tracked in LIT-4841. The
registry row stays, so llm.chat_completions.openai.prompt_cache_5m.nonstream.works
now reports as an uncovered gap instead of failing every run.
Refs LIT-4821, LIT-4841
* test(e2e): drop the responses metadata redis-ttl check
It failed on a Redis read timeout against the stage serverless cache
(berrie-litellm-stage-ieib2i.serverless.use1.cache.amazonaws.com:6379), a
reachability problem this suite has hit before rather than a proxy defect the
assertion can pin down.
The file held only this test. Its other cell,
llm.responses.openai.basic.nonstream.works, is still covered by
test_responses_e2e.py; other.config.responses.metadata_redis_ttl_bounded becomes
an uncovered registry row, taking headline coverage 314/431 -> 312/431.
Refs LIT-4821
* test(e2e): fix passthrough header propagation and openai body, drop the cost check
Three separate problems behind the two passthrough failures.
The header test 404'd because POST /config/pass_through_endpoint is a
control-plane write and the worker serving the route only registers it on its next
config reload; measured at ~18s on a live proxy. Wait for the route to stop 404ing
before calling it. The readiness probe reuses the master key and omits
anthropic-version so polling does not bill a completion per attempt.
The openai passthrough body sent max_tokens, which the gpt-5 family rejects
outright ("Unsupported parameter: 'max_tokens' is not supported with this model").
Confirmed against OpenAI directly: max_tokens 400s, max_completion_tokens 200s.
Passthrough forwards the body untouched by design, so the body was simply wrong.
test_openai_passthrough_nonstreaming_logs_cost still finds no SpendLogs row for
its call_id after the fix, so it is removed rather than left red; the gemini and
anthropic passthrough cost checks still cover that path.
Passthrough suite is 8/8 green.
Refs LIT-4821
300 lines
9.7 KiB
Python
300 lines
9.7 KiB
Python
"""Client for the MCP e2e suite: admin server registration plus the api_key tool
|
|
surface.
|
|
|
|
An admin registers an upstream MCP server through the management API
|
|
(`/v1/mcp/server`, persisted in the DB) and grants a virtual key access to it via
|
|
`object_permission.mcp_servers`. Keys then reach the server through the REST bridge
|
|
the proxy exposes for api_key auth (`/mcp-rest/tools/list`, `/mcp-rest/tools/call`),
|
|
which `user_api_key_auth` gates the same way the JSON-RPC `/mcp` surface does. The
|
|
request/response bodies are co-located here because only this suite speaks MCP.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import time
|
|
from collections.abc import Mapping
|
|
from dataclasses import dataclass
|
|
|
|
from pydantic import BaseModel, ConfigDict, Field, RootModel
|
|
|
|
from e2e_http import Headers, NoBody, Result, Success, unwrap
|
|
from models import KeyGenerateBody, ObjectPermission
|
|
from proxy_client import ProxyClient
|
|
|
|
McpToolArg = str | int | float | bool | list[str] | dict[str, str]
|
|
McpToolArguments = Mapping[str, McpToolArg]
|
|
|
|
|
|
class ApiKeyHeaders(Headers):
|
|
x_litellm_api_key: str = Field(serialization_alias="x-litellm-api-key")
|
|
|
|
|
|
class McpServerNewBody(BaseModel):
|
|
server_name: str
|
|
alias: str
|
|
url: str
|
|
transport: str = "http"
|
|
auth_type: str | None = None
|
|
static_headers: dict[str, str] | None = None
|
|
allowed_tools: list[str] | None = None
|
|
mcp_access_groups: list[str] | None = None
|
|
|
|
|
|
class McpServerNewResponse(BaseModel):
|
|
server_id: str
|
|
|
|
|
|
class McpServerRow(BaseModel):
|
|
server_id: str
|
|
alias: str | None = None
|
|
url: str | None = None
|
|
|
|
|
|
class McpServersListResponse(RootModel[list[McpServerRow]]):
|
|
pass
|
|
|
|
|
|
class McpToolMcpInfo(BaseModel):
|
|
server_id: str | None = None
|
|
alias: str | None = None
|
|
|
|
|
|
class McpToolEntry(BaseModel):
|
|
name: str
|
|
description: str | None = None
|
|
mcp_info: McpToolMcpInfo | None = None
|
|
|
|
|
|
class McpToolsListResponse(BaseModel):
|
|
tools: list[McpToolEntry] = []
|
|
error: str | None = None
|
|
message: str | None = None
|
|
|
|
def tool_names_for_server(self, server_id: str) -> frozenset[str]:
|
|
return frozenset(
|
|
tool.name
|
|
for tool in self.tools
|
|
if tool.mcp_info is not None and tool.mcp_info.server_id == server_id
|
|
)
|
|
|
|
def tool_name_containing(self, server_id: str, needle: str) -> str | None:
|
|
needle_l = needle.lower()
|
|
for tool in self.tools:
|
|
if tool.mcp_info is None or tool.mcp_info.server_id != server_id:
|
|
continue
|
|
if needle_l in tool.name.lower() or tool.name.lower().endswith(needle_l):
|
|
return tool.name
|
|
return None
|
|
|
|
|
|
class BlockedWordSpec(BaseModel):
|
|
keyword: str
|
|
action: str = "BLOCK"
|
|
|
|
|
|
class ContentFilterMcpParams(BaseModel):
|
|
"""litellm_content_filter params scoped to the MCP tool-call hook. mode is
|
|
pre_mcp_call because a pre_call config silently no-ops on the tools/call path
|
|
(the event type is rewritten to pre_mcp_call for call_mcp_tool), and default_on
|
|
is required there because per-key/request guardrail selection is dropped from
|
|
the synthetic MCP request the hook sees."""
|
|
|
|
guardrail: str = "litellm_content_filter"
|
|
mode: str = "pre_mcp_call"
|
|
default_on: bool = True
|
|
blocked_words: list[BlockedWordSpec]
|
|
|
|
|
|
class GuardrailSpecBody(BaseModel):
|
|
guardrail_name: str
|
|
litellm_params: ContentFilterMcpParams
|
|
|
|
|
|
class GuardrailCreateBody(BaseModel):
|
|
guardrail: GuardrailSpecBody
|
|
|
|
|
|
class GuardrailCreateResponse(BaseModel):
|
|
guardrail_id: str
|
|
|
|
|
|
class McpCallToolBody(BaseModel):
|
|
name: str
|
|
arguments: dict[str, McpToolArg]
|
|
server_id: str
|
|
|
|
|
|
class McpCallContent(BaseModel):
|
|
type: str | None = None
|
|
text: str | None = None
|
|
|
|
|
|
class McpCallToolResponse(BaseModel):
|
|
model_config = ConfigDict(populate_by_name=True)
|
|
content: list[McpCallContent] = []
|
|
is_error: bool | None = Field(default=None, alias="isError")
|
|
|
|
@property
|
|
def first_text(self) -> str | None:
|
|
return self.content[0].text if self.content else None
|
|
|
|
@property
|
|
def all_text(self) -> str:
|
|
return "\n".join(part.text for part in self.content if part.text)
|
|
|
|
|
|
@dataclass(frozen=True, slots=True)
|
|
class McpClient:
|
|
proxy: ProxyClient
|
|
|
|
def register_server(
|
|
self,
|
|
*,
|
|
server_name: str,
|
|
alias: str,
|
|
url: str,
|
|
transport: str = "http",
|
|
auth_type: str | None = None,
|
|
static_headers: dict[str, str] | None = None,
|
|
allowed_tools: list[str] | None = None,
|
|
mcp_access_groups: list[str] | None = None,
|
|
) -> str:
|
|
return unwrap(
|
|
self.proxy.transport.post(
|
|
"/v1/mcp/server",
|
|
headers=self.proxy.transport.master,
|
|
json=McpServerNewBody(
|
|
server_name=server_name,
|
|
alias=alias,
|
|
url=url,
|
|
transport=transport,
|
|
auth_type=auth_type,
|
|
static_headers=static_headers,
|
|
allowed_tools=allowed_tools,
|
|
mcp_access_groups=mcp_access_groups,
|
|
),
|
|
response_type=McpServerNewResponse,
|
|
)
|
|
).server_id
|
|
|
|
def delete_server(self, server_id: str) -> None:
|
|
_ = self.proxy.transport.delete(
|
|
f"/v1/mcp/server/{server_id}",
|
|
headers=self.proxy.transport.master,
|
|
json=NoBody(),
|
|
response_type=NoBody,
|
|
)
|
|
|
|
def registered_servers(self) -> list[McpServerRow]:
|
|
return unwrap(
|
|
self.proxy.transport.get(
|
|
"/v1/mcp/server",
|
|
headers=self.proxy.transport.master,
|
|
params=NoBody(),
|
|
response_type=McpServersListResponse,
|
|
)
|
|
).root
|
|
|
|
def generate_key(
|
|
self,
|
|
*,
|
|
user_id: str,
|
|
mcp_servers: list[str] | None,
|
|
mcp_access_groups: list[str] | None = None,
|
|
models: list[str] | None = None,
|
|
) -> str:
|
|
object_permission = (
|
|
ObjectPermission(mcp_servers=mcp_servers, mcp_access_groups=mcp_access_groups)
|
|
if mcp_servers is not None or mcp_access_groups is not None
|
|
else None
|
|
)
|
|
return self.proxy.generate_key(
|
|
KeyGenerateBody(
|
|
models=models if models is not None else [],
|
|
user_id=user_id,
|
|
object_permission=object_permission,
|
|
)
|
|
)
|
|
|
|
def list_tools(self, key: str) -> Result[McpToolsListResponse]:
|
|
return self.proxy.transport.get(
|
|
"/mcp-rest/tools/list",
|
|
headers=ApiKeyHeaders(x_litellm_api_key=key),
|
|
params=NoBody(),
|
|
response_type=McpToolsListResponse,
|
|
)
|
|
|
|
def await_tool(self, key: str, server_id: str, needle: str) -> str:
|
|
"""Poll tools/list until `server_id` serves a tool matching `needle`, and
|
|
return its fully-qualified name. Fails at poll_timeout.
|
|
|
|
/v1/mcp/server returns as soon as the DB row is written, but the gateway
|
|
runs the initialize + tools/list handshake against the upstream lazily on
|
|
the first request that needs it, and reports a server it has not
|
|
discovered yet exactly like a dead one: an empty tool list. Waiting is
|
|
what separates the two.
|
|
"""
|
|
deadline = time.monotonic() + self.proxy.poll_timeout
|
|
while True:
|
|
result = self.list_tools(key)
|
|
if isinstance(result, Success):
|
|
tool_name = result.data.tool_name_containing(server_id, needle)
|
|
if tool_name is not None:
|
|
return tool_name
|
|
if time.monotonic() >= deadline:
|
|
raise AssertionError(
|
|
f"server {server_id} never served a tool matching {needle!r} within "
|
|
f"{self.proxy.poll_timeout}s of registration (upstream unreachable, or "
|
|
f"the key's grant was not applied); last tools/list: {result}"
|
|
)
|
|
time.sleep(self.proxy.poll_interval)
|
|
|
|
def register_mcp_content_filter(self, *, name: str, blocked_keyword: str) -> str:
|
|
"""Register a default-on content-filter guardrail that runs on the MCP
|
|
tool-call hook (pre_mcp_call) and blocks a single keyword. The keyword is
|
|
unique per test, so default_on only ever intercepts this test's own
|
|
banned tool call on the shared proxy."""
|
|
return unwrap(
|
|
self.proxy.transport.post(
|
|
"/guardrails",
|
|
headers=self.proxy.transport.master,
|
|
json=GuardrailCreateBody(
|
|
guardrail=GuardrailSpecBody(
|
|
guardrail_name=name,
|
|
litellm_params=ContentFilterMcpParams(
|
|
blocked_words=[BlockedWordSpec(keyword=blocked_keyword)],
|
|
),
|
|
)
|
|
),
|
|
response_type=GuardrailCreateResponse,
|
|
)
|
|
).guardrail_id
|
|
|
|
def delete_guardrail(self, guardrail_id: str) -> None:
|
|
_ = self.proxy.transport.delete(
|
|
f"/guardrails/{guardrail_id}",
|
|
headers=self.proxy.transport.master,
|
|
json=NoBody(),
|
|
response_type=NoBody,
|
|
)
|
|
|
|
def call_tool(
|
|
self,
|
|
key: str,
|
|
*,
|
|
server_id: str,
|
|
name: str,
|
|
arguments: McpToolArguments,
|
|
) -> Result[McpCallToolResponse]:
|
|
return self.proxy.transport.post(
|
|
"/mcp-rest/tools/call",
|
|
headers=ApiKeyHeaders(x_litellm_api_key=key),
|
|
json=McpCallToolBody(
|
|
name=name, arguments=dict(arguments), server_id=server_id
|
|
),
|
|
response_type=McpCallToolResponse,
|
|
)
|
|
|
|
|
|
def build_client(proxy: ProxyClient) -> McpClient:
|
|
return McpClient(proxy=proxy)
|