ssex

package module
v0.2.2 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 3, 2026 License: MIT Imports: 16 Imported by: 0

README

SSE Package

github.com/gtkit/ssex 提供 Server-Sent Events 的服务端写入与上游解码能力。

入口只用标准库的 http.ResponseWriter*http.Request,net/http、gin、chi 都可直接使用;唯一的直接第三方依赖是 github.com/gtkit/json/v2(它自身另有传递依赖)。

1. 能力总览

能力 API
写响应头并立即刷给客户端 Stream.Start() / Writer.WriteHeaders()
命名事件 / 带 id 事件 / data-only 帧 Event / EventWithID / Data / Send
保活 Comment / Ping / Heartbeat
重连间隔建议 Retry
起流状态判断 Started
终止流并抑制客户端自动重连 Close
错误分类 ErrClientGone / ErrStreamClosed / ErrInvalidArgument / ErrWriteTimeout / ErrFrameTooLarge
单帧写超时 WithWriteTimeout
解析上游 SSE 流(转发大模型输出) Decode
在线连接注册表与带外推送 Hub
读取客户端重连位点 LastEventID

针对长连接做的事:起流时清零 http.Server.WriteTimeout(否则待支付订单、大模型长响应会在全局超时到期时被服务端 RST)、每帧设置独立写截止时间、下发 X-Accel-Buffering: no 关闭 nginx 缓冲、HTTP/2 下按 RFC 9113 省略 Connection 头、Flush 错误在当帧就暴露。

帧安全:id / event 字段值含换行或 NUL 时报错且不写字节;注释与原样透传载荷中的 CRLF / 孤立 CR 归一后逐行加前缀,内容无法逃出所在字段伪造帧。

2. 包结构

当前主要有两层:

2.1 Writer

文件:

职责:

  • 最底层 SSE 帧写入
  • 写响应头
  • event/data
  • comment
  • retry

适合:

  • 你需要完全自己控制输出顺序
  • 你明确知道什么时候开始响应
2.2 Stream

文件:

职责:

  • Writer 之上增加一层业务友好的能力
  • 首次输出时自动写 Header
  • 提供 Started() 状态
  • 提供 Ping() / Error() / Comment() / Retry() 这些标准辅助方法

适合:

  • 大多数业务 SSE 场景
  • 比如订单状态流、LLM 流式输出

当前建议:

  • 业务代码优先使用 Stream
  • 只有极少数需要完全手控底层帧顺序的场景,才直接用 Writer
