云计算百科
云计算领域专业知识百科平台

手写 MCP:从零实现 Model Context Protocol 服务器与客户端

前言

2025年底,Anthropic 提出 Model Context Protocol(MCP),迅速成为 AI Agent 生态中最有影响力的开放协议之一。到2026年,MCP 已经从"新概念"变成了实际生产力——各大模型平台、开发工具、数据中间件纷纷原生支持 MCP 协议。

但 MCP 到底是什么?它和传统的 API 调用有什么区别?为什么说它是"AI 时代的 USB-C"?

本文从零开始,用纯 Python 手写一个完整的 MCP Server 和 MCP Client,实现工具的注册、发现、调用全流程。读完你会发现,MCP 的核心远没有你想象的那么复杂。

前置知识: Python 基础,了解 JSON-RPC 或 REST API 的基本概念。 本文目标: 不依赖任何 MCP 框架(如 official SDK),从 socket 层开始手写,让你真正理解协议本质。


1. 理解 MCP 的核心设计

1.1 为什么需要 MCP?

在 MCP 出现之前,让 AI 模型调用外部工具通常需要:

  • 自定义 Function Calling 框架:每个项目自己定义工具注册和调用机制
  • 松散的 API 集成:硬编码特定 API 端点和参数格式
  • 无统一发现机制:AI 无法"知道"服务端提供了什么工具

