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

文章来源声明: 原文作者:吴佳浩Alben; 来源站点:掘金; 原文链接:https://juejin.cn/post/7687446787750543366; 本文基于上述来源整理/加工,觅优补充点评,仅供技术学习交流。版权归原作者所有。
觅优短评

面向已跑通 MCP Demo、准备落地企业内网的团队,文章把鉴权、多租户隔离、SSE 长连接治理讲得具体可抄,避坑清单尤其实用,适合后端与 Agent 平台工程师参考。

从零手写一个生产级 MCP Server:鉴权、流式传输与状态管理 ---------------------------------

请添加图片描述

作 者:吴佳浩(Alben)
公众号:全栈架构师笔记
系列专栏:《企业级 Agent 实战指南———MCP 与 Agent Tools 工程化落地实战》· 第 02 篇


导读
官方 Demo 的 Stdio 模式只能在本地玩具项目里跑一跑,真正走进企业内网,MCP Server 必须是一个具备鉴权、多租户隔离、连接池与健康检查的独立微服务。
很多人以为写 MCP Server 就是套个 FastMCP 装饰器;但在高并发生产环境下,连接泄露、未捕获异常导致子进程暴毙、以及多租户 Token 越权,才是最致命的隐形杀手。
从 Demo 走向生产,核心在于把无状态的 HTTP 请求包装为可控的 Long-Session 状态机。


在第 01 篇中,我们拆解了 MCP 协议的设计哲学。很多工程师在看完官方文档后,通常会用官方 Python SDK 写一个基于 Stdio 的脚本:

<span># 典型的本地 Demo 写法(无法直接用于企业生产)</span>
<span>from</span> mcp.server.fastmcp <span>import</span> FastMCP
mcp = FastMCP(<span>"DemoServer"</span>)

<span>@mcp.tool()</span>
<span>def</span> <span>query_db</span>(<span>sql: <span>str</span></span>) -> <span>str</span>:
    <span># 裸连数据库,无鉴权,无租户隔离,无超时控制</span>
    <span>return</span> db.execute(sql)

mcp.run() <span># 默认走 stdio</span>

这种写法在本地测试很顺畅,但一旦要把这个 MCP Server 部署在企业 Kubernetes 集群中,供全公司数十个 Agent 同时调用时,会立刻遭遇三大生产级灾难:

灾难现象具体表现架构根因
1. 跨租户越权穿透研发部的 Agent 查到了财务部的缺乏基于 Request 粒度的 Tenant
(Tenant Leakage)薪酬表,引发严重合规审查事故上下文注入与动态连接池隔离
2. 长连接雪崩多个 Agent 并发建立 SSE 长连接,缺乏连接心跳保活、会话复用与
(Connection Exhaust)导致服务器文件句柄耗尽崩溃僵尸会话(Ghost Session)淘汰
3. 异常引发全局暴毙某次 Tool 执行抛出未捕获异常,缺少全局错误屏障与 JSON-RPC
(Process Crash)整个 MCP 进程直接退出服务中断标准错误码封装

要构建企业级 MCP Server,必须支持 Streamable HTTP / SSE 通道,并具备完整的 鉴权拦截、租户动态隔离与生命周期管理


一、企业级 MCP Server 的微服务架构拓扑

在生产环境中,MCP Server 绝不是单机脚本,而是一个标准的云原生微服务:

mermaid diagram

  • 🔸 双通道架构(SSE + HTTP POST):客户端通过 /sse 端点建立 Server-Sent Events 长连接接收下行事件;通过 /messages?session_id=xxx 发送上行 JSON-RPC 请求;
  • 🔸 会话生命周期隔离:每一个 Agent 连接分配唯一的 session_id,保证多 Agent 之间状态互不干扰;
  • 🔸 多租户上下文透传:将网关鉴权后的 tenant_id 注入到异步协程上下文中(ContextVar),底层数据库连接池按租户动态路由。

一句话总结这一章的核心观点:
生产级 MCP Server 的本质,是把传统的 Web API 包装成符合 JSON-RPC 2.0 规范与长连接会话协议的微服务。


二、生产级代码实现:基于 FastAPI 的企业级 MCP Server

以下示例基于 Python 3.11+ 与 FastAPI 构建企业级 MCP Server,并采用原生 StreamingResponse 实现 Server-Sent Events(SSE) 数据传输。示例完整演示了企业级 MCP 服务端的核心能力,包括 多租户鉴权、工具注册、JSON-RPC 协议处理、SSE 双向通信、会话管理以及工具异常隔离 等关键模块。

需要说明的是,本示例主要用于展示企业级 MCP Server 的整体架构设计与核心实现思路。实际生产环境中,还应进一步结合 OAuth/JWT 鉴权、会话回收(Session GC)、限流熔断、日志审计、可观测性(OpenTelemetry)、高可用部署 等能力,构建完整的企业级 MCP 服务体系。由于示例使用的是 FastAPI 原生 StreamingResponse 手工实现 SSE,因此无需额外依赖第三方 SSE 库即可完成协议通信。FastAPI 也提供了更高层的 SSE 支持,可根据项目需求选择使用。

<span>"""
enterprise_mcp_server.py

生产级企业 MCP Server 实现

包含:
- Streamable HTTP / SSE 双通道
- JWT 鉴权
- 多租户上下文注入
- Tool 注册
- Tool 异常隔离
"""</span>

<span>import</span> asyncio
<span>import</span> json
<span>import</span> uuid
<span>from</span> contextvars <span>import</span> ContextVar
<span>from</span> typing <span>import</span> <span>Any</span>, <span>Dict</span>, <span>List</span>

<span>from</span> fastapi <span>import</span> Depends, FastAPI, Header, HTTPException, Request, status
<span>from</span> fastapi.responses <span>import</span> JSONResponse, StreamingResponse
<span>from</span> pydantic <span>import</span> BaseModel

app = FastAPI(title=<span>"Enterprise MCP Server"</span>, version=<span>"1.0.0"</span>)

current_tenant_id: ContextVar[<span>str</span>] = ContextVar(
    <span>"current_tenant_id"</span>, default=<span>"default"</span>
)


<span>class</span> <span>SessionContext</span>:
    <span>def</span> <span>__init__</span>(<span>self, session_id: <span>str</span>, tenant_id: <span>str</span></span>):
        self.session_id = session_id
        self.tenant_id = tenant_id
        self.queue: asyncio.Queue = asyncio.Queue()
        self.last_active = asyncio.get_event_loop().time()


active_sessions: <span>Dict</span>[<span>str</span>, SessionContext] = {}


<span>class</span> <span>ToolDefinition</span>(<span>BaseModel</span>):
    name: <span>str</span>
    description: <span>str</span>
    inputSchema: <span>Dict</span>[<span>str</span>, <span>Any</span>]


REGISTERED_TOOLS: <span>Dict</span>[<span>str</span>, <span>Any</span>] = {}
TOOL_SCHEMAS: <span>List</span>[ToolDefinition] = []


<span>def</span> <span>mcp_tool</span>(<span>name: <span>str</span>, description: <span>str</span>, schema: <span>Dict</span>[<span>str</span>, <span>Any</span>]</span>):
    <span>def</span> <span>decorator</span>(<span>fn</span>):
        REGISTERED_TOOLS[name] = fn
        TOOL_SCHEMAS.append(
            ToolDefinition(
                name=name,
                description=description,
                inputSchema=schema,
            )
        )
        <span>return</span> fn

    <span>return</span> decorator


<span>@mcp_tool(<span>
    name=<span>"query_tenant_metrics"</span>,
    description=<span>"Query production metrics for the authenticated tenant safely."</span>,
    schema={
        <span>"type"</span>: <span>"object"</span>,
        <span>"properties"</span>: {
            <span>"metric_name"</span>: {
                <span>"type"</span>: <span>"string"</span>,
                <span>"description"</span>: <span>"e.g. qps, error_rate, latency"</span>,
            },
            <span>"time_range"</span>: {
                <span>"type"</span>: <span>"string"</span>,
                <span>"description"</span>: <span>"e.g. 1h, 24h, 7d"</span>,
            },
        },
        <span>"required"</span>: [<span>"metric_name"</span>],
    },
</span>)</span>
<span>async</span> <span>def</span> <span>query_tenant_metrics</span>(<span>metric_name: <span>str</span>, time_range: <span>str</span> = <span>"1h"</span></span>) -> <span>str</span>:
    tenant = current_tenant_id.get()
    <span>return</span> json.dumps(
        {
            <span>"tenant_id"</span>: tenant,
            <span>"metric"</span>: metric_name,
            <span>"time_range"</span>: time_range,
            <span>"data"</span>: {
                <span>"avg_value"</span>: <span>142.5</span>,
                <span>"p99_latency_ms"</span>: <span>23.4</span>,
                <span>"status"</span>: <span>"healthy"</span>,
            },
        }
    )


