-
Notifications
You must be signed in to change notification settings - Fork 3.9k
Add extensible client and server transport APIs #3517
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Large diffs are not rendered by default.
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,27 @@ | ||
| import secrets | ||
| from dataclasses import dataclass | ||
|
|
||
| from mcp.server import ServerRequestContext | ||
| from mcp.server.mcpserver import MCPServer | ||
| from mcp.server.request_state import RequestStateSecurity | ||
| from mcp.shared.transport import TransportContext | ||
|
|
||
|
|
||
| @dataclass(kw_only=True, frozen=True) | ||
| class VerifiedPeer(TransportContext): | ||
| principal: str | ||
|
|
||
|
|
||
| def principal(ctx: ServerRequestContext) -> str: | ||
| if not isinstance(ctx.transport, VerifiedPeer): | ||
| raise ValueError("Verified transport identity is required") | ||
| return ctx.transport.principal | ||
|
|
||
|
|
||
| mcp = MCPServer( | ||
| "broker-service", | ||
| request_state_security=RequestStateSecurity( | ||
| keys=[secrets.token_bytes(32)], | ||
| bind_principal=principal, | ||
| ), | ||
| ) | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,57 @@ | ||
| from collections.abc import AsyncIterator | ||
| from contextlib import asynccontextmanager | ||
| from dataclasses import dataclass | ||
| from typing import Any | ||
|
|
||
| import anyio | ||
|
|
||
| from mcp import Client | ||
| from mcp.server.mcpserver import Context, MCPServer | ||
| from mcp.server.runtime import ServerRuntime | ||
| from mcp.shared.memory import create_client_server_memory_streams | ||
| from mcp.shared.transport import MessageMetadata, TransportContext, TransportStreams | ||
|
|
||
|
|
||
| @dataclass(kw_only=True, frozen=True) | ||
| class PeerContext(TransportContext): | ||
| peer: str | ||
|
|
||
|
|
||
| server = MCPServer("Custom transport") | ||
|
|
||
|
|
||
| @server.tool() | ||
| async def identify(ctx: Context) -> str: | ||
| transport = ctx.transport | ||
| assert isinstance(transport, PeerContext) | ||
| return transport.peer | ||
|
|
||
|
|
||
| @asynccontextmanager | ||
| async def memory_client(runtime: ServerRuntime[Any], peer: str) -> AsyncIterator[TransportStreams]: | ||
| async with create_client_server_memory_streams() as (client_streams, server_streams): | ||
|
|
||
| @asynccontextmanager | ||
| async def server_transport() -> AsyncIterator[TransportStreams]: | ||
| async with server_streams[0], server_streams[1]: | ||
| yield server_streams | ||
|
|
||
| def build_context(metadata: MessageMetadata) -> PeerContext: | ||
| return PeerContext(kind="memory", can_send_request=True, peer=peer) | ||
|
|
||
| await runtime.connect(server_transport(), transport_builder=build_context) | ||
| yield client_streams | ||
|
|
||
|
|
||
| async def main() -> None: | ||
| async with server.serve(max_connections=10) as runtime: | ||
| async with Client(memory_client(runtime, "alice")) as alice: | ||
| async with Client(memory_client(runtime, "bob")) as bob: | ||
| alice_result = await alice.call_tool("identify") | ||
| bob_result = await bob.call_tool("identify") | ||
| assert alice_result.structured_content == {"result": "alice"} | ||
| assert bob_result.structured_content == {"result": "bob"} | ||
|
|
||
|
|
||
| if __name__ == "__main__": | ||
| anyio.run(main) |
| Original file line number | Diff line number | Diff line change | ||||
|---|---|---|---|---|---|---|
| @@ -0,0 +1,54 @@ | ||||||
| from collections.abc import AsyncIterator | ||||||
| from contextlib import asynccontextmanager | ||||||
| from typing import Any | ||||||
|
|
||||||
| import anyio | ||||||
|
|
||||||
| from mcp import Client | ||||||
| from mcp.server.mcpserver import Context, MCPServer | ||||||
| from mcp.server.runtime import ServerRuntime | ||||||
| from mcp.shared.direct_dispatcher import create_direct_dispatcher_pair | ||||||
| from mcp.shared.dispatcher import Dispatcher | ||||||
| from mcp.shared.transport import DispatcherTransport, TransportContext | ||||||
|
|
||||||
| server = MCPServer("Dispatcher transport") | ||||||
|
|
||||||
|
|
||||||
| @server.tool() | ||||||
| async def greet(name: str, ctx: Context) -> str: | ||||||
| assert ctx.transport is not None | ||||||
| assert not ctx.transport.can_send_request | ||||||
| return f"Hello, {name}!" | ||||||
|
|
||||||
|
|
||||||
| def direct_client(runtime: ServerRuntime[Any]) -> DispatcherTransport: | ||||||
| @asynccontextmanager | ||||||
| async def connection() -> AsyncIterator[Dispatcher[TransportContext]]: | ||||||
| client_dispatcher, server_dispatcher = create_direct_dispatcher_pair() | ||||||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. P2: The example fails on its first tool call because Prompt for AI agents
Suggested change
|
||||||
|
|
||||||
| @asynccontextmanager | ||||||
| async def server_connection() -> AsyncIterator[Dispatcher[TransportContext]]: | ||||||
| try: | ||||||
| yield server_dispatcher | ||||||
| finally: | ||||||
| server_dispatcher.close() | ||||||
|
|
||||||
| try: | ||||||
| await runtime.connect(DispatcherTransport(server_connection())) | ||||||
| yield client_dispatcher | ||||||
| finally: | ||||||
| client_dispatcher.close() | ||||||
| server_dispatcher.close() | ||||||
|
|
||||||
| return DispatcherTransport(connection()) | ||||||
|
|
||||||
|
|
||||||
| async def main() -> None: | ||||||
| async with server.serve() as runtime: | ||||||
| async with Client(direct_client(runtime)) as client: | ||||||
| result = await client.call_tool("greet", {"name": "Alice"}) | ||||||
| assert result.structured_content == {"result": "Hello, Alice!"} | ||||||
|
|
||||||
|
|
||||||
| if __name__ == "__main__": | ||||||
| anyio.run(main) | ||||||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -1,21 +1,5 @@ | ||
| """Transport protocol for MCP clients.""" | ||
|
|
||
| from __future__ import annotations | ||
|
|
||
| from contextlib import AbstractAsyncContextManager | ||
| from typing import Protocol | ||
|
|
||
| from mcp.shared._stream_protocols import ReadStream, WriteStream | ||
| from mcp.shared.message import SessionMessage | ||
| from mcp.shared.transport import ReadStream, Transport, TransportStreams, WriteStream | ||
|
|
||
| __all__ = ["ReadStream", "WriteStream", "Transport", "TransportStreams"] | ||
|
|
||
| TransportStreams = tuple[ReadStream[SessionMessage | Exception], WriteStream[SessionMessage]] | ||
|
|
||
|
|
||
| class Transport(AbstractAsyncContextManager[TransportStreams], Protocol): | ||
| """Protocol for MCP transports. | ||
|
|
||
| A transport is an async context manager that yields read and write streams | ||
| for bidirectional communication with an MCP server. | ||
| """ |
Uh oh!
There was an error while loading. Please reload this page.