Coverage for agentos/mcp/server.py: 27%
220 statements
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-06 19:15 +0800
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-06 19:15 +0800
1"""MCP Server 实现 — 将 AgentOS 暴露为 MCP Server。
3支持 stdio JSON-RPC 2.0 传输,暴露 LLM 对话、工具调用、Agent 运行等能力。
4其他 MCP 客户端(如 Claude Desktop、Cursor)可直接连接使用。
6用法:
7 agentos mcp-server # 以 stdio 模式启动
8 agentos mcp-server --port 9000 # 以 HTTP SSE 模式启动(可选)
9"""
11from __future__ import annotations
13import asyncio
14import json
15import logging
16import os
17import sys
18from typing import Any
20logger = logging.getLogger(__name__)
22# ── MCP Server 核心 ─────────────────────────
25class MCPServer:
26 """MCP Server — stdio JSON-RPC 2.0 传输。
28 实现 MCP 协议的 server 端,暴露 AgentOS 能力。
29 客户端通过 stdio 发送 JSON-RPC 请求,服务器响应。
31 支持的操作:
32 - initialize: 协议握手,返回 capabilities
33 - tools/list: 列出可用工具
34 - tools/call: 调用工具
35 - resources/list: 列出可用资源
36 - prompts/list: 列出可用提示
37 """
39 def __init__(
40 self,
41 server_info=None,
42 *,
43 name: str = "agentos",
44 version: str = "1.5.2",
45 tools: list | None = None,
46 resources: list | None = None,
47 prompts: list | None = None,
48 ):
49 if server_info is not None:
50 if isinstance(server_info, ServerInfo):
51 self.name = server_info.name
52 self.version = server_info.version
53 else:
54 # backward compat: first arg is name
55 self.name = server_info
56 self.version = version
57 else:
58 self.name = name
59 self.version = version
60 self._tools: dict[str, Any] = {}
61 self._resources: dict[str, MCPResource] = {}
62 self._prompts: dict[str, MCPPromptDef] = {}
63 self._initialized = False
65 for t in tools or []:
66 self.register_tool(t)
67 for r in resources or []:
68 self._resources[r.uri] = r
69 for p in prompts or []:
70 self._prompts[p.name] = p
72 # 内置工具
73 self._register_builtin_tools()
75 @property
76 def info(self) -> ServerInfo:
77 return ServerInfo(name=self.name, version=self.version)
79 @property
80 def tools(self) -> dict[str, Any]:
81 return self._tools
83 def register_tool(self, tool):
84 """注册一个 MCP 工具。兼容 MCPToolDef 和 Tool dataclass。"""
85 if hasattr(tool, "name"):
86 self._tools[tool.name] = tool
87 else:
88 # MCPToolDef compat
89 self._tools[tool.name] = tool
91 async def list_tools(self) -> list:
92 """异步列出所有工具。"""
93 result = []
94 for t in self._tools.values():
95 result.append(
96 {
97 "name": t.name,
98 "description": t.description,
99 "inputSchema": getattr(t, "input_schema", {}),
100 }
101 )
102 return result
104 def run_stdio(self):
105 """以 stdio 模式运行 MCP Server(同步阻塞)。"""
106 asyncio.run(self._run_stdio_async())
108 async def _run_stdio_async(self):
109 """异步 stdio 循环。"""
110 loop = asyncio.get_event_loop()
111 reader = asyncio.StreamReader()
112 protocol = asyncio.StreamReaderProtocol(reader)
113 await loop.connect_read_pipe(lambda: protocol, sys.stdin)
115 writer_transport, writer_protocol = await loop.connect_write_pipe(
116 asyncio.streams.FlowControlMixin, os.fdopen(sys.stdout.fileno(), "wb")
117 )
118 writer = asyncio.StreamWriter(writer_transport, writer_protocol, reader, loop)
120 logger.info(f"MCP Server '{self.name}' v{self.version} started (stdio)")
122 while True:
123 try:
124 line = await reader.readline()
125 if not line:
126 break
127 line_str = line.decode("utf-8").strip()
128 if not line_str:
129 continue
131 try:
132 request = json.loads(line_str)
133 except json.JSONDecodeError:
134 continue
136 response = await self._handle_request(request)
137 if response is not None:
138 payload = json.dumps(response, ensure_ascii=False) + "\n"
139 writer.write(payload.encode("utf-8"))
140 await writer.drain()
141 except Exception as e:
142 logger.error(f"MCP Server error: {e}")
143 break
145 async def _handle_request(self, request: dict) -> dict | None:
146 """处理单个 JSON-RPC 请求。"""
147 req_id = request.get("id")
148 method = request.get("method", "")
149 params = request.get("params", {})
151 # 通知类(无 id),不回复
152 if req_id is None:
153 if method == "notifications/initialized":
154 self._initialized = True
155 return None
157 try:
158 result = await self._dispatch(method, params)
159 return {
160 "jsonrpc": "2.0",
161 "id": req_id,
162 "result": result,
163 }
164 except MCPError as e:
165 return {
166 "jsonrpc": "2.0",
167 "id": req_id,
168 "error": {"code": e.code, "message": e.message},
169 }
170 except Exception as e:
171 return {
172 "jsonrpc": "2.0",
173 "id": req_id,
174 "error": {"code": -32603, "message": str(e)},
175 }
177 async def _dispatch(self, method: str, params: dict) -> Any:
178 """路由到对应的处理器。"""
179 handlers = {
180 "initialize": self._handle_initialize,
181 "tools/list": self._handle_tools_list,
182 "tools/call": self._handle_tools_call,
183 "resources/list": self._handle_resources_list,
184 "resources/read": self._handle_resources_read,
185 "prompts/list": self._handle_prompts_list,
186 "prompts/get": self._handle_prompts_get,
187 }
188 handler = handlers.get(method)
189 if handler is None:
190 raise MCPError(-32601, f"Method not found: {method}")
191 return await handler(params)
193 # ── MCP 协议方法 ──────────────────────
195 async def _handle_initialize(self, params: dict) -> dict:
196 return {
197 "protocolVersion": "2024-11-05",
198 "serverInfo": {
199 "name": self.name,
200 "version": self.version,
201 },
202 "capabilities": {
203 "tools": {"listChanged": False},
204 "resources": {"subscribe": False, "listChanged": False},
205 "prompts": {"listChanged": False},
206 },
207 }
209 async def _handle_tools_list(self, params: dict) -> dict:
210 tools = []
211 for t in self._tools.values():
212 tools.append(
213 {
214 "name": t.name,
215 "description": t.description,
216 "inputSchema": t.input_schema,
217 }
218 )
219 return {"tools": tools}
221 async def _handle_tools_call(self, params: dict) -> dict:
222 tool_name = params.get("name", "")
223 arguments = params.get("arguments", {})
224 tool = self._tools.get(tool_name)
225 if tool is None:
226 raise MCPError(-32602, f"Unknown tool: {tool_name}")
228 try:
229 result = (
230 tool.handler(arguments)
231 if not asyncio.iscoroutinefunction(tool.handler)
232 else await tool.handler(arguments)
233 )
234 return {
235 "content": [
236 {"type": "text", "text": str(result) if not isinstance(result, str) else result}
237 ]
238 }
239 except Exception as e:
240 return {
241 "content": [{"type": "text", "text": f"Error: {str(e)}"}],
242 "isError": True,
243 }
245 async def _handle_resources_list(self, params: dict) -> dict:
246 resources = []
247 for r in self._resources.values():
248 resources.append(
249 {
250 "uri": r.uri,
251 "name": r.name,
252 "description": r.description,
253 "mimeType": r.mime_type,
254 }
255 )
256 return {"resources": resources}
258 async def _handle_resources_read(self, params: dict) -> dict:
259 uri = params.get("uri", "")
260 r = self._resources.get(uri)
261 if r is None:
262 raise MCPError(-32602, f"Unknown resource: {uri}")
263 text = r.content() if callable(r.content) else r.content
264 return {"contents": [{"uri": uri, "mimeType": r.mime_type, "text": str(text)}]}
266 async def _handle_prompts_list(self, params: dict) -> dict:
267 prompts = []
268 for p in self._prompts.values():
269 prompts.append(
270 {
271 "name": p.name,
272 "description": p.description,
273 "arguments": p.arguments,
274 }
275 )
276 return {"prompts": prompts}
278 async def _handle_prompts_get(self, params: dict) -> dict:
279 prompt_name = params.get("name", "")
280 prompt_args = params.get("arguments", {})
281 p = self._prompts.get(prompt_name)
282 if p is None:
283 raise MCPError(-32602, f"Unknown prompt: {prompt_name}")
284 template = p.template(prompt_args) if callable(p.template) else p.template
285 return {
286 "description": p.description,
287 "messages": [{"role": "user", "content": {"type": "text", "text": template}}],
288 }
290 # ── 内置工具 ──────────────────────────
292 def _register_builtin_tools(self):
293 """注册 AgentOS 内置 MCP 工具。"""
295 self.register_tool(
296 MCPToolDef(
297 name="agentos_chat",
298 description="使用 AgentOS LLM 进行对话(支持 OpenAI/DeepSeek/Anthropic/Claude/Ollama)",
299 input_schema={
300 "type": "object",
301 "properties": {
302 "messages": {
303 "type": "array",
304 "description": "对话消息列表",
305 "items": {
306 "type": "object",
307 "properties": {
308 "role": {
309 "type": "string",
310 "enum": ["system", "user", "assistant"],
311 },
312 "content": {"type": "string"},
313 },
314 "required": ["role", "content"],
315 },
316 },
317 "model": {"type": "string", "description": "模型名称,默认从配置读取"},
318 "temperature": {"type": "number", "description": "温度参数(0-2)"},
319 "max_tokens": {"type": "integer", "description": "最大输出 token 数"},
320 },
321 "required": ["messages"],
322 },
323 handler=self._tool_agentos_chat,
324 )
325 )
327 self.register_tool(
328 MCPToolDef(
329 name="agentos_list_tools",
330 description="列出 AgentOS 中所有可用的工具(含 MCP 工具)",
331 input_schema={
332 "type": "object",
333 "properties": {
334 "format": {
335 "type": "string",
336 "enum": ["openai", "anthropic"],
337 "description": "输出格式",
338 },
339 },
340 },
341 handler=self._tool_list_tools,
342 )
343 )
345 self.register_tool(
346 MCPToolDef(
347 name="agentos_version",
348 description="获取 AgentOS 版本信息",
349 input_schema={"type": "object", "properties": {}},
350 handler=self._tool_version,
351 )
352 )
354 async def _tool_agentos_chat(self, args: dict) -> str:
355 """调用 AgentOS LLM 对话。"""
356 try:
357 from agentos.llm import LLMClient, LLMMessage
358 except ImportError:
359 return "Error: AgentOS LLM 模块不可用。请确认已安装 nexus-agentos。"
361 messages_raw = args.get("messages", [])
362 model = args.get("model")
363 temperature = args.get("temperature", 0.7)
364 max_tokens = args.get("max_tokens", 4096)
366 messages = [LLMMessage(role=m["role"], content=m["content"]) for m in messages_raw]
368 client = LLMClient(model=model)
369 response = await client.chat(messages, temperature=temperature, max_tokens=max_tokens)
370 return response.content
372 def _tool_list_tools(self, args: dict) -> str:
373 """列出可用工具。"""
374 fmt = args.get("format", "openai")
375 tools_list = []
376 for name, tool in self._tools.items():
377 if fmt == "openai":
378 tools_list.append(
379 {
380 "type": "function",
381 "function": {
382 "name": name,
383 "description": tool.description,
384 "parameters": tool.input_schema,
385 },
386 }
387 )
388 else:
389 tools_list.append(
390 {
391 "name": name,
392 "description": tool.description,
393 "input_schema": tool.input_schema,
394 }
395 )
396 return json.dumps(tools_list, ensure_ascii=False, indent=2)
398 def _tool_version(self, args: dict) -> str:
399 """返回版本信息。"""
400 try:
401 from agentos import __version__
402 except ImportError:
403 __version__ = self.version
404 return json.dumps(
405 {
406 "name": self.name,
407 "version": self.version,
408 "agentos_version": __version__,
409 "tools_count": len(self._tools),
410 },
411 ensure_ascii=False,
412 )
415# ── 数据结构 ───────────────────────────────
418class MCPToolDef:
419 """MCP 工具定义。"""
421 def __init__(
422 self,
423 name: str,
424 description: str,
425 input_schema: dict,
426 handler,
427 ):
428 self.name = name
429 self.description = description
430 self.input_schema = input_schema
431 self.handler = handler
434class MCPResource:
435 """MCP 资源定义。"""
437 def __init__(
438 self,
439 uri: str,
440 name: str = "",
441 description: str = "",
442 mime_type: str = "text/plain",
443 content: Any = "",
444 ):
445 self.uri = uri
446 self.name = name
447 self.description = description
448 self.mime_type = mime_type
449 self.content = content
452class MCPPromptDef:
453 """MCP 提示模板定义。"""
455 def __init__(
456 self,
457 name: str,
458 description: str = "",
459 arguments: list = None,
460 template: Any = "",
461 ):
462 self.name = name
463 self.description = description
464 self.arguments = arguments or []
465 self.template = template
468class MCPError(Exception):
469 """MCP 协议错误(与服务端共用异常类)。"""
471 def __init__(self, code: int, message: str):
472 self.code = code
473 self.message = message
474 super().__init__(f"MCP Error [{code}]: {message}")
477# ── 便捷函数 ───────────────────────────────
480def create_default_server() -> MCPServer:
481 """创建预配置了 AgentOS 内置工具的 MCP Server。"""
482 return MCPServer(
483 name="agentos",
484 version="1.5.2",
485 )
488def start_mcp_server(port: int = 0):
489 """启动 MCP Server。
491 Args:
492 port: 0 表示 stdio 模式,>0 表示 HTTP SSE 模式(暂未实现)。
493 """
494 if port == 0:
495 server = create_default_server()
496 server.run_stdio()
497 else:
498 print("MCP HTTP SSE 模式暂未实现。请使用 stdio 模式(port=0)。")
499 sys.exit(1)
502# ── ServerInfo & Tool (test compatibility) ──
503from dataclasses import dataclass, field # noqa: E402
506@dataclass
507class ServerInfo:
508 name: str
509 version: str
510 description: str = ""
513@dataclass
514class Tool:
515 name: str
516 description: str
517 input_schema: dict
518 call: callable = field(default=lambda params: None)
521@dataclass
522class AgentCard:
523 agent_id: str
524 name: str
525 version: str
526 capabilities: list = field(default_factory=list)
527 endpoint: str = ""