返回文章列表

SSE 流式代理:为什么第一帧之后不能再返回 JSON 错误

从 ModelGate 的流式代理实现出发,解释 SSE 的响应边界、首帧前后错误处理、Flush、Context 取消与资源释放。

SSE 流式代理:为什么第一帧之后不能再返回 JSON 错误

在普通 HTTP 接口中,错误处理通常很直接:业务函数返回错误,Handler 选择合适的状态码,再向客户端返回一段 JSON。

{
  "error": {
    "code": "provider_unavailable",
    "message": "upstream provider is unavailable"
  }
}

但大模型的流式接口并不是等完整答案生成后再一次返回。Provider 会持续产生增量内容,网关收到一部分就向客户端转发一部分。在这条连接可能持续数十秒的情况下,错误随时都可能发生。

这时会出现一个看似简单、实际很关键的问题:

如果网关已经向客户端发送了一部分内容,后面又发生错误,还能不能改成 HTTP 500,并返回 JSON?

答案是不能。

这篇文章从 ModelGate 的 SSE 流式代理出发,说明这条边界为什么存在,以及网关应该怎样分别处理首帧之前和首帧之后的错误。

项目地址:theHerta27/ModelGate

SSE 解决了什么问题

SSE 的全称是 Server-Sent Events。它允许服务器在一个持续存在的 HTTP 响应中,不断向客户端发送事件。

一个最简单的 SSE 响应可能是:

data: {"content":"你"}

data: {"content":"好"}

data: [DONE]

每个事件由若干字段组成,并以空行结束。对于大模型接口,常见做法是把每次生成的增量内容放进 data: 字段,最后发送 [DONE] 表示正常结束。

因此,非流式和流式响应的差别不只是“返回得快一点”:

flowchart TD
    Request["客户端请求"] --> Mode{"响应方式"}
    Mode -->|非流式| Wait["等待完整结果"]
    Wait --> JSON["一次返回 JSON"]
    Mode -->|SSE| Stream["持续接收增量"]
    Stream --> Events["逐帧写入并 Flush"]

非流式响应在结果产生前还没有向客户端提交响应;SSE 则会很早发送响应头和第一帧数据,后续处理都发生在一条已经开始的响应里。

HTTP 响应存在一个提交点

一次 HTTP 响应由状态码、响应头和响应体组成。服务端可以先设置:

Content-Type: text/event-stream
Cache-Control: no-cache

然后写入状态码和响应体。

关键在于:一旦响应头或第一段响应体已经被发送,状态码就确定了。

在 Go 的 net/http 中,如果 Handler 没有主动调用 WriteHeader,第一次调用 Write 时会隐式发送 200 OK。之后再调用 WriteHeader(500),已经无法把客户端收到的状态改成 500。

可以把这个过程理解为:

stateDiagram-v2
    [*] --> Uncommitted: 尚未写入响应
    Uncommitted --> JSONError: 首帧前失败
    Uncommitted --> Streaming: 写入响应头或第一帧
    Streaming --> Streaming: 继续发送 SSE
    Streaming --> Closed: 完成、取消或中途失败
    JSONError --> Closed

真正的分界线不是“代码执行到哪一行”,而是响应是否已经提交给客户端。

首帧前失败:仍然可以返回结构化错误

假设网关刚刚完成参数校验,正在请求上游 Provider,但还没有发送任何 SSE 数据。此时如果 Provider 立即返回 429、503,或者根本无法建立连接,网关仍然拥有完整的 HTTP 响应控制权。

它可以返回:

HTTP/1.1 503 Service Unavailable
Content-Type: application/json

同时附带统一的 JSON 错误结构。

这也是 ModelGate 在发送 SSE 响应前先尝试取得第一帧的原因:如果第一帧都无法获得,当前请求本质上还没有进入流式阶段,可以沿用普通 HTTP 错误处理。

简化后的逻辑类似:

stream, err := provider.ChatStream(ctx, req)
if err != nil {
    writeJSONError(w, err)
    return
}
defer stream.Close()

first, err := stream.Recv()
if err != nil {
    writeJSONError(w, err)
    return
}

startSSE(w)
writeChunk(w, first)
flush(w)

这段代码不是为了展示完整实现,而是说明顺序:先确认流能够开始,再正式提交 SSE 响应。

首帧后失败:HTTP 状态已经无法重来

第一帧写出后,客户端可能已经显示了一部分模型回答。此时如果上游连接中断、事件格式损坏,或者网关读取下一帧失败,服务端不能把前面已经返回的 200 OK 改成 500

也不能在 SSE 数据后面直接拼接普通 JSON:

data: {"content":"已经生成的内容"}

{"error":{"code":"internal_error"}}

客户端正在按照 SSE 规则解析事件,突然出现一段不属于该协议的 JSON,只会让响应语义更加混乱。

理论上可以自行定义一种 SSE Error Event,但这要求所有客户端都理解同一套扩展协议。如果网关宣称兼容某个既有接口,却擅自插入新的事件格式,反而可能破坏兼容性。