<span>async</span> <span>def</span> <span>auth_middleware</span>(<span>
    x_tenant_id: <span>str</span> = Header(<span>..., alias=<span>"X-Tenant-ID"</span></span>),
    authorization: <span>str</span> = Header(<span>..., alias=<span>"Authorization"</span></span>),
</span>) -> <span>str</span>:
    <span>if</span> <span>not</span> authorization.startswith(<span>"Bearer "</span>) <span>or</span> <span>not</span> x_tenant_id:
        <span>raise</span> HTTPException(
            status_code=status.HTTP_401_UNAUTHORIZED,
            detail=<span>"Unauthorized: Missing valid Bearer token or Tenant Header"</span>,
        )

    current_tenant_id.<span>set</span>(x_tenant_id)
    <span>return</span> x_tenant_id


<span>@app.get(<span><span>"/sse"</span></span>)</span>
<span>async</span> <span>def</span> <span>sse_endpoint</span>(<span>request: Request, tenant_id: <span>str</span> = Depends(<span>auth_middleware</span>)</span>):
    session_id = <span>str</span>(uuid.uuid4())
    session_ctx = SessionContext(session_id=session_id, tenant_id=tenant_id)
    active_sessions[session_id] = session_ctx

    <span>async</span> <span>def</span> <span>event_generator</span>():
        <span>try</span>:
            <span>yield</span> (
                <span>"event: endpoint\n"</span>
                <span>f"data: /messages?session_id=<span>{session_id}</span>\n\n"</span>
            )

            <span>while</span> <span>True</span>:
                <span>if</span> <span>await</span> request.is_disconnected():
                    <span>break</span>

                <span>try</span>:
                    msg = <span>await</span> asyncio.wait_for(
                        session_ctx.queue.get(), timeout=<span>15.0</span>
                    )
                    <span>yield</span> (
                        <span>"event: message\n"</span>
                        <span>f"data: <span>{json.dumps(msg)}</span>\n\n"</span>
                    )
                <span>except</span> asyncio.TimeoutError:
                    <span>yield</span> <span>": ping\n\n"</span>
        <span>finally</span>:
            active_sessions.pop(session_id, <span>None</span>)

    <span>return</span> StreamingResponse(
        event_generator(),
        media_type=<span>"text/event-stream"</span>,
        headers={
            <span>"Cache-Control"</span>: <span>"no-cache"</span>,
            <span>"Connection"</span>: <span>"keep-alive"</span>,
            <span>"X-Accel-Buffering"</span>: <span>"no"</span>,
        },
    )


