一、引言

模型上下文协议(Model Context Protocol,MCP),是由Anthropic推出的开源协议,旨在实现大语言模型与外部数据源和工具的集成,用来在大模型和数据源之间建立安全双向的连接。该协议通过相同的协议同时处理本地资源(例如数据库、文件、服务等)和远程资源(例如Slack或GitHub等API)。

该协议是cs架构设计,客户端和服务端之间支持三种传输方式:stdio、SSE 与 Streamable HTTP,其中SSE逐渐会被Streamable HTTP代替,今天本文实现一个基于stdio传输的简易版mcp sdk。

二、实现

mcp客户端调用mcp服务端的过程大致可以概述为以下几个步骤:

  1. 客户端启动
  2. 客户端通过子进程启动服务端
  3. 客户端获取服务端的工具列表及其他资源
  4. 客户端将获取的工具列表添加到大模型上下文中
  5. 大模型在需要调用工具的时候调用某个工具
  6. 客户端接收到调用逻辑后向服务端请求对应工具的调用逻辑
  7. 服务端收到接收到调用请求后执行具体的工具函数,然后将结果返回给客户端
  8. 客户端将返回的结果添加到大模型上下文消息列表中让大模型继续后续的处理

stdio服务

首先定义一个stdio服务stdio.py,用来统一处理输入和输出:

这里需要先定义两个内存流(读和写),内存流的作用是在各模块间传输数据

  1. stdio服务接收父进程的输入,清洗数据和处理异常情况,将处理后的数据写入内存读流
  2. 其他模块通过内存读流接收输入,执行具体的逻辑,然后将结果写入内存写流
  3. stdio服务通过内存写流接收结果,然后输出给父进程
import sys
from contextlib import asynccontextmanager
from io import TextIOWrapper
import json
import sys

import anyio
from anyio.streams.memory import MemoryObjectReceiveStream, MemoryObjectSendStream

# 将io处理函数包装为上下文管理器,以便在with语句中使用
@asynccontextmanager
async def stdio_server():
    # 将标准输入流(sys.stdin)转换为一个异步文件对象,并使用UTF-8编码
    stdin = anyio.wrap_file(TextIOWrapper(sys.stdin.buffer, encoding="utf-8"))
    stdout = anyio.wrap_file(TextIOWrapper(sys.stdout.buffer, encoding="utf-8"))
    # 定义两个内存流用来在不同模块间传递消息,一个用来接收消息(接收请求),一个用来发送消息(发送响应)
    read_stream: MemoryObjectReceiveStream[str | Exception]
    read_stream_writer: MemoryObjectSendStream[str | Exception]

    write_stream: MemoryObjectSendStream[str]
    write_stream_reader: MemoryObjectReceiveStream[str]

    read_stream_writer, read_stream = anyio.create_memory_object_stream(0)
    write_stream, write_stream_reader = anyio.create_memory_object_stream(0)
    # 定义两个并发执行的异步函数,分别用来读取标准输入和写入标准输出
    async def stdin_reader():
        try:
            async with read_stream_writer:
                async for line in stdin:
                    line = line.strip()
                    if not line:
                        continue
                    try:
                        logger.info(f"Received line: {line}")
                        message = json.loads(line) # 将 JSON 格式的字符串 解析(反序列化)为 Python 对象
                    except Exception as exc:
                        logger.error(f"JSON decode error: {exc}, line: {line}")
                        await read_stream_writer.send(exc)
                        continue
                    await read_stream_writer.send(message) # 模块间传递消息
        except anyio.ClosedResourceError:
            await anyio.lowlevel.checkpoint()

    async def stdout_writer():
        try:
            async with write_stream_reader:
                # 其他模块可以通过write_stream发送消息到这个函数
                async for message in write_stream_reader:
                    logger.info(f"Sending response to client: {message}")
                    await stdout.write(json.dumps(message) + "\n") # 将消息转换为JSON格式并写入标准输出
                    await stdout.flush() # 刷新输出缓冲区,确保消息及时输出
        except anyio.ClosedResourceError:
            await anyio.lowlevel.checkpoint()

    # 并发执行上述两个函数,确保同时处理输入和输出
    async with anyio.create_task_group() as tg:
        tg.start_soon(stdin_reader)
        tg.start_soon(stdout_writer)
        yield read_stream, write_stream