因此 ModelGate 采用更保守的处理:

  1. 记录结构化错误日志;
  2. 结束当前流;
  3. 关闭上游资源;
  4. 不再尝试伪造一个新的 HTTP JSON 错误。

客户端最终看到的是一次不完整的流,而不是一个“先成功、后又变成500”的响应。

Flush 为什么重要

Go 的 HTTP 服务器和中间代理可能会缓冲输出。如果只调用 Write,一小段数据不一定立即到达客户端。

流式 Handler 通常需要在写入一帧后调用 Flush

func writeSSE(w http.ResponseWriter, payload []byte) error {
    if _, err := fmt.Fprintf(w, "data: %s\n\n", payload); err != nil {
        return err
    }

    flusher, ok := w.(http.Flusher)
    if !ok {
        return errors.New("streaming is not supported")
    }
    flusher.Flush()
    return nil
}

Flush 的作用是要求当前服务端尽快把缓冲数据写出去,但它并不能控制整条网络链路。反向代理、CDN或者客户端自身仍然可能继续缓冲,因此部署流式服务时还要检查代理层配置。

这也是为什么“本地测试可以逐字输出”并不能自动证明生产链路一定具有相同表现。

Context 取消:用户离开后,请求也应该结束

假设用户关闭页面,但网关仍然保持着到 Provider 的连接。上游可能继续生成内容,网关继续读取数据,Token和连接资源也继续被消耗。

正确的请求链路应该让取消信号向下传播:

sequenceDiagram
    participant C as Client
    participant G as ModelGate
    participant P as Provider
    C->>G: 发起 SSE 请求
    G->>P: 使用同一 Context 请求上游
    C--xG: 页面关闭 / 连接取消
    G--xP: Context Done
    P-->>G: 请求结束

在 Go 中,客户端断开后,请求的 Context 会被取消。网关向 Provider 创建请求时继续使用这个 Context,上游调用就能观察到取消信号。

Context 并不会强行“杀死”某段代码。具体的网络请求、阻塞读取和业务函数仍然需要主动监听或使用这个 Context。只有整条调用链都正确传递它,取消才真正有效。

结束一个流,不只是退出循环

流式处理通常还拥有响应体、解析器和上游连接等资源。无论正常完成、客户端取消还是中途失败,都应该进入统一的清理路径。

ModelGate 对流的 Close 采用幂等设计:即使不同清理路径重复调用,也只执行一次实际关闭。Recv 在正常结束后稳定返回 EOF,避免调用方得到相互矛盾的状态。

正常结束时,Provider 会产生 [DONE];调用方收到后停止读取并关闭资源。异常结束时虽然不一定存在 [DONE],但同样需要确保响应体被关闭。

这里的关键不是多写一个 defer,而是先明确资源所有权:

  • 谁创建 Stream?
  • 谁负责读取?
  • 谁负责关闭?
  • Context 取消后由谁退出?
  • 重复关闭是否安全?

只有这些问题有确定答案,异常路径才不容易留下连接和 goroutine。

事件大小也需要边界

SSE 是持续传输,但“持续”不等于单个事件可以无限大。

如果解析器一直读取却等不到事件结束的空行,内存占用可能不断增加。ModelGate 因此为每个完整 SSE 事件设置了 1 MiB 上限,超过限制就停止处理。

限制应该施加在“完整事件”上,而不是随意限制底层某一次网络读取。因为一个合法事件可能被 TCP 拆成多个小片段,也可能一次读取到多个事件。应用层边界和网络分片不是同一个概念。

这一点也解释了为什么流式解析不能简单地假设“一次 Read 就是一帧”。

怎样验证这些错误边界

只测试正常输出无法证明流式代理可靠。ModelGate 的测试重点包括:

  • 第一帧前失败时返回结构化 JSON;
  • 第一帧后失败时终止流,不再修改状态码;
  • 多个 data: 事件能够逐帧转发;
  • [DONE] 后稳定结束;
  • keep-alive、注释和非 data 字段不会被误当成模型内容;
  • 超过 1 MiB 的单个事件被拒绝;
  • 客户端取消能够传播到 Provider;
  • Close 可以重复调用;
  • EOF 行为保持一致。

HTTP Smoke Test 则从真实请求角度检查响应头、增量内容和结束标记,而不仅仅调用某个内部函数。

这些测试能证明已覆盖场景下的行为符合预期,但不能证明所有代理、网络故障和客户端实现都已经被覆盖。流式链路中的真实缓冲和断连问题,仍然需要在完整部署环境中验证。

我从这个问题中理解到的事情

开始实现 SSE 之前,我更容易把错误处理理解成:函数返回 error,Handler 再选择一个状态码。

但流式响应让我看到,错误处理能力会随着响应状态变化:

尚未发送响应
→ 可以选择状态码和 JSON 错误

已经发送第一帧
→ 状态码已确定,只能在流协议范围内处理或结束连接

因此,一个函数“返回了错误”,并不代表外层永远可以用同一种方式把错误交给客户端。系统已经对外做出了什么承诺,会决定后续还剩下哪些选择。

这也是 SSE 流式代理中最重要的边界。

延伸阅读