2.3 其他文件
  • decode.go:上游 SSE 解码器(Decode / Message
  • hub.go:在线连接注册表与带外推送(Hub
  • event.go:待推送的事件值(Event
  • errors.go:错误哨兵与分类(5 个哨兵,见 4.10)
  • options.go:Functional Options(WithWriteTimeout

3. 快速开始

3.1 最小示例
package demo

import (
    "time"

    "github.com/gin-gonic/gin"
    "github.com/gtkit/ssex"
)

func StreamDemo(c *gin.Context) {
    stream := ssex.NewStream(c.Writer, c.Request)

    // 首次 Event 会自动写 text/event-stream 头
    if err := stream.Event("status", gin.H{
        "status": "pending",
    }); err != nil {
        return
    }

    if err := stream.Ping(time.Now()); err != nil {
        return
    }

    if err := stream.Event("status", gin.H{
        "status": "done",
        "final":  true,
    }); err != nil {
        return
    }
}
3.2 net/http 原生
func StreamDemo(w http.ResponseWriter, r *http.Request) {
    stream := ssex.NewStream(w, r)

    if err := stream.Event("status", map[string]string{"status": "pending"}); err != nil {
        return
    }
    _ = stream.Close(map[string]string{"reason": "done"}) // 生产代码要处理它,见 4.9
}

4. API 说明

4.1 NewStream(w, r)

创建一个业务层友好的 SSE 输出器。wr 都必须非 nil——传 nil 是编程错误,库不做防御性降级,因为一个写不出字节的 Stream 比直接 panic 难定位得多。同理 Decode(r)rLastEventID(r)rHeartbeat(ctx, …)ctx 都必须非 nil。

这是本库唯一一处会 panic 的路径,且是经过评估后有意保留的:拿到 error 的调用方对"没有 ResponseWriter"无事可做,唯一正确处理就是改代码,而这个错误在任何一次请求里都会立刻暴露;让全部正确调用方多写一次判断换不到任何东西。取值合法性问题(字段含换行、retry 越界等)一律返回 error,不 panic。该行为由 TestNilArgumentsPanic 固定,PATCH 版本不会改动它。

stream := ssex.NewStream(c.Writer, c.Request)
4.2 stream.Event(name, payload)

发送具名事件。

_ = stream.Event("status", gin.H{
    "status": "pending",
})

输出类似:

event: status
data: {"status":"pending"}

4.3 stream.Ping(at)

发送标准保活注释帧。

_ = stream.Ping(time.Now())

输出类似:

: ping 2026-03-30T10:00:00Z

4.4 stream.Error(payload)

发送标准业务错误事件。

_ = stream.Error(gin.H{
    "error": "order not found",
})

输出类似:

event: error
data: {"error":"order not found"}

注意:

  • 这里是服务端业务事件
  • 不是浏览器 EventSource.onerror 那个网络层错误回调
4.5 stream.Data(payload)

发送 data-only 帧(无 event: 行),适合 OpenAI 风格流式接口。

普通 payload 会经 github.com/gtkit/json/v2 序列化:

_ = stream.Data(gin.H{"delta": "hello"})

需要字面输出 [DONE] 这类非 JSON 哨兵时,用 ssex.Raw

_ = stream.Data(ssex.Raw("[DONE]"))
4.6 stream.Comment(text)

发送一条 SSE 注释帧。

_ = stream.Comment("keepalive")

输出类似:

: keepalive

用途:

  • 保活
  • 调试
  • 某些代理环境下防止长时间空闲断开
4.7 stream.Retry(milliseconds)

发送 SSE 的 retry 指令,提示客户端后续重连间隔。

_ = stream.Retry(3000)

输出类似:

retry: 3000

用途:

  • 给浏览器原生 EventSource 提供重连节奏建议
4.8 stream.Started()

判断当前响应是否已经开始输出。

这个方法在“业务失败时还想回普通 JSON”的场景里很有用。

例如:

if err != nil {
    if !stream.Started() {
        // 还能回普通 JSON
        c.JSON(http.StatusBadRequest, ...)
        return
    }
    // 已经开始 SSE 输出,只能继续发 SSE error
    _ = stream.Error(gin.H{"error": "internal error"})
    return
}
4.9 stream.Close(payload)

终止本轮推送:发送一条 event: close 事件,并拒绝后续所有写入。

为什么必须有这一步:浏览器 EventSource 在服务端正常结束流后会按重连间隔自动重连。订单进入终态、LLM 输出结束后直接 return,前端会马上重连,形成无限循环。

// 服务端:终态后终止流
if err := stream.Close(gin.H{"reason": "delivered"}); err != nil {
    // 终止帧没送达前端会继续重连,除客户端已断开外都要能发现(见下)
    logger.Warn("发送终止事件失败", zap.Error(err))
}
// 前端:收到 close 事件后主动断开,否则仍会重连
es.addEventListener('close', () => es.close());

Close 只终结本流的写入许可,不关闭 HTTP 连接(连接在 handler 返回时结束)。 终止后任何写入返回 ssex.ErrStreamClosed,重复 Close 同样返回它。

它的错误不该一律忽略。 终止帧没送达时前端收不到 close,会继续按重连间隔重连,而服务端毫无信号。客户端已断开(ErrClientGone)属正常收尾,其余要能被发现:

func closeStream(stream *ssex.Stream, payload any) {
    err := stream.Close(payload)
    if err == nil || errors.Is(err, ssex.ErrClientGone) {
        return // 正常收尾
    }
    // 写超时或未知写错误:连接还在但帧没发出去,前端会继续重连
    logger.Warn("发送终止事件失败", zap.Error(err))
}
4.10 错误判定

写方法的错误分三类,用 errors.Is 区分:

判定 含义 建议处理
ErrClientGone 客户端已断开 静默结束,并取消上游请求(如正在进行的大模型调用)以停止计费
ErrStreamClosed 自己已 Close 调用顺序问题,检查业务逻辑
ErrInvalidArgument 参数非法:字段值含换行/NUL、retry 为负、心跳间隔非正 编程错误。仅当 Started() 为 false 时(即首帧就失败)响应头尚未提交、可改回普通 JSON;流已开始后它只表示这一帧被拒,此时不能再写普通响应体
ErrWriteTimeout 单帧写入超过 WithWriteTimeout 客户端读取过慢,连接仍活着;按需告警
ErrFrameTooLarge 解码时单帧的 Message.Data 超过 1048575 字节 上游异常或恶意,终止本次转发
其余 真实写失败 记录日志 / 告警

ErrClientGone 的错误链保留原因,因此 errors.Is(err, context.Canceled) 同样成立。 写超时不会被判定为 ErrClientGone——慢客户端不等于断开。

for chunk := range upstream {
    if err := stream.Data(chunk); err != nil {
        if errors.Is(err, ssex.ErrClientGone) {
            cancelUpstream() // 前端已关页面,别再为它烧 token
            return
        }
        logger.Error("sse 写入失败", err)
        return
    }
}

首帧失败仍可回 JSONStream 的写方法先构造并校验完整帧,只有构造成功才提交响应头。 因此首帧因 ErrInvalidArgument 或序列化失败而报错时 Started() 仍为 false, 调用方可以改用普通 JSON 响应回错。

起流错误要处理Start()WriteHeaders() 返回 error。纯推送型 handler 起流失败(客户端已断开)时应直接返回,不必再进消费循环白等。

4.11 Option
stream := ssex.NewStream(c.Writer, c.Request, ssex.WithWriteTimeout(30*time.Second))
Option 默认值 说明
WithWriteTimeout(d) 10s 单帧写入的截止时长;非正值忽略。弱网移动端可放大,内网可收紧以更快发现死连接

该截止时间只作用于单帧,长连接整体不受它限制——长连接本身靠起流时清零 http.Server.WriteTimeout 保活。

4.12 stream.Heartbeat(ctx, interval)

长时间无数据的流(等支付结果、等登录态)必须有心跳,否则代理层会按空闲连接断开。

// 错误结果走 hbErr,"已退出"走 hbDone——两者必须分开,理由见下
hbCtx, stopHeartbeat := context.WithCancel(c.Request.Context())
hbErr := make(chan error, 1)
hbDone := make(chan struct{})
go func() {
    defer close(hbDone)
    if err := stream.Heartbeat(hbCtx, 15*time.Second); err != nil {
        hbErr <- err
    }
}()
defer func() {
    stopHeartbeat()
    <-hbDone // 等它真的退出,再让 handler 返回
}()

阻塞运行,启动时机与所在 goroutine 由你控制。ctx 取消返回 ctx 错误,客户端断开返回 ErrClientGone,流已 Close 返回 ErrStreamClosedinterval 非正立即报错。

两条约束都不能省:

  • 必须等它退出:底层 ResponseWriter 只在 handler 执行期间有效(gin 会池化复用,见 9.1)。handler 返回后心跳若还在写,就写到别人的响应上了。
  • 错误信号与退出信号必须分开:若用同一个 channel 兼任,主循环消费掉唯一那次发送后,defer 里的第二次接收就永远等不到发送者——handler 永久阻塞,连排在后面的 release() 都不会执行。hbDoneclose,读多少次都不会阻塞。
4.13 ssex.Decode(r) — 解析上游 SSE

把上游大模型的 text/event-stream 解码成结构化帧,用于转发:

for msg, err := range ssex.Decode(resp.Body) {
    if err != nil {
        return err
    }
    if string(msg.Data) == "[DONE]" {
        break
    }
    if err := stream.Data(ssex.Raw(string(msg.Data))); err != nil {
        return err
    }
}

Message 字段:IDNameevent: 字段,空为默认事件)、Data []byteRetry

按 WHATWG event stream 规范解析,实际行为如下:

  • 行分隔符支持 CRLF / CR / LF
  • 字段值只剥掉冒号后的一个空格——大模型的增量 token 常以空格开头,多剥会损坏拼出的文本
  • Data 原样保留上游字节:[DONE] 这类非 JSON 哨兵可直接比对,增量 token 的前后空格完整保留
  • 多行 data:\n 拼接;流首 BOM、注释行与未知字段跳过
  • 帧是否产出以本帧是否出现过 data 字段为准(与浏览器 EventSource 一致),值为空串也照常产出
  • IDRetry 是连接级状态,按规范跨帧沿用(The buffer does not get reset),转发时 id 连续性不会断;Name 每帧独立
  • 流尾残留的不完整帧(缺结尾空行)按规范丢弃,与浏览器一致——连接被掐断时残留字节通常是截断的半条 JSON
  • 产出的 Message.Data 最多 1048575 字节,超限返回 ErrFrameTooLarge 而非静默截断;单行与多行同一口径(多行按拼接后的长度算),单行长度另有一个略高的硬上限防止超长行撑爆缓冲
  • retry 只接受纯 ASCII 数字,且限制在不会让 time.Duration 溢出的范围内,超范围时忽略该字段
  • 读到流尾正常结束
  • 与规范的一处差异:Data 原样保留上游字节,不执行 UTF-8 解码替换。转发链路上改写字节会让服务端交给前端的内容与上游不一致,而浏览器接收端本身会按 UTF-8 解码;需要严格校验时自行处理 Data
4.14 ssex.NewHub() — 带外推送

产生事件的 goroutine 不持有连接时(支付回调推订单状态、后台踢登录态)用它:

hub := ssex.NewHub()                        // 无后台 goroutine,不需要关闭

// 推送方(支付回调、MQ 消费者……)
// 顺序不能反:先在事务里持久化状态与 revision、提交成功,再发布事件。
delivered, dropped := hub.Push(uid, ssex.Event{ID: rev, Name: "status", Data: order})
metrics.Observe(delivered, dropped) // 只做监控,不参与任何业务判断

delivered 不是交付确认。 它只表示事件成功进入了多少个本机内存队列,不代表客户端收到、Flush 成功、浏览器处理完成、其他实例上的连接收到、连接没有随即断开,也不代表这条事件没有被后续溢出挤掉。因此绝不能写成"delivered == 0 才落库"——状态必须先持久化再发布,delivered / dropped 只用于监控。

连接侧的 handler 模板:

func Subscribe(c *gin.Context) {
    // key 必须由服务端身份计算,且非空:空 key 会让所有认证异常的请求
    // 共享同一个队列、互相收到对方的事件(见 4.15)。
    uid := c.GetString("uid")
    if uid == "" {
        c.AbortWithStatusJSON(http.StatusUnauthorized, gin.H{"error": "未登录"})
        return
    }

    // 先订阅、再起流:Start 最坏要等一个完整的 WithWriteTimeout,
    // 这段时间里的推送若无人订阅就永久丢失。
    events, release := hub.Subscribe(uid)
    defer release()

    stream := ssex.NewStream(c.Writer, c.Request)

    // 纯推送型 handler 在第一条事件到来前不写字节,不起流则前端一直等响应头、
    // 迟迟不触发 onopen。起流失败时响应头尚未提交,还能回普通 JSON。
    if err := stream.Start(); err != nil {
        return
    }

    // 心跳错误要能被主循环看到:慢客户端会让心跳返回 ErrWriteTimeout,
    // 而此时请求上下文可能还没结束,只丢弃错误会让 handler 继续空等。
    //
    // 错误走 hbErr,"已退出"走 hbDone,两者必须分开(见 4.12);
    // 并且必须等它退出——c.Writer 会被 gin 池化复用(见 9.1)。
    hbCtx, stopHeartbeat := context.WithCancel(c.Request.Context())
    hbErr := make(chan error, 1)
    hbDone := make(chan struct{})
    go func() {
        defer close(hbDone)
        if err := stream.Heartbeat(hbCtx, 15*time.Second); err != nil {
            hbErr <- err
        }
    }()
    defer func() {
        stopHeartbeat()
        <-hbDone
    }()

    for {
        select {
        case <-shutdownCtx.Done(): // 应用级停机信号,见 6.9
            closeStream(stream, gin.H{"reason": "server shutting down"}) // 见 4.9
            return
        case <-stream.Context().Done():
            return
        case <-hbErr: // 心跳写失败:连接已不可用
            return
        case e := <-events:
            if err := stream.Send(e); err != nil {
                return
            }
        }
    }
}
方法 说明
Subscribe(key) 注册连接,返回事件队列与注销函数(必须 defer 调用,重复调用安全)
Push(key, event) 定向投递,返回 (delivered, dropped)
Broadcast(event) 投给所有在线连接,返回 (delivered, dropped)
Online(key) 该 key 当前在线连接数
WithQueueSize(n) 每连接队列容量,默认 32,非正值忽略
返回值 Push / Broadcast 返回 (delivered, dropped)dropped 是被挤掉的事件数

七条使用约束:

  1. Start() 再进消费循环:纯推送型 handler 在第一条事件到来前不写字节,起流才能让前端及时拿到响应头并触发 onopenStart() 返回错误说明连接已不可用,直接返回。
  2. 投递与写出分离Push 只把事件放进队列,写出由持有连接的 handler 完成,因此写动作始终发生在 handler 自己的 goroutine、ResponseWriter 仍有效的窗口内。
  3. 队列满时保留最新(latest-wins):挤掉队首最旧的一条,让最新事件一定入队,dropped 是被挤掉的旧事件数。状态推送里最新一条描述的就是当前状态,必须送达。
  4. 顺序只在单连接内保证:多个 goroutine 并发 Push / Broadcast 时,它们之间的相对顺序由调度决定,不同连接可能观察到不同顺序。状态事件应携带单调递增的 revision(放进 Event.ID 或载荷),消费端按 revision 取大者、忽略更旧的值。
  5. 载荷投递后不可再改:同一份 Event.Data 会被该 key 下的多个连接各自序列化,且发生在它们各自的 goroutine 里;投递后再修改其中的 map、切片或指针内容会构成数据竞争。需要复用结构体就投递前拷贝一份。
  6. 容量与顺序的组合风险:latest-wins 按到达顺序保留最后一条,不比较业务 revision。多个并发推送方存在时,到达顺序可能与 revision 顺序相反——容量为 1 时后到的旧版本会挤掉先到的新版本,消费端再也看不到它,lastRev 过滤也救不回来。因此 WithQueueSize(1) 只在同一 key 的推送已串行化时才安全;否则至少满足一条:同一 key 经单点串行化后再 Push/容量留大于 1 并由消费端按 revision 过滤/终态后回读一次事实源做校准。
  7. 适用边界:Hub 面向状态推送——每条事件自带完整状态,丢掉中间过程无妨。AI token 流每一条都是文本的一部分,少一条就损坏输出,应在 handler 内直接用 stream.Data 写出,不经 Hub。

消费循环用 stream.Context().Done() 退出,队列由 GC 回收。

4.15 Hub 的两条硬边界

一、Hub 只管当前进程的连接。 多实例部署时,连接在实例 A、支付回调打到实例 B,实例 B 的 Push 找不到那条连接,返回 delivered == 0,前端收不到更新;Online 也只反映本机而非集群。滚动发布、负载均衡、EventSource 重连换实例都会触发。

单实例可以直接用。多实例必须让每个实例都收到状态事件,再各自推给本地 Hub:

数据库事务:写状态 + 单调递增 revision
        ↓ 提交成功
Outbox / MQ / Redis Pub/Sub
        ↓ 广播式消费(每个实例都收到)
各实例 hub.Push(key, event)
        ↓
本实例上的连接收到状态

也可以用一致性哈希把同一 key 的连接与事件路由到同一实例,但运维复杂度更高。无论哪种,存储始终是事实源,Hub 只是加速通道。

AI 对话流在 handler 内直连上游、不经 Hub,因此不受这条限制。

二、key 必须是服务端计算出的安全作用域。 Hub 的 key 就是一个字符串,两个不同租户的 orderID=100 会落进同一个 key,造成跨租户状态串流。约束:

  • key 只能由已认证身份 + 已授权资源计算,不能直接用 URL、query 或客户端传入的值
  • 多租户必须带 tenant 作用域,如 tenantID:userID:orderID,或用服务端生成的全局唯一不可猜测 ID
  • 禁止空 key:空 key 会让所有认证异常的请求共享同一个队列、互相收到对方的事件
  • 不要裸拼接"a:b" + ":" + "c""a" + ":" + "b:c" 会得到同一个 key。用长度前缀编码或服务端生成的全局唯一 ID,并统一封装成一个 key 构造函数,禁止业务代码各自拼接
  • 读资源时沿用授权阶段的同一标识,不要退回裸资源 ID——否则快照可能读到另一个租户的数据
  • 订阅前二次校验授权结果:授权实现若因 bug 返回零值标识与 nil error,key 会退化成一个固定值,多个异常请求就此落进同一队列。检查租户与资源标识都非空,可以把这类 bug 从静默串流降级成一次 403
  • Broadcast 只用于确实允许所有在线用户看到的内容

gincompat 里的 resource.scopeKey() 用长度前缀编码做了示范,并有 TestScopeKeyHasNoCollision 与走真实 handler 的 TestTenantCannotReachOtherTenantOrder 两个回归测试。

5. 常见模式

5.1 转发大模型输出

session,中间连续 chunk,结束 done,出错 error;上游用 Decode 读,下游用 Data / Send 写。

上游客户端要按阶段设超时,不要用 http.Client.Timeout——它覆盖整个响应体读取周期,会把正常的长流掐断:

// 全局复用一个客户端。Timeout 留零值,改为限制建连、TLS 与等响应头这三个阶段;
// "最长生成时长"由业务用 context.WithTimeout 单独控制。
//
// 必须从 DefaultTransport 克隆再覆盖,不要自己 new(http.Transport)——后者会丢掉
// ProxyFromEnvironment(企业代理失效)、ForceAttemptHTTP2(自定义 DialContext 时
// 不再尝试 HTTP/2)、MaxIdleConns / IdleConnTimeout / ExpectContinueTimeout 等
// 标准库调好的默认值。
var upstreamClient = &http.Client{Transport: newUpstreamTransport()}

func newUpstreamTransport() *http.Transport {
    tr := http.DefaultTransport.(*http.Transport).Clone()
    tr.TLSHandshakeTimeout = 5 * time.Second
    tr.ResponseHeaderTimeout = 30 * time.Second // 首个响应头就绪的上限,不含流式响应体
    // 只调建连超时,保留 KeepAlive
    tr.DialContext = (&net.Dialer{Timeout: 5 * time.Second, KeepAlive: 30 * time.Second}).DialContext

    return tr
}
// 上游请求绑定下游 ctx:客户端一断开,上游请求随之取消;
// 同时叠加一个业务级的最长生成时长。
upCtx, cancel := context.WithTimeout(c.Request.Context(), maxGenerationTime)
defer cancel()

req, err := http.NewRequestWithContext(upCtx, http.MethodPost, upstreamURL, body)
if err != nil {
    return err
}
resp, err := upstreamClient.Do(req)
if err != nil {
    return err
}
defer func() { _ = resp.Body.Close() }() // 始终关闭,否则连接与 goroutine 泄漏

if resp.StatusCode != http.StatusOK {
    return fmt.Errorf("上游状态码 %d", resp.StatusCode)
}

// 用 mime.ParseMediaType 精确比较,不要 strings.HasPrefix:后者会接受
// text/event-streaming、text/event-stream-error 这类非 SSE 类型,
// 又会拒绝合法的大小写与空白变体(media type 本身大小写不敏感)。
rawCT := resp.Header.Get("Content-Type")
mediaType, _, err := mime.ParseMediaType(rawCT)
if err != nil || mediaType != "text/event-stream" {
    return fmt.Errorf("上游 Content-Type 非 SSE: %q", rawCT) // 上游报错时常返回 JSON
}

// 上游校验通过、确定要开这条流了,显式起流把响应头发出去。
// 不起流的话,响应头要等到第一个 token(可能几十秒)或第一次心跳才提交,
// 前端迟迟不触发 onopen;而在这之前失败还能回普通 JSON。
if err := stream.Start(); err != nil {
    return err
}

// 起流之后启动心跳:模型思考、工具调用、上游暂时停顿期间下游一个字节
// 都没有,网关 / CDN / Ingress 的空闲超时会在模型还在工作时掐断连接。
//
// 心跳失败要连带取消上游——转发循环阻塞在读上游,没法同时 select hbErr,
// 靠 cancel 让 Decode 的读失败、循环自然退出。
hbCtx, stopHeartbeat := context.WithCancel(c.Request.Context())
hbErr := make(chan error, 1)
hbDone := make(chan struct{})
go func() {
    defer close(hbDone)
    if err := stream.Heartbeat(hbCtx, 15*time.Second); err != nil {
        hbErr <- err
        cancel() // 下游写不出去了,别再让上游继续生成
    }
}()
defer func() {
    stopHeartbeat()
    <-hbDone // 等心跳退出,见 9.1
}()

// completed 记录是否真的收到了结束哨兵。循环还可能因为上游正常 EOF、
// 代理提前结束响应、上游异常关闭但没返回读取错误而结束——那几种情况下
// 内容是被截断的,绝不能当成正常完成。
var completed bool
for msg, err := range ssex.Decode(resp.Body) {
    if err != nil {
        return finish(stream, hbErr, "incomplete", err)
    }
    if string(msg.Data) == "[DONE]" {
        completed = true

        break
    }
    if err := stream.Data(ssex.Raw(string(msg.Data))); err != nil {
        cancel() // 无论哪种失败都要停掉上游,别再为它烧 token

        // 按 9.3 的分级处理:断开是正常收尾,写超时与未知错误要能告警
        if errors.Is(err, ssex.ErrClientGone) || errors.Is(err, context.Canceled) {
            return nil
        }
        return err
    }
}

if !completed {
    // 让前端能区分"生成完成"与"被截断":终止原因不同,前端可以据此提示重试,
    // 而不是把半截回答当成最终答案。
    return finish(stream, hbErr, "incomplete", errors.New("上游未返回结束哨兵,输出可能被截断"))
}
closeStream(stream, gin.H{"reason": "done"}) // 终止帧失败要能发现,见 4.9

两个异常出口都走 finish,它同时解决两件事:报告哪个错误、还要不要写终止帧。

// finish 在转发异常结束时收尾。
//
// 先看心跳,有两个原因:
//
//  1. 心跳失败会 cancel 上游,因此上游返回的 context canceled 只是连带结果,
//     真正的原因是下游写不出去。不优先取心跳错误,ErrWriteTimeout 就会被上游
//     读错误遮蔽,生产上丢掉"这条连接写不动"的信号。
//  2. 心跳失败说明连接已经写不动了,此时再写终止帧只会白等一个 WithWriteTimeout
//     才失败,拖慢这条失败连接的释放,还会产生一次重复告警。所以直接跳过 closeStream。
func finish(stream *ssex.Stream, hbErr <-chan error, reason string, cause error) error {
    select {
    case hbFail := <-hbErr:
        if errors.Is(hbFail, ssex.ErrClientGone) || errors.Is(hbFail, context.Canceled) {
            return nil // 下游已断开,正常收尾
        }

        return hbFail
    default:
    }

    closeStream(stream, gin.H{"reason": reason})

    return cause
}

要点:

  • 上游请求用 http.NewRequestWithContext 绑定下游 r.Context(),并在下游写失败(ErrClientGone / ErrWriteTimeout)时立即 cancel()
  • 校验上游状态码,并用 mime.ParseMediaType 精确比较 Content-Type:上游出错时往往返回 JSON 而非事件流
  • 区分"收到结束哨兵"与"只是读到 EOF":后者意味着内容被截断(上游异常关闭、代理提前结束响应都会走到这里),必须用不同的终止原因告诉前端,否则半截回答会被当成完整答案
  • defer resp.Body.Close() 一定要有
  • token 流在 handler 内直接写出,不经 Hub:Hub 的队列会在满时丢弃事件,而 token 少一条就损坏文本
  • 校验通过后显式 Start():否则响应头要等第一个 token 或第一次心跳才提交
  • 长空闲期要有心跳:模型思考、工具调用期间下游没有字节,链路上任何一层的空闲超时都会掐断连接。只有当上游能保证最大空闲间隔小于链路最短 idle timeout 时才可以省掉心跳
  • 上游客户端按阶段设超时http.Client.Timeout 覆盖整个响应体读取,会掐断正常长流;改用 DialContext / TLSHandshakeTimeout / ResponseHeaderTimeout,最长生成时长用 context.WithTimeout 控制
  • 结束时用 Close 终止流,前端 es.close(),避免自动重连
  • 需要断线续传就用 EventWithID 带上 id,客户端重连时经 LastEventID(r) 读回起点,由业务决定从哪一条开始续推(见 6.7)

这段模板的可执行版本在 gincompat/relay_test.go:用 httptest 起假上游跑完整 handler,覆盖结束哨兵有/无、Content-Type 精确比较(含大小写变体与 text/event-streaming 这类前缀陷阱)、校验通过即起流、长空闲期心跳、下游断开取消上游。模板改坏了那里会挂,文档不会悄悄失真。

5.2 订单状态流

先发一条 status 快照,后续由支付回调经 Hub.Push 带外推送,定期注释心跳,终态后 Close

要点:

  • 快照先发:客户端可能在状态变更之后才连上来
  • 心跳必须有:等支付的连接可能几分钟没有数据
  • 终态后 Close,否则 EventSource 会一直重连
  • 状态先持久化再发布事件;delivered == 0 只说明本机此刻没有在等的连接,不能用它决定是否落库(见 4.14)

6. 推荐实践

6.1 业务层统一使用 Stream

推荐:

stream := ssex.NewStream(c.Writer, c.Request)
_ = stream.Event("status", payload)

不推荐每个模块都自己手写:

  • WriteHeaders
  • started 标志
  • 心跳格式
  • error 格式
6.2 事件名保持稳定

SSE 的事件名本质上就是客户端协议的一部分。

一旦对外开放,尽量不要频繁改:

  • status
  • ping
  • error
  • chunk
  • done
6.3 error 事件和网络错误分开处理

必须区分两种错误:

  1. 服务端发出的 event: error
  2. 浏览器/客户端自己的连接错误回调

这两者不是一个层级。

6.4 终态后主动结束流

如果业务天然有终态:

  • 订单 delivered/closed/...
  • LLM done

服务端在终态后调用 stream.Close(...),前端监听 close 事件并执行 es.close()。 只 return 而不发终止事件,EventSource 会自动重连,形成无限重连。

6.5 让 SSE 路由绕过压缩中间件

压缩中间件(如 gzip)会先把响应攒进自己的缓冲区,Flush 只刷到它的下一层,帧于是停在中间层不落地——表现是前端长时间收不到任何事件,直到缓冲被填满才一次性到达。

给 SSE 路由单独注册路由组、让它跳过压缩中间件即可。反向代理侧同理:本包已下发 X-Accel-Buffering: no(nginx 据此关闭该响应的缓冲),其他网关按各自方式关闭响应缓冲。

6.6 不要把双向交互硬塞进 SSE

SSE 只适合:

  • 服务端 -> 客户端

如果你的业务需要:

  • 客户端持续发消息
  • 双向实时交互
  • 房间/广播/多人会话

那应该考虑 WebSocket,而不是继续堆 SSE。

6.7 断线续传以存储为事实源

LastEventID(r) 只负责读回客户端的重连位点,重放数据从哪来由业务决定。可靠的做法:

  • 状态以数据库/缓存为唯一事实源,每次变更写入一个单调递增的 revision,并把它作为事件 id 下发(EventWithID
  • 客户端重连时带回 Last-Event-ID,服务端据此从存储里取 revision > LastEventID 的记录补发
  • 顺序上先订阅、后取快照:反过来(先快照再订阅)会漏掉两步之间发生的变更
  • 先订阅时,快照的 revision 可能比队列里已有的事件旧,按 revision 取大者即可,别让旧事件覆盖新快照
res, err := svc.Authorize(ctx, tenantID, uid, orderToken) // 先授权,拿到内部安全标识
if err != nil {
    return err
}

events, release := hub.Subscribe(res.ScopeKey())          // 再订阅
defer release()

snapshot, err := svc.Load(ctx, res)                       // 后取快照,沿用同一标识
if err != nil {
    return err
}
_ = stream.EventWithID(strconv.FormatInt(snapshot.Revision, 10), "status", snapshot)
6.8 浏览器接入:鉴权与 POST 流

原生 EventSource 的构造参数只有 URL 与 withCredentials,不能携带自定义请求头,也不能发 POST body。据此选接入方式:

场景 做法
登录态 / 订单状态(同源或可带 Cookie) new EventSource(url, { withCredentials: true }),鉴权走 Cookie;服务端从 Cookie 解出用户
Bearer Token 鉴权 fetch + ReadableStream 自行读流,把 token 放进 Authorization
AI 对话(请求体较大,需 POST) 同上用 fetchmethod: 'POST' 发 body,响应仍是 text/event-stream

fetch 时前端需要自己解析帧(按空行分帧、data: 行以 \n 拼接),并自行实现重连——EventSource 的自动重连与 Last-Event-ID 回传都不再由浏览器代劳。服务端侧本包的输出格式不变。

// Bearer Token / POST 流
const resp = await fetch('/api/chat', {
  method: 'POST',
  headers: { 'Authorization': `Bearer ${token}`, 'Content-Type': 'application/json' },
  body: JSON.stringify({ prompt }),
});
const reader = resp.body.pipeThrough(new TextDecoderStream()).getReader();
6.9 优雅停机要监听应用级信号

http.Server.Shutdown 会先关闭监听器、再关闭空闲连接,然后等活动连接变空闲。SSE handler 只要还在循环就一直是活动的——因此只调 Shutdown 而 handler 不监听停机信号,会一直等到 Shutdown 的 context 超时。

反过来,handler 一旦监听了应用级停机信号,Shutdown 自己就完成了等待:信号让 handler Close 并返回 → 连接变空闲 → Shutdown 返回。

// 应用启动时准备一个全局停机 context
shutdownCtx, shutdownDone := context.WithCancel(context.Background())

// 收到停机信号
shutdownDone() // 让所有 SSE handler 走 Close 分支主动收尾(见 9.2 模板的 appShutdown 分支)

// Shutdown 关掉监听器后不再有新请求,随后等已有连接变空闲——
// 正在收尾的 SSE handler 返回时连接就变空闲了。
if err := srv.Shutdown(ctx); err != nil {
    logger.Error("优雅停机未在期限内完成", zap.Error(err))
}

不要额外维护一个统计在线 handler 的 WaitGroupWait()Wait() 期间监听器还没关,新请求的 handler 仍会 Add(1),而 sync.WaitGroup 要求"计数器为零时开始的正数 Add 必须发生在 Wait 之前"——这既可能漏等新 handler,也是对 WaitGroup 的误用。把关闭监听器这件事交给 Shutdown 做,顺序才是安全的。

6.10 容量、限流与监控由应用提供

Hub 只维护注册表,不做全局连接数上限、单 key 连接上限或 IP 限流——这些策略属于应用与网关。库把做决策所需的原始信号交出来:

信号 来源
单 key 连接数(瞬时快照) Hub.Online(key)
被挤掉的旧事件数 Push / Broadcastdropped 返回值
客户端断开、写超时 ErrClientGone / ErrWriteTimeout
帧超限 ErrFrameTooLarge

生产上建议再自行记录:连接活跃时长与异常断开率、上游 AI 请求的取消率、每实例的文件描述符数量。

Online 不能用于严格限流。 它只是某个时刻单个 key 的快照,也不提供全局在线总数;if hub.Online(key) < limit { hub.Subscribe(key) } 是非原子的 check-then-act,并发请求照样会突破上限。它适合监控与近似判断,严格限流要用应用级 semaphore、原子计数器或网关。

7. 测试

当前公共层测试:

覆盖点包括:

  • 写头(含 HTTP/2 不下发 Connection)、起流即刷出响应头
  • 普通事件、带 id 事件、data-only 帧与 Raw 哨兵
  • 帧注入防护:字段值含换行 / NUL 被拒,注释与 raw data 中的孤立 \r 被归一
  • retry 值域校验
  • context cancel 判定为 ErrClientGone,序列化失败不误判
  • 写超时可配置
  • 自动起流、Close 终止后拒绝写入
  • ping 注释心跳、error、comment、retry、Heartbeat 间隔与退出条件
  • 上游解码:三种换行、只剥一个空格、多行 data、[DONE]、id/retry 合法与非法、末帧无空行、读取错误、超长帧
  • Hub:定向投递、同 key 多连接、注销幂等、广播、在线数、满队列丢弃计数、队列容量配置
  • 并发写入与 Hub 并发注册/推送(-race)、写入与解码热路径 benchmark
  • 交叉验证:写侧输出交由解码器解回(含两处帧注入交给解析器判决)、规范章节的官方示例向量、真实连接上客户端断开触发 ErrClientGone、Hub 与 Stream 真实连接端到端、真实 HTTP/2 连接
  • 可靠性:Hub latest-wins(最新状态送达、连续溢出只留最后一条、容量窗口)、帧构造失败时不提交响应头、起流错误可观察、整帧与单行大小上限、retry 溢出、非法 UTF-8 原样保留
  • Fuzz:FuzzDecodeFuzzRoundTrip(随机 CR/LF、非法 UTF-8、超大多行帧、超大 retry、提前停止迭代)
  • 起流写路径:刷新受 per-write deadline 约束、解除连接级截止时间失败时不提交响应头、nil Option 被跳过
  • gin 集成(gincompat 独立模块):鉴权在起流前拒绝、快照+推送+终态 Close、心跳 goroutine 在 handler 返回前收尾(-race 并发多轮)、长连接不被 WriteTimeout 截断、断开判定、Started() 分界、应用停机收尾、跨租户隔离与 scope key 无碰撞
  • AI 转发模板(5.1 的可执行版本,gincompat/relay_test.go):结束哨兵有/无(缺哨兵必须报 incomplete)、上游状态码与 Content-Type 精确校验、校验通过即起流、长空闲期心跳保活、下游断开取消上游

8. 依赖与版本

  • 唯一直接依赖:github.com/gtkit/json/v2go mod graph 里的其余条目均为它的传递依赖
  • 入口只用标准库 HTTP 接口,因此 net/http、gin、chi 共用同一套实现
  • 版本按 SemVer 走,破坏性变更在 CHANGELOG.md 以 ⚠ 标注并附迁移写法
8.1 默认构建实际链入什么

go.mod 里有十几条 // indirect,但默认构建一条都不会编进你的二进制

$ go list -deps github.com/gtkit/ssex | grep -E '^[^/]+\.[^/]+/'
github.com/gtkit/json/v2
github.com/gtkit/ssex

$ go list -f '{{join .Imports "\n"}}' github.com/gtkit/json/v2
bytes
encoding/json
io
unsafe

原因是 gtkit/json/v2 是 build tag 门控的 JSON 门面:默认那份实现走 encoding/json,sonic / json-iterator / goccy 三个引擎各自被 build tag 挡在门外,那些 // indirect 条目只为让 tag 可用而存在。

所以:

  • 默认构建的 JSON 行为等同于 encoding/json:默认那份实现把 API 设为委托 encoding/json 的后端
  • 需要更快的编解码时,下游用 -tags sonic(或 jsoniter / go_json)换引擎,不必改一行代码——SSE 高频转发 JSON 正是这类场景

代价是三项,都不为零:

  • 体积多约 5 KB。最小可执行程序实测(darwin/amd64、Go 1.26):直接用 encoding/json 为 2,956,128 字节,改用 gtkit/json/v2 为 2,961,072 字节,多 4,944 字节(+0.17%)。多出来的是门面代码与两个 init
  • 每次调用多一层接口分派Marshal 走的是 var API Core 这个接口变量,不是直接函数委托
  • 模块图会传播:那些 // indirect 条目经 MVS 进入下游的构建列表,可能顶高下游同名模块的版本

9. Gin 集成

库不依赖 gin,但 gin 是最常见的使用场景。这一节的三件事都是 gin 特有的坑。

9.1 c.Writer 的生命周期等于 handler 的生命周期

gin 的 c.WriterContext 的内部字段(c.Writer = &c.writermem)。handler 返回后 Context 被归还对象池,下一个请求会把这个 writer 重置到另一个响应上。

因此:在 handler goroutine 之外写入的 goroutine,必须在 handler 返回前退出。否则会写到别人的响应里。

c.Copy() 解决不了这个问题——它把副本的 writermem.ResponseWriter 置为 nil,副本只能用来读请求参数与头部,不能写响应。需要在后台 goroutine 里做的事,只有两种正确形态:

  • 要写响应:goroutine 必须在 handler 返回前收尾(如下面模板里的心跳)
  • 只读请求元数据:用 c.Copy(),或提前把需要的值取成局部变量
9.2 完整 handler 模板
func OrderEvents(c *gin.Context) {
    // 1. 认证:必须在 Start() 之前,此时还没提交响应头,可以正常返回 JSON 错误。
    uid, tenantID := c.GetString("uid"), c.GetString("tenant") // 由鉴权中间件写入
    if uid == "" || tenantID == "" {
        c.AbortWithStatusJSON(http.StatusUnauthorized, gin.H{"error": "未登录"})
        return
    }
    orderID := c.Param("id")

    // 2. 资源授权:与"读快照"分开做。授权是廉价的归属检查,先做掉可以避免
    //    产生未授权订阅(会让 Online 出现伪在线)。
    //    它返回一个内部安全标识,后续所有操作都用它——避免"授权带 tenantID、
    //    读快照只用 orderID"这种作用域漂移读到别的租户数据。
    res, err := svc.Authorize(c.Request.Context(), tenantID, uid, orderID)
    if err != nil || !res.Valid() {
        // 二次校验:授权实现若因 bug 返回零值标识与 nil error,key 会退化成一个
        // 固定值,多个异常请求就此落进同一队列、互相串流(见 4.15)。
        c.AbortWithStatusJSON(http.StatusForbidden, gin.H{"error": "无权访问"})
        return
    }

    // 3. key 由服务端从已授权资源计算。统一走一个构造函数,禁止业务代码各自拼接:
    //    裸拼接 tenantID + ":" + orderID 在任一段含分隔符时会碰撞
    //    ("a:b"+":"+"c" 与 "a"+":"+"b:c" 得到同一个 key),见 4.15。
    key := res.ScopeKey()

    // 4. 先订阅,再读快照。顺序反过来会永久漏事件:Load 与 Subscribe 之间发生的
    //    状态变更此时无人订阅,推送直接丢弃,客户端会永远停在旧快照上(见 6.7)。
    events, release := hub.Subscribe(key)
    defer release()

    // 5. 用同一个已验证标识读快照,不能退回裸 orderID。失败时还没起流,可以回 JSON。
    snapRev, snapshot, err := svc.Load(c.Request.Context(), res)
    if err != nil {
        c.AbortWithStatusJSON(http.StatusNotFound, gin.H{"error": "订单不存在"})
        return
    }

    stream := ssex.NewStream(c.Writer, c.Request, ssex.WithWriteTimeout(10*time.Second))

    // 6. 起流:失败时响应头尚未提交,仍可回 JSON。
    if err := stream.Start(); err != nil {
        handleStreamError(c, stream, err)
        return
    }

    // 7. 发快照,revision 放进事件 id。
    if err := stream.EventWithID(strconv.FormatInt(snapRev, 10), "status", snapshot); err != nil {
        handleStreamError(c, stream, err)
        return
    }

    // 8. 心跳:错误走 hbErr、"已退出"走 hbDone,并在 handler 返回前等它退出(见 4.12、9.1)。
    hbCtx, stopHeartbeat := context.WithCancel(c.Request.Context())
    hbErr := make(chan error, 1)
    hbDone := make(chan struct{})
    go func() {
        defer close(hbDone)
        if err := stream.Heartbeat(hbCtx, 15*time.Second); err != nil {
            hbErr <- err
        }
    }()
    defer func() {
        stopHeartbeat()
        <-hbDone
    }()

    // 9. 事件循环:四个退出口。lastRev 从快照起步并随成功发送推进——
    //    只比快照不够:Hub 不保证跨推送方的到达顺序,快照之后仍可能先到 rev 10、
    //    后到 rev 9,只比快照会把 9 也转发出去,前端状态回退。
    lastRev := snapRev
    for {
        select {
        case <-appShutdown.Done(): // 应用停机(见 6.9)
            closeStream(stream, gin.H{"reason": "server shutting down"})
            return

        case <-stream.Context().Done(): // 客户端断开
            return

        case err := <-hbErr: // 心跳写失败,连接已不可用
            handleStreamError(c, stream, err)
            return

        case e := <-events:
            rev, ok := revisionOf(e)
            if !ok {
                // 生产者没填 revision:告警而不是静默丢弃。若把解析失败当成 0,
                // 事件会因为 0 <= lastRev 被默默吃掉,排查时只看到"前端收不到状态"。
                // 只记可定位的元数据:Event.Data 可能含身份信息、token、
                // 订单详情或 AI 对话内容,不要整个序列化进日志。
                logger.Warn("事件缺少 revision",
                    zap.String("key", maskKey(key)),
                    zap.String("event_id", e.ID),
                    zap.String("event_name", e.Name),
                    zap.String("payload_type", fmt.Sprintf("%T", e.Data)),
                    zap.String("trace_id", traceID(c)))
                continue
            }
            if rev <= lastRev {
                continue // 积压的旧事件,或乱序到达的回退版本
            }
            if err := stream.Send(e); err != nil {
                handleStreamError(c, stream, err)
                return
            }
            lastRev = rev

            if isTerminal(e) { // 终态:主动结束,避免前端自动重连(见 6.4)
                closeStream(stream, gin.H{"reason": "final"})
                return
            }
        }
    }
}

这段模板的可执行版本在 gincompat/handler_test.go——模板改了那里会挂,文档不会悄悄失真。

9.3 错误处理规范

gin 场景最容易写错的是流已开始后又调用 c.JSON——响应头早已提交,再写普通响应体只会产生一个损坏的响应。用 Started() 分界:

func handleStreamError(c *gin.Context, stream *ssex.Stream, err error) {
    // 客户端断开与上下文取消是正常收尾,调试级记录即可
    if errors.Is(err, ssex.ErrClientGone) || errors.Is(err, context.Canceled) {
        logger.Debug("sse 客户端断开", zap.Error(err))
        c.Abort()
        return
    }

    // 写超时、序列化失败与未知写错误需要告警
    _ = c.Error(err)
    logger.Error("sse 写入失败", zap.Error(err))

    if !stream.Started() {
        // 响应头尚未提交,还能回普通 JSON
        c.AbortWithStatusJSON(http.StatusInternalServerError, gin.H{"error": "建立事件流失败"})
        return
    }

    // 已开始 SSE:只能记录并结束 handler
    c.Abort()
}

规则:

条件 允许的动作
stream.Started() == false 可以 c.JSON / AbortWithStatusJSON 返回普通 JSON
stream.Started() == true 只能记录错误 + c.Abort()禁止 c.JSON / c.String / AbortWithStatusJSON

错误分级:ErrClientGonecontext.Canceled 属正常收尾,调试级记录;ErrWriteTimeout、序列化错误与未知写错误需要告警。

返回什么状态码与响应体是业务策略,因此库只提供 Started() 这个分界信号,不封装响应内容。

日志不要序列化整个 Event.Data 或响应载荷。 它可能含身份信息、登录 token、订单详情、AI 对话内容、手机号或地址。只记可定位的元数据:key 的脱敏值、Event.IDEvent.Name、载荷类型、Trace ID、错误原因。

9.4 中间件兼容

SSE 路由要避开这几类中间件:

类型 后果
响应体缓存 / 统一包装 帧被攒在中间层,前端收不到
压缩(gzip 等) 同上,Flush 只刷到压缩层(见 6.5)
固定请求总时长的 timeout 长连接被按普通请求掐断
handler 返回后重写响应体 响应已提交,重写产生损坏响应
把所有错误转成 JSON 流已开始后再写 JSON 同样损坏响应

鉴权、Recovery、Tracing 可以正常使用,但鉴权必须在 Start() 之前完成——起流之后就没法再回 401 了。

替换 c.Writer 的自定义中间件必须实现 Unwrap() http.ResponseWriter 并正确透传 Flush 否则 http.ResponseController 沿不到底层连接,SetWriteDeadline 返回 http.ErrNotSupported,本库按约定静默降级——后果是逐帧写超时失效,且解除 http.Server.WriteTimeout 的能力一并失效,长连接仍会在全局超时到期时被掐断。这个降级不会报错,所以上线前要核对真实的中间件链,而不只是裸 gin。

type myWriter struct {
    gin.ResponseWriter
    // ...自己的字段
}

// 必须有:否则 SetWriteDeadline 与 WriteTimeout 清除都会降级
func (w *myWriter) Unwrap() http.ResponseWriter { return w.ResponseWriter }

自查办法:在真实中间件链下跑一个把 http.Server.WriteTimeout 设成几百毫秒的用例,持续写若干秒并断言最后一帧仍能收到——gincompat 里的 TestSurvivesServerWriteTimeout 就是这个形状,把它挪到你自己的路由与中间件上即可。

用独立路由组隔离:

sseGroup := router.Group("/events")
sseGroup.Use(AuthMiddleware(), gin.Recovery()) // 不挂压缩、缓存、timeout
sseGroup.GET("/orders/:id", OrderEvents)
9.5 兼容性测试

gin 的集成由仓库内的独立模块 gincompat 持续验证(自带 go.mod,因此主模块的依赖清单里没有 gin):

cd gincompat && go test -race ./...

覆盖:长连接不被 http.Server.WriteTimeout 截断、客户端断开判定为 ErrClientGone、Hub 端到端推送并由客户端 Decode 解回、心跳 goroutine 在 handler 返回前收尾(-race 下并发多轮)、心跳写错误不阻塞 handler、起流前后错误处理分界、与 Recovery / Logger 中间件共存。

状态推送的正确性契约也在这里固定:读快照期间的变更不丢、积压与乱序的旧 revision 被过滤、无效 revision 被上报而非静默丢弃、空身份与零值授权结果不进入 Hub、跨租户隔离(两个租户用同一 orderID,断言各自真实连接的完整事件序列)、终止帧失败被上报。其中隔离与终止帧两类断言做过反向验证——故意改坏实现后测试确实会失败。

Documentation

Overview

Package ssex 提供 Server-Sent Events 写入器:解决 SSE 长连接被 http.Server.WriteTimeout 杀死的问题,并为每帧写入设置 per-write deadline。

入口只用标准库的 http.ResponseWriter 与 *http.Request,net/http、gin、chi 等都可直接使用;gin 里传 c.Writer 与 c.Request 即可。

Example (LLMRelay)

Example_lLMRelay 演示把大模型的流式输出转发给前端: data-only 帧 + [DONE] 哨兵,并按错误类别决定是否取消上游请求。

package main

import (
	"errors"
	"fmt"
	"net/http"
	"net/http/httptest"

	"github.com/gtkit/ssex"
)

// newExampleStream 构造一个可打印输出的 Stream,仅供 Example 使用。
// 真实代码里传 handler 的 w 与 r(gin 里是 c.Writer 与 c.Request)。
func newExampleStream(opts ...ssex.Option) (*ssex.Stream, *httptest.ResponseRecorder) {
	recorder := httptest.NewRecorder()
	req := httptest.NewRequest(http.MethodGet, "/sse", nil)

	return ssex.NewStream(recorder, req, opts...), recorder
}

func main() {
	stream, recorder := newExampleStream()

	for _, delta := range []string{"Hel", "lo"} {
		if err := stream.Data(map[string]string{"delta": delta}); err != nil {
			if errors.Is(err, ssex.ErrClientGone) {
				return // 前端已断开:停止转发并取消上游大模型请求,避免继续计费
			}
			fmt.Println("write failed:", err)

			return
		}
	}
	_ = stream.Data(ssex.Raw("[DONE]"))

	fmt.Print(recorder.Body.String())
}
Output:
data: {"delta":"Hel"}

data: {"delta":"lo"}

data: [DONE]
Example (OrderStatus)

Example_orderStatus 演示订单状态推送:状态事件 + 终态后终止流。 前端须监听 close 事件并调用 EventSource.close(),否则流结束后会自动重连。

package main

import (
	"fmt"
	"net/http"
	"net/http/httptest"

	"github.com/gtkit/ssex"
)

// newExampleStream 构造一个可打印输出的 Stream,仅供 Example 使用。
// 真实代码里传 handler 的 w 与 r(gin 里是 c.Writer 与 c.Request)。
func newExampleStream(opts ...ssex.Option) (*ssex.Stream, *httptest.ResponseRecorder) {
	recorder := httptest.NewRecorder()
	req := httptest.NewRequest(http.MethodGet, "/sse", nil)

	return ssex.NewStream(recorder, req, opts...), recorder
}

func main() {
	stream, recorder := newExampleStream()

	_ = stream.Event("status", map[string]string{"status": "pending"})
	_ = stream.Event("status", map[string]string{"status": "delivered"})
	_ = stream.Close(map[string]string{"reason": "delivered"})

	fmt.Print(recorder.Body.String())
}
Output:
event: status
data: {"status":"pending"}

event: status
data: {"status":"delivered"}

event: close
data: {"reason":"delivered"}

Index

Examples

Constants

View Source
const Version = "v0.2.2"

Version 是本包的当前版本号,与 git 附注标签保持一致。

发版脚本(make release-patch / release-minor)会原地自增下面这行的版本号并据此 打标签,因此这一行的形状——Version 常量赋值为带 v 前缀的三段版本号——不能改动。

Variables

View Source
var ErrClientGone = errors.New("ssex: client gone")

ErrClientGone 表示客户端已断开,本次写入不可能送达;用 errors.Is 判定。

典型处理:静默结束本次流并取消上游请求(如正在进行的大模型调用),无需告警。 判定依据是请求上下文已被取消(net/http 在客户端断开时取消它)或底层连接已关闭; 断开发生在上下文取消之前的竞态窗口内,错误归类为普通写失败。 客户端读取过慢导致的写超时属于 ErrWriteTimeout,不属于本类。

错误链保留原因,因此 errors.Is(err, context.Canceled) 同样可判定。

View Source
var ErrFrameTooLarge = errors.New("ssex: frame too large")

ErrFrameTooLarge 表示解码时单帧的 data 超过上限,用 errors.Is 判定。

上限按产出的 Message.Data 大小计:最多 1048575 字节,单行与多行口径一致。 单行长度另有一个略高的硬上限(防止超长行撑爆缓冲),越过它同样归入本类。

View Source
var ErrInvalidArgument = errors.New("ssex: invalid argument")

ErrInvalidArgument 表示调用方传入了非法参数:字段值含换行或 NUL、retry 为负、 心跳间隔非正等。这类错误在帧构造阶段就返回,不写出任何字节。

若它发生在**首帧**(流尚未开始,Started() 为 false),响应头还没提交, 调用方可以改用普通 JSON 响应回错;流已开始之后出现的这类错误只表示该帧被拒绝。

View Source
var ErrStreamClosed = errors.New("ssex: stream closed")

ErrStreamClosed 表示流已由服务端显式终止(见 Stream.Close),不再接受写入;用 errors.Is 判定。

View Source
var ErrWriteTimeout = errors.New("ssex: write timeout")

ErrWriteTimeout 表示单帧写入超过了写截止时间(见 WithWriteTimeout), 通常意味着客户端读取过慢。它不等于客户端断开,因此不会被判定为 ErrClientGone。

Functions

func Decode

func Decode(r io.Reader) iter.Seq2[Message, error]

Decode 按 SSE 规范(WHATWG event stream)解码 r,逐帧产出 Message。

典型用途是把上游大模型的 text/event-stream 转发给前端:

for msg, err := range ssex.Decode(resp.Body) {
    if err != nil {
        return err
    }
    if string(msg.Data) == "[DONE]" {
        break
    }
    if err := stream.Data(ssex.Raw(string(msg.Data))); err != nil {
        return err
    }
}

行分隔符支持 CRLF / CR / LF;流首 BOM、注释行(以 : 开头)与未知字段被忽略; 字段值只剥掉冒号后的一个空格。按规范,从未出现 data 字段的帧不产出 (与浏览器 EventSource 一致),已出现但值为空串的帧照常产出。 读取出错时产出一次该错误并终止;读到流尾正常结束,不产出 io.EOF。 流尾残留的不完整帧(缺结尾空行)按规范丢弃,与浏览器一致——连接被中途掐断时, 残留字节通常是截断的半条 JSON,交出去只会让前端解析失败。

产出的 Message.Data 最多 1048575 字节(maxDataSize);超过即产出可判定为 ErrFrameTooLarge 的错误,不静默截断。这个口径对单行与多行一致——多行时按 各 data 行拼接后的长度算。单行长度另有一个略高的硬上限,防止超长行撑爆缓冲。

r 必须非 nil。nil 是调用方的编程错误,迭代时会 panic 而不是静默产出空流—— 静默降级会把"上游连接没建起来"伪装成"上游没有输出"。

与规范的一处差异:Data 原样保留上游字节,不执行 UTF-8 解码替换。 转发链路上改写字节会让服务端交给前端的内容与上游不一致,而浏览器接收端 本身会按 UTF-8 解码;需要严格校验的调用方可自行处理 Data。

Example

ExampleDecode 演示解析上游大模型的 SSE 流(转发链路的读侧)。

package main

import (
	"fmt"
	"strings"

	"github.com/gtkit/ssex"
)

func main() {
	upstream := strings.NewReader(
		"data: {\"delta\":\"Hel\"}\n\n" +
			"data: {\"delta\":\"lo\"}\n\n" +
			"data: [DONE]\n\n",
	)

	for msg, err := range ssex.Decode(upstream) {
		if err != nil {
			fmt.Println("decode failed:", err)

			return
		}
		if string(msg.Data) == "[DONE]" {
			break // 真实代码里此处 stream.Close(...) 收尾
		}
		fmt.Printf("%s\n", msg.Data)
	}
}
Output:
{"delta":"Hel"}
{"delta":"lo"}

func LastEventID

func LastEventID(r *http.Request) string

LastEventID 返回 EventSource 自动重连时携带的 `Last-Event-ID` 请求头 (即客户端最后收到的 EventWithID 的 id),无则返回空串。

r 必须非 nil。 服务端据此决定断线续推的起点,由业务选择从哪一条开始重放。

func Raw

func Raw(data string) any

Raw 将 data 标记为 Writer.Data 或 Stream.Data 的已编码 SSE data payload。

仅在确实需要绕过 JSON 序列化的 data-only 帧中使用,例如 OpenAI 风格哨兵: Data(ssex.Raw("[DONE]"))。裸换行仍会被拆成多条 data 行,不能注入 event 或 id 字段。

Types

type Event

type Event struct {
	// ID 非空时写出 `id:` 行,供客户端重连时通过 Last-Event-ID 回传。
	// 状态类推送建议填单调递增的 revision,消费端按 revision 取大者(见 Hub 的顺序契约)。
	ID string
	// Name 是事件名;为空时写出 data-only 帧,前端经 EventSource 的 onmessage 接收。
	Name string
	// Data 是事件载荷,自动 JSON 序列化;Raw(...) 原样透传。
	//
	// 所有权:交给 Hub.Push / Hub.Broadcast 之后即视为交出所有权。同一份载荷会被
	// 该 key 下的多个连接各自序列化,且发生在它们各自的 handler goroutine 里——
	// 投递后再修改其中的 map、切片或指针内容会构成数据竞争。需要复用结构体时,
	// 投递前拷贝一份。
	Data any
}

Event 是一条待推送给前端的 SSE 事件,用于 Hub 投递与 Stream.Send。

与解码方向的 Message 相对:Message 描述"从上游读到的一帧",Data 是原始字节; Event 描述"要写给前端的一帧",Data 由写入层序列化。

type Hub

type Hub struct {
	// contains filtered or unexported fields
}

Hub 按 key 管理在线 SSE 连接,支持从任意 goroutine 定向推送或广播。 典型用途是带外推送:支付回调推订单状态、后台踢下线推登录态—— 产生事件的 goroutine 并不持有连接。

进程边界:Hub 只管**当前进程**的连接。多实例部署时,连接在实例 A 而支付回调 打到实例 B,实例 B 的 Push 找不到那条连接、返回 delivered == 0,前端收不到更新; Online 也只反映本机而非集群。单实例可直接用;多实例必须让每个实例都收到状态 事件(Outbox / MQ / Redis Pub/Sub 广播式消费),再各自 Push 给本地 Hub, 或用一致性哈希把同一 key 的连接与事件路由到同一实例。存储始终是事实源, Hub 只是加速通道。AI 对话流在 handler 内直连上游、不经 Hub,不受此限制。

key 的安全作用域:**Hub 不校验 key**,它只把 key 当成一个不透明字符串—— 包括空串在内的任何值都能订阅成功。因此作用域治理完全由调用方负责:

  • key 必须由已认证身份与已授权资源计算,不能直接用 URL、query 或客户端传入的值;
  • 必须非空:空 key 会让所有认证异常的请求共享同一个队列、互相收到对方的事件, 调用方应在 Subscribe 之前拒掉空身份;
  • 多租户必须带 tenant 作用域,且不要用裸拼接——"a:b"+":"+"c" 与 "a"+":"+"b:c" 会得到同一个 key。用长度前缀编码或服务端生成的全局唯一 ID,并统一封装成 一个 key 构造函数,禁止业务代码各自拼接;
  • 读取资源时要沿用授权阶段的同一作用域标识,不要退回裸资源 ID,否则快照可能 读到另一个租户的数据;
  • Broadcast 只用于确实允许所有在线用户看到的内容。

README 第 9 节给出了完整写法,gincompat 里有对应的可执行版本与跨租户回归测试。

适用边界:Hub 面向**状态推送**——每条事件自带完整状态,队列满时挤掉旧的、 保留最新的(latest-wins)即可。AI token 流每一条都是文本的一部分, 少一条就损坏输出,应在 handler 内直接用 Stream.Data 写出,不经 Hub。

顺序契约:事件顺序在**单个连接内**按入队顺序保证。多个 goroutine 并发 Push / Broadcast 时,它们之间的相对顺序由调度决定,不同连接可能观察到不同顺序。 因此订单、登录这类状态事件应携带单调递增的 revision(放进 Event.ID 或载荷), 消费端按 revision 取大者、忽略更旧的值。要提供跨推送方的全局顺序,就必须把所有 投递串行化到单点,广播吞吐会退化为单 goroutine——而网络与客户端处理本身也不保证顺序。

容量与顺序的组合风险:latest-wins 按**到达顺序**保留最后一条,不比较业务 revision。若同一 key 有多个并发推送方,到达顺序可能与 revision 顺序相反—— 容量为 1 时,后到的旧版本会挤掉先到的新版本,消费端再也看不到那条新版本, 按 revision 过滤也救不回来(它已经不在队列里了)。

因此 WithQueueSize(1) 只在"同一 key 的推送已经串行化"时才安全。否则至少满足一条:

  • 同一 key 的事件经单点串行化后再 Push,使到达顺序与 revision 顺序一致;
  • 容量留大于 1,并由消费端维护 lastRev 按 revision 过滤;
  • 终态事件到达后回读一次事实源做校准,不完全相信内存队列。

容量与限流:Hub 不做全局连接数上限、单 key 上限或 IP 限流——这些策略属于应用 与网关。Online 与 Push 的返回值可以用于监控与近似判断,但**不能用来做严格限流**: Online 只是某个时刻单个 key 的快照,也不提供全局在线总数, "if Online(key) < limit { Subscribe(key) }" 是非原子的 check-then-act, 并发请求照样会突破上限。严格限流请用应用级 semaphore、原子计数器或网关。

并发安全。Hub 自身不写连接:Push 只把事件投进目标连接的有界队列, 由持有连接的 handler 取出后写出。反过来(Hub 直接写)在 handler 返回后 会踩到已失效的 ResponseWriter。

Hub 没有后台 goroutine,因此不需要关闭。

func NewHub

func NewHub(opts ...HubOption) *Hub

NewHub 创建一个连接注册表;nil Option 跳过。

func (*Hub) Broadcast

func (h *Hub) Broadcast(e Event) (delivered, dropped int)

Broadcast 把 e 投给所有在线连接,返回投递成功与被丢弃的连接数。

func (*Hub) Online

func (h *Hub) Online(key string) int

Online 返回 key 当前的在线连接数,可用于判断"还有没有人在等这个结果"。

func (*Hub) Push

func (h *Hub) Push(key string, e Event) (delivered, dropped int)

Push 把 e 投给 key 下所有在线连接,返回入队成功的连接数与被挤掉的旧事件数。 key 不在线时返回 (0, 0),不报错、不阻塞。

delivered 不是交付确认:它只表示事件进入了多少个**本机内存队列**,不代表 客户端收到、Flush 成功、浏览器处理完成、其他实例上的连接收到、连接没有随即 断开,也不代表这条事件没有被后续溢出挤掉。因此 delivered 与 dropped 只用于 监控,MUST NOT 用来判断业务是否成功、或决定要不要持久化。

正确顺序是:在事务里写入状态与单调递增 revision → 提交成功 → 再 Push。 反过来(先 Push、delivered == 0 才落库)会在事件丢失时同时丢掉状态。

目标队列已满时挤掉该连接队首最旧的一条,让 e 一定入队(latest-wins), 既不阻塞也不重试:状态推送里最新一条描述的就是当前状态,必须送达; 而让一个卡住的浏览器拖住支付回调这类关键路径是更坏的结果。 dropped 是被挤掉的旧事件数,交给调用方决定是否告警。

Example

ExampleHub_Push 演示带外推送:支付回调把订单状态推给正在等待的连接。

package main

import (
	"fmt"

	"github.com/gtkit/ssex"
)

func main() {
	hub := ssex.NewHub()

	// 连接侧:handler 注册自己,退出时注销
	events, release := hub.Subscribe("u1")
	defer release()

	// 推送侧:支付回调,与持有连接的 handler 不在同一个 goroutine
	delivered, dropped := hub.Push("u1", ssex.Event{Name: "status", Data: "paid"})
	fmt.Printf("delivered=%d dropped=%d online=%d\n", delivered, dropped, hub.Online("u1"))

	// 连接侧:取出后写给客户端(真实代码里是 stream.Send(e))
	e := <-events
	fmt.Printf("%s %v\n", e.Name, e.Data)
}
Output:
delivered=1 dropped=0 online=1
status paid

func (*Hub) Subscribe

func (h *Hub) Subscribe(key string) (events <-chan Event, release func())

Subscribe 以 key 注册当前连接,返回事件队列与注销函数。 同一 key 可注册多个连接(同一用户多标签页 / 多端),各自收到完整副本。

调用方必须 defer 调用 release(重复调用安全)。事件队列不会被 Hub 关闭: 关闭它会与并发的 Push 构成"向已关闭 channel 发送"的 panic 竞态。 用连接上下文结束消费循环:

// key 由服务端从已授权资源计算,统一走构造函数、不要裸拼接
events, release := hub.Subscribe(res.ScopeKey())
defer release()

// 订阅之后再读快照,否则两步之间的状态变更无人接收、会永久丢失;
// 用授权阶段返回的同一标识,不要退回裸资源 ID
snapRev, snapshot, err := svc.Load(r.Context(), res)
if err != nil {
    // 尚未起流,这里还能回普通 JSON
    return
}

stream := ssex.NewStream(w, r) // gin 里传 c.Writer, c.Request
if err := stream.Start(); err != nil {
    return // 必须先起流,否则响应头要等到第一条推送才发出,前端迟迟不触发 onopen
}
if err := stream.EventWithID(strconv.FormatInt(snapRev, 10), "status", snapshot); err != nil {
    return
}

lastRev := snapRev // 随成功发送推进,不能只与快照比
for {
    select {
    case <-stream.Context().Done():
        return
    case e := <-events:
        rev, ok := revisionOf(e)
        if !ok {
            // 生产者没填 revision:告警而不是静默丢弃。日志只记 key 脱敏值、
            // Event.ID / Name 与载荷类型,不要序列化整个 Data。
            continue
        }
        if rev <= lastRev { // 积压的旧事件,或乱序到达的回退版本
            continue
        }
        if err := stream.Send(e); err != nil {
            return
        }
        lastRev = rev
    }
}

lastRev 必须随成功发送推进:只与快照比较挡不住快照之后乱序到达的回退版本 (Hub 不保证跨并发推送方的到达顺序,先到 rev 10、后到 rev 9 时 9 也会被转发, 前端状态回退)。

完整模板(含心跳收尾、应用停机、错误分界、租户作用域 key)见 README 第 9 节。

type HubOption

type HubOption func(*hubOptions)

HubOption 配置 Hub。

func WithQueueSize

func WithQueueSize(n int) HubOption

WithQueueSize 设置每个连接事件队列的容量,默认 32;非正值被忽略。 队列满时 Push / Broadcast 挤掉该连接队首最旧的一条,不阻塞推送方。

type Message

type Message struct {
	// ID 是当前生效的 last event ID。按 SSE 规范它是连接级状态:某帧带了 id 之后,
	// 后续未带 id 的帧沿用该值(规范原文 "The buffer does not get reset"),
	// 只有上游显式发送空 id 才重置。转发时据此保持 id 连续,断线续传起点才不会错。
	ID string
	// Name 是 `event:` 字段值,无该字段时为空串(默认事件)。每帧独立,不跨帧延续。
	Name string
	// Data 是该帧所有 `data:` 行以 \n 拼接的结果,空白原样保留。
	Data []byte
	// Retry 是当前生效的建议重连间隔,同为连接级状态;上游从未发送过合法值时为 0。
	Retry time.Duration
}

Message 是从上游 SSE 流解码出的一帧。

与推送方向的 Event 相对:Data 是原始字节,不做 JSON 解析也不做 trim—— OpenAI 风格的 `[DONE]` 哨兵不是合法 JSON,而大模型的增量 token 常带前后空格, 任何 trim 都会损坏拼出的文本。

type Option

type Option func(*options)

Option 配置 Writer / Stream 的可选行为。

func WithWriteTimeout

func WithWriteTimeout(d time.Duration) Option

WithWriteTimeout 设置单帧写入的截止时长,默认 10s;非正值被忽略。

该截止时间只作用于单帧:客户端读取过慢导致某一帧写入超时时, 这一次写入返回错误(不会被判定为 ErrClientGone,慢客户端不等于断开), 长连接整体不受此值限制。弱网移动端可适当放大,内网可收紧以更快发现死连接。

Example

ExampleWithWriteTimeout 演示为弱网客户端放宽单帧写超时(默认 10s)。

package main

import (
	"fmt"
	"net/http"
	"net/http/httptest"
	"time"

	"github.com/gtkit/ssex"
)

// newExampleStream 构造一个可打印输出的 Stream,仅供 Example 使用。
// 真实代码里传 handler 的 w 与 r(gin 里是 c.Writer 与 c.Request)。
func newExampleStream(opts ...ssex.Option) (*ssex.Stream, *httptest.ResponseRecorder) {
	recorder := httptest.NewRecorder()
	req := httptest.NewRequest(http.MethodGet, "/sse", nil)

	return ssex.NewStream(recorder, req, opts...), recorder
}

func main() {
	stream, recorder := newExampleStream(ssex.WithWriteTimeout(30 * time.Second))

	_ = stream.Event("status", map[string]string{"status": "pending"})

	fmt.Print(recorder.Body.String())
}
Output:
event: status
data: {"status":"pending"}

type Stream

type Stream struct {
	// contains filtered or unexported fields
}

Stream 是面向业务的 SSE 写入器。相比底层 Writer,Stream 额外:

  1. 首个事件自动提交 SSE 响应头,且帧构造失败时不提交(调用方仍可回普通 JSON);
  2. 跟踪响应是否已开始;
  3. 提供统一的 ping / error / heartbeat 辅助方法;
  4. 支持显式终止流(Close),终止后拒绝后续写入。

并发安全:Stream 用互斥锁串行化所有写方法,可从不同 goroutine (如心跳 goroutine + 业务 goroutine)并发调用。

生命周期:底层 http.ResponseWriter 只在 handler 执行期间有效,框架可能池化 复用它——gin 的 c.Writer 就是 Context 的内部字段,handler 返回后归还对象池, 下个请求会把它重置到另一个响应上。因此在 handler goroutine 之外写入的 goroutine 必须在 handler 返回前退出,否则会写到别人的响应上。 心跳这类后台写入的正确收尾写法见 Heartbeat。

func NewStream

func NewStream(w http.ResponseWriter, r *http.Request, opts ...Option) *Stream

NewStream 创建一个 SSE Stream;可用 Option 调整写入行为。

gin 里这样调用:ssex.NewStream(c.Writer, c.Request)。

w 与 r 都必须非 nil,理由见 New。

func (*Stream) Close

func (s *Stream) Close(payload any) error

Close 终止本流:发送一条名为 close 的事件,并拒绝后续所有写入(返回 ErrStreamClosed)。

为什么需要它:EventSource 在服务端正常结束流后会按重连间隔自动重连, 订单进入终态、LLM 输出结束后直接 return 会让前端反复重连。约定前端监听 close 事件后调用 EventSource.close(),这轮推送才真正结束。

Close 只终结本流的写入许可,不关闭 HTTP 连接——连接在 handler 返回时结束。 即使终止事件写入失败(客户端已断开),流同样标记为已终止。重复调用返回 ErrStreamClosed。

func (*Stream) Comment

func (s *Stream) Comment(text string) error

Comment 发送一条注释帧。

func (*Stream) Context

func (s *Stream) Context() context.Context

Context 返回绑定到本 SSE 连接的请求上下文。

func (*Stream) Data

func (s *Stream) Data(payload any) error

Data 发送一条 data-only 帧(OpenAI 风格,见 Writer.Data); 响应尚未开始时自动先提交 SSE 响应头。

func (*Stream) Error

func (s *Stream) Error(payload any) error

Error 发送一条标准业务 error 事件。

func (*Stream) Event

func (s *Stream) Event(name string, payload any) error

Event 发送一条命名 SSE 事件;响应尚未开始时自动先提交 SSE 响应头。

func (*Stream) EventWithID

func (s *Stream) EventWithID(id, name string, payload any) error

EventWithID 发送一条带 `id:` 字段的命名 SSE 事件(断线续传,见 Writer.EventWithID); 响应尚未开始时自动先提交 SSE 响应头。

func (*Stream) Heartbeat

func (s *Stream) Heartbeat(ctx context.Context, interval time.Duration) error

Heartbeat 每隔 interval 发送一条保活注释帧,直到 ctx 取消或写入失败。

阻塞运行,启动时机与所在 goroutine 由调用方控制。它在 handler 之外的 goroutine 里写入,因此必须在 handler 返回前退出——底层 ResponseWriter 会被框架池化复用 (见 Stream 的生命周期说明)。用独立 ctx 加等待,而不是只发个停止信号:

hbCtx, stop := context.WithCancel(c.Request.Context())
hbErr := make(chan error, 1)   // 传错误
hbDone := make(chan struct{})  // 表示"已退出"
go func() {
    defer close(hbDone)
    if err := stream.Heartbeat(hbCtx, 15*time.Second); err != nil {
        hbErr <- err
    }
}()
defer func() {
    stop()
    <-hbDone // 等它真的退出,再让 handler 返回
}()

独立 ctx 的作用是:handler 因终态、写失败或应用停机而返回时(此时请求上下文 可能还没取消)也能停掉心跳。把 hbErr 接进主 select 就能同时感知心跳写失败。

错误信号与退出信号必须是两个 channel。用同一个兼任时,主循环的 case 消费掉 唯一那次发送后,defer 里的第二次接收就永远等不到发送者——handler 永久阻塞, 连按 LIFO 排在后面的注销逻辑(如 Hub 的 release)都不会执行。 hbDone 用 close,读多少次都立即返回。

长时间无数据的流(等支付结果、等登录态)必须有心跳,否则代理层会按空闲连接断开。 ctx 取消时返回 ctx 错误;客户端断开返回 ErrClientGone;流已终止返回 ErrStreamClosed。 interval 非正时立即返回 ErrInvalidArgument,不写出任何字节。 ctx 必须非 nil:传 nil 会在 select 时 panic,而不是变成一个永不退出的心跳。

func (*Stream) Ping

func (s *Stream) Ping(at time.Time) error

Ping 发送一条标准保活注释帧。

func (*Stream) Retry

func (s *Stream) Retry(milliseconds int) error

Retry 告知客户端建议的重连间隔(毫秒);负值返回 ErrInvalidArgument。

func (*Stream) Send

func (s *Stream) Send(e Event) error

Send 写出一条 Event,语义与逐字段调用一致:Name 为空时写 data-only 帧, ID 为空时省略 `id:` 行。供 Hub 的消费循环直接写出投递来的事件。

func (*Stream) Start

func (s *Stream) Start() error

Start 显式提交 SSE 响应头并立即刷给客户端。

纯推送型 handler(如从 Hub 消费事件)必须先调用它:否则在第一条事件到来前 连接上一个字节都没有,前端迟迟不触发 onopen,空闲连接还可能被代理层掐断。 返回错误表示响应头没能送达,客户端已断开时可判定为 ErrClientGone。

func (*Stream) Started

func (s *Stream) Started() bool

Started 返回 SSE 响应是否已开始写入。

首帧因参数非法或序列化失败而报错时它仍为 false——响应头尚未提交, 调用方可以改用普通 JSON 回错。

type Writer

type Writer struct {
	// contains filtered or unexported fields
}

Writer 是低层 SSE 写入器,负责设置 SSE 响应头与逐帧写入。

只依赖标准库 HTTP 接口,因此 net/http、gin、chi 等都可直接使用; gin 里传 c.Writer 与 c.Request 即可。

并发安全:Writer **非并发安全**——它直接写底层 http.ResponseWriter, 不做任何串行化。若需从多个 goroutine(如心跳 + 业务推送)写入同一连接, 请改用 Stream(它用互斥锁串行化所有写方法)。

func New

func New(w http.ResponseWriter, r *http.Request, opts ...Option) *Writer

New 创建一个 SSE Writer;可用 Option 调整写入行为。

gin 里这样调用:ssex.New(c.Writer, c.Request)。

w 与 r 都必须非 nil。传 nil 是调用方的编程错误,后续写入会 panic—— 库不为此做防御性降级:一个写不出任何字节的 Writer 比直接 panic 更难定位。

func (*Writer) Comment

func (w *Writer) Comment(text string) error

Comment 写入一条 SSE 注释帧。 注释帧不会触发前端的业务事件回调,常用于链路保活、调试标记或代理层防空闲断开。 多行文本按行拆成多条注释行,任何换行形式都无法逃出注释语义。

func (*Writer) Context

func (w *Writer) Context() context.Context

Context 返回绑定到本 SSE 连接的请求上下文。

func (*Writer) Data

func (w *Writer) Data(payload any) error

Data 写入一条 data-only 帧(仅 `data:` 行,无事件名),即 OpenAI 风格的 流式块格式;payload 自动 JSON 序列化,Raw(...) 原样透传—— 终止哨兵可写作 Data(ssex.Raw("[DONE]")),输出字面 `data: [DONE]`。 前端经 EventSource 的 onmessage(默认事件)接收。

func (*Writer) Event

func (w *Writer) Event(name string, payload any) error

Event 写入一条命名 SSE 事件,payload 自动 JSON 序列化;写入带 per-write deadline。 name 含 \r / \n / NUL 时返回 ErrInvalidArgument,不写出任何字节。

func (*Writer) EventWithID

func (w *Writer) EventWithID(id, name string, payload any) error

EventWithID 写入一条带 `id:` 字段的命名 SSE 事件,用于断线续传: EventSource 自动重连时会把最后收到的 id 放进 `Last-Event-ID` 头回传 (服务端用 LastEventID 读取,自行决定从哪续推)。 id 为空串时不输出 `id:` 行,行为等同 Event。 id / name 含 \r / \n / NUL 时返回 ErrInvalidArgument,不写出任何字节。

func (*Writer) Retry

func (w *Writer) Retry(milliseconds int) error

Retry 写入 SSE 的 retry 指令,提示客户端后续重连间隔(毫秒)。 这是 SSE 协议的一部分,浏览器/EventSource 客户端会把它作为建议重连时间使用。 milliseconds 为负时返回 ErrInvalidArgument 且不写出任何字节:SSE 规范只接受 ASCII 数字,客户端会静默忽略非法值,静默失败比显式报错更难排查。

func (*Writer) WriteHeaders

func (w *Writer) WriteHeaders() error

WriteHeaders 提交 SSE 响应头并立即刷给客户端,同时解除 http.Server.WriteTimeout 对本长连接的写截止时间。

返回错误表示响应头没能送达:客户端已断开时可判定为 ErrClientGone。 底层不支持刷新时返回 nil(静默降级)。 响应头一旦提交就不应重复调用,否则标准库会报 superfluous response.WriteHeader。

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL