MCP构建指南

利用AI,结合官方文档收集整理,作为了解

MCP 开发指南:构建模型上下文协议服务器

Date: October 23, 2025
Version: 1.0
Specification Version: 2024-11-05

目录

  1. 快速开始

  2. 协议概述

  3. 核心架构

  4. 协议基础

  5. 生命周期管理

  6. 传输机制

  7. 服务端开发

  8. 客户端开发

  9. 实战示例

  10. 调试与测试

  11. 部署指南

  12. 最佳实践


快速开始

环境准备

系统要求

  • Python 3.8+ 或 Node.js 16+

  • 网络连接

  • 1GB 可用内存

安装 MCP SDK

1
2
3
4
5
# Python SDK
pip install "mcp[cli]"

# TypeScript SDK
npm install @modelcontextprotocol/sdk

5分钟创建第一个 MCP 服务器

Step 1: 创建项目目录

1
mkdir mcp-demo && cd mcp-demo

Step 2: 编写基础服务器代码

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
# server.py
from mcp import MCPServer
from mcp.types import ToolRegistration

class HelloWorldServer(MCPServer):
def __init__(self):
super().__init__(name="HelloWorldServer", version="1.0.0")
self.register_tools()

def register_tools(self):
"""注册工具"""
hello_tool: ToolRegistration = {
"name": "say_hello",
"description": "向指定名称的用户打招呼",
"parameters": {
"name": {
"type": "string",
"description": "用户名称",
"required": True
}
},
"handler": self.handle_say_hello
}
self.register_tool(hello_tool)

async def handle_say_hello(self, parameters):
"""处理打招呼工具调用"""
name = parameters.get("name", "World")
return {"message": f"Hello, {name}!"}

if __name__ == "__main__":
server = HelloWorldServer()
server.run(host="0.0.0.0", port=8080)

Step 3: 启动服务器

1
python server.py

Step 4: 测试服务器

使用 MCP 客户端测试:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
# client.py
import asyncio
from mcp import MCPClient

async def test_server():
client = MCPClient(host="localhost", port=8080)
session = await client.connect()

result = await session.call_tool(
tool_name="say_hello",
parameters={"name": "MCP Developer"}
)

print(f"服务器响应: {result}")
# 输出: 服务器响应: {'message': 'Hello, MCP Developer!'}

if __name__ == "__main__":
asyncio.run(test_server())

协议概述

什么是 MCP

MCP (Model Context Protocol) 是一种专门为大语言模型(LLM)设计的标准化通信协议,旨在建立模型与外部工具、数据源和服务之间的统一交互接口。

核心价值

MCP 解决了 AI 应用开发中的关键问题:

  • 标准化接口:为不同工具和服务提供统一调用方式

  • 上下文管理:支持复杂对话场景中的状态保持

  • 安全隔离:通过三层架构实现严格的安全边界

  • 扩展性:允许动态添加新功能而不影响现有系统

协议特点

  • 基于 JSON-RPC 2.0:成熟的远程过程调用规范

  • 有状态会话:维护上下文信息的会话机制

  • 双向通信:支持客户端和服务器双向消息交换

  • 多传输支持:stdio 和 HTTP/SSE 传输机制


核心架构

三层架构模式

MCP 采用客户端-主机-服务器三层架构:

1
2
3
4
┌─────────────────┐     ┌─────────────────┐     ┌─────────────────┐
│ 客户端 │ │ 主机 │ │ 服务器 │
│ (Client) │────▶│ (Host) │────▶│ (Server) │
└─────────────────┘ └─────────────────┘ └─────────────────┘

组件详情

主机 (Host)

职责

  • 管理客户端实例的生命周期

  • 控制连接权限和执行安全策略

  • 协调 AI/LLM 集成

  • 确保系统稳定运行

核心功能

  • 连接管理和路由

  • 权限验证和安全控制

  • 负载均衡和故障恢复

  • 监控和日志记录

客户端 (Client)

职责

  • 维护与服务器的独立连接

  • 建立有状态会话

  • 处理协议协商

  • 管理消息路由

核心功能

  • 会话状态管理

  • 消息序列化和解析

  • 错误处理和重试

  • 传输协议适配

服务器 (Server)

职责

  • 公开特定的资源和工具

  • 独立运行和管理

  • 通过客户端处理请求

  • 支持本地和远程服务

核心功能

  • 工具注册和执行

  • 资源管理和访问

  • 提示词管理

  • 事件通知


协议基础

消息类型

MCP 定义了三种基本消息类型,基于 JSON-RPC 2.0 规范。

1. 请求 (Request)

特点

  • 双向消息,可在客户端和服务器间双向发送

  • 必须包含唯一 ID

  • 支持参数传递

  • 需要响应

格式示例

1
2
3
4
5
6
7
8
9
10
11
{
"jsonrpc": "2.0",
"id": "request-123",
"method": "tool.call",
"params": {
"toolName": "say_hello",
"arguments": {
"name": "Alice"
}
}
}

2. 响应 (Response)

特点

  • 作为对请求的回复

  • 必须包含对应请求的 ID

  • 必须设置 result 或 error

  • 错误码必须是整数

成功响应示例

1
2
3
4
5
6
7
{
"jsonrpc": "2.0",
"id": "request-123",
"result": {
"message": "Hello, Alice!"
}
}

错误响应示例

1
2
3
4
5
6
7
8
9
10
11
12
{
"jsonrpc": "2.0",
"id": "request-123",
"error": {
"code": -32602,
"message": "参数验证失败",
"data": {
"parameter": "name",
"message": "名称不能为空"
}
}
}

3. 通知 (Notification)

特点

  • 不需要响应的单向消息

  • 不能包含 ID 字段

  • 用于状态更新和事件通知

  • 减少通信开销

格式示例

1
2
3
4
5
6
7
8
{
"jsonrpc": "2.0",
"method": "status.update",
"params": {
"status": "processing",
"progress": 50
}
}

