package mcp // GET /mcp —— 服务端→客户端的事件流(Streamable HTTP 的 SSE 通道)。 // // # 这半边解决什么 // // 工具面(POST /mcp)只解决「接入方**问**」。这一条解决「服务端**说**」: // 邮件投递时把 new_mail / session_update 推给接入方,让它**拉起对话** —— // 与各桥靠 /api/v1/events/stream 收信是同一件事,只是方言不同: // // 桥: id: 7\nevent: new_mail\ndata: {…}\n\n // MCP: {"jsonrpc":"2.0","method":"notifications/message","params":{…}} // // # 为什么复用 sse.Manager 而不是另起一套 // // Manager 里那些东西**都是踩过坑才对的**:writeMu 串行化(2026-09-28 -race // 实测 http.ResponseWriter 并发写会把 JSON 劈成半截,800 帧只切出 459 个完整)、 // Last-Event-ID 回放(宁可重复也不丢失)、心跳(反代按空闲 30-58s 掐连接)、 // 环形缓冲上限、断线清理。复制一份等于把那些坑再踩一遍, // 而两边的修复从此各走各的。 // // 代价是 sse.Client 多了一个 Frame 钩子(见 sse/manager.go)。 // 默认 nil = AgentMail 原格式,各桥与 WebUI 行为一字未变。 import ( "fmt" "net/http" "strings" "github.com/agentmail/gateway/internal/middleware" "github.com/agentmail/gateway/internal/sse" ) // Frame 把一条 AgentMail 事件渲染成 MCP 方言(JSON-RPC 通知)。 // // 独立成函数是为了能被单测直接调用、逐字节断言。 // // data 已是 JSON 字节(SendWithID marshal 过;回放时本来就是原始字节), // 嵌进 params 时**原样拼接**而不是再 marshal 一次 —— 再 marshal 会把已序列化的 // JSON 转义成字符串,客户端得解析两次,且长度翻倍。 func Frame(id, eventType string, data []byte) string { return fmt.Sprintf(`{"jsonrpc":"2.0","method":"notifications/message","params":{`+ `"event_id":%q,"event_type":%q,"payload":%s}}`, id, eventType, string(data)) + "\n\n" } // HandleGET 处理 GET /mcp:开一条 SSE 长连,把投递事件以 MCP 方言推给接入方。 // // 认证在**外层**(与 POST 同一条 AgentAuth 路由):身份已在 context 里, // 未带凭证的请求根本到不了这里。 func (s *Server) HandleGET(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodGet { w.Header().Set("Allow", "GET, POST") http.Error(w, "MCP GET 通道只接受 GET(工具面走 POST)", http.StatusMethodNotAllowed) return } agent := middleware.GetAgentName(r) if agent == "" { // AgentAuth 已拦过;能到这里说明路由挂错了位置。 // 报出来而不是默默接受 —— 匿名订阅会变成一处静默的越权读。 http.Error(w, "MCP 事件流需要 Agent 凭证", http.StatusUnauthorized) return } // 必须声明接受 SSE:否则客户端拿到一堆自己解析不了的文本流。 // 这条不是形式主义 —— 规范里 POST 与 GET 的语义不同,接错的客户端 // 会静默地什么都不收到(连接开着,但内容它不是当 SSE 读的)。 if !acceptsEventStream(r) { http.Error(w, "MCP 事件流要求 Accept: text/event-stream", http.StatusNotAcceptable) return } // ★ 同一把 Manager、同一个缓冲区、同一套回放 —— 与各桥唯一的差别是方言。 // // Frame 在**注册时**传入(而不是注册后设):回放发生在"先写响应、再注册" // 的前半段,事后设会让重连那一批走默认格式 —— 同一个连接两种方言, // 且只在真实断线时现形。见 sse.AddClientWithFrame 的注释。 client := sse.Default.AddClientWithFrame(w, r, agent, "", Frame) if client == nil { http.Error(w, "SSE not supported", http.StatusInternalServerError) return } // 连接还活着就阻塞在这里;断开(客户端关、反代掐、服务关停)即返回。 <-r.Context().Done() sse.Default.RemoveClient(client.ID) } // acceptsEventStream 判 Accept 头里有没有 text/event-stream。 // // 手写而不是用 mime.ParseMediaType + 循环:Accept 是**列表**且带 q 值, // 完整解析要处理 `text/event-stream;q=0.9, application/json;q=0.5`。 // 这里只关心「有没有声明」,不关心优先级 —— 多解析的那部分没有消费者。 func acceptsEventStream(r *http.Request) bool { for _, part := range strings.Split(r.Header.Get("Accept"), ",") { media := strings.TrimSpace(strings.SplitN(part, ";", 2)[0]) if strings.EqualFold(media, "text/event-stream") { return true } } return false }