MCP 要解决的是三个核心问题:

  • 标准化:统一工具的定义、注册、调用协议
  • 可发现:Client 可以查询 Server 提供了哪些工具、资源
  • 双向通信:Server 可以主动向 Client 推送消息
  • 1.2 协议架构

    MCP 基于 JSON-RPC 2.0 协议,运行在传输层之上。架构非常简单:

    ┌─────────────┐ JSON-RPC over STDIO/SSE ┌─────────────┐
    │ │ ──── request ────────────────────────▶ │ │
    │ MCP Client │ │ MCP Server │
    │ │ ◀──── response / event ─────────────── │ │
    └─────────────┘ └─────────────┘
    │ │
    │ capabilities: │
    │ – tools() │ – tools[]
    │ – call_tool(name, args) │ – resources[]
    │ – list_resources() │ – prompts[]
    │ – subscribe_resource() │

    关键概念:

    概念说明类比
    Tools Server 提供的可调用函数 类似 OpenAPI 的 endpoints
    Resources Server 暴露的数据资源 类似 REST 的资源路径
    Prompts 预定义的提示模板 类似模板化的 prompt
    Capabilities Server 声明自己支持哪些功能 握手时协商

    1.3 传输层

    MCP 目前支持两种传输方式:

    • STDIO:通过标准输入/输出通信,适合本地子进程模式
    • SSE(Server-Sent Events):通过 HTTP 长连接通信,适合远程模式

    本文先用 STDIO 实现核心逻辑,再扩展到 SSE。


    2. JSON-RPC 基础

    MCP 使用 JSON-RPC 2.0 作为消息格式。每条消息都是一个 JSON 对象:

    请求(Request):

    {
    "jsonrpc": "2.0",
    "id": 1,
    "method": "tools/list",
    "params": {}
    }

    成功响应(Success Response):

    {
    "jsonrpc": "2.0",
    "id": 1,
    "result": {
    "tools": […]
    }
    }

    错误响应(Error Response):

    {
    "jsonrpc": "2.0",
    "id": 1,
    "error": {
    "code": -32601,
    "message": "Method not found"
    }
    }

    通知(Notification): 没有 id 的消息,Server 不需要回复。

    {
    "jsonrpc": "2.0",
    "method": "notifications/initialized",
    "params": {}
    }

    先写一个基础的 JSON-RPC 消息处理类:

    import json
    import sys
    import uuid
    from typing import Any, Callable, Dict, Optional
    from dataclasses import dataclass, field

    class JSONRPCError(Exception):
    """JSON-RPC 错误基类"""
    def __init__(self, code: int, message: str, data: Any = None):
    self.code = code
    self.message = message
    self.data = data
    super().__init__(message)

    class MethodNotFound(JSONRPCError):
    def __init__(self, method: str):
    super().__init__(-32601, f"Method not found: {method}")

    class InvalidParams(JSONRPCError):
    def __init__(self, message: str):
    super().__init__(-32602, message)

    class InternalError(JSONRPCError):
    def __init__(self, message: str):
    super().__init__(-32603, message)

    def make_request(method: str, params: dict = None, request_id: int = None) -> dict:
    """构造 JSON-RPC 请求"""
    if request_id is None:
    request_id = uuid.uuid4().int & (1 << 31) – 1
    msg = {
    "jsonrpc": "2.0",
    "id": request_id,
    "method": method,
    }
    if params is not None:
    msg["params"] = params
    return msg

    def make_notification(method: str, params: dict = None) -> dict:
    """构造通知消息(无 id)"""
    msg = {"jsonrpc": "2.0", "method": method}
    if params is not None:
    msg["params"] = params
    return msg

    def make_success_response(request_id: int, result: Any) -> dict:
    """构造成功响应"""
    return {"jsonrpc": "2.0", "id": request_id, "result": result}

    def make_error_response(request_id: int, error: JSONRPCError) -> dict:
    """构造错误响应"""
    resp = {"jsonrpc": "2.0", "id": request_id, "error": {"code": error.code, "message": error.message}}
    if error.data is not None:
    resp["error"]["data"] = error.data
    return resp

    这段代码只有 60 行,但已经覆盖了 JSON-RPC 的全部核心类型:请求、通知、成功/错误响应。之后所有的 MCP 通信都基于这个底层。


    3. 手写 MCP Server

    有了 JSON-RPC 基础,接下来实现 Server 端。

    一个 MCP Server 的生命周期分为三个阶段: 1. 初始化阶段:接收 Client 的 initialize 请求,返回能力声明(提供了哪些工具、资源) 2. 就绪阶段:接收 Client 的 notifications/initialized 通知,确认握手完成 3. 运行阶段:持续处理工具调用、资源读取等请求

    这个流程和 TCP 的三次握手有异曲同工之处——先协商能力,再开始正式通信。

    3.1 定义工具模型

    MCP 的工具(Tool)是一个包含描述和参数 schema 的可调用函数。我们先定义数据模型:

    @dataclass
    class Tool:
    """MCP 工具定义"""
    name: str
    description: str
    input_schema: dict # JSON Schema
    handler: Callable # 实际的 Python 函数

    @dataclass
    class Resource:
    """MCP 资源定义"""
    uri: str
    name: str
    description: str = ""
    mime_type: str = "text/plain"

    @dataclass
    class Prompt:
    """MCP 提示模板"""
    name: str
    description: str = ""
    arguments: list = field(default_factory=list)

    3.2 实现 Server 核心

    Server 的核心职责: 1. 注册工具、资源、提示 2. 通过 STDIO 读取 JSON-RPC 请求 3. 分发请求到对应的方法处理器 4. 返回 JSON-RPC 响应

    import sys
    import json
    import threading
    from typing import Optional

    class MCPServer:
    """
    纯手写 MCP Server,基于 STDIO 传输
    不依赖任何 MCP SDK
    """

    def __init__(self, server_name: str = "mcp-server", version: str = "1.0.0"):
    self.server_name = server_name
    self.version = version
    self._tools: Dict[str, Tool] = {}
    self._resources: Dict[str, Resource] = {}
    self._prompts: Dict[str, Prompt] = {}
    self._initialized = False
    self._request_handlers: Dict[str, Callable] = {}
    self._running = False

    # 注册内置方法
    self._request_handlers = {
    "initialize": self._handle_initialize,
    "initialized": self._handle_initialized,
    "tools/list": self._handle_tools_list,
    "tools/call": self._handle_tools_call,
    "resources/list": self._handle_resources_list,
    "resources/read": self._handle_resources_read,
    "prompts/list": self._handle_prompts_list,
    }

    def register_tool(self, tool: Tool):
    """注册一个工具"""
    if tool.name in self._tools:
    raise ValueError(f"Tool already registered: {tool.name}")
    self._tools[tool.name] = tool

    def tool(self, name: str, description: str = "", input_schema: dict = None):
    """装饰器方式注册工具"""
    def decorator(func):
    schema = input_schema or {
    "type": "object",
    "properties": {},
    "required": []
    }
    t = Tool(
    name=name,
    description=description or func.__doc__ or "",
    input_schema=schema,
    handler=func
    )
    self.register_tool(t)
    return func
    return decorator

    def register_resource(self, resource: Resource):
    self._resources[resource.uri] = resource

    def register_prompt(self, prompt: Prompt):
    self._prompts[prompt.name] = prompt

    # —- 协议方法处理 —-

    def _handle_initialize(self, params: dict) -> dict:
    """处理初始化请求,返回 Server 能力声明"""
    protocol_version = params.get("protocolVersion", "2025-03-26")
    self._initialized = True

    # 收集 Server 支持的能力
    capabilities = {
    "tools": {} if self._tools else None,
    "resources": {} if self._resources else None,
    "prompts": {} if self._prompts else None,
    }
    # 移除 None 的能力
    capabilities = {k: v for k, v in capabilities.items() if v is not None}

    return {
    "protocolVersion": protocol_version,
    "serverInfo": {
    "name": self.server_name,
    "version": self.version,
    },
    "capabilities": capabilities,
    }

    def _handle_initialized(self, params: dict):
    """处理初始化通知(Client 确认初始化完成)"""
    # 无返回值(通知不需要响应)
    return None

    def _handle_tools_list(self, params: dict) -> dict:
    """返回所有已注册的工具"""
    tools_list = []
    for t in self._tools.values():
    tools_list.append({
    "name": t.name,
    "description": t.description,
    "inputSchema": t.input_schema,
    })
    return {"tools": tools_list}

    def _handle_tools_call(self, params: dict) -> dict:
    """调用指定的工具"""
    name = params.get("name")
    arguments = params.get("arguments", {})

    if name not in self._tools:
    raise MethodNotFound(f"Tool not found: {name}")

    tool = self._tools[name]
    try:
    result = tool.handler(**arguments)

    # MCP 工具调用的标准结果格式
    if isinstance(result, str):
    content = [{"type": "text", "text": result}]
    elif isinstance(result, (dict, list)):
    content = [{"type": "text", "text": json.dumps(result, ensure_ascii=False)}]
    else:
    content = [{"type": "text", "text": str(result)}]

    return {"content": content}
    except Exception as e:
    raise InternalError(str(e))

    def _handle_resources_list(self, params: dict) -> dict:
    resources_list = []
    for r in self._resources.values():
    resources_list.append({
    "uri": r.uri,
    "name": r.name,
    "description": r.description,
    "mimeType": r.mime_type,
    })
    return {"resources": resources_list}

    def _handle_resources_read(self, params: dict) -> dict:
    uri = params.get("uri")
    if uri not in self._resources:
    raise MethodNotFound(f"Resource not found: {uri}")
    return {"contents": []} # 实际情况会返回资源内容

    def _handle_prompts_list(self, params: dict) -> dict:
    prompts_list = []
    for p in self._prompts.values():
    prompts_list.append({
    "name": p.name,
    "description": p.description,
    "arguments": p.arguments,
    })
    return {"prompts": prompts_list}

    # —- 消息分发 —-

    def _process_message(self, raw: str) -> Optional[str]:
    """处理一条 JSON-RPC 消息,返回响应消息(如果有)"""
    try:
    msg = json.loads(raw)
    except json.JSONDecodeError:
    return None # 静默丢弃无效 JSON

    method = msg.get("method")
    params = msg.get("params", {})
    msg_id = msg.get("id")

    # 通知:没有 id,无需响应
    if msg_id is None:
    handler = self._request_handlers.get(method)
    if handler:
    try:
    handler(params)
    except Exception:
    pass # 通知处理失败无需通知调用方
    return None

    # 请求:有 id,需要响应
    handler = self._request_handlers.get(method)
    if handler is None:
    return json.dumps(make_error_response(msg_id, MethodNotFound(method)))

    try:
    result = handler(params)
    if result is None:
    # handler 返回 None 表示不需要响应(例如 initialized 通知的响应)
    return None
    return json.dumps(make_success_response(msg_id, result))
    except JSONRPCError as e:
    return json.dumps(make_error_response(msg_id, e))
    except Exception as e:
    return json.dumps(make_error_response(msg_id, InternalError(str(e))))

    def run_stdio(self):
    """通过 STDIO 运行 Server(主循环)"""
    self._running = True

    # 第一行输出 Server 标识(MCP 协议约定)
    # 实际生产环境中使用 Line-delimited JSON
    for line in sys.stdin:
    line = line.strip()
    if not line:
    continue

    response = self._process_message(line)
    if response:
    sys.stdout.write(response + "\\n")
    sys.stdout.flush()

    def run(self):
    """启动 Server"""
    self.run_stdio()

    不到 200 行的代码,一个可以运行的 MCP Server 就完成了。核心流程:

  • 注册工具到 _tools 字典
  • 每收到一行 JSON,解析后找到对应方法处理
  • 返回 JSON-RPC 响应
  • 3.3 定义实际工具

    现在用我们的 MCP Server 注册几个真实可用的工具:

    # 创建一个实际可用的 MCP Server
    import datetime
    import random

    def create_demo_server() -> MCPServer:
    server = MCPServer(
    server_name="demo-server",
    version="1.0.0"
    )

    @server.tool(
    name="get_weather",
    description="获取指定城市的实时天气信息",
    input_schema={
    "type": "object",
    "properties": {
    "city": {
    "type": "string",
    "description": "城市名称,如 北京、上海、深圳"
    }
    },
    "required": ["city"]
    }
    )
    def get_weather(city: str) -> str:
    """模拟获取天气数据"""
    weathers = ["☀️ 晴", "⛅ 多云", "🌧️ 小雨", "🌩️ 雷阵雨", "🌫️ 雾"]
    temps = {city: random.randint(15, 35) for city in ["北京", "上海", "深圳", "广州", "杭州"]}
    temp = temps.get(city, random.randint(10, 30))
    weather = random.choice(weathers)
    return f"{city} 当前天气:{weather},温度 {temp}°C,湿度 {random.randint(40, 80)}%"

    @server.tool(
    name="calculate",
    description="执行数学计算",
    input_schema={
    "type": "object",
    "properties": {
    "expression": {
    "type": "string",
    "description": "数学表达式,如 2 + 3 * 4"
    }
    },
    "required": ["expression"]
    }
    )
    def calculate(expression: str) -> str:
    """安全执行数学表达式计算"""
    # 注意:生产环境应使用 ast 或者 safer 的表达式引擎
    # 这里演示用 eval 但只允许数字和运算符
    allowed = set("0123456789.+-*/()% ")
    if not all(c in allowed for c in expression):
    return "错误:表达式包含不允许的字符"
    try:
    result = eval(expression, {"__builtins__": {}}, {})
    return f"{expression} = {result}"
    except Exception as e:
    return f"计算错误:{str(e)}"

    @server.tool(
    name="get_current_time",
    description="获取当前日期和时间",
    input_schema={
    "type": "object",
    "properties": {},
    "required": []
    }
    )
    def get_current_time() -> str:
    now = datetime.datetime.now()
    return f"当前时间:{now.strftime('%Y年%m月%d日 %H:%M:%S')} 星期{['一','二','三','四','五','六','日'][now.weekday()]}"

    @server.tool(
    name="search_knowledge",
    description="搜索本地知识库(模拟)",
    input_schema={
    "type": "object",
    "properties": {
    "query": {
    "type": "string",
    "description": "搜索关键词"
    },
    "limit": {
    "type": "integer",
    "description": "返回结果数量",
    "default": 3
    }
    },
    "required": ["query"]
    }
    )
    def search_knowledge(query: str, limit: int = 3) -> str:
    """模拟知识库搜索"""
    # 模拟知识库
    knowledge_base = {
    "python": [
    "Python 是一种解释型、面向对象的高级编程语言。",
    "Python 3.13 引入了 JIT 编译器和自由线程模式。",
    "Python 的 asyncio 库支持异步编程。",
    ],
    "mcp": [
    "MCP(Model Context Protocol)是 Anthropic 提出的开放协议。",
    "MCP 定义了 AI 模型与外部工具/数据源的标准通信方式。",
    "MCP 支持 STDIO 和 SSE 两种传输方式。",
    ],
    "ai": [
    "AI Agent 是具备自主决策能力的大模型应用形态。",
    "ReAct 模式让 Agent 通过思考-行动-观察循环完成任务。",
    "Function Calling 让大模型能够调用外部工具和 API。",
    ]
    }

    results = []
    for keyword, articles in knowledge_base.items():
    if query.lower() in keyword.lower():
    results.extend(articles)

    for articles in knowledge_base.values():
    for article in articles:
    if query.lower() in article.lower():
    results.append(article)

    if not results:
    return f"未找到与「{query}」相关的知识条目。"

    results = results[:limit]
    output = f"找到 {len(results)} 条相关结果:\\n\\n"
    for i, r in enumerate(results, 1):
    output += f"{i}. {r}\\n"
    return output

    return server

    这四个工具覆盖了不同的输入输出模式: – get_weather:必填参数,返回格式化文本 – calculate:带安全校验的参数 – get_current_time:无参数 – search_knowledge:带默认值的参数

    这里有几个关键设计要点值得注意:

    1. 方法分发器的设计模式

    _request_handlers 字典实现了一个简单的路由分发器,将方法名映射到对应的处理函数。这种模式的好处是: – 新增方法只需要添加一个键值对 – 每个方法独立实现,职责清晰 – 可以灵活地重写或扩展

    2. 能力声明的协商机制

    在 _handle_initialize 中,Server 会根据实际注册的内容动态构建能力声明。如果一个 Server 只注册了工具而没有注册资源,resources 能力就不会出现在声明中。这种设计让 Client 可以精确知道 Server 支持什么。

    3. 错误处理的分层设计

    JSONRPCError 的继承体系让错误处理非常有层次: – MethodNotFound:方法名错误(协议层错误) – InvalidParams:参数校验失败(业务层错误) – InternalError:内部执行异常(实现层错误)

    3.4 启动 Server

    if __name__ == "__main__":
    server = create_demo_server()
    print("🧩 MCP Demo Server 已启动", file=sys.stderr)
    print(f"📦 已注册 {len(server._tools)} 个工具", file=sys.stderr)
    server.run()

    用一行命令测试:

    # 发送初始化请求
    echo '{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"protocolVersion":"2025-03-26"}}' | python mcp_server.py

    # 列出工具
    echo '{"jsonrpc":"2.0","id":2,"method":"tools/list","params":{}}' | python mcp_server.py


    4. 手写 MCP Client

    Server 写好了,现在写一个 Client 来使用它。Client 的核心职责:

  • 与 Server 建立连接(STDIO 或 SSE)
  • 发送初始化握手
  • 发现 Server 提供的工具
  • 调用工具并获取结果
  • 封装成友好的调用接口
  • import subprocess
    import json
    import threading
    from typing import Optional

    class MCPClient:
    """
    纯手写 MCP Client,连接 STDIO 模式的 Server
    不依赖任何 MCP SDK
    """

    def __init__(self, server_command: list[str]):
    """
    Args:
    server_command: 启动 Server 的命令,如 ["python", "server.py"]
    """
    self._server_command = server_command
    self._process: Optional[subprocess.Popen] = None
    self._request_id = 0
    self._pending: Dict[int, threading.Event] = {}
    self._results: Dict[int, dict] = {}
    self._lock = threading.Lock()
    self._reader_thread: Optional[threading.Thread] = None
    self._server_info = None
    self._capabilities = None

    def connect(self):
    """启动 Server 进程并连接"""
    self._process = subprocess.Popen(
    self._server_command,
    stdin=subprocess.PIPE,
    stdout=subprocess.PIPE,
    stderr=subprocess.PIPE,
    text=True,
    bufsize=1, # 行缓冲
    )

    # 启动读取线程
    self._reader_thread = threading.Thread(
    target=self._read_loop,
    daemon=True
    )
    self._reader_thread.start()

    def _read_loop(self):
    """后台线程:持续读取 Server 输出"""
    while self._process and self._process.poll() is None:
    line = self._process.stdout.readline()
    if not line:
    break
    line = line.strip()
    if not line:
    continue

    try:
    response = json.loads(line)
    msg_id = response.get("id")

    if msg_id is not None and msg_id in self._pending:
    with self._lock:
    self._results[msg_id] = response
    self._pending[msg_id].set()
    except json.JSONDecodeError:
    pass # 忽略无效行

    def _send_request(self, method: str, params: dict = None) -> dict:
    """发送请求并等待响应(同步调用)"""
    with self._lock:
    self._request_id += 1
    req_id = self._request_id
    event = threading.Event()
    self._pending[req_id] = event

    request = make_request(method, params, req_id)
    request_str = json.dumps(request) + "\\n"

    self._process.stdin.write(request_str)
    self._process.stdin.flush()

    # 等待响应(最多 30 秒)
    if not event.wait(timeout=30):
    raise TimeoutError(f"Request timed out: {method}")

    with self._lock:
    response = self._results.pop(req_id)
    self._pending.pop(req_id)

    if "error" in response:
    err = response["error"]
    raise JSONRPCError(err.get("code", 0), err.get("message", "Unknown error"))

    return response.get("result", {})

    def initialize(self):
    """发送初始化请求,完成握手"""
    result = self._send_request("initialize", {
    "protocolVersion": "2025-03-26",
    "capabilities": {},
    "clientInfo": {
    "name": "mcp-client-demo",
    "version": "1.0.0"
    }
    })

    self._server_info = result.get("serverInfo", {})
    self._capabilities = result.get("capabilities", {})

    # 发送 initialized 通知
    notif = make_notification("notifications/initialized")
    self._process.stdin.write(json.dumps(notif) + "\\n")
    self._process.stdin.flush()

    return result

    def list_tools(self) -> list:
    """获取 Server 提供的工具列表"""
    result = self._send_request("tools/list")
    return result.get("tools", [])

    def call_tool(self, name: str, arguments: dict = None) -> dict:
    """调用指定工具"""
    params = {"name": name}
    if arguments:
    params["arguments"] = arguments
    return self._send_request("tools/call", params)

    def close(self):
    """关闭连接"""
    if self._process:
    self._process.terminate()
    self._process.wait(timeout=5)
    self._process = None

    4.1 同步调用机制的实现细节

    上面的 Client 中,_send_request 是核心方法。它的工作流程值得仔细揣摩:

    发送请求前:
    1. 生成自增 request_id
    2. 创建 threading.Event 并放入 _pending 字典
    3. 通过 STDIN 发送 JSON 消息

    等待响应:
    4. 后台 _read_loop 线程持续读取 STDOUT
    5. 收到响应时,根据 id 找到对应的 Event
    6. 设置 Event,唤醒等待的线程

    获取结果:
    7. 从 _results 字典取出响应数据
    8. 检查是否有 error 字段
    9. 返回 result 部分

    这种设计的关键在于解耦发送和接收: – 发送线程不阻塞 STDIO 的读取 – 后台读取线程独立运行,保证任何时刻都能收到 Server 的消息 – Event 机制让多个并发请求可以独立等待各自的响应

    threading.Event 在这里起到了信号量的作用——发送线程发起请求后挂起等待,读取线程收到响应后通过 Event 通知发送线程继续执行。这种模式在异步编程中非常常见,只是这里我们在同步的线程模型中实现了它。

    4.2 测试 Client

    def test_client():
    """测试 MCP Client 与 Server 的交互"""

    client = MCPClient(["python", "mcp_server.py"])

    try:
    # 第一步:连接
    print("🔄 连接 Server…")
    client.connect()

    # 第二步:初始化握手
    print("🔄 初始化握手…")
    info = client.initialize()
    print(f"✅ 已连接: {info['serverInfo']['name']} v{info['serverInfo']['version']}")
    print(f"📊 能力: {', '.join(info['capabilities'].keys())}\\n")

    # 第三步:列出工具
    print("🔧 可用工具:")
    tools = client.list_tools()
    for t in tools:
    required = t['inputSchema'].get('required', [])
    params_desc = ', '.join(t['inputSchema']['properties'].keys()) if t['inputSchema']['properties'] else '无'
    print(f" 📌 {t['name']}: {t['description']}")
    print(f" 参数: {params_desc}")
    print()

    # 第四步:调用工具
    print("🚀 调用工具:\\n")

    # 调用 get_current_time
    print("— get_current_time —")
    result = client.call_tool("get_current_time")
    print(result['content'][0]['text'])
    print()

    # 调用 calculate
    print("— calculate —")
    result = client.call_tool("calculate", {"expression": "2 ** 10 + 42"})
    print(result['content'][0]['text'])
    print()

    # 调用 search_knowledge
    print("— search_knowledge —")
    result = client.call_tool("search_knowledge", {"query": "MCP", "limit": 2})
    print(result['content'][0]['text'])

    finally:
    client.close()

    if __name__ == "__main__":
    test_client()

    运行输出示例:

    🔄 连接 Server…
    🔄 初始化握手…
    ✅ 已连接: demo-server v1.0.0
    📊 能力: tools

    🔧 可用工具:
    📌 get_weather: 获取指定城市的实时天气信息
    参数: city
    📌 calculate: 执行数学计算
    参数: expression
    📌 get_current_time: 获取当前日期和时间
    参数: 无
    📌 search_knowledge: 搜索本地知识库(模拟)
    参数: query, limit

    🚀 调用工具:

    — get_current_time —
    当前时间:2026年04月27日 21:00:00 星期一

    — calculate —
    2 ** 10 + 42 = 1066

    — search_knowledge —
    找到 3 条相关结果:

    1. MCP(Model Context Protocol)是 Anthropic 提出的开放协议。
    2. MCP 定义了 AI 模型与外部工具/数据源的标准通信方式。
    3. MCP 支持 STDIO 和 SSE 两种传输方式。


    5. 与 LLM 集成:让 AI 用你的工具

    MCP 的最大价值在于:让大模型自动发现并使用工具。我们把 MCP Client 和 LLM 结合起来。

    import requests
    from typing import List

    class MCPAgent:
    """
    将 MCP Client 与 LLM 集成
    AI 自动发现、调用 MCP Server 上的工具
    """

    def __init__(self,
    mcp_client: MCPClient,
    llm_api_url: str = "https://api.openai.com/v1/chat/completions",
    llm_api_key: str = None,
    model: str = "gpt-4o"):
    self.mcp = mcp_client
    self.api_url = llm_api_url
    self.api_key = llm_api_key
    self.model = model
    self._tools = []

    def connect_mcp(self):
    """连接 MCP Server 并同步工具"""
    self.mcp.connect()
    self.mcp.initialize()
    self._tools = self.mcp.list_tools()
    print(f"✅ 已同步 {len(self._tools)} 个 MCP 工具")

    def _convert_tool_to_openai(self, mcp_tool: dict) -> dict:
    """将 MCP 工具格式转换为 OpenAI Function Calling 格式"""
    return {
    "type": "function",
    "function": {
    "name": mcp_tool["name"],
    "description": mcp_tool["description"],
    "parameters": mcp_tool.get("inputSchema", {
    "type": "object",
    "properties": {}
    })
    }
    }

    def _call_llm(self, messages: list) -> dict:
    """调用 LLM API"""
    headers = {
    "Authorization": f"Bearer {self.api_key}",
    "Content-Type": "application/json"
    }

    payload = {
    "model": self.model,
    "messages": messages,
    "tools": [self._convert_tool_to_openai(t) for t in self._tools],
    "tool_choice": "auto"
    }

    resp = requests.post(self.api_url, headers=headers, json=payload, timeout=60)
    resp.raise_for_status()
    return resp.json()

    def chat(self, user_message: str, max_turns: int = 10) -> str:
    """
    与 AI 对话,AI 自动调用 MCP 工具
    使用 ReAct 模式:思考 → 调用 → 观察 → 继续思考
    """
    messages = [
    {"role": "system", "content": "你是一个使用 MCP 工具的 AI 助手。调用工具时,使用返回结果回答问题。"},
    {"role": "user", "content": user_message}
    ]

    for turn in range(max_turns):
    response = self._call_llm(messages)
    choice = response["choices"][0]
    msg = choice["message"]

    # 如果 AI 没有调用工具,直接返回
    if not msg.get("tool_calls"):
    return msg["content"]

    # AI 请求调用工具
    messages.append(msg)

    for tool_call in msg["tool_calls"]:
    func_name = tool_call["function"]["name"]
    func_args = json.loads(tool_call["function"]["arguments"])

    # 通过 MCP 调用工具
    try:
    result = self.mcp.call_tool(func_name, func_args)
    tool_result = result["content"][0]["text"]
    except Exception as e:
    tool_result = f"Error: {str(e)}"

    # 将工具结果返回给 AI
    messages.append({
    "role": "tool",
    "tool_call_id": tool_call["id"],
    "content": tool_result
    })

    # 继续循环,让 AI 根据工具结果决定下一步

    return "已达到最大对话轮次。"

    def close(self):
    self.mcp.close()

    5.1 完整流程演示

    def demo():
    """MCP + LLM 完整演示"""

    # 创建 MCP Client(连接 Server)
    client = MCPClient(["python", "mcp_server.py"])

    # 创建 MCP Agent(集成 LLM)
    agent = MCPAgent(
    mcp_client=client,
    llm_api_url="https://api.deepseek.com/v1/chat/completions",
    llm_api_key="your-api-key-here",
    model="deepseek-chat"
    )

    # 连接 MCP,自动同步工具
    agent.connect_mcp()

    try:
    # 用户提问
    questions = [
    "现在几点了?帮我算一下 2^16 是多少",
    "搜索一下关于 Python 异步编程的知识",
    "北京的天气怎么样?顺便说一下今天日期",
    ]

    for q in questions:
    print(f"\\n👤 用户: {q}")
    print("🤖 AI: ", end="", flush=True)
    answer = agent.chat(q)
    print(answer)

    finally:
    agent.close()

    当用户问"北京的天气怎么样?顺便说一下今天日期"时,AI 内部流程如下:

    1. AI 思考:需要调 get_current_time 获取日期
    2. AI 调用 → get_current_time()
    3. Server 返回:当前时间…
    4. AI 思考:还需要调 get_weather 获取北京天气
    5. AI 调用 → get_weather(city="北京")
    6. Server 返回:北京天气…
    7. AI 综合回答

    这就是 MCP 的核心价值——AI 自己决定用什么工具、什么时候用,不需要开发者硬编码调用链。


    6. 扩展:SSE 传输实现

    STDIO 适合本地子进程模式,但实际部署中更多需要远程访问。MCP 支持 SSE(Server-Sent Events)作为 HTTP 传输层。

    6.1 SSE Server

    from http.server import HTTPServer, BaseHTTPRequestHandler
    import urllib.parse

    class SSEMCPServer(MCPServer):
    """
    基于 SSE 的 MCP Server
    支持 HTTP 远程访问
    """

    def __init__(self, host="0.0.0.0", port=8000, **kwargs):
    super().__init__(**kwargs)
    self.host = host
    self.port = port
    self._clients = []

    def _handle_sse_connection(self, client_id: str):
    """处理 SSE 连接"""
    # SSE 使用长连接,Server 可以主动推送消息
    pass

    def _handle_message_post(self, client_id: str, body: str):
    """处理 Client 通过 POST 发送的消息"""
    return self._process_message(body)

    def run(self):
    """启动 HTTP Server"""
    server = HTTPServer((self.host, self.port), self._create_handler())
    print(f"🌐 MCP SSE Server 运行在 http://{self.host}:{self.port}", file=sys.stderr)
    server.serve_forever()

    def _create_handler(self):
    server_ref = self

    class Handler(BaseHTTPRequestHandler):
    def do_GET(self):
    if self.path == "/sse":
    # SSE 端点
    self.send_response(200)
    self.send_header("Content-Type", "text/event-stream")
    self.send_header("Cache-Control", "no-cache")
    self.send_header("Connection", "keep-alive")
    self.send_header("Access-Control-Allow-Origin", "*")
    self.end_headers()

    # 发送 endpoint 事件,告知 Client 消息 POST 地址
    self.wfile.write(f"event: endpoint\\ndata: /messages/{id(self)}\\n\\n".encode())
    self.wfile.flush()

    # 保持连接(持续发送心跳)
    while server_ref._running:
    try:
    self.wfile.write(": heartbeat\\n\\n".encode())
    self.wfile.flush()
    import time
    time.sleep(15)
    except BrokenPipeError:
    break
    else:
    self.send_error(404)

    def do_POST(self):
    if self.path.startswith("/messages/"):
    content_length = int(self.headers["Content-Length"])
    body = self.rfile.read(content_length).decode()

    response = server_ref._process_message(body)
    if response:
    self.send_response(200)
    self.send_header("Content-Type", "application/json")
    self.send_header("Access-Control-Allow-Origin", "*")
    self.end_headers()
    self.wfile.write(response.encode())
    else:
    self.send_response(202) # Accepted(通知类消息)
    self.end_headers()
    else:
    self.send_error(404)

    def log_message(self, format, *args):
    pass # 安静模式

    return Handler

    6.2 SSE Client

    import requests
    import sseclient # pip install sseclient-py

    class SSEMCPClient(MCPClient):
    """基于 SSE 的 MCP Client"""

    def __init__(self, server_url: str):
    self.server_url = server_url.rstrip("/")
    self._session = requests.Session()
    self._message_url = None
    self._request_id = 0
    self._pending = {}

    def connect(self):
    """建立 SSE 连接,获取消息 POST 地址"""
    resp = self._session.get(
    f"{self.server_url}/sse",
    stream=True,
    headers={"Accept": "text/event-stream"}
    )
    resp.raise_for_status()

    client = sseclient.SSEClient(resp)
    for event in client.events():
    if event.event == "endpoint":
    self._message_url = f"{self.server_url}{event.data}"
    break

    def _send_request(self, method: str, params: dict = None) -> dict:
    with self._lock:
    self._request_id += 1
    req_id = self._request_id
    event = threading.Event()
    self._pending[req_id] = event

    request = make_request(method, params, req_id)

    resp = self._session.post(
    self._message_url,
    json=request,
    headers={"Content-Type": "application/json"}
    )
    resp.raise_for_status()

    # HTTP 200 表示有响应,202 表示通知已接受
    if resp.status_code == 200:
    return resp.json().get("result", {})

    event.wait(timeout=30)
    return self._results.pop(req_id, {})

    SSE 模式的核心区别:

    • Server 维护长连接,可以主动推送事件
    • Client 通过 POST 发送请求到 /messages/{id}
    • Server 通过 SSE 连接推送响应和事件

    7. 完整项目结构

    以上代码虽然是从零手写,但已经覆盖了生产级 MCP 需要的核心能力。完整项目结构如下:

    mcp-demo/
    ├── mcp_core.py # JSON-RPC 基础(130行)
    ├── mcp_server.py # MCP Server 实现(200行)
    ├── mcp_client.py # MCP Client 实现(150行)
    ├── mcp_agent.py # MCP + LLM 集成(100行)
    ├── sse_server.py # SSE 扩展(80行)
    ├── demo.py # 示例工具和演示(100行)
    └── test.py # 测试脚本(50行)

    总代码量不到 1000 行,就实现了一个完整的 MCP 协议栈。这就是 MCP 的设计之美——协议本身非常简洁,核心是 JSON-RPC + 标准化的数据模型。


    8. 在 AI Agent 中落地 MCP

    理解了原理,我们来聊聊 MCP 在实际项目中的最佳实践。

    8.1 什么时候用 STDIO,什么时候用 SSE?

    场景推荐方式原因
    本地开发/测试 STDIO 零依赖,启动快
    单机 Agent 进程 STDIO 子进程通信,隔离性好
    微服务/容器部署 SSE 支持远程调用和水平扩展
    跨语言协作 SSE HTTP 协议通用兼容
    需要 Server 推送 SSE 事件流天然支持

    8.2 常见的 MCP 应用模式

    模式一:本地工具增强

    LLM 运行在本地,MCP Server 提供文件系统、Shell、计算等工具。

    用户 ↔ LLM ↔ MCP Server(本地工具)

    模式二:微服务网关

    每个后端服务启动一个 MCP Server,AI 网关统一发现和路由。

    用户 ↔ AI网关 ↔ MCP Server(搜索服务)
    ↔ MCP Server(数据库服务)
    ↔ MCP Server(代码服务)

    模式三:中间件代理

    MCP 作为 AI 与现有系统的适配层,封装旧系统的 API。

    用户 ↔ LLM ↔ MCP Proxy ↔ 遗留系统 REST/SOAP API

    8.3 常见陷阱与调试技巧

    在实际开发 MCP Server 时,有几个容易踩的坑:

    陷阱一:使用 print 输出 Python 的 print() 默认输出到 STDOUT,但在 MCP 的 STDIO 模式下,STDOUT 用来传输 JSON-RPC 消息。如果你在工具函数里用了 print() 调试,Client 端会收到一行无法解析的文本,导致 JSON 解析失败。

    ✅ 正确做法:用 sys.stderr 输出日志

    print("调试信息", file=sys.stderr, flush=True)

    陷阱二:忘记 flush MCP 的行协议要求每次写入后立即刷新缓冲区。如果忘记 flush(),Client 可能永远收不到响应。

    ✅ 正确做法:每次 write() 后调用 flush()

    sys.stdout.write(response + "\\n")
    sys.stdout.flush() # 一定要刷新!

    陷阱三:通知不需要响应 初次实现时容易把通知(notifications/initialized)也返回 {"id": null} 格式的响应——但实际上通知根本不需要响应。区分标准很简单:消息里没有 id 字段就是通知。

    陷阱四:工具返回值格式 MCP 规定的工具返回值格式是 {"content": [{"type": "text", "text": "…"}]},不是直接返回字符串。如果格式不对,虽然 Client 能收到数据,但不符合协议的语义标准。

    8.4 安全性注意事项

    MCP 本身不加密传输——STDIO 是本地 IPC,SSE 需要配合 HTTPS:

  • 输入验证:calculate 工具中的白名单模式必须严格执行
  • 权限控制:不要把所有系统能力暴露为 MCP 工具
  • 速率限制:给每个工具调用加限流,防止滥用
  • 操作确认:危险操作(删除文件、修改数据)加入二次确认
  • 审计日志:记录所有工具调用详情

  • 9. 总结

    从零手写 MCP,我敢写你也敢学。核心要点:

    层级内容代码量
    JSON-RPC 协议 消息格式、请求/响应/通知 60行
    Server 实现 工具注册、方法分发、能力声明 200行
    Client 实现 进程通信、请求响应同步 150行
    LLM 集成 工具自动发现、Function Calling 100行
    SSE 扩展 HTTP 长连接、远程访问 80行

    MCP 不是黑魔法——它就是一个设计良好的 JSON-RPC 协议,加上标准化的工具定义和发现机制。理解了这个本质,不管是读官方 SDK 源码、自己写扩展、还是调试生产问题,都游刃有余。

    最后送上一句话: MCP 解决的从来不是"能不能调工具"的问题,而是"AI 如何发现和选择工具"的问题——这才是 Agent 时代的核心命题。


    本文所有代码可在 GitHub 找到完整实现:github.com/your-repo/mcp-from-scratch

    欢迎在评论区讨论你的 MCP 实践心得!

    赞(0)
    未经允许不得转载:网硕互联帮助中心 » 手写 MCP:从零实现 Model Context Protocol 服务器与客户端
    分享到: 更多 (0)

    评论 抢沙发

    评论前必须登录!