错误码规范

错误码范围 含义 示例
-32768 到 -32000 JSON-RPC 保留错误码 -32600: 无效请求
-32099 到 -32000 服务器端错误 -32001: 工具不存在
-32099 到 -32000 自定义错误码 -32002: 权限不足

生命周期管理

会话生命周期

MCP 会话包含三个主要阶段:初始化、操作和关闭。

1. 初始化阶段

初始化是客户端和服务器的第一次交互,建立通信基础。

初始化流程

1
2
3
4
5
6
7
8
9
10
11
12
13
客户端                  服务器
│ │
│ initialize 请求 │
│─────────────────────▶│
│ │
│ │ 验证版本和能力
│ │
│ initialize 响应 │
│◀─────────────────────│
│ │
│ initialized 通知 │
│─────────────────────▶│
│ │

初始化请求

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
{
"jsonrpc": "2.0",
"id": 1,
"method": "initialize",
"params": {
"protocolVersion": "2024-11-05",
"capabilities": {
"roots": {
"listChanged": true
},
"sampling": {}
},
"clientInfo": {
"name": "MyClient",
"version": "1.0.0"
}
}
}

初始化响应

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
{
"jsonrpc": "2.0",
"id": 1,
"result": {
"protocolVersion": "2024-11-05",
"capabilities": {
"tools": {
"listChanged": true
},
"resources": {
"subscribe": true
}
},
"serverInfo": {
"name": "MyServer",
"version": "1.0.0"
}
}
}

2. 操作阶段

操作阶段是会话的核心,处理所有工具调用和资源访问。

主要操作

  • 工具调用:tool.call

  • 资源读取:resource.read

  • 资源订阅:resource.subscribe

  • 提示词获取:prompt.get

工具调用示例

1
2
3
4
5
6
7
8
9
10
11
{
"jsonrpc": "2.0",
"id": "call-456",
"method": "tool.call",
"params": {
"toolName": "calculate",
"arguments": {
"expression": "2 + 3 * 4"
}
}
}

3. 关闭阶段

优雅地终止会话并清理资源。

关闭流程

1
2
3
4
5
6
7
8
9
10
11
12
13
客户端                  服务器
│ │
│ shutdown 通知 │
│─────────────────────▶│
│ │
│ │ 清理资源
│ │
│ shutdown 确认 │
│◀─────────────────────│
│ │
│ 连接关闭 │
│─────────────────────▶│
│ │

版本协商

版本协商策略

  • 客户端发送支持的最新版本

  • 服务器响应相同版本或支持的其他版本

  • 版本格式:YYYY-MM-DD

  • 不兼容时断开连接

能力协商

客户端能力

  • roots:提供文件系统根目录

  • sampling:支持 LLM 采样请求

  • experimental:实验性功能

服务器能力

  • tools:提供可调用工具

  • resources:提供可读资源

  • prompts:提供提示模板

  • logging:结构化日志

  • experimental:实验性功能


传输机制

标准输入输出 (stdio)

适合本地集成和命令行工具。

特点

  • 客户端将服务器作为子进程启动

  • 通过 stdin/stdout 传输 JSON-RPC 消息

  • 消息以换行符分隔

  • stderr 用于日志记录

Python 实现示例

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
# stdio_server.py
import sys
import json
import asyncio

async def handle_stdin():
for line in sys.stdin:
line = line.strip()
if not line:
continue

try:
message = json.loads(line)
# 处理消息
response = await process_message(message)
if response:
print(json.dumps(response))
sys.stdout.flush()
except Exception as e:
error = {
"jsonrpc": "2.0",
"error": {
"code": -32603,
"message": str(e)
}
}
print(json.dumps(error))
sys.stdout.flush()

async def process_message(message):
"""处理 MCP 消息"""
if message["method"] == "initialize":
return {
"jsonrpc": "2.0",
"id": message["id"],
"result": {
"protocolVersion": "2024-11-05",
"capabilities": {"tools": {}}
}
}
return None

if __name__ == "__main__":
asyncio.run(handle_stdin())

基于 SSE 的 HTTP

适合远程通信和 Web 应用。

服务器端点

  • SSE 端点:接收服务器消息

  • HTTP POST 端点:发送客户端消息

工作流程

1
2
3
4
1. 客户端连接到 SSE 端点
2. 服务器发送 endpoint 事件
3. 客户端使用 POST 端点发送消息
4. 服务器通过 SSE 发送响应

TypeScript 实现示例

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
// http_server.ts
import { createServer } from 'http';
import { SSE } from 'sse';

const server = createServer((req, res) => {
if (req.url === '/mcp/sse') {
// SSE 连接
const sse = new SSE(req, res);

// 发送 endpoint 事件
sse.send({
event: 'endpoint',
data: JSON.stringify({
endpoint: '/mcp/message'
})
});

} else if (req.url === '/mcp/message' && req.method === 'POST') {
// 处理消息
let body = '';
req.on('data', chunk => body += chunk);
req.on('end', () => {
const message = JSON.parse(body);
handleMessage(message, res);
});
}
});

function handleMessage(message: any, res: any) {
// 处理 MCP 消息
const response = processMessage(message);
res.writeHead(200, { 'Content-Type': 'application/json' });
res.end(JSON.stringify(response));
}

server.listen(8080, () => {
console.log('MCP HTTP server running on port 8080');
});

服务端开发

核心组件

1. 服务器类

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
from mcp import MCPServer
from typing import Dict, Any, Coroutine

class MyServer(MCPServer):
def __init__(self, name: str, version: str):
super().__init__(name=name, version=version)
self.register_components()

def register_components(self):
"""注册工具、资源和提示词"""
self.register_tools()
self.register_resources()
self.register_prompts()

def register_tools(self):
"""注册工具"""
pass

def register_resources(self):
"""注册资源"""
pass

def register_prompts(self):
"""注册提示词"""
pass

