diff --git a/pyproject.toml b/pyproject.toml index df95c1a3a..dea5d446f 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -7,9 +7,7 @@ dependencies = [ "python-dotenv>=1.1.0", "exceptiongroup>=1.2.2", "httpx>=0.28.1", - # "mcp>=1.7.1,<2.0.0", - # use git commit until 1.7.2 is released - "mcp", + "mcp>=1.8.0,<2.0.0", "openapi-pydantic>=0.5.1", "rich>=13.9.4", "typer>=0.15.2", @@ -70,11 +68,6 @@ Documentation = "https://gofastmcp.com" requires = ["hatchling", "uv-dynamic-versioning>=0.7.0"] build-backend = "hatchling.build" -[tool.uv] - -[tool.uv.sources] -mcp = { git = "https://github.com/modelcontextprotocol/python-sdk", rev = "a027d75f609000378522c5873c2a16aa1963d487" } - [tool.hatch.version] source = "uv-dynamic-versioning" diff --git a/src/fastmcp/server/http.py b/src/fastmcp/server/http.py index 3db6138f5..385ea5280 100644 --- a/src/fastmcp/server/http.py +++ b/src/fastmcp/server/http.py @@ -14,6 +14,7 @@ from mcp.server.auth.provider import OAuthAuthorizationServerProvider from mcp.server.auth.routes import create_auth_routes from mcp.server.auth.settings import AuthSettings from mcp.server.sse import SseServerTransport +from mcp.server.streamable_http_manager import StreamableHTTPSessionManager from starlette.applications import Starlette from starlette.middleware import Middleware from starlette.middleware.authentication import AuthenticationMiddleware @@ -24,9 +25,6 @@ from starlette.types import Receive, Scope, Send from fastmcp.utilities.logging import get_logger -# This import is vendored until it is finalized in the upstream SDK -from fastmcp.vendor.streamable_http_manager import StreamableHTTPSessionManager - if TYPE_CHECKING: from fastmcp.server.server import FastMCP diff --git a/src/fastmcp/vendor/streamable_http_manager.py b/src/fastmcp/vendor/streamable_http_manager.py deleted file mode 100644 index 37cdc23ec..000000000 --- a/src/fastmcp/vendor/streamable_http_manager.py +++ /dev/null @@ -1,241 +0,0 @@ -"""StreamableHTTP Session Manager for MCP servers.""" - -# follows https://github.com/modelcontextprotocol/python-sdk/blob/ihrpr/shttp/src/mcp/server/streamable_http_manager.py -# and can be removed once that spec is finalized - -from __future__ import annotations - -import contextlib -import logging -from collections.abc import AsyncIterator -from http import HTTPStatus -from typing import Any -from uuid import uuid4 - -import anyio -from anyio.abc import TaskStatus -from mcp.server.lowlevel.server import Server as MCPServer -from mcp.server.streamable_http import ( - MCP_SESSION_ID_HEADER, - EventStore, - StreamableHTTPServerTransport, -) -from starlette.requests import Request -from starlette.responses import Response -from starlette.types import Receive, Scope, Send - -logger = logging.getLogger(__name__) - - -class StreamableHTTPSessionManager: - """ - Manages StreamableHTTP sessions with optional resumability via event store. - - This class abstracts away the complexity of session management, event storage, - and request handling for StreamableHTTP transports. It handles: - - 1. Session tracking for clients - 2. Resumability via an optional event store - 3. Connection management and lifecycle - 4. Request handling and transport setup - - Args: - app: The MCP server instance - event_store: Optional event store for resumability support. - If provided, enables resumable connections where clients - can reconnect and receive missed events. - If None, sessions are still tracked but not resumable. - json_response: Whether to use JSON responses instead of SSE streams - stateless: If True, creates a completely fresh transport for each request - with no session tracking or state persistence between requests. - - """ - - def __init__( - self, - app: MCPServer[Any], - event_store: EventStore | None = None, - json_response: bool = False, - stateless: bool = False, - ): - self.app = app - self.event_store = event_store - self.json_response = json_response - self.stateless = stateless - - # Session tracking (only used if not stateless) - self._session_creation_lock = anyio.Lock() - self._server_instances: dict[str, StreamableHTTPServerTransport] = {} - - # The task group will be set during lifespan - self._task_group = None - - @contextlib.asynccontextmanager - async def run(self) -> AsyncIterator[None]: - """ - Run the session manager with proper lifecycle management. - - This creates and manages the task group for all session operations. - - Use this in the lifespan context manager of your Starlette app: - - @contextlib.asynccontextmanager - async def lifespan(app: Starlette) -> AsyncIterator[None]: - async with session_manager.run(): - yield - """ - async with anyio.create_task_group() as tg: - # Store the task group for later use - self._task_group = tg - logger.info("StreamableHTTP session manager started") - try: - yield # Let the application run - finally: - logger.info("StreamableHTTP session manager shutting down") - # Cancel task group to stop all spawned tasks - tg.cancel_scope.cancel() - self._task_group = None - # Clear any remaining server instances - self._server_instances.clear() - - async def handle_request( - self, - scope: Scope, - receive: Receive, - send: Send, - ) -> None: - """ - Process ASGI request with proper session handling and transport setup. - - Dispatches to the appropriate handler based on stateless mode. - - Args: - scope: ASGI scope - receive: ASGI receive function - send: ASGI send function - """ - if self._task_group is None: - raise RuntimeError( - "Task group is not initialized. Make sure to use the run()." - ) - - # Dispatch to the appropriate handler - if self.stateless: - await self._handle_stateless_request(scope, receive, send) - else: - await self._handle_stateful_request(scope, receive, send) - - async def _handle_stateless_request( - self, - scope: Scope, - receive: Receive, - send: Send, - ) -> None: - """ - Process request in stateless mode - creating a new transport for each request. - - Args: - scope: ASGI scope - receive: ASGI receive function - send: ASGI send function - """ - logger.debug("Stateless mode: Creating new transport for this request") - # No session ID needed in stateless mode - http_transport = StreamableHTTPServerTransport( - mcp_session_id=None, # No session tracking in stateless mode - is_json_response_enabled=self.json_response, - event_store=None, # No event store in stateless mode - ) - - # Start server in a new task - async def run_stateless_server( - *, task_status: TaskStatus[None] = anyio.TASK_STATUS_IGNORED - ): - async with http_transport.connect() as streams: - read_stream, write_stream = streams - task_status.started() - await self.app.run( - read_stream, - write_stream, - self.app.create_initialization_options(), - stateless=True, - ) - - # Assert task group is not None for type checking - assert self._task_group is not None - # Start the server task - await self._task_group.start(run_stateless_server) - - # Handle the HTTP request and return the response - await http_transport.handle_request(scope, receive, send) - - async def _handle_stateful_request( - self, - scope: Scope, - receive: Receive, - send: Send, - ) -> None: - """ - Process request in stateful mode - maintaining session state between requests. - - Args: - scope: ASGI scope - receive: ASGI receive function - send: ASGI send function - """ - request = Request(scope, receive) - request_mcp_session_id = request.headers.get(MCP_SESSION_ID_HEADER) - - # Existing session case - if ( - request_mcp_session_id is not None - and request_mcp_session_id in self._server_instances - ): - transport = self._server_instances[request_mcp_session_id] - logger.debug("Session already exists, handling request directly") - await transport.handle_request(scope, receive, send) - return - - if request_mcp_session_id is None: - # New session case - logger.debug("Creating new transport") - async with self._session_creation_lock: - new_session_id = uuid4().hex - http_transport = StreamableHTTPServerTransport( - mcp_session_id=new_session_id, - is_json_response_enabled=self.json_response, - event_store=self.event_store, # May be None (no resumability) - ) - - assert http_transport.mcp_session_id is not None - self._server_instances[http_transport.mcp_session_id] = http_transport - logger.info(f"Created new transport with session ID: {new_session_id}") - - # Define the server runner - async def run_server( - *, task_status: TaskStatus[None] = anyio.TASK_STATUS_IGNORED - ) -> None: - async with http_transport.connect() as streams: - read_stream, write_stream = streams - task_status.started() - await self.app.run( - read_stream, - write_stream, - self.app.create_initialization_options(), - stateless=False, # Stateful mode - ) - - # Assert task group is not None for type checking - assert self._task_group is not None - # Start the server task - await self._task_group.start(run_server) - - # Handle the HTTP request and return the response - await http_transport.handle_request(scope, receive, send) - else: - # Invalid session ID - response = Response( - "Bad Request: No valid session ID provided", - status_code=HTTPStatus.BAD_REQUEST, - ) - await response(scope, receive, send) diff --git a/uv.lock b/uv.lock index df71b5da9..18220c0af 100644 --- a/uv.lock +++ b/uv.lock @@ -331,6 +331,7 @@ dev = [ { name = "pytest-cov" }, { name = "pytest-flakefinder" }, { name = "pytest-report" }, + { name = "pytest-timeout" }, { name = "pytest-xdist" }, { name = "ruff" }, ] @@ -339,7 +340,7 @@ dev = [ requires-dist = [ { name = "exceptiongroup", specifier = ">=1.2.2" }, { name = "httpx", specifier = ">=0.28.1" }, - { name = "mcp", git = "https://github.com/modelcontextprotocol/python-sdk.git?rev=a027d75f609000378522c5873c2a16aa1963d487" }, + { name = "mcp", specifier = ">=1.8.0,<2.0.0" }, { name = "openapi-pydantic", specifier = ">=0.5.1" }, { name = "python-dotenv", specifier = ">=1.1.0" }, { name = "rich", specifier = ">=13.9.4" }, @@ -361,6 +362,7 @@ dev = [ { name = "pytest-cov", specifier = ">=6.1.1" }, { name = "pytest-flakefinder" }, { name = "pytest-report", specifier = ">=0.2.1" }, + { name = "pytest-timeout", specifier = ">=2.4.0" }, { name = "pytest-xdist", specifier = ">=3.6.1" }, { name = "ruff" }, ] @@ -571,8 +573,8 @@ wheels = [ [[package]] name = "mcp" -version = "1.7.1.dev17+a027d75" -source = { git = "https://github.com/modelcontextprotocol/python-sdk.git?rev=a027d75f609000378522c5873c2a16aa1963d487#a027d75f609000378522c5873c2a16aa1963d487" } +version = "1.8.0" +source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "anyio" }, { name = "httpx" }, @@ -584,6 +586,10 @@ dependencies = [ { name = "starlette" }, { name = "uvicorn", marker = "sys_platform != 'emscripten'" }, ] +sdist = { url = "https://files.pythonhosted.org/packages/ff/97/0a3e08559557b0ac5799f9fb535fbe5a4e4dcdd66ce9d32e7a74b4d0534d/mcp-1.8.0.tar.gz", hash = "sha256:263dfb700540b726c093f0c3e043f66aded0730d0b51f04eb0a3eb90055fe49b", size = 264641, upload-time = "2025-05-08T20:09:06.255Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/b2/b2/4ac3bd17b1fdd65658f18de4eb0c703517ee0b483dc5f56467802a9197e0/mcp-1.8.0-py3-none-any.whl", hash = "sha256:889d9d3b4f12b7da59e7a3933a0acadae1fce498bfcd220defb590aa291a1334", size = 119544, upload-time = "2025-05-08T20:09:04.458Z" }, +] [[package]] name = "mdurl" @@ -956,6 +962,18 @@ dependencies = [ ] sdist = { url = "https://files.pythonhosted.org/packages/3b/82/e141da085de0b6dac3f047ae009e136bcedbcfca4ada082a55359d6f735e/pytest-report-0.2.1.tar.gz", hash = "sha256:d382e8db4c52a815d39dae5f21ee5edc0da3ae8ec19a22e55e9be5c60714a39d", size = 3517, upload-time = "2016-05-11T02:08:04.665Z" } +[[package]] +name = "pytest-timeout" +version = "2.4.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "pytest" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/ac/82/4c9ecabab13363e72d880f2fb504c5f750433b2b6f16e99f4ec21ada284c/pytest_timeout-2.4.0.tar.gz", hash = "sha256:7e68e90b01f9eff71332b25001f85c75495fc4e3a836701876183c4bcfd0540a", size = 17973, upload-time = "2025-05-05T19:44:34.99Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/fa/b6/3127540ecdf1464a00e5a01ee60a1b09175f6913f0644ac748494d9c4b21/pytest_timeout-2.4.0-py3-none-any.whl", hash = "sha256:c42667e5cdadb151aeb5b26d114aff6bdf5a907f176a007a30b940d3d865b5c2", size = 14382, upload-time = "2025-05-05T19:44:33.502Z" }, +] + [[package]] name = "pytest-xdist" version = "3.6.1"