<span>@app.post(<span><span>"/messages"</span></span>)</span>
<span>async</span> <span>def</span> <span>message_endpoint</span>(<span>
    request: Request,
    session_id: <span>str</span>,
    tenant_id: <span>str</span> = Depends(<span>auth_middleware</span>),
</span>):
    session_ctx = active_sessions.get(session_id)
    <span>if</span> <span>not</span> session_ctx:
        <span>raise</span> HTTPException(status_code=<span>404</span>, detail=<span>"Session not found or expired"</span>)

    current_tenant_id.<span>set</span>(tenant_id)

    payload = <span>await</span> request.json()
    method = payload.get(<span>"method"</span>)
    req_id = payload.get(<span>"id"</span>)

    <span>if</span> method == <span>"initialize"</span>:
        <span>await</span> session_ctx.queue.put(
            {
                <span>"jsonrpc"</span>: <span>"2.0"</span>,
                <span>"id"</span>: req_id,
                <span>"result"</span>: {
                    <span>"protocolVersion"</span>: <span>"2024-11-05"</span>,
                    <span>"capabilities"</span>: {<span>"tools"</span>: {}},
                    <span>"serverInfo"</span>: {
                        <span>"name"</span>: <span>"EnterpriseProductionMCPServer"</span>,
                        <span>"version"</span>: <span>"1.0.0"</span>,
                    },
                },
            }
        )
        <span>return</span> JSONResponse({<span>"status"</span>: <span>"accepted"</span>})

    <span>if</span> method == <span>"tools/list"</span>:
        <span>await</span> session_ctx.queue.put(
            {
                <span>"jsonrpc"</span>: <span>"2.0"</span>,
                <span>"id"</span>: req_id,
                <span>"result"</span>: {
                    <span>"tools"</span>: [t.model_dump() <span>for</span> t <span>in</span> TOOL_SCHEMAS],
                },
            }
        )
        <span>return</span> JSONResponse({<span>"status"</span>: <span>"accepted"</span>})

    <span>if</span> method == <span>"tools/call"</span>:
        params = payload.get(<span>"params"</span>, {})
        tool_name = params.get(<span>"name"</span>)
        arguments = params.get(<span>"arguments"</span>, {})

        fn = REGISTERED_TOOLS.get(tool_name)

        <span>if</span> fn <span>is</span> <span>None</span>:
            response = {
                <span>"jsonrpc"</span>: <span>"2.0"</span>,
                <span>"id"</span>: req_id,
                <span>"error"</span>: {
                    <span>"code"</span>: -<span>32601</span>,
                    <span>"message"</span>: <span>f"Tool '<span>{tool_name}</span>' not found"</span>,
                },
            }
        <span>else</span>:
            <span>try</span>:
                result = <span>await</span> fn(**arguments)
                response = {
                    <span>"jsonrpc"</span>: <span>"2.0"</span>,
                    <span>"id"</span>: req_id,
                    <span>"result"</span>: {
                        <span>"content"</span>: [
                            {<span>"type"</span>: <span>"text"</span>, <span>"text"</span>: <span>str</span>(result)}
                        ]
                    },
                }
            <span>except</span> Exception <span>as</span> e:
                response = {
                    <span>"jsonrpc"</span>: <span>"2.0"</span>,
                    <span>"id"</span>: req_id,
                    <span>"error"</span>: {
                        <span>"code"</span>: -<span>32000</span>,
                        <span>"message"</span>: <span>f"Execution Error: <span>{e}</span>"</span>,
                    },
                }

        <span>await</span> session_ctx.queue.put(response)
        <span>return</span> JSONResponse({<span>"status"</span>: <span>"accepted"</span>})

    <span>return</span> JSONResponse({<span>"status"</span>: <span>"ignored"</span>})


三、生产级 MCP Server 的四大避坑指南

在将 MCP Server 部署到企业私有云时,团队必须严防以下四个典型深水区问题:

  • 🔸 NGINX / 网关的 Buffer 缓冲导致 SSE 断流:反向代理服务器默认会开启响应缓冲(Buffering),导致客户端无法实时收到长连接事件。必须在响应头中明确添加 X-Accel-Buffering: no
  • 🔸 协程并发下的租户数据混淆:严禁在全局变量中暂存 tenant_id,在异步 Python 环境中必须使用 contextvars.ContextVar,确保异步调度切换时租户边界严密隔离;
  • 🔸 僵尸连接导致的内存泄漏:必须设计基于 asyncio.TimeoutError 的心跳机制(Ping),当客户端异常断网未发送关闭信号时,及时清理 active_sessions 字典;
  • 🔸 工具超时的硬熔断机制:任何工具调用必须套上超时装饰器(如 30s 熔断),防止某个卡死在后端的 SQL 查询把整个工作线程池拖垮。

本篇总结

  • 🔸 Stdio 只用于本地极客调试,企业中台必须上 Streamable HTTP/SSE 双通道架构
  • 🔸 通过 ContextVar 与依赖注入实现多租户上下文的绝对物理/逻辑隔离
  • 🔸 构筑全局异常屏障与 JSON-RPC 标准错误包装,防止进程意外崩溃
  • 🔸 配置 NGINX 零缓冲与心跳保活,杜绝生产环境长连接雪崩

掌握了如何编写高可用的 MCP Server 之后,下一个核心问题是:当 Agent 拥有了强大的工具调用能力,如何保证它在执行敏感命令时不破坏生产环境?

在下一篇中,我们将深入探讨:《Tool 的安全性与执行沙箱:从 Docker 到 gVisor 的防御架构》

筒子们本篇为《企业级 Agent 实战指南》· 第二章的第 2篇,后续续会更新完整的agent的开发的全部过程,如果你对Agent开发感兴趣不妨关注一下本合集。