2. 工具开发

工具注册格式

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
from mcp.types import ToolRegistration

def register_my_tool(server: MCPServer):
"""注册自定义工具"""
tool: ToolRegistration = {
"name": "tool_name",
"description": "工具描述",
"parameters": {
"param1": {
"type": "string",
"description": "参数1描述",
"required": True
},
"param2": {
"type": "number",
"description": "参数2描述",
"required": False,
"default": 0
}
},
"handler": tool_handler
}
server.register_tool(tool)

async def tool_handler(parameters: Dict[str, Any]) -> Dict[str, Any]:
"""工具处理函数"""
param1 = parameters.get("param1")
param2 = parameters.get("param2", 0)

# 业务逻辑处理
result = process_data(param1, param2)

return {
"success": True,
"result": result
}

工具类型

  1. 同步工具:立即返回结果

  2. 异步工具:后台执行,需要轮询状态

  3. 事件工具:执行后发送事件通知

3. 资源管理

资源注册示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
from mcp.types import ResourceRegistration

def register_config_resource(server: MCPServer):
"""注册配置资源"""
resource: ResourceRegistration = {
"name": "config",
"uri": "config://app",
"description": "应用配置资源",
"handler": config_resource_handler
}
server.register_resource(resource)

async def config_resource_handler(request: Dict[str, Any]) -> Dict[str, Any]:
"""资源处理函数"""
config_data = {
"app_name": "My MCP Server",
"version": "1.0.0",
"features": ["tools", "resources", "prompts"]
}

return {
"contents": [{"text": json.dumps(config_data, indent=2)}]
}

4. 提示词管理

提示词注册示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
from mcp.types import PromptRegistration

def register_prompts(server: MCPServer):
"""注册提示词"""
code_reviewer_prompt: PromptRegistration = {
"name": "code_reviewer",
"description": "代码审查专家提示词",
"parameters": {
"language": {
"type": "string",
"description": "编程语言",
"required": True
}
},
"template": """你是一位专业的{{language}}代码审查专家。
请仔细审查以下代码,指出潜在的问题、性能优化点和最佳实践建议。
重点关注:
1. 代码可读性和维护性
2. 潜在的bug和安全问题
3. 性能优化机会
4. 编码规范符合性"""
}
server.register_prompt(code_reviewer_prompt)

事件系统

事件类型

  • tool.executed:工具执行完成

  • resource.updated:资源更新

  • session.closed:会话关闭

  • error.occurred:错误发生

事件发送示例

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
async def long_running_tool_handler(parameters: Dict[str, Any]):
"""长时间运行的工具"""
task_id = str(uuid.uuid4())

# 启动后台任务
asyncio.create_task(
background_task(task_id, parameters)
)

return {
"task_id": task_id,
"status": "processing"
}

async def background_task(task_id: str, parameters: Dict[str, Any]):
"""后台任务"""
try:
# 模拟长时间处理
await asyncio.sleep(30)

# 任务完成,发送事件
await self.send_event({
"event_type": "task.completed",
"task_id": task_id,
"result": {"status": "success"}
})

except Exception as e:
await self.send_event({
"event_type": "task.failed",
"task_id": task_id,
"error": str(e)
})

客户端开发

客户端基础

连接管理

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
from mcp import MCPClient
from mcp.types import Session

class MyMCP客户端:
def __init__(self, host: str, port: int):
self.client = MCPClient(host=host, port=port)
self.session: Session = None

async def connect(self) -> bool:
"""建立连接"""
try:
self.session = await self.client.connect()
return True
except Exception as e:
print(f"连接失败: {e}")
return False

async def disconnect(self):
"""断开连接"""
if self.session:
await self.session.close()
self.session = None

async def call_tool(self, tool_name: str, parameters: Dict[str, Any]) -> Dict[str, Any]:
"""调用工具"""
if not self.session:
raise Exception("未建立连接")

return await self.session.call_tool(
tool_name=tool_name,
parameters=parameters
)

异步操作处理

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
async def handle_async_tool_call(client: MyMCP客户端, tool_name: str, parameters: Dict[str, Any]):
"""处理异步工具调用"""
try:
# 调用异步工具
result = await client.call_tool(tool_name, parameters)

if result.get("status") == "processing":
task_id = result.get("task_id")

# 轮询任务状态
while True:
status = await client.call_tool("get_task_status", {"task_id": task_id})

if status.get("status") == "completed":
return status.get("result")
elif status.get("status") == "failed":
raise Exception(status.get("error"))

await asyncio.sleep(2)

return result

except Exception as e:
print(f"工具调用失败: {e}")
raise

事件订阅

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
async def subscribe_to_events(client: MyMCP客户端):
"""订阅服务器事件"""
async for event in client.session.events():
try:
event_data = json.loads(event.data)
event_type = event_data.get("event_type")

if event_type == "task.completed":
handle_task_completed(event_data)
elif event_type == "resource.updated":
handle_resource_updated(event_data)

except Exception as e:
print(f"处理事件失败: {e}")

def handle_task_completed(event_data: Dict[str, Any]):
"""处理任务完成事件"""
task_id = event_data.get("task_id")
result = event_data.get("result")
print(f"任务 {task_id} 完成: {result}")

实战示例

示例 1: 文件管理器服务器

功能需求

  • 列出目录内容

  • 读取文件内容

  • 创建和修改文件

  • 删除文件和目录

实现代码

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
# file_manager_server.py
import os
import json
from mcp import MCPServer
from mcp.types import ToolRegistration, ResourceRegistration

class FileManagerServer(MCPServer):
def __init__(self):
super().__init__(name="FileManagerServer", version="1.0.0")
self.register_tools()
self.register_resources()

def register_tools(self):
"""注册文件管理工具"""
# 列出目录工具
list_dir_tool: ToolRegistration = {
"name": "list_directory",
"description": "列出指定目录的内容",
"parameters": {
"path": {
"type": "string",
"description": "目录路径",
"required": True
}
},
"handler": self.handle_list_directory
}
self.register_tool(list_dir_tool)

