-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtest_mcp.py
More file actions
280 lines (226 loc) · 12.1 KB
/
Copy pathtest_mcp.py
File metadata and controls
280 lines (226 loc) · 12.1 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
from __future__ import annotations
from pathlib import Path
from typing import Any, cast
import anyio
import httpx
import pytest
from mcp.client.session import ClientSession
from mcp.client.streamable_http import streamablehttp_client
from mcp.shared.memory import create_connected_server_and_client_session
from vesper import a2a, mcp_server
from vesper.a2a import create_http_app
from vesper.config import Settings
from vesper.service import CiderAgentService
from vesper.storage import PreferenceStore
from tests.conftest import StubResolver, StubRpcClient
def _make_service(tmp_path: Path) -> tuple[Settings, CiderAgentService]:
settings = Settings(
http_host="127.0.0.1",
http_port=8766,
public_base_url="http://127.0.0.1:8766",
cider_base_url="http://localhost:10767",
cider_api_token="secret-token",
default_search_source="catalog",
resolver_backend="fallback",
resolver_base_url="https://api.openai.com/v1",
resolver_model=None,
resolver_api_key=None,
resolver_include_reasoning=False,
resolver_include_raw_output=False,
request_timeout_seconds=10.0,
verify_tls=True,
log_level="INFO",
database_path=tmp_path / "mcp-test.db",
config_path=None,
)
service = CiderAgentService(
settings,
rpc_client=StubRpcClient(),
preference_store=PreferenceStore(settings.database_path),
resolver=StubResolver(),
)
return settings, service
def test_mcp_lists_only_transport_and_ask_tools(monkeypatch, tmp_path: Path) -> None:
settings, service = _make_service(tmp_path)
monkeypatch.setattr(mcp_server, "get_settings", lambda: settings)
monkeypatch.setattr(mcp_server, "get_service", lambda: service)
async def _exercise() -> None:
async with create_connected_server_and_client_session(mcp_server.create_mcp_server()) as session:
tools = await session.list_tools()
resources = await session.list_resources()
assert [tool.name for tool in tools.tools] == ["play", "pause", "next", "previous", "ask"]
assert resources.resources == []
anyio.run(_exercise)
def test_mcp_transport_tools_delegate_to_service(monkeypatch, tmp_path: Path) -> None:
settings, service = _make_service(tmp_path)
monkeypatch.setattr(mcp_server, "get_settings", lambda: settings)
monkeypatch.setattr(mcp_server, "get_service", lambda: service)
async def _exercise() -> None:
async with create_connected_server_and_client_session(mcp_server.create_mcp_server()) as session:
play = await session.call_tool("play", {})
pause = await session.call_tool("pause", {})
next_track = await session.call_tool("next", {})
previous = await session.call_tool("previous", {})
assert play.structuredContent == {"status": "ok", "result": {"path": "/play", "body": None}}
assert pause.structuredContent == {"status": "ok", "result": {"path": "/pause", "body": None}}
next_content = next_track.structuredContent
assert next_content is not None
assert next_content["result"]["path"] == "/next"
prev_content = previous.structuredContent
assert prev_content is not None
assert prev_content["result"]["path"] == "/previous"
anyio.run(_exercise)
def test_mcp_ask_returns_text_request_payload(monkeypatch, tmp_path: Path) -> None:
settings, service = _make_service(tmp_path)
monkeypatch.setattr(mcp_server, "get_settings", lambda: settings)
monkeypatch.setattr(mcp_server, "get_service", lambda: service)
async def _exercise() -> None:
async with create_connected_server_and_client_session(mcp_server.create_mcp_server()) as session:
result = await session.call_tool("ask", {"text": "play some kep1er"})
assert result.isError is False
content = result.structuredContent
assert content is not None
assert content["status"] == "ok"
assert content["input"] == "play some kep1er"
assert content["execution"]["action"] == "search"
anyio.run(_exercise)
def test_mcp_ask_rejects_empty_text(monkeypatch, tmp_path: Path) -> None:
settings, service = _make_service(tmp_path)
monkeypatch.setattr(mcp_server, "get_settings", lambda: settings)
monkeypatch.setattr(mcp_server, "get_service", lambda: service)
async def _exercise() -> None:
async with create_connected_server_and_client_session(mcp_server.create_mcp_server()) as session:
result = await session.call_tool("ask", {"text": ""})
assert result.isError is True
assert "text cannot be empty" in cast(Any, result.content[0]).text
anyio.run(_exercise)
def test_http_app_enables_requested_transports(monkeypatch, tmp_path: Path) -> None:
settings, service = _make_service(tmp_path)
monkeypatch.setattr(a2a, "get_settings", lambda: settings)
monkeypatch.setattr(a2a, "get_service", lambda: service)
monkeypatch.setattr(mcp_server, "get_settings", lambda: settings)
monkeypatch.setattr(mcp_server, "get_service", lambda: service)
async def _exercise() -> None:
a2a_only = create_http_app(include_a2a=True)
mcp_only = create_http_app(include_mcp=True)
both = create_http_app(include_a2a=True, include_mcp=True)
async with a2a_only.router.lifespan_context(a2a_only):
transport_without = httpx.ASGITransport(app=a2a_only)
async with httpx.AsyncClient(
transport=transport_without,
base_url="http://127.0.0.1:8766",
follow_redirects=True,
) as client:
missing = await client.post("/mcp", json={})
agent_card = await client.get("/.well-known/agent-card")
assert missing.status_code == 404
assert agent_card.status_code == 200
async with mcp_only.router.lifespan_context(mcp_only):
transport_with = httpx.ASGITransport(app=mcp_only)
async with httpx.AsyncClient(transport=transport_with, base_url="http://127.0.0.1:8766") as client:
health = await client.get("/healthz")
mcp_response = await client.post("/mcp", json={})
agent_card = await client.get("/.well-known/agent-card")
assert health.status_code == 200
assert mcp_response.status_code in {200, 400, 406}
assert agent_card.status_code == 404
async with both.router.lifespan_context(both):
transport_both = httpx.ASGITransport(app=both)
async with httpx.AsyncClient(transport=transport_both, base_url="http://127.0.0.1:8766") as client:
agent_card = await client.get("/.well-known/agent-card")
mcp_response = await client.post("/mcp", json={})
assert agent_card.status_code == 200
assert mcp_response.status_code in {200, 400, 406}
anyio.run(_exercise)
def test_streamable_http_client_can_call_mounted_mcp(monkeypatch, tmp_path: Path) -> None:
settings, service = _make_service(tmp_path)
monkeypatch.setattr(a2a, "get_settings", lambda: settings)
monkeypatch.setattr(a2a, "get_service", lambda: service)
monkeypatch.setattr(mcp_server, "get_settings", lambda: settings)
monkeypatch.setattr(mcp_server, "get_service", lambda: service)
app = create_http_app(include_mcp=True)
def client_factory(headers=None, timeout=None, auth=None):
return httpx.AsyncClient(
transport=httpx.ASGITransport(app=app),
base_url="http://127.0.0.1:8766",
headers=headers,
timeout=timeout,
auth=auth,
follow_redirects=True,
)
async def _exercise() -> None:
async with app.router.lifespan_context(app):
async with streamablehttp_client("http://127.0.0.1:8766/mcp", httpx_client_factory=client_factory) as streams:
read_stream, write_stream, _ = streams
async with ClientSession(read_stream, write_stream) as session:
await session.initialize()
tools = await session.list_tools()
result = await session.call_tool("play", {})
assert [tool.name for tool in tools.tools] == ["play", "pause", "next", "previous", "ask"]
play_content = result.structuredContent
assert play_content is not None
assert play_content["result"]["path"] == "/play"
anyio.run(_exercise)
def test_mounted_mcp_requests_do_not_stop_parent_session_worker(monkeypatch, tmp_path: Path) -> None:
settings, service = _make_service(tmp_path)
monkeypatch.setattr(a2a, "get_settings", lambda: settings)
monkeypatch.setattr(a2a, "get_service", lambda: service)
monkeypatch.setattr(mcp_server, "get_settings", lambda: settings)
monkeypatch.setattr(mcp_server, "get_service", lambda: service)
app = create_http_app(include_mcp=True)
def client_factory(headers=None, timeout=None, auth=None):
return httpx.AsyncClient(
transport=httpx.ASGITransport(app=app),
base_url="http://127.0.0.1:8766",
headers=headers,
timeout=timeout,
auth=auth,
follow_redirects=True,
)
async def _exercise() -> None:
async with app.router.lifespan_context(app):
worker = service._session_worker_thread
assert worker is not None
assert worker.is_alive()
async with streamablehttp_client("http://127.0.0.1:8766/mcp", httpx_client_factory=client_factory) as streams:
read_stream, write_stream, _ = streams
async with ClientSession(read_stream, write_stream) as session:
await session.initialize()
await session.list_tools()
assert service._session_worker_thread is worker
assert worker.is_alive()
anyio.run(_exercise)
@pytest.mark.parametrize("include_a2a, include_mcp", [(True, False), (False, True), (True, True)])
def test_worker_does_not_run_after_http_lifespan_shutdown(monkeypatch, tmp_path: Path, include_a2a: bool, include_mcp: bool) -> None:
"""The worker must be stopped once the HTTP app lifespan exits (#45).
All three HTTP transport combinations delegate worker start/stop to the
single Application.worker_lifespan, so the worker must not be alive after
the lifespan context closes.
"""
settings, service = _make_service(tmp_path)
monkeypatch.setattr(a2a, "get_settings", lambda: settings)
monkeypatch.setattr(a2a, "get_service", lambda: service)
monkeypatch.setattr(mcp_server, "get_settings", lambda: settings)
monkeypatch.setattr(mcp_server, "get_service", lambda: service)
app = create_http_app(include_a2a=include_a2a, include_mcp=include_mcp)
async def _exercise() -> None:
async with app.router.lifespan_context(app):
assert service._session_worker_thread is not None
assert service._session_worker_thread.is_alive()
# After the lifespan exits, the worker must have been stopped.
assert service._session_worker_thread is None
anyio.run(_exercise)
def test_worker_does_not_run_after_standalone_mcp_lifespan_shutdown(monkeypatch, tmp_path: Path) -> None:
"""Standalone MCP (e.g. stdio) must stop the worker on lifespan exit (#45)."""
settings, service = _make_service(tmp_path)
monkeypatch.setattr(mcp_server, "get_settings", lambda: settings)
monkeypatch.setattr(mcp_server, "get_service", lambda: service)
server = mcp_server.create_mcp_server()
async def _exercise() -> None:
async with create_connected_server_and_client_session(server) as session:
await session.list_tools()
assert service._session_worker_thread is not None
assert service._session_worker_thread.is_alive()
# After the standalone MCP session closes, the worker must have stopped.
assert service._session_worker_thread is None
anyio.run(_exercise)