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

1"""MCP Server 实现 — 将 AgentOS 暴露为 MCP Server。 

2 

3支持 stdio JSON-RPC 2.0 传输,暴露 LLM 对话、工具调用、Agent 运行等能力。 

4其他 MCP 客户端(如 Claude Desktop、Cursor)可直接连接使用。 

5 

6用法: 

7 agentos mcp-server # 以 stdio 模式启动 

8 agentos mcp-server --port 9000 # 以 HTTP SSE 模式启动(可选) 

9""" 

10 

11from __future__ import annotations 

12 

13import asyncio 

14import json 

15import logging 

16import os 

17import sys 

18from typing import Any, Dict, List, Optional 

19 

20logger = logging.getLogger(__name__) 

21 

22# ── MCP Server 核心 ───────────────────────── 

23 

24 

25class MCPServer: 

26 """MCP Server — stdio JSON-RPC 2.0 传输。 

27 

28 实现 MCP 协议的 server 端,暴露 AgentOS 能力。 

29 客户端通过 stdio 发送 JSON-RPC 请求,服务器响应。 

30 

31 支持的操作: 

32 - initialize: 协议握手,返回 capabilities 

33 - tools/list: 列出可用工具 

34 - tools/call: 调用工具 

35 - resources/list: 列出可用资源 

36 - prompts/list: 列出可用提示 

37 """ 

38 

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 

64 

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 

71 

72 # 内置工具 

73 self._register_builtin_tools() 

74 

75 @property 

76 def info(self) -> "ServerInfo": 

77 return ServerInfo(name=self.name, version=self.version) 

78 

79 @property 

80 def tools(self) -> Dict[str, Any]: 

81 return self._tools 

82 

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 

90 

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 

101 

102 def run_stdio(self): 

103 """以 stdio 模式运行 MCP Server(同步阻塞)。""" 

104 asyncio.run(self._run_stdio_async()) 

105 

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) 

112 

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) 

117 

118 logger.info(f"MCP Server '{self.name}' v{self.version} started (stdio)") 

119 

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 

128 

129 try: 

130 request = json.loads(line_str) 

131 except json.JSONDecodeError: 

132 continue 

133 

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 

142 

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", {}) 

148 

149 # 通知类(无 id),不回复 

150 if req_id is None: 

151 if method == "notifications/initialized": 

152 self._initialized = True 

153 return None 

154 

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 } 

174 

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) 

190 

191 # ── MCP 协议方法 ────────────────────── 

192 

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 } 

206 

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} 

216 

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}") 

223 

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 } 

238 

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} 

249 

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 } 

261 

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} 

271 

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 } 

285 

286 # ── 内置工具 ────────────────────────── 

287 

288 def _register_builtin_tools(self): 

289 """注册 AgentOS 内置 MCP 工具。""" 

290 

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 )) 

317 

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 )) 

329 

330 self.register_tool(MCPToolDef( 

331 name="agentos_version", 

332 description="获取 AgentOS 版本信息", 

333 input_schema={"type": "object", "properties": {}}, 

334 handler=self._tool_version, 

335 )) 

336 

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。" 

343 

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) 

348 

349 messages = [LLMMessage(role=m["role"], content=m["content"]) for m in messages_raw] 

350 

351 client = LLMClient(model=model) 

352 response = await client.chat(messages, temperature=temperature, max_tokens=max_tokens) 

353 return response.content 

354 

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) 

368 

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) 

381 

382 

383# ── 数据结构 ─────────────────────────────── 

384 

385 

386class MCPToolDef: 

387 """MCP 工具定义。""" 

388 

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 

400 

401 

402class MCPResource: 

403 """MCP 资源定义。""" 

404 

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 

418 

419 

420class MCPPromptDef: 

421 """MCP 提示模板定义。""" 

422 

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 

434 

435 

436class MCPError(Exception): 

437 """MCP 协议错误(与服务端共用异常类)。""" 

438 

439 def __init__(self, code: int, message: str): 

440 self.code = code 

441 self.message = message 

442 super().__init__(f"MCP Error [{code}]: {message}") 

443 

444 

445# ── 便捷函数 ─────────────────────────────── 

446 

447 

448def create_default_server() -> MCPServer: 

449 """创建预配置了 AgentOS 内置工具的 MCP Server。""" 

450 return MCPServer( 

451 name="agentos", 

452 version="1.5.2", 

453 ) 

454 

455 

456def start_mcp_server(port: int = 0): 

457 """启动 MCP Server。 

458 

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) 

468 

469# ── ServerInfo & Tool (test compatibility) ── 

470from dataclasses import dataclass, field 

471 

472@dataclass 

473class ServerInfo: 

474 name: str 

475 version: str 

476 description: str = "" 

477 

478@dataclass 

479class Tool: 

480 name: str 

481 description: str 

482 input_schema: dict 

483 call: callable = field(default=lambda params: None) 

484 

485@dataclass 

486class AgentCard: 

487 agent_id: str 

488 name: str 

489 version: str 

490 capabilities: list = field(default_factory=list) 

491 endpoint: str = ""