feat(server): Add support for streamable-http transport type
This commit is contained in:
+1
-1
@@ -46,7 +46,7 @@ services:
|
|||||||
|
|
||||||
mcp:
|
mcp:
|
||||||
build: .
|
build: .
|
||||||
command: ["--host", "0.0.0.0"]
|
command: ["--host", "0.0.0.0", "--transport", "streamable-http"]
|
||||||
ports:
|
ports:
|
||||||
- 8000:8000
|
- 8000:8000
|
||||||
environment:
|
environment:
|
||||||
|
|||||||
@@ -2,7 +2,7 @@ import click
|
|||||||
import logging
|
import logging
|
||||||
import uvicorn
|
import uvicorn
|
||||||
from collections.abc import AsyncIterator
|
from collections.abc import AsyncIterator
|
||||||
from contextlib import asynccontextmanager
|
from contextlib import asynccontextmanager, AsyncExitStack
|
||||||
from dataclasses import dataclass
|
from dataclasses import dataclass
|
||||||
|
|
||||||
from starlette.applications import Starlette
|
from starlette.applications import Starlette
|
||||||
@@ -83,9 +83,19 @@ def get_app(transport: str = "sse", enabled_apps: list[str] | None = None):
|
|||||||
f"Unknown app: {app_name}. Available apps: {list(available_apps.keys())}"
|
f"Unknown app: {app_name}. Available apps: {list(available_apps.keys())}"
|
||||||
)
|
)
|
||||||
|
|
||||||
mcp_app = mcp.sse_app() if transport == "sse" else mcp.streamable_http_app()
|
if transport == "sse":
|
||||||
|
mcp_app = mcp.sse_app()
|
||||||
|
lifespan = None
|
||||||
|
else:
|
||||||
|
mcp_app = mcp.streamable_http_app()
|
||||||
|
|
||||||
app = Starlette(routes=[Mount("/", app=mcp_app)])
|
@asynccontextmanager
|
||||||
|
async def lifespan(app: Starlette):
|
||||||
|
async with AsyncExitStack() as stack:
|
||||||
|
await stack.enter_async_context(mcp.session_manager.run())
|
||||||
|
yield
|
||||||
|
|
||||||
|
app = Starlette(routes=[Mount("/", app=mcp_app)], lifespan=lifespan)
|
||||||
|
|
||||||
return app
|
return app
|
||||||
|
|
||||||
|
|||||||
+10
-10
@@ -6,7 +6,7 @@ from typing import Any, AsyncGenerator
|
|||||||
import pytest
|
import pytest
|
||||||
from httpx import HTTPStatusError
|
from httpx import HTTPStatusError
|
||||||
from mcp import ClientSession
|
from mcp import ClientSession
|
||||||
from mcp.client.sse import sse_client
|
from mcp.client.streamable_http import streamablehttp_client
|
||||||
|
|
||||||
from nextcloud_mcp_server.client import NextcloudClient
|
from nextcloud_mcp_server.client import NextcloudClient
|
||||||
|
|
||||||
@@ -39,18 +39,18 @@ async def nc_client() -> AsyncGenerator[NextcloudClient, Any]:
|
|||||||
await client.close()
|
await client.close()
|
||||||
|
|
||||||
|
|
||||||
@pytest.fixture
|
@pytest.fixture(scope="session")
|
||||||
async def nc_mcp_client() -> AsyncGenerator[ClientSession, Any]:
|
async def nc_mcp_client() -> AsyncGenerator[ClientSession, Any]:
|
||||||
"""
|
"""
|
||||||
Fixture to create an MCP client session for integration tests.
|
Fixture to create an MCP client session for integration tests using streamable-http.
|
||||||
"""
|
"""
|
||||||
logger.info("Creating SSE client")
|
logger.info("Creating Streamable HTTP client")
|
||||||
sse_context = sse_client(url="http://127.0.0.1:8000/sse")
|
streamable_context = streamablehttp_client("http://127.0.0.1:8000/mcp")
|
||||||
session_context = None
|
session_context = None
|
||||||
|
|
||||||
try:
|
try:
|
||||||
read, write = await sse_context.__aenter__()
|
read_stream, write_stream, _ = await streamable_context.__aenter__()
|
||||||
session_context = ClientSession(read, write)
|
session_context = ClientSession(read_stream, write_stream)
|
||||||
session = await session_context.__aenter__()
|
session = await session_context.__aenter__()
|
||||||
await session.initialize()
|
await session.initialize()
|
||||||
logger.info("MCP client session initialized successfully")
|
logger.info("MCP client session initialized successfully")
|
||||||
@@ -71,14 +71,14 @@ async def nc_mcp_client() -> AsyncGenerator[ClientSession, Any]:
|
|||||||
logger.warning(f"Error closing session: {e}")
|
logger.warning(f"Error closing session: {e}")
|
||||||
|
|
||||||
try:
|
try:
|
||||||
await sse_context.__aexit__(None, None, None)
|
await streamable_context.__aexit__(None, None, None)
|
||||||
except RuntimeError as e:
|
except RuntimeError as e:
|
||||||
if "cancel scope" in str(e):
|
if "cancel scope" in str(e):
|
||||||
logger.debug(f"Ignoring cancel scope teardown issue: {e}")
|
logger.debug(f"Ignoring cancel scope teardown issue: {e}")
|
||||||
else:
|
else:
|
||||||
logger.warning(f"Error closing SSE client: {e}")
|
logger.warning(f"Error closing streamable HTTP client: {e}")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning(f"Error closing SSE client: {e}")
|
logger.warning(f"Error closing streamable HTTP client: {e}")
|
||||||
|
|
||||||
|
|
||||||
@pytest.fixture
|
@pytest.fixture
|
||||||
|
|||||||
Reference in New Issue
Block a user