## 这半边解决什么
工具面(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 格全绿
362 lines
12 KiB
Go
362 lines
12 KiB
Go
// Package mcp 把 MCP(Model Context Protocol)实现进网关本身。
|
||
//
|
||
// # 为什么在服务端而不是独立进程
|
||
//
|
||
// 早先的形状是 `plugins/zcode-mail-bridge/mcp/server.mjs`:一个独立进程,
|
||
// 用 HTTP 调本网关。这条路有四个真实成本:
|
||
//
|
||
// 1. **工具语义有两份**。桥里的 read_inbox / send_mail 是**手抄**网关的语义,
|
||
// 抄错就是行为分叉。实测已经踩到过一次:`connect_to_server` 只发
|
||
// `X-Agent-Secret` 头,而 `/agent/register` 只认 Bearer 或 body 里的
|
||
// secret ⇒ secret-only 的 Agent 调它必然 400。
|
||
// 2. **鉴权与作用域要再实现一遍**。工作区收窄、会话收窄、冷静期、配额
|
||
// 这些规则住在服务端;独立进程拿不到,只能靠 HTTP 重走一遍。
|
||
// 3. **多一跳 + 多一个故障点**。宿主 → 桥进程 → HTTP → 网关。
|
||
// 4. **接入端仍要装东西**。本机装 node、装桥、配环境变量。
|
||
//
|
||
// 进服务端之后:工具**包装现有 handler**(见 tools.go),同一条代码路径、
|
||
// 同一套鉴权与收窄;宿主只需填一个 URL。
|
||
//
|
||
// # 传输:Streamable HTTP
|
||
//
|
||
// MCP 规范 2025-06-18 的传输:客户端 POST 一个 JSON-RPC 消息到单一端点,
|
||
// 服务端回 202(无输出)或一条 SSE 流。单条请求-响应场景最简单的是
|
||
// **直接回 JSON**(POST 一次拿到一个 JSON-RPC 响应),本实现这样做;
|
||
// 会把 Accept 头里的 `text/event-stream` 也认下来,返回 `Content-Type:
|
||
// application/json`(规范允许服务端在无待推送消息时如此)。
|
||
//
|
||
// 之所以不引 `github.com/modelcontextprotocol/go-sdk`:协议面只有四个方法,
|
||
// 而引 SDK 会带来一条依赖链;与本仓其余部分零依赖的取向一致(见
|
||
// server/go.mod)。手写让这一层成为可单测的纯函数。
|
||
package mcp
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"fmt"
|
||
"io"
|
||
"log"
|
||
"net/http"
|
||
"strings"
|
||
)
|
||
|
||
// 协议版本。客户端报的版本原样回显(见 handleMessage 的说明)。
|
||
const (
|
||
ProtocolVersion = "2025-06-18"
|
||
FallbackVersion = "2024-11-05"
|
||
|
||
ServerName = "agentmail"
|
||
ServerVersion = "1.0.0"
|
||
)
|
||
|
||
// JSON-RPC 错误码(只列本实现真的会返回的)。
|
||
const (
|
||
rpcParse = -32700
|
||
rpcInvalidRequest = -32600
|
||
rpcMethodNotFound = -32601
|
||
rpcInvalidParams = -32602
|
||
rpcInternal = -32603
|
||
)
|
||
|
||
// rpcMessage 是 JSON-RPC 消息。请求与响应共用(协议本身如此),故不分类型。
|
||
type rpcMessage struct {
|
||
JSONRPC string `json:"jsonrpc"`
|
||
// ID 是指针而不是 RawMessage:协议规定**每个响应都必须带 id**,
|
||
// 且解析失败时 id 必须是 JSON null。用 RawMessage 配 `omitempty` 时,
|
||
// nil 会让整个字段消失 —— 实测过一次(TestMalformedJSONGetsParseError
|
||
// 抓到的):响应里没有 id,客户端会一直等这条的响应。
|
||
//
|
||
// 指针的取舍:非空指针指向 RawMessage(可能是 `null`、数字、字符串);
|
||
// nil 指针表示「无 id」(通知)。
|
||
ID *json.RawMessage `json:"id"`
|
||
Method string `json:"method,omitempty"`
|
||
Params json.RawMessage `json:"params,omitempty"`
|
||
Result any `json:"result,omitempty"`
|
||
Error *rpcError `json:"error,omitempty"`
|
||
}
|
||
|
||
// nullID 是 JSON-RPC 规定的「id 为 null」(解析失败时用)。
|
||
var nullID = json.RawMessage("null")
|
||
|
||
// rawID 把指针化的 id 还原成 RawMessage(nil ⇒ nil)。
|
||
func rawID(p *json.RawMessage) json.RawMessage {
|
||
if p == nil {
|
||
return nil
|
||
}
|
||
return *p
|
||
}
|
||
|
||
type rpcError struct {
|
||
Code int `json:"code"`
|
||
Message string `json:"message"`
|
||
}
|
||
|
||
// toolCallParams 是 `tools/call` 的参数。
|
||
type toolCallParams struct {
|
||
Name string `json:"name"`
|
||
Arguments map[string]any `json:"arguments"`
|
||
}
|
||
|
||
// toolResult 是 `tools/call` 的结果。
|
||
//
|
||
// isError 的存在是**功能性的**:MCP 的约定是工具执行失败回 result +
|
||
// isError:true,而不是 JSON-RPC error —— 后者模型只看到"协议错误",
|
||
// 拿不到失败原因就没法改道(换个 attachment_id 重试之类)。
|
||
type toolResult struct {
|
||
Content []content `json:"content"`
|
||
IsError bool `json:"isError,omitempty"`
|
||
}
|
||
|
||
type content struct {
|
||
Type string `json:"type"`
|
||
Text string `json:"text,omitempty"`
|
||
}
|
||
|
||
// textResult 造一条成功结果。
|
||
func textResult(s string) toolResult {
|
||
return toolResult{Content: []content{{Type: "text", Text: s}}}
|
||
}
|
||
|
||
// errorResult 造一条失败结果(**不是** JSON-RPC error)。
|
||
func errorResult(format string, a ...any) toolResult {
|
||
return toolResult{
|
||
Content: []content{{Type: "text", Text: fmt.Sprintf(format, a...)}},
|
||
IsError: true,
|
||
}
|
||
}
|
||
|
||
// idPtr 把 nil 归一成「显式的 JSON null」—— 响应**必须**带 id 字段。
|
||
func idPtr(id json.RawMessage) *json.RawMessage {
|
||
if id == nil {
|
||
return &nullID
|
||
}
|
||
return &id
|
||
}
|
||
|
||
func result(id json.RawMessage, v any) *rpcMessage {
|
||
return &rpcMessage{JSONRPC: "2.0", ID: idPtr(id), Result: v}
|
||
}
|
||
|
||
func failure(id json.RawMessage, code int, format string, a ...any) *rpcMessage {
|
||
return &rpcMessage{
|
||
JSONRPC: "2.0",
|
||
ID: idPtr(id),
|
||
Error: &rpcError{Code: code, Message: fmt.Sprintf(format, a...)},
|
||
}
|
||
}
|
||
|
||
// Tool 是一次 MCP 工具调用。
|
||
type Tool interface {
|
||
// Schema 声明工具名、说明与入参 JSON Schema。
|
||
Schema() ToolSchema
|
||
// Run 执行。ctx 是**当前请求的 context**,里面带着 AgentAuth 放进来的
|
||
// 身份(middleware.AgentNameKey)。
|
||
//
|
||
// ★ 为什么 ctx 必须显式传进来(而不是在实现里 context.Background()):
|
||
// 身份就住在这个 ctx 里。丢掉它 ⇒ 每个工具调用都 Unauthorized,
|
||
// 而症状是「模型说连上了但读不到任何信」—— 很难当场归因。
|
||
// 这条是被判据逼出来的(TestReadInboxRequiresWorkspaceSameAsHTTP
|
||
// 先是报 Unauthorized 才暴露出来)。
|
||
Run(ctx context.Context, args map[string]any) (string, error)
|
||
}
|
||
|
||
// ToolSchema 是 tools/list 里每个条目的形状。
|
||
//
|
||
// Annotations 是 MCP 规范里的提示字段(readOnlyHint / destructiveHint 等)。
|
||
// 本实现**透传**它:部分宿主据此算风险等级并在 plan 档下放行非破坏性工具,
|
||
// 漏传的后果不是"少个提示"而是工具在该档下全被拒。
|
||
type ToolSchema struct {
|
||
Name string `json:"name"`
|
||
Description string `json:"description"`
|
||
InputSchema map[string]any `json:"inputSchema"`
|
||
Annotations map[string]any `json:"annotations,omitempty"`
|
||
}
|
||
|
||
// Server 是 MCP 端点。
|
||
type Server struct {
|
||
tools map[string]Tool
|
||
log *log.Logger
|
||
}
|
||
|
||
// NewServer 造一个端点。
|
||
func NewServer(logger *log.Logger) *Server {
|
||
if logger == nil {
|
||
logger = log.Default()
|
||
}
|
||
return &Server{tools: map[string]Tool{}, log: logger}
|
||
}
|
||
|
||
// Register 注册一个工具。同名时后者覆盖前者(测试里常用)。
|
||
func (s *Server) Register(t Tool) {
|
||
s.tools[t.Schema().Name] = t
|
||
}
|
||
|
||
// RegisterAll 批量注册。
|
||
func (s *Server) RegisterAll(ts ...Tool) {
|
||
for _, t := range ts {
|
||
s.Register(t)
|
||
}
|
||
}
|
||
|
||
// ToolCount 供测试与 /mcp 自述用。
|
||
func (s *Server) ToolCount() int { return len(s.tools) }
|
||
|
||
// HandleHTTP 处理一次 POST。
|
||
//
|
||
// 认证在**外层**(main.go 把 middleware.AgentAuth 挂在这条路由上)——
|
||
// 与普通 Agent 端点同一套凭证(Bearer 密钥或 name/secret),
|
||
// 所以 MCP 不能成为绕过既有鉴权与收窄的后门。
|
||
func (s *Server) HandleHTTP(w http.ResponseWriter, r *http.Request) {
|
||
if r.Method != http.MethodPost {
|
||
w.Header().Set("Allow", "GET, POST")
|
||
writeJSON(w, http.StatusMethodNotAllowed, failure(nil, rpcInvalidRequest,
|
||
"POST 是工具面;事件流用 GET /api/v1/mcp(需 Accept: text/event-stream)"))
|
||
return
|
||
}
|
||
// 限制请求体:工具调用都是小 JSON,附件走独立的 /attachments 端点。
|
||
// 1MB 足够,且挡住"把整个文件塞进 JSON"的用法。
|
||
const maxBody = 1 << 20
|
||
body, err := io.ReadAll(http.MaxBytesReader(w, r.Body, maxBody))
|
||
if err != nil {
|
||
writeJSON(w, http.StatusBadRequest, failure(nil, rpcParse, "读请求体失败:%v", err))
|
||
return
|
||
}
|
||
|
||
var msg rpcMessage
|
||
if err := json.Unmarshal(body, &msg); err != nil {
|
||
writeJSON(w, http.StatusBadRequest, failure(nil, rpcParse, "不是合法的 JSON"))
|
||
return
|
||
}
|
||
if msg.JSONRPC != "2.0" {
|
||
writeJSON(w, http.StatusBadRequest, failure(rawID(msg.ID), rpcInvalidRequest, `jsonrpc 字段必须是 "2.0"`))
|
||
return
|
||
}
|
||
|
||
out := s.handleMessage(&msg, r)
|
||
if out == nil {
|
||
// 通知(无 id):没有响应体。按规范回 202。
|
||
w.WriteHeader(http.StatusAccepted)
|
||
return
|
||
}
|
||
writeJSON(w, http.StatusOK, out)
|
||
}
|
||
|
||
// handleMessage 分发一条消息,返回要写回的响应;通知返回 nil。
|
||
//
|
||
// 抽出来是为了能**不经过 HTTP** 单测(httptest 之外也能穷举协议分支)。
|
||
func (s *Server) handleMessage(msg *rpcMessage, r *http.Request) *rpcMessage {
|
||
// 通知没有 id。回了响应,客户端会把响应与请求错配,后续调用全乱。
|
||
isNotification := msg.ID == nil
|
||
id := rawID(msg.ID)
|
||
|
||
switch msg.Method {
|
||
case "initialize":
|
||
if isNotification {
|
||
return nil
|
||
}
|
||
var p struct {
|
||
ProtocolVersion string `json:"protocolVersion"`
|
||
}
|
||
_ = json.Unmarshal(msg.Params, &p)
|
||
// 回显客户端给的版本:不认识的也回显,交由客户端决定是否降级。
|
||
// 自作主张改成我们的版本会让客户端以为协商成功而按新语义调用。
|
||
version := p.ProtocolVersion
|
||
if version == "" {
|
||
version = FallbackVersion
|
||
}
|
||
return result(id, map[string]any{
|
||
"protocolVersion": version,
|
||
"capabilities": map[string]any{"tools": map[string]any{"listChanged": false}},
|
||
"serverInfo": map[string]any{"name": ServerName, "version": ServerVersion},
|
||
})
|
||
|
||
case "notifications/initialized", "initialized":
|
||
return nil // 纯通知
|
||
|
||
case "ping":
|
||
if isNotification {
|
||
return nil
|
||
}
|
||
return result(id, map[string]any{})
|
||
|
||
case "tools/list":
|
||
if isNotification {
|
||
return nil
|
||
}
|
||
list := make([]ToolSchema, 0, len(s.tools))
|
||
for _, t := range s.tools {
|
||
list = append(list, t.Schema())
|
||
}
|
||
// 稳定顺序:map 迭代随机会让客户端每次刷新看到不同排列。
|
||
sortTools(list)
|
||
return result(id, map[string]any{"tools": list})
|
||
|
||
case "tools/call":
|
||
if isNotification {
|
||
return nil
|
||
}
|
||
var p toolCallParams
|
||
if err := json.Unmarshal(msg.Params, &p); err != nil {
|
||
return failure(id, rpcInvalidParams, "tools/call 参数不是合法 JSON:%v", err)
|
||
}
|
||
if p.Name == "" {
|
||
return failure(id, rpcInvalidParams, "tools/call 缺少 name")
|
||
}
|
||
t, ok := s.tools[p.Name]
|
||
if !ok {
|
||
// 未知工具名:回 INVALID_PARAMS 而不是「执行失败」——
|
||
// 前者说"你叫错了",后者说"我试了但失败",模型的反应不同。
|
||
return failure(id, rpcInvalidParams, "没有名为 %s 的工具(可用:%s)", p.Name, strings.Join(s.toolNames(), ", "))
|
||
}
|
||
args := p.Arguments
|
||
if args == nil {
|
||
args = map[string]any{} // 缺 arguments 当空对象,不抛错
|
||
}
|
||
text, err := t.Run(r.Context(), args)
|
||
if err != nil {
|
||
s.log.Printf("[mcp] 工具 %s 失败:%v", p.Name, err)
|
||
// 失败走 result + isError,**不是** JSON-RPC error(见 toolResult 注释)。
|
||
return result(id, errorResult("工具 %s 执行失败:%s", p.Name, err.Error()))
|
||
}
|
||
return result(id, textResult(text))
|
||
|
||
default:
|
||
if isNotification {
|
||
return nil
|
||
}
|
||
return failure(id, rpcMethodNotFound, "不支持的方法 %q", msg.Method)
|
||
}
|
||
}
|
||
|
||
// toolNames 返回已注册工具名(错误文案里提示模型可用集合)。
|
||
func (s *Server) toolNames() []string {
|
||
names := make([]string, 0, len(s.tools))
|
||
for n := range s.tools {
|
||
names = append(names, n)
|
||
}
|
||
sortStrings(names)
|
||
return names
|
||
}
|
||
|
||
func sortTools(list []ToolSchema) {
|
||
for i := 1; i < len(list); i++ {
|
||
for j := i; j > 0 && list[j].Name < list[j-1].Name; j-- {
|
||
list[j], list[j-1] = list[j-1], list[j]
|
||
}
|
||
}
|
||
}
|
||
|
||
func sortStrings(s []string) {
|
||
for i := 1; i < len(s); i++ {
|
||
for j := i; j > 0 && s[j] < s[j-1]; j-- {
|
||
s[j], s[j-1] = s[j-1], s[j]
|
||
}
|
||
}
|
||
}
|
||
|
||
func writeJSON(w http.ResponseWriter, status int, v any) {
|
||
w.Header().Set("Content-Type", "application/json; charset=utf-8")
|
||
w.WriteHeader(status)
|
||
_ = json.NewEncoder(w).Encode(v)
|
||
}
|