# 读取文件工具
read_file_tool: ToolRegistration = {
"name": "read_file",
"description": "读取文件内容",
"parameters": {
"path": {
"type": "string",
"description": "文件路径",
"required": True
}
},
"handler": self.handle_read_file
}
self.register_tool(read_file_tool)

# 写入文件工具
write_file_tool: ToolRegistration = {
"name": "write_file",
"description": "写入文件内容",
"parameters": {
"path": {
"type": "string",
"description": "文件路径",
"required": True
},
"content": {
"type": "string",
"description": "文件内容",
"required": True
}
},
"handler": self.handle_write_file
}
self.register_tool(write_file_tool)

def register_resources(self):
"""注册文件系统资源"""
file_resource: ResourceRegistration = {
"name": "file",
"uri": "file://",
"description": "文件系统资源",
"handler": self.handle_file_resource
}
self.register_resource(file_resource)

async def handle_list_directory(self, parameters):
"""处理列出目录请求"""
path = parameters.get("path")

try:
if not os.path.isdir(path):
return {"error": f"路径 {path} 不是目录"}

items = os.listdir(path)
result = []

for item in items:
item_path = os.path.join(path, item)
result.append({
"name": item,
"type": "directory" if os.path.isdir(item_path) else "file",
"size": os.path.getsize(item_path) if os.path.isfile(item_path) else None
})

return {"contents": result}

except Exception as e:
return {"error": str(e)}

async def handle_read_file(self, parameters):
"""处理读取文件请求"""
path = parameters.get("path")

try:
if not os.path.isfile(path):
return {"error": f"文件 {path} 不存在"}

with open(path, 'r', encoding='utf-8') as f:
content = f.read()

return {"content": content}

except Exception as e:
return {"error": str(e)}

async def handle_write_file(self, parameters):
"""处理写入文件请求"""
path = parameters.get("path")
content = parameters.get("content")

try:
with open(path, 'w', encoding='utf-8') as f:
f.write(content)

return {"success": True, "message": f"文件 {path} 写入成功"}

except Exception as e:
return {"error": str(e)}

async def handle_file_resource(self, request):
"""处理文件资源请求"""
uri = request.get("uri")
path = uri.replace("file://", "")

try:
if os.path.isfile(path):
with open(path, 'r', encoding='utf-8') as f:
content = f.read()

return {
"contents": [{"text": content}]
}
elif os.path.isdir(path):
items = os.listdir(path)
directory_info = json.dumps({
"type": "directory",
"items": items
}, indent=2)

return {
"contents": [{"text": directory_info}]
}
else:
return {
"error": f"路径 {path} 不存在"
}

except Exception as e:
return {
"error": str(e)
}

if __name__ == "__main__":
server = FileManagerServer()
server.run(host="0.0.0.0", port=8080)

示例 2: 天气查询服务器

功能需求

  • 获取实时天气信息

  • 获取天气预报

  • 支持多个城市

  • 缓存查询结果

实现代码

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
# weather_server.py
import asyncio
import aiohttp
import time
from mcp import MCPServer
from mcp.types import ToolRegistration

class WeatherServer(MCPServer):
def __init__(self):
super().__init__(name="WeatherServer", version="1.0.0")
self.api_key = "your_api_key" # 替换为实际的API密钥
self.cache = {}
self.cache_ttl = 3600 # 缓存有效期1小时
self.register_tools()

def register_tools(self):
"""注册天气查询工具"""
current_weather_tool: ToolRegistration = {
"name": "get_current_weather",
"description": "获取指定城市的实时天气",
"parameters": {
"city": {
"type": "string",
"description": "城市名称",
"required": True
}
},
"handler": self.handle_get_current_weather
}
self.register_tool(current_weather_tool)

forecast_tool: ToolRegistration = {
"name": "get_weather_forecast",
"description": "获取指定城市的天气预报",
"parameters": {
"city": {
"type": "string",
"description": "城市名称",
"required": True
},
"days": {
"type": "integer",
"description": "预报天数",
"required": False,
"default": 3
}
},
"handler": self.handle_get_weather_forecast
}
self.register_tool(forecast_tool)

async def get_weather_from_api(self, city: str, forecast: bool = False, days: int = 3):
"""从API获取天气数据"""
cache_key = f"{city}:{forecast}:{days}"

# 检查缓存
if cache_key in self.cache:
cached_data, timestamp = self.cache[cache_key]
if time.time() - timestamp < self.cache_ttl:
return cached_data

# 构建API请求
base_url = "https://api.weatherapi.com/v1/"
endpoint = "forecast.json" if forecast else "current.json"

params = {
"key": self.api_key,
"q": city,
"days": days
}

async with aiohttp.ClientSession() as session:
async with session.get(f"{base_url}{endpoint}", params=params) as response:
if response.status != 200:
raise Exception(f"API请求失败: {response.status}")

data = await response.json()

# 更新缓存
self.cache[cache_key] = (data, time.time())

return data

async def handle_get_current_weather(self, parameters):
"""处理实时天气查询"""
city = parameters.get("city")

try:
weather_data = await self.get_weather_from_api(city)

return {
"city": city,
"temperature": weather_data["current"]["temp_c"],
"condition": weather_data["current"]["condition"]["text"],
"humidity": weather_data["current"]["humidity"],
"wind_speed": weather_data["current"]["wind_kph"],
"last_updated": weather_data["current"]["last_updated"]
}

except Exception as e:
return {"error": str(e)}

async def handle_get_weather_forecast(self, parameters):
"""处理天气预报查询"""
city = parameters.get("city")
days = parameters.get("days", 3)

try:
weather_data = await self.get_weather_from_api(city, forecast=True, days=days)

