SSE Gateway

SSE Gateway 是 Go 后端作为 SSE 流式事件的中间代理层,负责从 Python AI Service 接收事件、验证格式、附加元数据、持久化关键事件并转发到前端。

#type / concept #status / evergreen #tech / backend #tech / network

[!info] related notes

SSE Gateway

一句话定义

SSE Gateway 是 Go 后端作为 SSE 流式事件的中间代理层。它不只是”转发”,而是负责事件验证、元数据附加、关键事件持久化、错误处理和连接管理。

它解决什么问题

Python AI Service 产生的 SSE 事件不能直接推送到前端:

  1. 需要鉴权: 前端的连接需要验证身份
  2. 需要持久化: assistant 的回复需要保存到数据库
  3. 需要审计: 每次调用需要记录日志
  4. 需要协议转换: Python 的事件格式可能和前端期望的不同
  5. 需要错误处理: Python 崩溃时需要优雅降级

核心原理

Gateway 的位置

Python AI Service ──SSE──→ Go SSE Gateway ──SSE──→ React 前端

                          ├─ 验证事件格式
                          ├─ 附加元数据
                          ├─ 持久化关键事件
                          ├─ 记录审计日志
                          └─ 错误处理

事件处理流程

func (g *SSEGateway) Proxy(w http.ResponseWriter, r *http.Request, aiStream <-chan Event) {
    // 设置 SSE 响应头
    w.Header().Set("Content-Type", "text/event-stream")
    w.Header().Set("Cache-Control", "no-cache")
    w.Header().Set("Connection", "keep-alive")

    flusher := w.(http.Flusher)
    var assistantBuffer string

    for event := range aiStream {
        // 1. 验证事件格式
        if !event.IsValid() {
            continue
        }

        // 2. 附加元数据
        event.Meta.SessionID = sessionID
        event.Meta.Timestamp = time.Now().Unix()

        // 3. 持久化关键事件
        switch event.Type {
        case "text_delta":
            assistantBuffer += event.Data["delta"].(string)
        case "done":
            g.messageRepo.SaveAssistantMessage(sessionID, assistantBuffer)
        case "error":
            g.messageRepo.MarkMessageFailed(sessionID, event.Data["message"].(string))
        }

        // 4. 审计
        g.auditLog.Record(event)

        // 5. 转发到前端
        fmt.Fprintf(w, "event: %s\ndata: %s\n\n", event.Type, event.JSON())
        flusher.Flush()

        // 6. 检查客户端断开
        if r.Context().Err() != nil {
            g.aiService.CancelRun(runID)
            return
        }
    }
}

心跳保活

func (g *SSEGateway) heartbeat(w http.ResponseWriter, flusher http.Flusher, done <-chan struct{}) {
    ticker := time.NewTicker(15 * time.Second)
    defer ticker.Stop()

    for {
        select {
        case <-ticker.C:
            fmt.Fprintf(w, "event: heartbeat\ndata: {}\n\n")
            flusher.Flush()
        case <-done:
            return
        }
    }
}

典型工程实现

多服务 SSE 管道

Python AI Service ──SSE──→ Go Gateway ──SSE──→ React

                          ├─ 附加 user_id
                          ├─ 附加 session_id
                          ├─ 持久化 text_delta
                          ├─ 记录 tool_call
                          └─ 心跳保活

常见坑

  1. 不做事件验证: Python 输出了格式错误的事件,前端崩溃
  2. 不持久化: assistant 回复没有保存,刷新后丢失
  3. 不做心跳: 连接被 Nginx/CDN 断开
  4. 不做客户端断开检测: 前端断了,后端还在转发
  5. 不做背压: Python 产生事件太快,Go 来不及转发

参考资料

创建于 2026/6/30 更新于 2026/7/15