Coverage for agentos/mcp/server.py: 27%
220 statements
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-06 10:59 +0800
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-06 10:59 +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, Dict, List, Optional
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: Optional[List] = None,
46 resources: Optional[List] = None,
47 prompts: Optional[List] = 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 "name": t.name,
97 "description": t.description,
98 "inputSchema": getattr(t, 'input_schema', {}),
99 })
100 return result
102 def run_stdio(self):
103 """以 stdio 模式运行 MCP Server(同步阻塞)。"""
104 asyncio.run(self._run_stdio_async())
106 async def _run_stdio_async(self):
107 """异步 stdio 循环。"""
108 loop = asyncio.get_event_loop()
109 reader = asyncio.StreamReader()
110 protocol = asyncio.StreamReaderProtocol(reader)
111 await loop.connect_read_pipe(lambda: protocol, sys.stdin)
113 writer_transport, writer_protocol = await loop.connect_write_pipe(
114 asyncio.streams.FlowControlMixin, os.fdopen(sys.stdout.fileno(), "wb")
115 )
116 writer = asyncio.StreamWriter(writer_transport, writer_protocol, reader, loop)
118 logger.info(f"MCP Server '{self.name}' v{self.version} started (stdio)")
120 while True:
121 try:
122 line = await reader.readline()
123 if not line:
124 break
125 line_str = line.decode("utf-8").strip()
126 if not line_str:
127 continue
129 try:
130 request = json.loads(line_str)
131 except json.JSONDecodeError:
132 continue
134 response = await self._handle_request(request)
135 if response is not None:
136 payload = json.dumps(response, ensure_ascii=False) + "\n"
137 writer.write(payload.encode("utf-8"))
138 await writer.drain()
139 except Exception as e:
140 logger.error(f"MCP Server error: {e}")
141 break
143 async def _handle_request(self, request: dict) -> dict | None:
144 """处理单个 JSON-RPC 请求。"""
145 req_id = request.get("id")
146 method = request.get("method", "")
147 params = request.get("params", {})
149 # 通知类(无 id),不回复
150 if req_id is None:
151 if method == "notifications/initialized":
152 self._initialized = True
153 return None
155 try:
156 result = await self._dispatch(method, params)
157 return {
158 "jsonrpc": "2.0",
159 "id": req_id,
160 "result": result,
161 }
162 except MCPError as e:
163 return {
164 "jsonrpc": "2.0",
165 "id": req_id,
166 "error": {"code": e.code, "message": e.message},
167 }
168 except Exception as e:
169 return {
170 "jsonrpc": "2.0",
171 "id": req_id,
172 "error": {"code": -32603, "message": str(e)},
173 }
175 async def _dispatch(self, method: str, params: dict) -> Any:
176 """路由到对应的处理器。"""
177 handlers = {
178 "initialize": self._handle_initialize,
179 "tools/list": self._handle_tools_list,
180 "tools/call": self._handle_tools_call,
181 "resources/list": self._handle_resources_list,
182 "resources/read": self._handle_resources_read,
183 "prompts/list": self._handle_prompts_list,
184 "prompts/get": self._handle_prompts_get,
185 }
186 handler = handlers.get(method)
187 if handler is None:
188 raise MCPError(-32601, f"Method not found: {method}")
189 return await handler(params)
191 # ── MCP 协议方法 ──────────────────────
193 async def _handle_initialize(self, params: dict) -> dict:
194 return {
195 "protocolVersion": "2024-11-05",
196 "serverInfo": {
197 "name": self.name,
198 "version": self.version,
199 },
200 "capabilities": {
201 "tools": {"listChanged": False},
202 "resources": {"subscribe": False, "listChanged": False},
203 "prompts": {"listChanged": False},
204 },
205 }
207 async def _handle_tools_list(self, params: dict) -> dict:
208 tools = []
209 for t in self._tools.values():
210 tools.append({
211 "name": t.name,
212 "description": t.description,
213 "inputSchema": t.input_schema,
214 })
215 return {"tools": tools}
217 async def _handle_tools_call(self, params: dict) -> dict:
218 tool_name = params.get("name", "")
219 arguments = params.get("arguments", {})
220 tool = self._tools.get(tool_name)
221 if tool is None:
222 raise MCPError(-32602, f"Unknown tool: {tool_name}")
224 try:
225 result = tool.handler(arguments) if not asyncio.iscoroutinefunction(tool.handler) else await tool.handler(arguments)
226 return {
227 "content": [
228 {"type": "text", "text": str(result) if not isinstance(result, str) else result}
229 ]
230 }
231 except Exception as e:
232 return {
233 "content": [
234 {"type": "text", "text": f"Error: {str(e)}"}
235 ],
236 "isError": True,
237 }
239 async def _handle_resources_list(self, params: dict) -> dict:
240 resources = []
241 for r in self._resources.values():
242 resources.append({
243 "uri": r.uri,
244 "name": r.name,
245 "description": r.description,
246 "mimeType": r.mime_type,
247 })
248 return {"resources": resources}
250 async def _handle_resources_read(self, params: dict) -> dict:
251 uri = params.get("uri", "")
252 r = self._resources.get(uri)
253 if r is None:
254 raise MCPError(-32602, f"Unknown resource: {uri}")
255 text = r.content() if callable(r.content) else r.content
256 return {
257 "contents": [
258 {"uri": uri, "mimeType": r.mime_type, "text": str(text)}
259 ]
260 }
262 async def _handle_prompts_list(self, params: dict) -> dict:
263 prompts = []
264 for p in self._prompts.values():
265 prompts.append({
266 "name": p.name,
267 "description": p.description,
268 "arguments": p.arguments,
269 })
270 return {"prompts": prompts}
272 async def _handle_prompts_get(self, params: dict) -> dict:
273 prompt_name = params.get("name", "")
274 prompt_args = params.get("arguments", {})
275 p = self._prompts.get(prompt_name)
276 if p is None:
277 raise MCPError(-32602, f"Unknown prompt: {prompt_name}")
278 template = p.template(prompt_args) if callable(p.template) else p.template
279 return {
280 "description": p.description,
281 "messages": [
282 {"role": "user", "content": {"type": "text", "text": template}}
283 ]
284 }
286 # ── 内置工具 ──────────────────────────
288 def _register_builtin_tools(self):
289 """注册 AgentOS 内置 MCP 工具。"""
291 self.register_tool(MCPToolDef(
292 name="agentos_chat",
293 description="使用 AgentOS LLM 进行对话(支持 OpenAI/DeepSeek/Anthropic/Claude/Ollama)",
294 input_schema={
295 "type": "object",
296 "properties": {
297 "messages": {
298 "type": "array",
299 "description": "对话消息列表",
300 "items": {
301 "type": "object",
302 "properties": {
303 "role": {"type": "string", "enum": ["system", "user", "assistant"]},
304 "content": {"type": "string"},
305 },
306 "required": ["role", "content"],
307 },
308 },
309 "model": {"type": "string", "description": "模型名称,默认从配置读取"},
310 "temperature": {"type": "number", "description": "温度参数(0-2)"},
311 "max_tokens": {"type": "integer", "description": "最大输出 token 数"},
312 },
313 "required": ["messages"],
314 },
315 handler=self._tool_agentos_chat,
316 ))
318 self.register_tool(MCPToolDef(
319 name="agentos_list_tools",
320 description="列出 AgentOS 中所有可用的工具(含 MCP 工具)",
321 input_schema={
322 "type": "object",
323 "properties": {
324 "format": {"type": "string", "enum": ["openai", "anthropic"], "description": "输出格式"},
325 },
326 },
327 handler=self._tool_list_tools,
328 ))
330 self.register_tool(MCPToolDef(
331 name="agentos_version",
332 description="获取 AgentOS 版本信息",
333 input_schema={"type": "object", "properties": {}},
334 handler=self._tool_version,
335 ))
337 async def _tool_agentos_chat(self, args: dict) -> str:
338 """调用 AgentOS LLM 对话。"""
339 try:
340 from agentos.llm import LLMClient, LLMMessage
341 except ImportError:
342 return "Error: AgentOS LLM 模块不可用。请确认已安装 nexus-agentos。"
344 messages_raw = args.get("messages", [])
345 model = args.get("model")
346 temperature = args.get("temperature", 0.7)
347 max_tokens = args.get("max_tokens", 4096)
349 messages = [LLMMessage(role=m["role"], content=m["content"]) for m in messages_raw]
351 client = LLMClient(model=model)
352 response = await client.chat(messages, temperature=temperature, max_tokens=max_tokens)
353 return response.content
355 def _tool_list_tools(self, args: dict) -> str:
356 """列出可用工具。"""
357 fmt = args.get("format", "openai")
358 tools_list = []
359 for name, tool in self._tools.items():
360 if fmt == "openai":
361 tools_list.append({
362 "type": "function",
363 "function": {"name": name, "description": tool.description, "parameters": tool.input_schema},
364 })
365 else:
366 tools_list.append({"name": name, "description": tool.description, "input_schema": tool.input_schema})
367 return json.dumps(tools_list, ensure_ascii=False, indent=2)
369 def _tool_version(self, args: dict) -> str:
370 """返回版本信息。"""
371 try:
372 from agentos import __version__
373 except ImportError:
374 __version__ = self.version
375 return json.dumps({
376 "name": self.name,
377 "version": self.version,
378 "agentos_version": __version__,
379 "tools_count": len(self._tools),
380 }, ensure_ascii=False)
383# ── 数据结构 ───────────────────────────────
386class MCPToolDef:
387 """MCP 工具定义。"""
389 def __init__(
390 self,
391 name: str,
392 description: str,
393 input_schema: dict,
394 handler,
395 ):
396 self.name = name
397 self.description = description
398 self.input_schema = input_schema
399 self.handler = handler
402class MCPResource:
403 """MCP 资源定义。"""
405 def __init__(
406 self,
407 uri: str,
408 name: str = "",
409 description: str = "",
410 mime_type: str = "text/plain",
411 content: Any = "",
412 ):
413 self.uri = uri
414 self.name = name
415 self.description = description
416 self.mime_type = mime_type
417 self.content = content
420class MCPPromptDef:
421 """MCP 提示模板定义。"""
423 def __init__(
424 self,
425 name: str,
426 description: str = "",
427 arguments: list = None,
428 template: Any = "",
429 ):
430 self.name = name
431 self.description = description
432 self.arguments = arguments or []
433 self.template = template
436class MCPError(Exception):
437 """MCP 协议错误(与服务端共用异常类)。"""
439 def __init__(self, code: int, message: str):
440 self.code = code
441 self.message = message
442 super().__init__(f"MCP Error [{code}]: {message}")
445# ── 便捷函数 ───────────────────────────────
448def create_default_server() -> MCPServer:
449 """创建预配置了 AgentOS 内置工具的 MCP Server。"""
450 return MCPServer(
451 name="agentos",
452 version="1.5.2",
453 )
456def start_mcp_server(port: int = 0):
457 """启动 MCP Server。
459 Args:
460 port: 0 表示 stdio 模式,>0 表示 HTTP SSE 模式(暂未实现)。
461 """
462 if port == 0:
463 server = create_default_server()
464 server.run_stdio()
465 else:
466 print("MCP HTTP SSE 模式暂未实现。请使用 stdio 模式(port=0)。")
467 sys.exit(1)
469# ── ServerInfo & Tool (test compatibility) ──
470from dataclasses import dataclass, field
472@dataclass
473class ServerInfo:
474 name: str
475 version: str
476 description: str = ""
478@dataclass
479class Tool:
480 name: str
481 description: str
482 input_schema: dict
483 call: callable = field(default=lambda params: None)
485@dataclass
486class AgentCard:
487 agent_id: str
488 name: str
489 version: str
490 capabilities: list = field(default_factory=list)
491 endpoint: str = ""