核心摘要

Go 适合实现流式 Transport,但 Goroutine 和 http.Flusher 并不等于 MCP Server。Transport Adapter 必须遵循固定的 MCP 版本和 Client Profile,分离身份与 Session,关联 JSON-RPC 消息,限制内存,传播取消,并在异常时 Fail Safe。

常见的 GET /ssePOST /message 是旧式兼容模式,不应被描述成当前 MCP 远程 Transport 的通用默认值。新系统应先评估当前 Streamable HTTP Profile;如果必须保留旧式路径,应隔离、标注版本并单独测试。

先选择 Protocol Profile

写 Handler 前记录:

决策 问题
MCP 版本 固定了哪个规范和 SDK 版本?
Transport 使用 stdio、Streamable HTTP 还是旧式兼容 Profile?
Session 状态是连接本地、Server 生成还是外部 Registry?
Identity Bearer Token 在哪里校验,Tenant 如何得出?
Delivery 如何处理顺序、取消、重连和重复请求?
Limits 请求、结果、队列、并发、时间和成本预算是多少?

不要把旧式 SSE 示例直接复制到新部署。端点名称和事件名称不能替代选定的协议契约。

旧式 SSE 模式包含什么

需要兼容时,其形状大致是:

text
Client -> GET /sse
Server -> 返回 Server 生成的 endpoint 信息
Client -> 使用 JSON-RPC 请求 POST 到消息端点
Server -> 通过事件流返回带有关联 ID 的响应

这会产生两个独立 HTTP 表面。Server 必须保证:

  • endpoint 信息不泄露凭证或权限证明;
  • Session 由 Server 生成并拥有;
  • POST 请求独立执行认证与授权;
  • 响应通过 JSON-RPC ID 关联,而不是只依赖到达顺序;
  • 队列溢出、断连、取消和重复请求都有明确行为。

事件流不是权限通道。已经认证过的 SSE 连接不能在 Token 过期或策略撤销后无限期保留权限。

Session 状态与身份

使用不透明的 Server 生成 ID,只保存当前 Profile 需要的状态:

go
package transport

import (
	"context"
	"crypto/rand"
	"encoding/hex"
	"errors"
	"sync"
	"time"
)

type Principal struct {
	Subject string
	Tenant  string
}

type Session struct {
	ID        string
	Principal Principal
	CreatedAt time.Time
	ExpiresAt time.Time
	Events    chan []byte
	Context   context.Context
	Cancel    context.CancelFunc
}

type Registry struct {
	mu       sync.RWMutex
	sessions map[string]*Session
}

func NewRegistry() *Registry {
	return &Registry{sessions: make(map[string]*Session)}
}

func newSessionID() (string, error) {
	var raw [24]byte
	if _, err := rand.Read(raw[:]); err != nil {
		return "", err
	}
	return hex.EncodeToString(raw[:]), nil
}

func (r *Registry) Create(parent context.Context, principal Principal, ttl time.Duration) (*Session, error) {
	id, err := newSessionID()
	if err != nil {
		return nil, err
	}
	ctx, cancel := context.WithCancel(parent)
	session := &Session{
		ID: id, Principal: principal, CreatedAt: time.Now(),
		ExpiresAt: time.Now().Add(ttl),
		// 示例容量:应根据事件大小和内存预算推导。
		Events: make(chan []byte, 64),
		Context: ctx, Cancel: cancel,
	}
	r.mu.Lock()
	r.sessions[id] = session
	r.mu.Unlock()
	return session, nil
}

func (r *Registry) Get(id string, principal Principal) (*Session, error) {
	r.mu.RLock()
	session := r.sessions[id]
	r.mu.RUnlock()
	if session == nil || time.Now().After(session.ExpiresAt) {
		return nil, errors.New("session_not_found")
	}
	if session.Principal != principal {
		return nil, errors.New("session_principal_mismatch")
	}
	return session, nil
}