forecast = []
for day in weather_data["forecast"]["forecastday"]:
forecast.append({
"date": day["date"],
"max_temp": day["day"]["maxtemp_c"],
"min_temp": day["day"]["mintemp_c"],
"condition": day["day"]["condition"]["text"],
"precipitation": day["day"]["totalprecip_mm"]
})

return {
"city": city,
"forecast": forecast,
"last_updated": weather_data["current"]["last_updated"]
}

except Exception as e:
return {"error": str(e)}

if __name__ == "__main__":
server = WeatherServer()
server.run(host="0.0.0.0", port=8081)

调试与测试

MCP 检查器

使用官方 MCP Inspector 工具进行调试:

1
2
3
4
5
# 安装 MCP Inspector
npm install -g @modelcontextprotocol/inspector

# 启动检查器
mcp-inspector --server http://localhost:8080

日志配置

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
import logging
import os
from datetime import datetime

def setup_logging():
"""设置日志配置"""
log_dir = "logs"
os.makedirs(log_dir, exist_ok=True)

log_file = os.path.join(
log_dir,
f"mcp-server-{datetime.now().strftime('%Y%m%d')}.log"
)

logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(name)s:%(lineno)d - %(message)s",
handlers=[
logging.FileHandler(log_file, encoding='utf-8'),
logging.StreamHandler()
]
)

return logging.getLogger("mcp-server")

# 在服务器中使用
logger = setup_logging()

class MyServer(MCPServer):
async def handle_say_hello(self, parameters):
logger.info(f"处理打招呼请求: {parameters}")
try:
name = parameters.get("name", "World")
result = {"message": f"Hello, {name}!"}
logger.info(f"请求处理成功: {result}")
return result
except Exception as e:
logger.error(f"请求处理失败: {e}")
raise

单元测试

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
# tests/test_server.py
import pytest
import asyncio
from unittest.mock import Mock, patch
from my_server import MyServer

@pytest.mark.asyncio
async def test_say_hello_tool():
"""测试打招呼工具"""
server = MyServer()

# 测试正常情况
result = await server.handle_say_hello({"name": "Test User"})
assert result == {"message": "Hello, Test User!"}

# 测试默认值
result = await server.handle_say_hello({})
assert result == {"message": "Hello, World!"}

@pytest.mark.asyncio
async def test_file_manager_list_directory():
"""测试文件管理器的列出目录功能"""
server = FileManagerServer()

# 测试存在的目录
result = await server.handle_list_directory({"path": "."})
assert "contents" in result
assert isinstance(result["contents"], list)

# 测试不存在的目录
result = await server.handle_list_directory({"path": "/nonexistent"})
assert "error" in result

集成测试

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
# tests/test_integration.py
import pytest
import asyncio
from mcp import MCPClient

@pytest.mark.asyncio
async def test_server_integration():
"""测试服务器集成功能"""
client = MCPClient(host="localhost", port=8080)

try:
# 连接服务器
session = await client.connect()
assert session is not None

# 测试工具调用
result = await session.call_tool(
tool_name="say_hello",
parameters={"name": "Integration Test"}
)
assert result == {"message": "Hello, Integration Test!"}

# 测试资源访问
resource_result = await session.get_resource("config://app")
assert "contents" in resource_result

finally:
await client.disconnect()

部署指南

Docker 容器化

Dockerfile

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
FROM python:3.12-slim

# 设置工作目录
WORKDIR /app

# 安装系统依赖
RUN apt-get update && apt-get install -y \
build-essential \
&& rm -rf /var/lib/apt/lists/*

# 设置 Python 环境
ENV PYTHONDONTWRITEBYTECODE=1
ENV PYTHONUNBUFFERED=1

# 安装依赖
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt

# 复制项目文件
COPY src/ ./src/

# 创建非 root 用户
RUN useradd -m appuser
USER appuser

# 暴露端口
EXPOSE 8080

# 健康检查
HEALTHCHECK --interval=30s --timeout=10s --start-period=5s --retries=3 \
CMD curl -f http://localhost:8080/health || exit 1

# 启动命令
CMD ["python", "src/server.py", "--host", "0.0.0.0", "--port", "8080"]

Docker Compose

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
version: '3.8'

services:
mcp-server:
build: .
container_name: mcp-server
restart: unless-stopped
ports:
- "8080:8080"
environment:
- NODE_ENV=production
- MCP_PORT=8080
- API_KEY=${API_KEY}
volumes:
- ./data:/app/data
- ./logs:/app/logs
networks:
- mcp-network

nginx-proxy:
image: nginx:alpine
container_name: mcp-nginx
restart: unless-stopped
ports:
- "80:80"
- "443:443"
volumes:
- ./nginx/conf.d:/etc/nginx/conf.d
- ./nginx/ssl:/etc/nginx/ssl
depends_on:
- mcp-server
networks:
- mcp-network

networks:
mcp-network:
driver: bridge

Kubernetes 部署

Deployment

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
apiVersion: apps/v1
kind: Deployment
metadata:
name: mcp-server
namespace: mcp-system
spec:
replicas: 3
selector:
matchLabels:
app: mcp-server
template:
metadata:
labels:
app: mcp-server
spec:
containers:
- name: mcp-server
image: your-registry/mcp-server:latest
ports:
- containerPort: 8080
resources:
requests:
memory: "256Mi"
cpu: "250m"
limits:
memory: "512Mi"
cpu: "500m"
env:
- name: NODE_ENV
value: "production"
- name: MCP_PORT
value: "8080"
- name: API_KEY
valueFrom:
secretKeyRef:
name: mcp-secrets
key: api-key
readinessProbe:
httpGet:
path: /health
port: 8080
initialDelaySeconds: 5
periodSeconds: 10
livenessProbe:
httpGet:
path: /health
port: 8080
initialDelaySeconds: 15
periodSeconds: 20

Service

1
2
3
4
5
6
7
8
9
10
11
12
13
apiVersion: v1
kind: Service
metadata:
name: mcp-server-service
namespace: mcp-system
spec:
selector:
app: mcp-server
ports:
- protocol: TCP
port: 80
targetPort: 8080
type: ClusterIP

Ingress

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
apiVersion: networking.k8s.io/v1
kind: Ingress
metadata:
name: mcp-server-ingress
namespace: mcp-system
annotations:
nginx.ingress.kubernetes.io/rewrite-target: /
nginx.ingress.kubernetes.io/ssl-redirect: "true"
spec:
tls:
- hosts:
- mcp.yourdomain.com
secretName: mcp-tls-secret
rules:
- host: mcp.yourdomain.com
http:
paths:
- path: /
pathType: Prefix
backend:
service:
name: mcp-server-service
port:
number: 80

云服务部署

AWS ECS 部署

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
# task-definition.json
{
"family": "mcp-server",
"networkMode": "awsvpc",
"requiresCompatibilities": ["FARGATE"],
"cpu": "256",
"memory": "512",
"executionRoleArn": "arn:aws:iam::your-account:role/ecs-execution-role",
"containerDefinitions": [
{
"name": "mcp-server",
"image": "your-account.dkr.ecr.region.amazonaws.com/mcp-server:latest",
"portMappings": [
{
"containerPort": 8080,
"protocol": "tcp"
}
],
"environment": [
{
"name": "NODE_ENV",
"value": "production"
},
{
"name": "MCP_PORT",
"value": "8080"
}
],
"logConfiguration": {
"logDriver": "awslogs",
"options": {
"awslogs-group": "/ecs/mcp-server",
"awslogs-region": "us-east-1",
"awslogs-stream-prefix": "ecs"
}
}
}
]
}

最佳实践

安全最佳实践

1. 身份认证与授权

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
from mcp import MCPServer
import jwt
from datetime import datetime, timedelta

class SecureServer(MCPServer):
def __init__(self, secret_key: str):
super().__init__(name="SecureServer", version="1.0.0")
self.secret_key = secret_key
self.register_middleware(self.auth_middleware)

async def auth_middleware(self, request, next_handler):
"""认证中间件"""
# 跳过初始化请求的认证
if request.method == "initialize":
return await next_handler(request)

# 获取认证令牌
auth_header = request.headers.get("Authorization")
if not auth_header or not auth_header.startswith("Bearer "):
return {
"jsonrpc": "2.0",
"error": {
"code": -32602,
"message": "缺少认证令牌"
}
}

token = auth_header.split(" ")[1]

try:
# 验证令牌
payload = jwt.decode(
token,
self.secret_key,
algorithms=["HS256"]
)

# 检查令牌过期时间
if payload.get("exp") < datetime.utcnow().timestamp():
return {
"jsonrpc": "2.0",
"error": {
"code": -32602,
"message": "令牌已过期"
}
}

# 将用户信息添加到请求
request.user = payload.get("user")

return await next_handler(request)

except jwt.InvalidTokenError:
return {
"jsonrpc": "2.0",
"error": {
"code": -32602,
"message": "无效的认证令牌"
}
}

2. 输入验证

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
from typing import Dict, Any
import re

class InputValidator:
@staticmethod
def validate_tool_parameters(parameters: Dict[str, Any], schema: Dict[str, Any]) -> Dict[str, Any]:
"""验证工具参数"""
validated = {}

for param_name, param_schema in schema.items():
if param_schema.get("required") and param_name not in parameters:
raise ValueError(f"参数 {param_name} 是必填项")

if param_name not in parameters:
if "default" in param_schema:
validated[param_name] = param_schema["default"]
continue

value = parameters[param_name]

# 类型验证
expected_type = param_schema.get("type")
if expected_type and not InputValidator._check_type(value, expected_type):
raise ValueError(f"参数 {param_name} 类型应为 {expected_type}")

# 格式验证
if "pattern" in param_schema and not re.match(param_schema["pattern"], str(value)):
raise ValueError(f"参数 {param_name} 格式不符合要求")

# 范围验证
if "min" in param_schema and value < param_schema["min"]:
raise ValueError(f"参数 {param_name} 最小值为 {param_schema['min']}")

if "max" in param_schema and value > param_schema["max"]:
raise ValueError(f"参数 {param_name} 最大值为 {param_schema['max']}")

validated[param_name] = value

return validated

@staticmethod
def _check_type(value, expected_type: str) -> bool:
"""检查值的类型"""
if expected_type == "string":
return isinstance(value, str)
elif expected_type == "number":
return isinstance(value, (int, float))
elif expected_type == "boolean":
return isinstance(value, bool)
elif expected_type == "array":
return isinstance(value, list)
elif expected_type == "object":
return isinstance(value, dict)
return True

3. 安全通信

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
import ssl
import asyncio
from aiohttp import web

async def create_secure_server(app, host, port, ssl_context):
"""创建安全的 HTTPS 服务器"""
runner = web.AppRunner(app)
await runner.setup()

site = web.TCPSite(
runner,
host=host,
port=port,
ssl_context=ssl_context
)

await site.start()
return runner

def create_ssl_context(certfile, keyfile):
"""创建 SSL 上下文"""
ssl_context = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH)
ssl_context.load_cert_chain(certfile=certfile, keyfile=keyfile)
return ssl_context

性能最佳实践

1. 连接池管理

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
import asyncio
from typing import Generic, TypeVar, List, Tuple

T = TypeVar('T')

class ConnectionPool(Generic[T]):
def __init__(self, create_connection, max_size=10, max_idle_time=300):
self.create_connection = create_connection
self.max_size = max_size
self.max_idle_time = max_idle_time
self.pool: List[Tuple[T, float]] = []
self.lock = asyncio.Lock()

async def get_connection(self) -> T:
"""获取连接"""
async with self.lock:
# 清理空闲连接
now = asyncio.get_event_loop().time()
self.pool = [
(conn, timestamp) for conn, timestamp in self.pool
if now - timestamp < self.max_idle_time
]

# 如果有可用连接,返回最旧的
if self.pool:
conn, _ = self.pool.pop(0)
return conn

# 如果池未满,创建新连接
if len(self.pool) < self.max_size:
conn = await self.create_connection()
return conn

# 池已满,等待其他连接释放
await asyncio.sleep(0.1)
return await self.get_connection()

async def release_connection(self, conn: T):
"""释放连接"""
async with self.lock:
if len(self.pool) < self.max_size:
self.pool.append((conn, asyncio.get_event_loop().time()))

2. 缓存策略

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
import time
from typing import Dict, Any, Optional
from functools import wraps

class CacheManager:
def __init__(self, default_ttl=300):
self.cache: Dict[str, Tuple[Any, float]] = {}
self.default_ttl = default_ttl

def get(self, key: str) -> Optional[Any]:
"""获取缓存"""
if key in self.cache:
value, expire_time = self.cache[key]
if time.time() < expire_time:
return value
# 过期缓存清理
del self.cache[key]
return None

def set(self, key: str, value: Any, ttl: Optional[int] = None):
"""设置缓存"""
expire_time = time.time() + (ttl or self.default_ttl)
self.cache[key] = (value, expire_time)

def delete(self, key: str):
"""删除缓存"""
if key in self.cache:
del self.cache[key]

def clear(self):
"""清空缓存"""
self.cache.clear()

# 缓存装饰器
def cache_result(cache_manager: CacheManager, ttl: Optional[int] = None):
"""缓存结果装饰器"""
def decorator(func):
@wraps(func)
async def wrapper(*args, **kwargs):
# 生成缓存键
cache_key = f"{func.__name__}:{str(args)}:{str(kwargs)}"

# 尝试从缓存获取
cached_result = cache_manager.get(cache_key)
if cached_result is not None:
return cached_result

# 执行函数
result = await func(*args, **kwargs)

# 缓存结果
cache_manager.set(cache_key, result, ttl)

return result
return wrapper
return decorator

3. 异步处理

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
import asyncio
import uuid
from typing import Dict, Any

class AsyncTaskManager:
def __init__(self):
self.tasks: Dict[str, asyncio.Task] = {}
self.task_results: Dict[str, Any] = {}
self.task_errors: Dict[str, Exception] = {}

async def create_task(self, coro, task_id: Optional[str] = None) -> str:
"""创建异步任务"""
task_id = task_id or str(uuid.uuid4())

async def task_wrapper():
try:
result = await coro
self.task_results[task_id] = result
except Exception as e:
self.task_errors[task_id] = e
finally:
if task_id in self.tasks:
del self.tasks[task_id]

task = asyncio.create_task(task_wrapper())
self.tasks[task_id] = task

return task_id

def get_task_status(self, task_id: str) -> Dict[str, Any]:
"""获取任务状态"""
if task_id in self.task_results:
return {
"status": "completed",
"result": self.task_results[task_id]
}
elif task_id in self.task_errors:
return {
"status": "failed",
"error": str(self.task_errors[task_id])
}
elif task_id in self.tasks:
return {"status": "processing"}
else:
return {"status": "unknown"}

async def wait_for_task(self, task_id: str, timeout: Optional[float] = None) -> Any:
"""等待任务完成"""
if task_id not in self.tasks:
if task_id in self.task_results:
return self.task_results[task_id]
elif task_id in self.task_errors:
raise self.task_errors[task_id]
else:
raise ValueError(f"任务 {task_id} 不存在")

task = self.tasks[task_id]

try:
await asyncio.wait_for(task, timeout)

if task_id in self.task_results:
return self.task_results[task_id]
elif task_id in self.task_errors:
raise self.task_errors[task_id]
else:
raise ValueError(f"任务 {task_id} 没有结果")

except asyncio.TimeoutError:
raise TimeoutError(f"任务 {task_id} 超时")

开发最佳实践

1. 代码组织

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
mcp-project/
├── src/
│ ├── server/ # 服务器核心代码
│ │ ├── __init__.py
│ │ ├── server.py # 服务器主类
│ │ ├── middleware.py # 中间件
│ │ └── lifecycle.py # 生命周期管理
│ ├── tools/ # 工具模块
│ │ ├── __init__.py
│ │ ├── weather.py # 天气工具
│ │ ├── file_manager.py # 文件管理工具
│ │ └── calculator.py # 计算器工具
│ ├── resources/ # 资源模块
│ │ ├── __init__.py
│ │ ├── config.py # 配置资源
│ │ └── status.py # 状态资源
│ ├── prompts/ # 提示词模块
│ │ ├── __init__.py
│ │ └── templates.py # 提示词模板
│ ├── utils/ # 工具函数
│ │ ├── __init__.py
│ │ ├── logger.py # 日志工具
│ │ ├── validator.py # 验证工具
│ │ └── cache.py # 缓存工具
│ └── main.py # 应用入口
├── tests/ # 测试代码
├── config/ # 配置文件
├── docker/ # Docker 配置
├── docs/ # 文档
└── requirements.txt # 依赖管理

2. 配置管理

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
import os
import yaml
from typing import Dict, Any, Optional

class ConfigManager:
def __init__(self, config_file: str = "config/config.yaml"):
self.config_file = config_file
self.config: Dict[str, Any] = {}
self.load_config()

def load_config(self):
"""加载配置文件"""
if os.path.exists(self.config_file):
with open(self.config_file, 'r', encoding='utf-8') as f:
self.config = yaml.safe_load(f) or {}

# 加载环境变量配置
self.load_env_config()

def load_env_config(self):
"""从环境变量加载配置"""
env_mapping = {
"MCP_HOST": ("server", "host"),
"MCP_PORT": ("server", "port"),
"MCP_DEBUG": ("server", "debug"),
"API_KEY": ("services", "api_key"),
"LOG_LEVEL": ("logging", "level")
}

for env_var, config_path in env_mapping.items():
if env_var in os.environ:
self.set_config_value(config_path, os.environ[env_var])

def get_config_value(self, path: tuple, default: Any = None) -> Any:
"""获取配置值"""
value = self.config
for key in path:
if isinstance(value, dict) and key in value:
value = value[key]
else:
return default
return value

def set_config_value(self, path: tuple, value: Any):
"""设置配置值"""
current = self.config
for i, key in enumerate(path[:-1]):
if key not in current or not isinstance(current[key], dict):
current[key] = {}
current = current[key]

current[path[-1]] = value

def get_server_config(self) -> Dict[str, Any]:
"""获取服务器配置"""
return {
"host": self.get_config_value(("server", "host"), "0.0.0.0"),
"port": self.get_config_value(("server", "port"), 8080),
"debug": self.get_config_value(("server", "debug"), False)
}

def get_logging_config(self) -> Dict[str, Any]:
"""获取日志配置"""
return {
"level": self.get_config_value(("logging", "level"), "INFO"),
"file": self.get_config_value(("logging", "file"), "logs/mcp-server.log")
}

3. 错误处理

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
import logging
from typing import Dict, Any, Callable, Coroutine

class ErrorHandler:
def __init__(self, logger: logging.Logger):
self.logger = logger

def handle_errors(self, func: Callable[..., Coroutine[Any, Any, Any]]) -> Callable:
"""错误处理装饰器"""
async def wrapper(*args, **kwargs):
try:
return await func(*args, **kwargs)
except ValueError as e:
self.logger.warning(f"参数错误: {e}")
return self._create_error_response(-32602, str(e))
except PermissionError as e:
self.logger.warning(f"权限错误: {e}")
return self._create_error_response(-32603, str(e))
except Exception as e:
self.logger.error(f"服务器错误: {e}", exc_info=True)
return self._create_error_response(-32603, "服务器内部错误")

return wrapper

def _create_error_response(self, code: int, message: str, data: Any = None) -> Dict[str, Any]:
"""创建错误响应"""
error = {
"code": code,
"message": message
}

if data is not None:
error["data"] = data

return {
"jsonrpc": "2.0",
"error": error
}

# 使用示例
class MyServer(MCPServer):
def __init__(self):
super().__init__(name="MyServer", version="1.0.0")
self.error_handler = ErrorHandler(self.logger)
self.register_tools()

def register_tools(self):
tool: ToolRegistration = {
"name": "my_tool",
"description": "我的工具",
"parameters": {"param": {"type": "string", "required": True}},
"handler": self.error_handler.handle_errors(self.handle_my_tool)
}
self.register_tool(tool)

运维最佳实践

1. 监控指标

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
from prometheus_client import Counter, Gauge, Histogram, start_http_server
import time

class MetricsCollector:
def __init__(self):
# 会话相关指标
self.active_sessions = Gauge(
"mcp_active_sessions",
"当前活跃会话数"
)
self.session_total = Counter(
"mcp_session_total",
"总会话数"
)

# 工具调用指标
self.tool_calls = Counter(
"mcp_tool_calls_total",
"工具调用总数",
["tool_name", "status"]
)
self.tool_duration = Histogram(
"mcp_tool_duration_seconds",
"工具调用耗时",
["tool_name"]
)

# 错误指标
self.error_total = Counter(
"mcp_errors_total",
"错误总数",
["error_type"]
)

def record_session_start(self):
"""记录会话开始"""
self.session_total.inc()
self.active_sessions.inc()

def record_session_end(self):
"""记录会话结束"""
self.active_sessions.dec()

def record_tool_call(self, tool_name, duration, success=True):
"""记录工具调用"""
status = "success" if success else "failure"
self.tool_calls.labels(tool_name=tool_name, status=status).inc()
self.tool_duration.labels(tool_name=tool_name).observe(duration)

def record_error(self, error_type):
"""记录错误"""
self.error_total.labels(error_type=error_type).inc()

2. 健康检查

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
from typing import Dict, Any
import asyncio

class HealthChecker:
def __init__(self, server):
self.server = server
self.services = []

def register_service_check(self, name: str, check_func):
"""注册服务检查"""
self.services.append((name, check_func))

async def check_health(self) -> Dict[str, Any]:
"""执行健康检查"""
status = {
"status": "healthy",
"timestamp": time.time(),
"server": {
"name": self.server.name,
"version": self.server.version,
"uptime": time.time() - self.server.start_time
},
"services": {}
}

# 检查各个服务
for name, check_func in self.services:
try:
service_status = await check_func()
status["services"][name] = service_status

if service_status.get("status") != "healthy":
status["status"] = "unhealthy"

except Exception as e:
status["services"][name] = {
"status": "error",
"error": str(e)
}
status["status"] = "unhealthy"

return status

# 使用示例
class MyServer(MCPServer):
def __init__(self):
super().__init__(name="MyServer", version="1.0.0")
self.start_time = time.time()
self.health_checker = HealthChecker(self)
self.setup_health_checks()

def setup_health_checks(self):
"""设置健康检查"""
# 检查数据库连接
self.health_checker.register_service_check(
"database",
self.check_database
)

# 检查外部API
self.health_checker.register_service_check(
"external_api",
self.check_external_api
)

async def check_database(self):
"""检查数据库连接"""
try:
# 数据库连接检查逻辑
return {"status": "healthy", "latency": 0.1}
except Exception as e:
return {"status": "unhealthy", "error": str(e)}

async def check_external_api(self):
"""检查外部API"""
try:
# 外部API检查逻辑
return {"status": "healthy", "latency": 0.5}
except Exception as e:
return {"status": "unhealthy", "error": str(e)}

async def handle_health_request(self, request):
"""处理健康检查请求"""
return await self.health_checker.check_health()

官方资源