Files
MailUI4Agents/server/internal/mcp/stream.go
JianFeeeee 560c462768 feat(mcp): GET /api/v1/mcp —— 投递侧事件流(让接入方被动收信,不用轮询)
## 这半边解决什么

工具面(POST)只解决「接入方**问**」。这一条解决「服务端**说**」:
邮件投递时把 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` 钩子:
**nil = AgentMail 原格式,各桥与 WebUI 行为一字未变**(默认值即历史行为)。

## ★ 回放是第三条写路径,漏了就只在断线时现形

`Send` / `SendWithID` / `replay` 是三条写 Res 的路径。原先**三条都把格式写死**,
只改前两条的话:MCP 客户端**平时**一切正常,只有带 `Last-Event-ID` 重连时
才会收到一批自己解不开的帧 —— 同一个连接上两种方言。

判据 `TestCustomFrameAppliesToReplayToo` 专门钉这条,并带反向对照
(nil 帧必须回落 AgentMail 格式)。

`Frame` 必须在**注册时**传入(`AddClientWithFrame`),不能事后设 ——
回放发生在「先写响应、再注册」的前半段,事后设只影响之后推来的事件。
原先 `AddClient` 保留为薄封装,各桥与 WebUI 调用点一字未改。

## 判据(6 格)

    Frame 是 JSON-RPC 2.0 通知 + 帧完整性(单事件、\n\n 结尾)
    payload 原样嵌入(不是 JSON 字符串)—— 再 marshal 会让客户端解析两次
    event_id / event_type 必带(前者是 Last-Event-ID 续传的依据)
    Accept 判定(含 q 值、大小写)
    匿名 GET → 401(不能变成静默的匿名订阅)
    缺 Accept → 406(接错的客户端会静默收不到东西)

## 顺带修:TestAdvanceRecurrenceLunar 的时区缺陷(★ 今天第三次假红)

全量测试红了,查下来是**我今天早些时候改判据时引入的**,与本次改动无关。

农历换算必须按**本地公历日**算(`AdvanceRecurrence` 里那句
`eventTime.In(time.Local)` 就是这条规则)。库里读回的 EventTime 是 **UTC**
(DSN 用 `_timezone=UTC`),UTC 比本地晚 8 小时,跨零点时农历日差一天:

    start    (Local) = 2026-10-04        农历日 24
    after    (UTC)   = 2026-11-01 16:00   农历日 23   ← 断言没换算时区(错)
    after.In(Local)  = 2026-11-02 00:00   农历日 24   ← 正确

服务端代码一直是对的,是判据没照做。失败信息里现在打印时区,
免得下次要重新推导一遍。变异验证:去掉 `.In(time.Local)` → 红 1 ✓

(这条判据是农历的第三次假红了:3459605「断言要求不存在的农历日」、
今天早些「起点写死日期 + advanceToFuture 跳过过期月份」、现在「没换算时区」——
三次都是判据自己写错,代码三次都对。它依赖 Local 时区与「今天」,
天生脆弱,值得记着。)

## 验证

    go test ./...              14 包全绿
    go test ./internal/sse/    含新判据绿
    go test ./internal/mcp/    6 格新判据 + 原 19 格全绿
2026-10-02 15:37:34 +08:00

103 lines
4.5 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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
}