拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

从零手写MCP Server:鉴权、流式传输与状态管理的生产级实践

从零手写MCP Server:鉴权、流式传输与状态管理的生产级实践 1. 为什么我要从零手写一个 MCP Server市面上关于 MCP Server 的教程绝大多数停留在“跑通官方 SDK 的 echo demo”这个层面。你照着文档敲一遍确实能在本地跑起来一个能响应tools/list和tools/call的服务但只要往生产环境挪一步问题就全冒出来了多个客户端同时连上来会话状态互相串工具调用返回一个几百 KB 的 JSON客户端直接超时服务裸奔在公网上谁都能调你的工具。这三个问题——鉴权、流式传输、状态管理——恰好是 demo 和 production 之间的分水岭。我这次做的项目是一个面向内部工具链的 MCP Server它需要对接若干后端服务把数据库查询、文件操作、外部 API 调用这些能力以 MCP 工具的形式暴露给 AI 客户端。选型上我没有直接用某个框架一把梭而是基于官方 SDK 从零搭建原因很简单只有自己把鉴权、传输、状态这三层拆开写一遍你才知道每一层该在哪里设边界。框架帮你省掉的那部分代码恰恰是出问题时你最需要看懂的部分。这篇文章适合两类人一是已经跑通过 MCP 基础 demo、想把它推到生产环境的开发者二是正在设计 AI 工具调用网关、需要理解 MCP 协议工程化细节的架构同学。我会把整个实现过程拆成设计思路、核心细节、实操落地、问题排查四个部分每一段都尽量给出可直接抄的代码和参数。全文基于我实际项目中的取舍来写不是文档翻译踩过的坑我会明确标出来。2. 整体架构设计与技术选型拆解2.1 三层职责划分传输层、会话层、工具层在动手写代码之前我先把整个 Server 拆成三层这个划分决定了后面所有代码的组织方式。传输层负责和客户端之间的字节流搬运它不关心你传的是 JSON-RPC 还是别的什么只负责把数据可靠地送出去、收进来。这一层要处理的是连接建立、协议协商、流式分块、背压控制。会话层负责维护“谁在跟我说话”这件事每个连接对应一个会话会话里存着鉴权结果、上下文状态、订阅关系。工具层才是真正干活的每个工具是一个独立的处理单元接收参数、执行逻辑、返回结果。为什么这么分因为这三层的生命周期完全不同。传输层随连接生灭会话层随会话生灭工具层是无状态的、可复用的。如果你把鉴权逻辑写在工具函数里那么每加一个工具就要重复一遍鉴权如果你把会话状态塞在传输层的连接对象上那么一旦引入流式传输、连接复用状态就会错乱。分层不是为了好看是为了让每一层的变化不影响其他层。2.2 为什么选 SSE 而不是 WebSocket 做流式传输MCP 协议支持多种传输方式我最终选了SSEServer-Sent Events作为流式通道配合 HTTP POST 做请求入口。这个选择值得展开说。WebSocket 是全双工看起来更适合“实时”场景但它带来两个额外成本一是连接管理复杂你需要自己实现心跳、重连、消息分片二是很多企业网关、反向代理对 WebSocket 的支持并不透明部署时容易卡在基础设施上。而 SSE 本质上是“长连接的 HTTP 响应”服务端持续往一个text/event-stream的响应里写数据客户端用标准 EventSource 或 fetch 流式读取即可。它天然单向服务端到客户端而 MCP 的请求-响应模型恰好是客户端发起、服务端流式返回方向完全吻合。更关键的是SSE 走的是标准 HTTP意味着你现有的鉴权中间件、限流、日志、链路追踪全都能直接复用不需要为 WebSocket 单独搭一套。代价是 SSE 不支持客户端往同一个连接里推数据但 MCP 的请求本来就通过独立的 POST 发送所以这个“缺点”对我们不成立。注意SSE 连接容易被中间层缓冲。Nginx 默认会缓冲响应导致流式数据被攒成一坨才发出去。必须在 location 里显式关闭缓冲后面实操部分会给配置。2.3 状态管理的核心矛盾有状态会话 vs 无状态工具MCP 的会话是有状态的——客户端初始化时要协商协议版本、能力集后续调用要复用这个上下文。但工具执行本身应该是无状态的同一个工具被不同会话调用不应该互相影响。这两者的矛盾就是状态管理的核心。我的方案是会话状态集中存在一个 SessionStore 里用 sessionId 索引工具执行时只从会话里读取必要的上下文比如鉴权主体、用户偏好执行完不往会话里写业务数据。会话里存的是“身份和配置”不是“业务中间结果”。这样即使会话被回收、重建工具行为也不会变。如果你把业务状态写进会话那么一旦客户端重连换了 sessionId之前的状态就丢了这类 bug 极难排查。3. 鉴权体系的核心细节与实现要点3.1 鉴权放在哪一层传输层拦截还是工具层校验这是我在设计时纠结最久的一个点。把鉴权放在传输层连接建立时校验一次性能最好但粒度粗放在工具层每次调用校验粒度细但重复开销大。最终我采用的是两层结合传输层做身份认证Authentication你是谁工具层做权限校验Authorization你能不能干这件事。身份认证在 SSE 连接建立时完成校验通过后把主体信息subject、角色、权限范围绑定到会话上。之后每次工具调用工具层从会话里取出主体检查这个主体是否有权调用该工具、是否有权访问该工具声明的资源。这样认证只做一次授权按需做既省开销又不失粒度。具体实现上认证我用的是Bearer Token 签名校验。客户端在建立 SSE 连接时通过Authorization头带上 token服务端校验签名和有效期解析出主体信息。这里有个细节SSE 的 EventSource API 原生不支持自定义请求头所以如果你的客户端是浏览器环境得改用 fetch ReadableStream 的方式或者把 token 放在 query 参数里不推荐会进日志。服务端对两种方式都要兼容。3.2 Token 校验的完整实现与防绕过设计Token 校验这块我踩过一个坑最初只校验了签名没校验aud受众和iss签发者结果一个为其他服务签发的 token 也能通过。签名只证明 token 没被篡改不证明这个 token 是发给你的。完整的校验必须包含签名验证、过期时间exp、生效时间nbf、受众aud、签发者iss以及一个自定义的scope字段。import jwt from jwt import PyJWKClient class TokenVerifier: def __init__(self, jwks_url, audience, issuer): self.jwks_client PyJWKClient(jwks_url) self.audience audience self.issuer issuer def verify(self, token: str) - dict: signing_key self.jwks_client.get_signing_key_from_jwt(token) payload jwt.decode( token, signing_key.key, algorithms[RS256], audienceself.audience, issuerself.issuer, options{require: [exp, nbf, aud, iss, sub]}, ) return payload关于“鉴权绕过”我在测试阶段专门做了一轮对抗性验证发现几个必须堵的口子。第一算法混淆攻击如果服务端允许alg: none或者同时接受 HS256 和 RS256攻击者可以把 RS256 的 token 改成 HS256 用公钥当密钥签名。解决办法是algorithms参数写死成单一算法绝不动态读取 token 头里的alg。第二JWKS 缓存投毒如果 JWKS 拉取没有做缓存和校验攻击者可能诱导服务端拉取恶意 JWKS。解决办法是固定 JWKS 来源、设置合理缓存时间、校验 kid 存在性。第三时序攻击字符串比较用会泄露信息涉及密钥比较的地方一律用hmac.compare_digest。3.3 权限模型从粗粒度角色到细粒度 scope认证解决“你是谁”授权解决“你能干什么”。我一开始用的是简单的角色模型admin / user很快发现不够用同样是 userA 只能查自己的数据B 能查全表。于是改成scope 模型每个 token 携带一组 scope比如db:read、db:write、file:read:/home/user/*工具在调用前声明自己需要的 scope会话层做匹配。匹配逻辑我写了一个小函数支持通配符和路径前缀def has_scope(granted: list[str], required: str) - bool: for g in granted: if g required: return True if g.endswith(*) and required.startswith(g[:-1]): return True return False这里有个容易忽略的点scope 的匹配必须是“服务端声明、客户端无法伪造”。也就是说工具需要什么 scope 是服务端代码里写死的不能由客户端在请求里指定。否则客户端可以声称“我只需要 db:read”然后调用一个实际需要 db:write 的工具。工具声明和校验都在服务端完成客户端只提供 token。4. 流式传输的工程化落地4.1 SSE 通道的建立与消息格式设计SSE 的响应头必须包含Content-Type: text/event-stream、Cache-Control: no-cache、Connection: keep-alive。每条消息的格式是event: 类型\ndata: JSON\n\n注意末尾必须有两个换行否则客户端不会触发消息事件。MCP 的流式返回我设计成三种事件类型message传正常的 JSON-RPC 响应progress传工具执行的进度长任务用error传错误。客户端根据 event 类型分发处理。这里的关键是每条 data 必须是单行 JSON如果你的 JSON 里有换行必须转义否则 SSE 解析会断。async def sse_stream(session): yield event: message\ndata: json.dumps({jsonrpc: 2.0, id: 1, result: {}}) \n\n async for chunk in session.progress_channel: yield fevent: progress\ndata: {json.dumps(chunk)}\n\n4.2 背压控制当客户端读得比服务端写得慢流式传输最容易被忽视的问题是背压。如果服务端拼命往连接里写客户端消费不过来数据会在内核缓冲区堆积最终要么内存爆掉要么连接被重置。SSE 场景下服务端的写入是异步的你需要一个有界队列作为缓冲区队列满了就阻塞生产端而不是无限堆积。我用的是asyncio.Queue(maxsize100)生产者在put时如果队列满会自动 await形成天然的背压。同时给队列消费端设置超时如果客户端长时间不读比如超过 30 秒主动断开连接释放资源。这个超时值要根据你的工具执行时长来定长任务工具可以放宽到几分钟。实操心得不要用无界队列。我第一版用了默认的asyncio.Queue()无界压测时一个慢客户端直接把服务端内存吃到了 2GB。改成有界队列后内存稳定在几十 MB。4.3 大结果分块与心跳保活工具返回大结果时不能一次性json.dumps再发那样会有一个明显的卡顿。我的做法是分块序列化把结果按数组元素或对象字段切分每块单独发一个message事件客户端按顺序拼接。分块大小我定在 32KB 左右太小了事件数量爆炸太大了失去流式意义。心跳保活是另一个必须做的。SSE 连接如果长时间没有数据中间的反向代理或负载均衡会认为连接空闲而切断。解决办法是每隔 15 到 30 秒发一个注释行: keepalive\n\n注释行以冒号开头客户端会忽略但能保持连接活跃。心跳间隔要小于你基础设施里最短的空闲超时我一般设 15 秒比较保险。5. 状态管理的完整实现方案5.1 SessionStore 的数据结构与并发安全SessionStore 我用的是一个内存字典加锁的方案key 是 sessionIdvalue 是 Session 对象。Session 对象里存主体信息、创建时间、最后活跃时间、订阅的工具列表、进度通道。并发安全上Python 的asyncio是单线程事件循环字典的读写本身不会被打断但**“读-改-写”这种复合操作必须加锁**比如更新最后活跃时间、往订阅列表里加元素。class SessionStore: def __init__(self): self._sessions: dict[str, Session] {} self._lock asyncio.Lock() async def get(self, sid: str) - Session | None: async with self._lock: s self._sessions.get(sid) if s: s.last_active time.time() return s这里有个细节get操作我顺手更新了last_active这样清理任务只需要扫一遍字典就能找出过期会话。但更新操作在锁内做避免并发下的竞态。5.2 会话生命周期创建、续期与回收会话的生命周期管理有三个动作创建、续期、回收。创建发生在 SSE 连接建立、鉴权通过之后生成一个 UUID 作为 sessionId返回给客户端通过 SSE 的第一个事件。续期发生在每次工具调用或心跳时更新last_active。回收由一个后台任务定期执行扫描超过 TTL我设的是 30 分钟没有活跃的会话关闭其连接、清理资源。回收这里有个坑不能直接删字典里的条目就完事还得通知所有持有该会话引用的协程停止工作。我的做法是给 Session 加一个closed事件asyncio.Event回收时先 set 这个事件让相关协程感知到并自行退出再删条目。否则会出现“会话已删但协程还在往已关闭的连接写数据”的报错。5.3 多客户端并发下的状态隔离验证状态隔离是生产环境的硬要求。我写了一个并发测试同时建立 50 个会话每个会话调用同一个工具但传不同的参数验证返回结果没有串。测试中发现一个隐蔽的 bug工具函数里用了一个模块级的缓存字典key 只用了参数的一部分导致不同会话的调用命中了同一个缓存条目。工具函数必须是无状态的任何缓存都要把会话标识纳入 key或者干脆用请求级的作用域。修复方式是把缓存改成以(session_id, params_hash)为 key或者更彻底一点工具函数不持有任何跨请求状态需要缓存就下沉到专门的缓存层由缓存层负责隔离。我选了后者工具函数保持纯粹缓存逻辑独立出来。6. 实操过程与核心环节实现6.1 环境准备与依赖安装整个项目基于 Python 3.11核心依赖就三个mcp官方 SDK、pyjwttoken 校验、uvicornASGI 服务器。我不用 FastAPI 是因为 MCP 的 SSE 传输对响应流的控制要求比较细直接用 Starlette 更贴近底层。python -m venv venv source venv/bin/activate pip install mcp pyjwt[crypto] uvicorn starlette版本上有个注意点mcpSDK 迭代很快不同小版本的 API 有差异。我锁的是mcp1.2,1.3避免自动升级带来的破坏性变更。生产项目一定要锁版本别用pip install mcp裸装。6.2 服务端骨架搭建与路由注册服务端入口是一个 Starlette 应用注册两个路由GET /sse建立流式连接POST /messages接收客户端请求。这两个路由共享同一个 SessionStore。from starlette.applications import Starlette from starlette.routing import Route app Starlette(routes[ Route(/sse, endpointhandle_sse, methods[GET]), Route(/messages, endpointhandle_message, methods[POST]), ])handle_sse里做三件事校验 Authorization 头、创建会话、返回 StreamingResponse。handle_message里做两件事根据 sessionId 找到会话、把请求投递到会话的处理队列。注意handle_message必须校验请求里的 sessionId 和 token 主体是否匹配否则 A 会话的 token 可以往 B 会话发请求。6.3 鉴权中间件的接入与参数配置鉴权我封装成一个 Starlette 中间件对所有路由生效但/sse和/messages的校验逻辑略有不同/sse校验后要创建会话/messages校验后要匹配已有会话。中间件里我只做 token 的密码学校验把解析出的 payload 挂到request.state.auth具体的会话逻辑留给路由处理。配置项我抽到一个config.py里JWKS_URL、AUDIENCE、ISSUER、SESSION_TTL、HEARTBEAT_INTERVAL、QUEUE_MAXSIZE。这些值在不同环境开发、预发、生产不一样用环境变量注入代码里只留默认值。生产环境的 JWKS URL 必须是内网地址避免每次校验都走公网。6.4 流式响应的分块发送与客户端对接服务端分块发送的代码前面给过片段这里补全客户端的对接方式。客户端用 fetch 读取流const resp await fetch(/sse, { headers: { Authorization: Bearer ${token} } }); const reader resp.body.getReader(); const decoder new TextDecoder(); let buffer ; while (true) { const { done, value } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); const events buffer.split(\n\n); buffer events.pop(); for (const evt of events) { // 解析 event: 和 data: 行 } }这里的关键是按\n\n切分事件最后一段不完整的留在 buffer 里。很多人在这里踩坑直接按行读结果一个事件被拆成两半解析失败。SSE 的协议边界是空行不是换行。6.5 状态持久化与重启恢复的取舍会话状态要不要持久化我的结论是默认不持久化但提供可选的持久化钩子。原因MCP 会话本质上是短生命周期的交互上下文客户端断线重连后重新初始化即可持久化的收益不大反而引入 Redis 依赖和序列化复杂度。但如果你的场景要求“断线后恢复未完成的工具调用”那就需要把会话状态和任务状态分开持久化任务状态存 Redis会话状态仍可重建。我留了一个SessionPersistence接口默认实现是空操作需要时替换成 Redis 实现。这样既不强制依赖又保留了扩展点。7. 常见问题与排查技巧实录7.1 鉴权相关的高频问题速查现象可能原因排查方向401 但 token 看起来没过期时钟偏移检查服务端和签发方 NTP 同步401 且日志显示签名失败JWKS 缓存过期或 kid 不匹配手动拉一次 JWKS 对比 kid403 但 token 有效scope 不匹配打印 token 的 scope 和工具要求的 scope偶发 401JWKS 拉取超时加本地缓存和超时重试时钟偏移这个问题特别隐蔽。我遇到过一次服务端比签发方慢了 90 秒导致刚签发的 token 因为nbf生效时间还没到而被拒。解决办法是给nbf校验留一点容差比如 30 秒同时确保所有机器时间同步。7.2 流式传输断连与缓冲问题排查SSE 断连最常见的原因是中间层缓冲。排查步骤先在服务端本地直连测试如果本地正常、经过网关就断那基本是网关问题。Nginx 的配置要加location /sse { proxy_pass http://backend; proxy_buffering off; proxy_cache off; proxy_read_timeout 3600s; proxy_set_header Connection ; proxy_http_version 1.1; }proxy_buffering off是必须的proxy_read_timeout要大于你的最长工具执行时间。另外Connection 清空是为了让 Nginx 用 HTTP/1.1 的长连接避免默认的短连接行为。7.3 会话状态错乱的定位方法状态错乱的典型表现是“A 用户看到了 B 用户的数据”。定位方法在会话创建、每次工具调用、会话回收三个点打结构化日志日志里带上 sessionId 和主体标识。一旦出现错乱按 sessionId 串起日志就能看出是哪一步串了。我遇到过的两次错乱一次是缓存 key 没带 sessionId一次是异步任务里捕获了错误的会话变量闭包陷阱都是靠日志定位的。避坑技巧异步代码里绝对不要在循环里用闭包捕获循环变量。for s in sessions: asyncio.create_task(handle(s))这种写法如果handle里延迟使用了s可能拿到的是最后一个值。正确写法是asyncio.create_task(handle(s))时用默认参数绑定lambda ss: handle(s)。7.4 性能压测中的典型瓶颈压测时我发现的瓶颈排序第一是 JWKS 拉取每次校验都拉就完蛋必须缓存第二是 JSON 序列化大结果用orjson替代标准库能快 3 到 5 倍第三是日志写入同步写日志会阻塞事件循环必须用异步 handler。这三个优化做完单实例的并发会话数从 200 提到了 2000 以上。8. 我在这个项目里踩过的坑与经验总结第一个坑是过度设计。我一开始想做一个通用的插件系统工具可以热加载、动态注册结果复杂度爆炸调试困难。后来砍掉工具就是代码里静态注册的加工具就改代码重新部署。生产环境里静态注册的可预测性远比热加载的灵活性重要。第二个坑是低估了 SSE 的兼容性。某些企业代理会强制把text/event-stream改成text/plain导致客户端解析失败。解决办法是在客户端做容错不依赖 Content-Type直接按 SSE 格式解析。同时服务端可以加一个X-Accel-Buffering: no头兼容部分 Nginx 场景。第三个坑是会话回收的时机。我最初设的 TTL 是 5 分钟结果长任务工具执行到一半会话就被回收了。后来改成“有活跃任务时不回收”在 Session 里加一个active_tasks计数回收任务跳过计数大于零的会话。这个细节文档里不会写但生产环境一定会遇到。如果让我给准备做类似项目的同学一句建议先把鉴权、传输、状态这三层的边界画清楚再动手写代码。这三层里任何一层偷懒后面都要用加倍的调试时间还回来。我见过太多项目把鉴权塞在工具函数里、把状态挂在连接对象上最后要么加不了新工具要么一压测就崩。分层不是教条是让你在出问题时能快速定位到是哪一层的责任。
返回列表