一、为什么需要通知前六讲的交互模式全是Client 主动拉取PollingClient 调用tools/list获取工具列表Client 调用resources/read读取日志Client 调用tools/call执行动作这种模式在大多数场景下够用但有三个典型痛点场景Polling 的问题订单状态从“支付中”变为“已支付”Agent 不知道除非每隔几秒轮询一次运维工具新增了一个restart_podAgent 还在用旧的工具列表缓存直到缓存过期日志文件追加了新内容Agent 必须反复读取整个文件浪费带宽MCP 的通知机制notifications/*解决了这个问题Server 在状态变化时主动推送消息给 ClientClient 收到后决定是否刷新缓存或采取行动。二、MCP 通知的三类场景MCP 协议定义了三种标准通知通知方法触发时机Client 的响应notifications/tools/list_changedServer 的工具列表发生变化清除工具缓存下次tools/list重新拉取notifications/resources/list_changedServer 的资源列表或模板发生变化清除资源缓存notifications/initializedClient 初始化完成后通知 ServerServer 开始发送其他通知此外MCP 允许自定义通知notifications/*例如notifications/order/status_changed订单状态变更notifications/log/appended日志文件追加三、Go 实现带通知机制的 MCP Server下面扩展第4讲的日志 Resource Server增加工具列表变更通知模拟管理员新增工具日志文件追加通知模拟 tail -f通过 SSE 通道推送给 Clientpackage main import ( encoding/json fmt log net/http os path/filepath strings sync time ) // ---- MCP 消息结构 ---- type JSONRPCReq struct { JSONRPC string json:jsonrpc ID int json:id Method string json:method Params json.RawMessage json:params,omitempty } type JSONRPCResp struct { JSONRPC string json:jsonrpc ID int json:id Result interface{} json:result,omitempty Error *RPCError json:error,omitempty } type RPCError struct { Code int json:code Message string json:message } // ---- 通知消息JSON-RPC 通知格式无 id---- type Notification struct { JSONRPC string json:jsonrpc Method string json:method Params interface{} json:params,omitempty } // ---- 工具定义 ---- type ToolSpec struct { Name string json:name Description string json:description InputSchema any json:inputSchema } // ---- 资源定义 ---- type ResourceTemplate struct { URITemplate string json:uriTemplate Name string json:name Description string json:description MimeType string json:mimeType } // ---- SSE 客户端管理 ---- type SSEClientManager struct { mu sync.RWMutex clients map[chan Notification]bool } func NewSSEClientManager() *SSEClientManager { return SSEClientManager{ clients: make(map[chan Notification]bool), } } func (m *SSEClientManager) Subscribe() chan Notification { ch : make(chan Notification, 100) m.mu.Lock() m.clients[ch] true m.mu.Unlock() return ch } func (m *SSEClientManager) Unsubscribe(ch chan Notification) { m.mu.Lock() delete(m.clients, ch) m.mu.Unlock() close(ch) } func (m *SSEClientManager) Broadcast(notif Notification) { m.mu.RLock() defer m.mu.RUnlock() for ch : range m.clients { select { case ch - notif: default: // 客户端消费太慢丢弃消息 log.Println(通知队列满丢弃消息:, notif.Method) } } } // ---- MCP Server带通知能力---- type NotifyingMCPServer struct { sseManager *SSEClientManager tools map[string]ToolSpec mu sync.RWMutex } func NewNotifyingMCPServer(sseManager *SSEClientManager) *NotifyingMCPServer { return NotifyingMCPServer{ sseManager: sseManager, tools: map[string]ToolSpec{ current_time: { Name: current_time, Description: 返回当前 UTC 与北京时间, InputSchema: map[string]interface{}{ type: object, properties: map[string]interface{}{}, }, }, }, } } // 模拟管理员动态添加工具 func (s *NotifyingMCPServer) AddTool(name, desc string) { s.mu.Lock() s.tools[name] ToolSpec{ Name: name, Description: desc, InputSchema: map[string]interface{}{ type: object, properties: map[string]interface{}{}, }, } s.mu.Unlock() // 广播工具列表变更通知 s.sseManager.Broadcast(Notification{ JSONRPC: 2.0, Method: notifications/tools/list_changed, Params: map[string]interface{}{new_tool: name}, }) log.Printf(已添加工具 %s 并广播通知, name) } func (s *NotifyingMCPServer) Handle(req JSONRPCReq) JSONRPCResp { switch req.Method { case initialize: return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Result: map[string]interface{}{ protocolVersion: 2026-07-28, capabilities: map[string]interface{}{ tools: map[string]interface{}{}, resources: map[string]interface{}{}, }, }, } case tools/list: s.mu.RLock() specs : make([]ToolSpec, 0, len(s.tools)) for _, t : range s.tools { specs append(specs, t) } s.mu.RUnlock() return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Result: map[string]interface{}{tools: specs}, } case tools/call: var params struct { Name string json:name Arguments json.RawMessage json:arguments } json.Unmarshal(req.Params, params) s.mu.RLock() _, exists : s.tools[params.Name] s.mu.RUnlock() if !exists { return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Error: RPCError{Code: -32601, Message: 未知工具}, } } return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Result: map[string]interface{}{ content: []map[string]interface{}{ {type: text, text: fmt.Sprintf(工具 %s 执行成功, params.Name)}, }, }, } case resources/list: return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Result: map[string]interface{}{ resources: []interface{}{}, resourceTemplates: []ResourceTemplate{ { URITemplate: logs://{namespace}/{pod}, Name: Pod 日志, Description: 按 namespace 和 pod 名称读取容器日志, MimeType: text/plain, }, }, }, } case resources/read: var params struct { URI string json:uri } json.Unmarshal(req.Params, params) // 简化实现返回模拟日志 if !strings.HasPrefix(params.URI, logs://) { return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Error: RPCError{Code: -32603, Message: 不支持的 URI}, } } return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Result: map[string]interface{}{ contents: []map[string]interface{}{ {uri: params.URI, mimeType: text/plain, text: fmt.Sprintf([%s] 模拟日志内容, time.Now().Format(time.RFC3339))}, }, }, } default: return JSONRPCResp{ JSONRPC: 2.0, ID: req.ID, Error: RPCError{Code: -32601, Message: 不支持的方法}, } } } // ---- HTTP Transport双端点/mcp 用于 JSON-RPC/events 用于 SSE 通知---- func main() { sseManager : NewSSEClientManager() server : NewNotifyingMCPServer(sseManager) mux : http.NewServeMux() // JSON-RPC 端点 mux.HandleFunc(/mcp, func(w http.ResponseWriter, r *http.Request) { var req JSONRPCReq json.NewDecoder(r.Body).Decode(req) resp : server.Handle(req) w.Header().Set(Content-Type, application/json) json.NewEncoder(w).Encode(resp) }) // SSE 端点Client 订阅通知 mux.HandleFunc(/events, func(w http.ResponseWriter, r *http.Request) { flusher, ok : w.(http.Flusher) if !ok { http.Error(w, Streaming unsupported, http.StatusInternalServerError) return } w.Header().Set(Content-Type, text/event-stream) w.Header().Set(Cache-Control, no-cache) w.Header().Set(Connection, keep-alive) ch : sseManager.Subscribe() defer sseManager.Unsubscribe(ch) // 发送初始连接成功消息 fmt.Fprintf(w, event: connected\ndata: {\status\:\ok\}\n\n) flusher.Flush() ctx : r.Context() for { select { case -ctx.Done(): log.Println(SSE 客户端断开连接) return case notif : -ch: data, _ : json.Marshal(notif) fmt.Fprintf(w, event: message\ndata: %s\n\n, data) flusher.Flush() } } }) // 模拟每 15 秒自动添加一个新工具演示通知 go func() { toolNames : []string{ping_host, check_disk, restart_service, tail_log} idx : 0 for { time.Sleep(15 * time.Second) if idx len(toolNames) { server.AddTool(toolNames[idx], fmt.Sprintf(模拟工具 #%d, idx1)) idx } } }() // 模拟每 10 秒广播一条日志追加通知 go func() { for { time.Sleep(10 * time.Second) sseManager.Broadcast(Notification{ JSONRPC: 2.0, Method: notifications/log/appended, Params: map[string]interface{}{ uri: logs://production/web-abc, timestamp: time.Now().Format(time.RFC3339), line: fmt.Sprintf(新日志条目 at %s, time.Now().Format(time.RFC3339Nano)), }, }) } }() log.Println(MCP 通知 Server 启动:) log.Println( JSON-RPC: http://localhost:8083/mcp) log.Println( SSE: http://localhost:8083/events) log.Fatal(http.ListenAndServe(:8083, mux)) }四、Client 端对接通知Client 需要同时维护两个连接HTTP POST 到/mcp发送 JSON-RPC 请求SSE 连接到/events接收 Server 推送的通知// ---- MCP Client 带通知订阅 ---- type NotifyingMCPClient struct { serverURL string token string client *http.Client mu sync.RWMutex toolsCache []ToolSpec cacheExpiry time.Time } func (c *NotifyingMCPClient) StartEventLoop(ctx context.Context) { // 建立 SSE 连接 req, _ : http.NewRequestWithContext(ctx, GET, c.serverURL/events, nil) req.Header.Set(Authorization, Bearer c.token) resp, err : c.client.Do(req) if err ! nil { log.Printf(SSE 连接失败: %v, err) return } defer resp.Body.Close() reader : bufio.NewReader(resp.Body) for { line, err : reader.ReadString(\n) if err ! nil { log.Printf(SSE 读取失败: %v, err) return } if strings.HasPrefix(line, data: ) { var notif Notification json.Unmarshal([]byte(line[6:]), notif) switch notif.Method { case notifications/tools/list_changed: log.Println(收到工具列表变更通知清除缓存) c.mu.Lock() c.toolsCache nil c.cacheExpiry time.Time{} c.mu.Unlock() case notifications/log/appended: log.Printf(收到日志追加通知: %v, notif.Params) // 可以选择主动刷新相关 Resource } } } }五、通知 vs Polling 的选择场景推荐方式原因工具列表变更通知notifications/tools/list_changed变更频率极低通知即可无需轮询日志文件追加通知 增量读取通知告知有新内容Client 再用resources/read读增量订单状态变更通知实时性要求高轮询延迟大配置项变更通知变更频率低通知即可指标监控数据Polling定期拉取数据持续变化通知会过于频繁经验法则状态变化频率低于每分钟1次的用通知高于每分钟1次的用 Polling 本地缓存。六、安全分层L6本讲在 L1-L5 基础上补充通知机制的安全关注点L6通知安全SSE 端点同样需要认证Bearer Token不能开放给未经授权的 Client通知内容不应包含敏感信息Token、密钥、个人数据Server 应限制每个 Client 的通知队列大小防止慢 Client 导致 Server 内存泄漏通知速率限制每秒不超过 10 条通知防止恶意 Server 淹没 ClientClient 应验证通知的来源确认来自可信的 Server防止伪造通知触发缓存清空七、延伸阅读MCP Specification – Notifications标准通知方法的定义和格式SSE (Server-Sent Events) 规范W3CEventSource API 和事件流格式WebSocket vs SSE 对比为什么 MCP 选择 SSE 而非 WebSocket单向通知 vs 双向通信八、下一讲预告第8讲MCP 生态——Hub、网关与联邦单个 MCP Server 好写但生产环境面对的是几十上百个 Server谁来管理它们的注册与发现谁来统一鉴权与限流不同团队的 Server 如何互操作下一讲介绍 MCP Hub、MCP Gateway 的概念以及如何用 Go 实现一个轻量级的 MCP 网关统一管理多个下游 Server。开发之余的小工具推荐处理 Base64、JSON 格式化、JWT 解析、Crontab 计算、PDF 合并压缩这些碎片需求我常用一个纯前端本地工具箱zz365.top。所有计算在浏览器完成文件不上服务器关页即清。免费、无登录、无广告适合开发者当常驻标签页。