这是结构片段。真实实现还需要删除、过期清理、安全取消、水平扩展时的所有权和 Token 撤销策略。Session ID 只是 Server 状态的索引,不是 Token 校验的替代品。

流式 Handler 的责任

兼容 SSE 的 Handler 应:

  1. 认证初始请求;
  2. 创建 Server 拥有的 Session;
  3. 在写入前设置流式 Header;
  4. 只发送固定 Profile 所需的 endpoint 信息;
  5. 在队列事件、心跳、请求取消和 Session 过期之间进行选择;
  6. 只执行一次 Session 删除和资源关闭。
go
func writeEvent(w io.Writer, flush func(), event string, data []byte) error {
	if _, err := io.WriteString(w, "event: "+event+"\n"); err != nil {
		return err
	}
	if _, err := io.WriteString(w, "data: "+string(data)+"\n\n"); err != nil {
		return err
	}
	flush()
	return nil
}

该片段省略了 Handler、Registry 删除和 HTTP 错误辅助函数,不能复制后直接作为完整 Server。心跳间隔是部署参数,必须短于相关代理的空闲超时,但也不能制造不必要的流量。应通过真实 Load Balancer 和 CDN 测试,而不是只用本地 curl

消息端点与 JSON-RPC

消息端点校验的不应只有 JSON 语法:

  • HTTP Method、Content-Type 和请求大小;
  • JSON-RPC 版本、ID、Method 和参数 Schema;
  • Session 存在性、过期、Principal、Tenant 和 Transport Profile;
  • 请求取消、重复和幂等行为;
  • Tool 授权与 Resource 所有权;
  • 队列容量和下游超时。

使用有界 Decoder,并返回符合协议的错误。不要把任意 Client 输入直接索引到 Tool Dispatcher,也不要从模型参数推导对象所有权。

go
type Request struct {
	JSONRPC string          `json:"jsonrpc"`
	ID      any             `json:"id,omitempty"`
	Method  string          `json:"method"`
	Params  json.RawMessage `json:"params,omitempty"`
}

func decodeRequest(body io.Reader, maxBytes int64) (Request, error) {
	limited := io.LimitReader(body, maxBytes+1)
	var request Request
	decoder := json.NewDecoder(limited)
	if err := decoder.Decode(&request); err != nil {
		return Request{}, errors.New("invalid_json")
	}
	if request.JSONRPC != "2.0" || request.Method == "" {
		return Request{}, errors.New("invalid_jsonrpc_request")
	}
	return request, nil
}

片段需要标准库 encoding/jsonerrorsio 导入,并把按 Method 的参数校验留给应用。JSON-RPC Framing 不会授予 Tool 调用权限。

背压与顺序

事件 Channel 就是内存预算。队列满时需要明确策略:

策略 适用场景 风险
阻塞生产者 每个事件都必须保留 可能耗尽 Worker
拒绝新工作 Caller 可以安全重试 需要明确重试响应
关闭 Session 旧 Client 可以重新连接 必须支持恢复
写入持久队列 需要回放 增加顺序与删除复杂度

不要使用无界 Channel。保留 JSON-RPC 关联,并定义同一 Session 内响应是否有序。非幂等 Tool 在网络超时后不能自动重复执行,除非有幂等键或对账步骤。

认证与授权

认证应在创建 Session 前执行,并在每次消息请求上重新校验。实现应使用可信 Provider 配置校验 Issuer、Audience/Resource、签名算法、过期、Key Rotation 和 Scope。

授权仍是应用决策:

  • Principal 和 Tenant 来自可信上下文;
  • 对象所有权来自权威 Resource Service;
  • Tool Annotation 是提示,不是权限;
  • 删除、外发等高影响操作需要独立确认或工作流;
  • Token 过期和撤销可以使已有 Session 失效;
  • 下游 Tool Result 是不可信输入,不能改变策略决策。

暴露到互联网前,还需要 TLS、请求限制、限流、结构化审计,以及跨租户测试。

Proxy 与关闭行为