IO服务

定义一个IO服务类server.py作为sdk的入口,在内部调用stdio服务拿到读写内存流,然后分发到具体的工具中使用

import anyio
from stdio import stdio_server
from tools import add


class IOServer:
    async def runIO(self, read_stream: anyio.abc.ObjectReceiveStream[str | Exception], write_stream: anyio.abc.ObjectSendStream[str]) -> None:
        # 接收客户端的消息
        async for message in read_stream:
            if message["method"] == "add":
                result = add(message["args"])
                await write_stream.send({"type": "result", "result": result})
                continue
            # 发送响应给客户端
            await write_stream.send(
                {"status": "success", "message": "Message received"}
            )

    async def run_stdio_async(self) -> None:
        logger.info("Starting stdio server...")
        async with stdio_server() as (read_stream, write_stream):
            await self.runIO(read_stream, write_stream)
           
    def run(self):
        anyio.run(self.run_stdio_async)

工具集

真正的工具是在开发mcp服务时定义的,然后通过装饰器注入到IO服务中来调用,这里简单实现,直接定义一个工具集文件tool.py,里面有一个计算加法运算的函数add

def add(args:list[int]) -> int:
    """一个简单的加法函数,用于测试"""
    return args[0] + args[1]

三、测试

mcp服务端

编写test-server.py

from server import IOServer

mcp = IOServer()

if __name__ == '__main__':
    mcp.run()

mcp客户端

编写test-client.py

import asyncio
import argparse
import json
import sys
from pathlib import Path

process = None
async def communicate_with_server(command):
    # 启动服务器作为子进程
    global process
    # 使用传入的命令启动子进程
    process = await asyncio.create_subprocess_exec(
        *command,
        stdin=asyncio.subprocess.PIPE,
        stdout=asyncio.subprocess.PIPE,
        stderr=asyncio.subprocess.PIPE
    )

    # 创建任务来读取服务器的stderr输出
    async def read_stderr():
        while True:
            line = await process.stderr.readline()
            if not line:
                break
            print(f"Server stderr: {line.decode().strip()}")

    # 启动读取stderr的任务
    asyncio.create_task(read_stderr())
    # 等待服务器启动完成
    await asyncio.sleep(1)


    # # 发送消息给服务器
    await test_quest()
    # 读取服务器的响应,添加超时处理
    try:
        response_line = await asyncio.wait_for(process.stdout.readline(), timeout=5.0)
        if response_line:
            response_text = response_line.decode().strip()
            try:
                response = json.loads(response_text)
                print(f"Received response from server: {response}")
            except json.JSONDecodeError as e:
                print(f"JSON decode error: {e}, response: {response_text}")
    except asyncio.TimeoutError:
        print("Timeout: No response received from server within 5 seconds")
    # 关闭子进程
    process.stdin.close()
    await process.wait()

async def test_quest():
    # 每隔1s给服务器发送一条消息,发三条
    for i in range(3):
        await asyncio.sleep(1)
        message = {"method": "add", "args": [i, i + 1]}
        print("Sending message to server:", message)
        message_json = json.dumps(message) + "\n"
        process.stdin.write(message_json.encode())
        await process.stdin.drain()

if __name__ == "__main__":
    # 解析命令行参数
    parser = argparse.ArgumentParser(description='Test MCP server client')
    parser.add_argument('command', type=str, nargs='+',
                        help='Command to start the server (e.g., "python server.py" or "node server.js")')
    args = parser.parse_args()

    # 运行客户端
    asyncio.run(communicate_with_server(args.command))

运行客户端

python test-client.py python test-server.py

运行结果

可以看到客户端请求add方法后服务端计算后将结果返回给了客户端

四、总结

本文通过简单的示例代码展示了mcp通过标准IO传输方式如何从客户端到服务端的调用,通过本文可以让初学者更容易理解mcp开发的底层原理

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