Coverage for agentos/mcp/server.py: 27%

220 statements  

« prev     ^ index     » next       coverage.py v7.14.3, created at 2026-07-06 23:17 +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 

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

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 { 

97 "name": t.name, 

98 "description": t.description, 

99 "inputSchema": getattr(t, "input_schema", {}), 

100 } 

101 ) 

102 return result 

103 

104 def run_stdio(self): 

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

106 asyncio.run(self._run_stdio_async()) 

107 

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) 

114 

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) 

119 

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

121 

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 

130 

131 try: 

132 request = json.loads(line_str) 

133 except json.JSONDecodeError: 

134 continue 

135 

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 

144 

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

150 

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

152 if req_id is None: 

153 if method == "notifications/initialized": 

154 self._initialized = True 

155 return None 

156 

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 } 

176 

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) 

192 

193 # ── MCP 协议方法 ────────────────────── 

194 

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 } 

208 

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} 

220 

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

227 

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 } 

244 

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} 

257 

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

265 

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} 

277 

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 } 

289 

290 # ── 内置工具 ────────────────────────── 

291 

292 def _register_builtin_tools(self): 

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

294 

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 ) 

326 

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 ) 

344 

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 ) 

353 

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

360 

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) 

365 

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

367 

368 client = LLMClient(model=model) 

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

370 return response.content 

371 

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) 

397 

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 ) 

413 

414 

415# ── 数据结构 ─────────────────────────────── 

416 

417 

418class MCPToolDef: 

419 """MCP 工具定义。""" 

420 

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 

432 

433 

434class MCPResource: 

435 """MCP 资源定义。""" 

436 

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 

450 

451 

452class MCPPromptDef: 

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

454 

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 

466 

467 

468class MCPError(Exception): 

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

470 

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

472 self.code = code 

473 self.message = message 

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

475 

476 

477# ── 便捷函数 ─────────────────────────────── 

478 

479 

480def create_default_server() -> MCPServer: 

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

482 return MCPServer( 

483 name="agentos", 

484 version="1.5.2", 

485 ) 

486 

487 

488def start_mcp_server(port: int = 0): 

489 """启动 MCP Server。 

490 

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) 

500 

501 

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

503from dataclasses import dataclass, field # noqa: E402 

504 

505 

506@dataclass 

507class ServerInfo: 

508 name: str 

509 version: str 

510 description: str = "" 

511 

512 

513@dataclass 

514class Tool: 

515 name: str 

516 description: str 

517 input_schema: dict 

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

519 

520 

521@dataclass 

522class AgentCard: 

523 agent_id: str 

524 name: str 

525 version: str 

526 capabilities: list = field(default_factory=list) 

527 endpoint: str = ""