mirror of
https://github.com/PrefectHQ/fastmcp.git
synced 2026-08-21 04:54:17 +02:00
Update mcp pin to 1.8
This commit is contained in:
parent
27efa1c4c1
commit
01a348fe3a
4 changed files with 23 additions and 255 deletions
|
|
@ -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"
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
241
src/fastmcp/vendor/streamable_http_manager.py
vendored
241
src/fastmcp/vendor/streamable_http_manager.py
vendored
|
|
@ -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)
|
||||
24
uv.lock
generated
24
uv.lock
generated
|
|
@ -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"
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue