WebSocket 会话网关:支持双向打断与长流式推流的状态机设计 WebSocket 会话网关支持双向打断与长流式推流的状态机设计在智能体Agent系统的实时交互体验中虽然标准的 SSEServer-Sent Events能够极好地满足“服务端向客户端单向吐字”的流式场景但当业务演进到**“全双工多模态语音交互、实时协作画板、以及用户随时主动打断User Interruption”**的高级交互场景时单向的 SSE 协议便显得捉襟见肘。用户期望获得如同与真人交谈一般的敏捷体验当大模型正在长篇大论推流一段长文本时用户突然说了一句“停我不想听这个请直接告诉我价格”如果系统基于 SSE由于无法反向发送指令用户只能等待几十秒或者暴力刷新页面而在WebSocket 全双工长连接架构下客户端可以瞬间向服务端发送一个INTERRUPT控制帧服务端立即在毫秒级中断大模型生成协程、取消下游工具执行、释放显存 KV Cache并无缝无卡顿地开启新一轮推理如何在生产环境中设计并实现一个具备“实时流式下发 双向主动打断 状态机强一致性”的工业级 WebSocket 会话网关一、支持实时打断的 WebSocket 会话状态机模型[ 初始连接建立 ] ──► ( 状态 1: IDLE 空闲就绪 ) │ ▼ (接收到客户端 USER_QUERY 指令) ┌────────────────────────────────────────────────────────────────────────┐ │ 状态 2: GENERATING (正在流式推流与工具执行中) │ │ 动作: 协程 A 持续向下发数据帧: {type: CHUNK, text: ...} │ └─────────────────────────────────────┬──────────────────────────────────┘ │ ┌─────────────────────────────┴─────────────────────────────┐ ▼ (大模型正常吐完最后一个 Token) ▼ (客户端突然反向发送 INTERRUPT 指令) ┌───────────────────────┐ ┌───────────────────────┐ │ 状态 3: COMPLETED │ │ 状态 4: INTERRUPTING │ │ 下发: {type: DONE}│ │ 1. 触发 Context Cancel │ │ 回到 IDLE 等待下一轮 │ │ 2. 终止下游大模型推流 │ └───────────────────────┘ │ 3. 下发 {type: STOP} │ 4. 毫秒级复位至 IDLE │ └───────────────────────┘二、生产级 Go 语言 WebSocket 双向打断网关实现实操基于gorilla/websocket与 Go Context 原语构建具备绝对状态隔离与并发安全的会话控制器package wsgateway import ( context encoding/json log net/http sync time github.com/gorilla/websocket ) var upgrader websocket.Upgrader{ CheckOrigin: func(r *http.Request) bool { return true }, // 生产需配置真实域名白名单 } type WSClientMessage struct { Type string json:type // QUERY, INTERRUPT, PING Payload string json:payload // 用户自然语言或参数 } type WSServerMessage struct { Type string json:type // CHUNK, STOPPED, DONE, PONG Text string json:text,omitempty } type ActiveSession struct { conn *websocket.Conn mu sync.Mutex cancelFunc context.CancelFunc // 用于打断当前生成协程的取消句柄 isWorking bool } func (s *ActiveSession) HandleClientMessages() { defer s.conn.Close() for { _, msgBytes, err : s.conn.ReadMessage() if err ! nil { log.Printf(【WS 断开】客户端连接关闭: %v, err) s.InterruptCurrentGeneration() return } var clientMsg WSClientMessage if err : json.Unmarshal(msgBytes, clientMsg); err ! nil { continue } switch clientMsg.Type { case INTERRUPT: // 1. 核心打断指令瞬间掐断正在生成的协程 log.Printf(【主动打断 ⚡】接收到用户打断信号正在终止生成...) s.InterruptCurrentGeneration() s.sendJSON(WSServerMessage{Type: STOPPED}) case QUERY: // 2. 新提问到达若前序任务还在跑先打断旧任务再开启新任务 s.InterruptCurrentGeneration() // 创建可取消的专属 Context ctx, cancel : context.WithCancel(context.Background()) s.mu.Lock() s.cancelFunc cancel s.isWorking true s.mu.Unlock() // 异步启动新一轮大模型推流工作协程 go s.runStreamingGeneration(ctx, clientMsg.Payload) case PING: s.sendJSON(WSServerMessage{Type: PONG}) } } } func (s *ActiveSession) InterruptCurrentGeneration() { s.mu.Lock() defer s.mu.Unlock() if s.cancelFunc ! nil s.isWorking { s.cancelFunc() // 广播 Context 取消信号 s.cancelFunc nil s.isWorking false } } func (s *ActiveSession) runStreamingGeneration(ctx context.Context, prompt string) { defer func() { s.mu.Lock() s.isWorking false s.mu.Unlock() }() // 模拟从下游大模型通道拉取流式 Token tokens : []string{根据, 您的, 订单, 记录, 当前, 商品, 处于, 配货中} for _, token : range tokens { // 每次推流前必须检查是否被用户打断 select { case -ctx.Done(): log.Printf(【生成协程退出】检测到 Context 取消已安全释放资源。) return default: } // 物理推流下发 s.sendJSON(WSServerMessage{Type: CHUNK, Text: token}) time.Sleep(100 * time.Millisecond) // 模拟大模型生成间隔 } s.sendJSON(WSServerMessage{Type: DONE}) } func (s *ActiveSession) sendJSON(msg *WSServerMessage) { s.mu.Lock() defer s.mu.Unlock() _ s.conn.WriteJSON(msg) }三、生产治理的三大关键避坑要点写并发锁Write Mutex Lock不可或缺Gorilla WebSocket 的底层连接在同一时刻绝对不允许并发写Concurrent Write如果心跳协程和数据推流协程同时调用WriteJSON进程会直接抛出致命 panic。必须通过sync.Mutex严格保护所有写操作打断信号必须级联透传至大模型接口在 ContextCancel()触发后后端与 OpenAI 或自建 vLLM 的底层 HTTP/TCP 连接必须同步关闭通知显卡立即停止该请求的 KV Cache 计算避免被用户打断后大模型还在后台继续烧钱空转前端渲染防乱序与版本戳Generation Version Epoch每次触发打断后客户端与服务端递增会话版本号epoch epoch 1前端自动丢弃所有滞后的旧版本残余数据帧彻底消灭文字乱序闪烁。支持随时打断的双向 WebSocket 网关赋予了用户前所未有的控制感与交互灵敏度是打造下一代拟人化实时智能体的必备核心工程利器。