每个中间层都会影响流式行为:

  • 只在所用 Proxy 需要时关闭 Buffering;
  • 根据实测心跳配置读取和空闲超时;
  • 安全传播 Trace Context 和认证 Header;
  • 限制并发连接和响应字节;
  • 部署时在 Drain Deadline 内关闭 Session;
  • 将 Context Cancellation 传播到 Tool 和下游调用。

WriteTimeout: 0 不是普遍安全配置。某些流式 Server 可能需要它,但也可能让卡住的 Handler 永久存活。应使用心跳、取消、连接限制和经过测试的 Drain 策略管理生命周期。

测试矩阵

上线兼容 Transport 前测试:

场景 预期结果
畸形 JSON 或超大 Body 有界协议错误,不执行 Handler
未知或过期 Session 拒绝,不泄露其他 Tenant
有效 Session 但 Principal 不同 拒绝
Stream 存续期间 Token 过期 按策略重新认证或关闭
事件队列已满 有界拒绝、断开或持久化
Tool 执行中 Client 断开 取消传播到下游
重复 Request ID 或幂等键 确定性对账
节点失败后重连 不恢复未授权 Session
Proxy 缓冲或关闭空闲流 集成测试发现
Tool Result 含指令 按不可信数据处理

先针对固定 SDK 运行协议一致性测试,再用空闲 Session、活跃调用、重连风暴、队列压力和依赖故障做压测。本地成功的 curl 不是容量或安全证据。

使用 Go SDK 还是从零实现

如果维护中的 SDK 覆盖目标 MCP 版本、Transport、取消、授权 Hook 和生命周期,优先使用它。标准库实现适合学习或受控兼容 Adapter,但也要自行承担:

  • 协议一致性和版本变化;
  • JSON-RPC 边界与取消;
  • 认证和对象级授权;
  • Session 过期与分布式所有权;
  • 背压和重复投递;
  • 可观测性、安全更新和事故响应。

“没有第三方依赖”不等于“运维风险更低”。

上线检查清单

  • [ ] 固定 MCP 版本和 Transport Profile。
  • [ ] 将旧式 SSE 隔离并标记为兼容行为。
  • [ ] Session 由 Server 生成,并绑定 Principal、Tenant 和 Profile。
  • [ ] 每个消息请求重新执行认证和授权。
  • [ ] 执行请求、结果、队列、并发、超时和成本限制。
  • [ ] 测试取消、重连、重复投递和关闭。
  • [ ] 将 Tool Result 视为不可信数据。
  • [ ] Telemetry 不包含原始 Token、Secret 和无界 Payload。
  • [ ] SDK 或自定义 Transport 有一致性和压测证据。

常见问题

新的远程 MCP Server 应该使用旧式 SSE 吗?

不应默认使用。先固定当前支持的 Transport Profile;只有存在兼容需求时才使用旧式 SSE,并单独测试它的双通道 Session 行为。

Go 提供了什么?

它提供 HTTP 流式响应、取消 Context 和并发原语,但不会自动提供 MCP 授权、Session 策略、JSON-RPC 校验、背压或安全关闭。

Client 可以选择 Session ID 吗?

不能。ID 由 Server 创建并绑定。Client 只能把它作为索引提交,Bearer Token 和 Resource Policy 仍然是权威。

单个进程能支撑多少连接?

没有通用数字。应在目标环境测量 Transport、事件大小、心跳、Proxy、文件描述符、内存、下游工作和重连行为。

暴露到互联网前要补什么?

可信认证、逐请求授权、租户隔离、限制、取消、限流、Telemetry 脱敏以及故障和一致性测试。

结语

用 Go 实现流式 Adapter 可以帮助理解 HTTP Flush、Channel、取消和 JSON-RPC 关联的机制,但生产 MCP Transport 远不止这些。选定的 Protocol Profile、Server 生成的 Session、授权、有界内存、重连和生命周期同样重要。把旧式 SSE 视为兼容边界,在真实部署中测量,并在 SDK 覆盖和治理能力足够时优先使用维护中的实现。

一手来源