This commit is contained in:
devin-ai-integration[bot] 2026-10-03 20:10:30 +08:00 • committed by GitHub
commit 81ee9a57c6
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
9 changed files with 28 additions and 38 deletions

View file

@ -241,7 +241,7 @@ env -i PATH="$PATH" HOME="$HOME" PYTHONPATH="$PYTHONPATH" \
if [ "${INTEGRATION_COVERAGE:-0}" = 1 ]; then
for covered_pid in "$proxy_pid" "$peer_pid"; do
[ -n "$covered_pid" ] || continue
kill -TERM -- "-$covered_pid"
kill -TERM -- "-$covered_pid" 2>/dev/null || true
for _ in {1..300}; do
kill -0 "$covered_pid" 2>/dev/null || break
sleep 0.1

View file

@ -145,8 +145,6 @@ def _drain_sse_streams() -> None:
def _draining_sse_watcher(app: Callable[[Scope, Receive, Send], object]):
"""sse_starlette parks a per-loop watcher that only stops once AppStatus.should_exit flips."""
async def lifespan(scope: Scope, receive: Receive, send: Send) -> None:
while True:
message: Final = await receive()
@ -230,7 +228,6 @@ def jsonrpc_error(identity: object, code: int, message: str) -> Reply:
@contextmanager
def scripted_peer(*tools: ScriptedTool) -> Iterator[McpPeer]:
"""Raw JSON-RPC peer for shapes the SDK server cannot produce: half-written bodies, stalls, wire errors."""
observed: Final[queue.Queue[dict[str, object]]] = queue.Queue()
by_name: Final = {tool.name: tool for tool in tools}
@ -290,7 +287,6 @@ def echo_tool(name: str) -> ScriptedTool:
@contextmanager
def openapi_peer() -> Iterator[McpPeer]:
"""OpenAPI-described HTTP service plus the spec file the proxy turns into MCP tools."""
observed: Final[queue.Queue[dict[str, object]]] = queue.Queue()
def provider(request: Request) -> Reply:
@ -375,16 +371,6 @@ def peer_of(kind: PeerKind, *, rich: bool = False) -> Iterator[McpPeer]:
yield candidate
def register_mcp(scenario: Scenario, peer: McpPeer, alias: str, **fields: object) -> str:
response: Final = scenario.gateway.request(
"POST", "/v1/mcp/server", {"server_name": alias, "alias": alias, **peer.registration(), **fields}
)
identity: Final = response.json()["server_id"]
scenario.cleanups.callback(forget_mcp, scenario.gateway, identity)
assert response.status_code == 201, response.text
return identity
def forget_mcp(gateway: Gateway, identity: str) -> None:
response: Final = gateway.request("DELETE", f"/v1/mcp/server/{identity}")
assert response.status_code in (202, 404), response.text
@ -396,6 +382,23 @@ def delete_mcp(gateway: Gateway, identity: str) -> None:
assert read_rows('SELECT server_id FROM "LiteLLM_MCPServerTable" WHERE server_id = %s', (identity,)) == []
def register_mcp(
scenario: Scenario,
peer: McpPeer,
alias: str,
*,
cleanup: Callable[[Gateway, str], None] = delete_mcp,
**fields: object,
) -> str:
response: Final = scenario.gateway.request(
"POST", "/v1/mcp/server", {"server_name": alias, "alias": alias, **peer.registration(), **fields}
)
identity: Final = response.json()["server_id"]
scenario.cleanups.callback(cleanup, scenario.gateway, identity)
assert response.status_code == 201, response.text
return identity
def listed_tools(gateway: Gateway, key: str, identity: str | None = None) -> dict[str, dict[str, object]]:
response: Final = gateway.client.get(
"/mcp-rest/tools/list",
@ -438,8 +441,6 @@ INITIALIZE: Final = {
@dataclass(frozen=True, slots=True)
class Outcome:
"""What a caller saw from one MCP operation, normalised across entry points."""
status: int
error: str | None
tools: tuple[str, ...] = ()
@ -493,8 +494,6 @@ def _outcome_from_rest(response: httpx.Response) -> Outcome:
@dataclass(frozen=True, slots=True)
class McpCaller:
"""One caller's view of the gateway through a specific entry point."""
gateway: Gateway
key: str | None
entry: EntryPoint
@ -558,7 +557,6 @@ class McpCaller:
def _legacy_sse_rpc(
gateway: Gateway, headers: Mapping[str, str], method: str, params: JsonRpc | None
) -> httpx.Response:
"""Drive the legacy GET /mcp/sse + POST /mcp/sse/messages pair for one request and synthesise a JSON response."""
with gateway.client.stream("GET", "/mcp/sse", headers=headers, timeout=15) as stream:
if stream.status_code != 200:
stream.read()
@ -586,7 +584,6 @@ def _legacy_sse_rpc(
def official_client_outcomes(
gateway: Gateway, key: str, path: str, name: str, arguments: JsonRpc, *, legacy_sse: bool = False
) -> tuple[Outcome, Outcome]:
"""List then call through the official MCP client session, returning both outcomes."""
url: Final = str(gateway.client.base_url).rstrip("/") + path
headers: Final = {"x-litellm-api-key": key}

View file

@ -21,8 +21,6 @@ SUBJECTS: Final[tuple[Subject, ...]] = (
@dataclass(frozen=True, slots=True)
class Caller:
"""A key plus the request headers that make the proxy resolve the granted subject."""
key: str
headers: Mapping[str, str]
@ -75,11 +73,6 @@ def grant(
access_group: str | None = None,
allowed_tools: Mapping[str, tuple[str, ...]] | None = None,
) -> Caller:
"""Build a caller whose ``subject`` level grants exactly ``granted`` out of ``ceiling``.
``ceiling`` is what the key itself can reach before the subject narrows it; the key subject grants
``granted`` directly. Access groups take the group name that the granted servers were registered with,
and ``allowed_tools`` maps server id to the tools the key may call on it."""
gateway: Final = scenario.gateway
match subject:
case "key":

View file

@ -1,5 +1,3 @@
"""Stdio MCP peer the proxy spawns; every inbound JSON-RPC line is appended to the record file."""
import sys
from pathlib import Path

View file

@ -1,6 +1,3 @@
"""OAuth 2.1 authorization-server double: metadata, DCR, PKCE authorization code, refresh, client credentials,
token exchange and revocation, every request recorded."""
import base64
import hashlib
import json

View file

@ -17,6 +17,7 @@ from integration._support.mcp import (
McpCaller,
Outcome,
call_tool,
forget_mcp,
mcp_peer,
register_mcp,
tool_calls,
@ -364,7 +365,11 @@ def test_generated_create_edit_grant_revoke_delete_call_keeps_grants_and_tool_li
return
alias: Final = "fleet" + uuid.uuid4().hex[:8]
identity: Final = register_mcp(
self.scenario, upstream, alias, static_headers={"X-Integration-Server": alias}
self.scenario,
upstream,
alias,
cleanup=forget_mcp,
static_headers={"X-Integration-Server": alias},
)
self.servers[alias] = identity

View file

@ -299,7 +299,7 @@ def test_ungranted_key_gets_no_gateway_tools_and_the_peer_is_never_reached(gatew
assert _peer_add_calls(rig.peer) == (), "denied caller reached the peer"
requests: Final = rig.upstream_tools()
assert requests and all(rig.tool not in names for names in requests), requests
assert response.status_code in (200, 400, 401, 403), response.text
assert response.status_code == 200, response.text
@pytest.mark.parametrize("surface", ("chat", "responses", "messages"))

View file

@ -112,7 +112,7 @@ def test_edit_url_moves_calls_to_the_new_peer_without_touching_grants(gateway: G
def test_delete_removes_listing_calls_and_database_row(gateway: Gateway) -> None:
with mcp_peer() as peer, gateway.scenario() as scenario:
alias: Final = "mgmt" + uuid.uuid4().hex[:8]
identity: Final = register_mcp(scenario, peer, alias)
identity: Final = register_mcp(scenario, peer, alias, cleanup=forget_mcp)
key: Final = scenario.key(object_permission={"mcp_servers": [identity]})
name: Final = tool_names(gateway, key, identity)["add"]
delete_mcp(gateway, identity)
@ -286,7 +286,7 @@ def test_access_group_membership_follows_edits(gateway: Gateway) -> None:
def test_peer_worker_observes_create_edit_and_delete_without_restart(gateway: Gateway, peer: Gateway) -> None:
with mcp_peer() as first, mcp_peer() as second, gateway.scenario() as scenario:
alias: Final = "mgmt" + uuid.uuid4().hex[:8]
identity: Final = register_mcp(scenario, first, alias)
identity: Final = register_mcp(scenario, first, alias, cleanup=forget_mcp)
key: Final = scenario.key(object_permission={"mcp_servers": [identity]})
eventually(
lambda: peer.client.get(

View file

@ -63,7 +63,7 @@ def test_unreachable_peer_errors_while_a_healthy_sibling_keeps_serving(gateway:
listing: Final = caller.list_tools(good_id if entry == "rest" else None)
assert listing.error is None, listing.raw
assert {f"{good}-add", "add"} & set(listing.tools), listing.raw
assert not {f"{bad}-add"} & set(listing.tools) or entry != "rest", listing.raw
assert not [tool for tool in listing.tools if tool.startswith(bad)], listing.raw
healthy.drain()
served: Final = _call(caller, f"{good}-add", {"a": 2, "b": 3}, entry, good_id)
assert served.text == "5", served.raw