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 ¶
- Constants
- Variables
- func Decode(r io.Reader) iter.Seq2[Message, error]
- func LastEventID(r *http.Request) string
- func Raw(data string) any
- type Event
- type Hub
- type HubOption
- type Message
- type Option
- type Stream
- func (s *Stream) Close(payload any) error
- func (s *Stream) Comment(text string) error
- func (s *Stream) Context() context.Context
- func (s *Stream) Data(payload any) error
- func (s *Stream) Error(payload any) error
- func (s *Stream) Event(name string, payload any) error
- func (s *Stream) EventWithID(id, name string, payload any) error
- func (s *Stream) Heartbeat(ctx context.Context, interval time.Duration) error
- func (s *Stream) Ping(at time.Time) error
- func (s *Stream) Retry(milliseconds int) error
- func (s *Stream) Send(e Event) error
- func (s *Stream) Start() error
- func (s *Stream) Started() bool
- type Writer
- func (w *Writer) Comment(text string) error
- func (w *Writer) Context() context.Context
- func (w *Writer) Data(payload any) error
- func (w *Writer) Event(name string, payload any) error
- func (w *Writer) EventWithID(id, name string, payload any) error
- func (w *Writer) Retry(milliseconds int) error
- func (w *Writer) WriteHeaders() error
Examples ¶
Constants ¶
const Version = "v0.2.2"
Version 是本包的当前版本号,与 git 附注标签保持一致。
发版脚本(make release-patch / release-minor)会原地自增下面这行的版本号并据此 打标签,因此这一行的形状——Version 常量赋值为带 v 前缀的三段版本号——不能改动。
Variables ¶
var ErrClientGone = errors.New("ssex: client gone")
ErrClientGone 表示客户端已断开,本次写入不可能送达;用 errors.Is 判定。
典型处理:静默结束本次流并取消上游请求(如正在进行的大模型调用),无需告警。 判定依据是请求上下文已被取消(net/http 在客户端断开时取消它)或底层连接已关闭; 断开发生在上下文取消之前的竞态窗口内,错误归类为普通写失败。 客户端读取过慢导致的写超时属于 ErrWriteTimeout,不属于本类。
错误链保留原因,因此 errors.Is(err, context.Canceled) 同样可判定。
var ErrFrameTooLarge = errors.New("ssex: frame too large")
ErrFrameTooLarge 表示解码时单帧的 data 超过上限,用 errors.Is 判定。
上限按产出的 Message.Data 大小计:最多 1048575 字节,单行与多行口径一致。 单行长度另有一个略高的硬上限(防止超长行撑爆缓冲),越过它同样归入本类。
var ErrInvalidArgument = errors.New("ssex: invalid argument")
ErrInvalidArgument 表示调用方传入了非法参数:字段值含换行或 NUL、retry 为负、 心跳间隔非正等。这类错误在帧构造阶段就返回,不写出任何字节。
若它发生在**首帧**(流尚未开始,Started() 为 false),响应头还没提交, 调用方可以改用普通 JSON 响应回错;流已开始之后出现的这类错误只表示该帧被拒绝。
var ErrStreamClosed = errors.New("ssex: stream closed")
ErrStreamClosed 表示流已由服务端显式终止(见 Stream.Close),不再接受写入;用 errors.Is 判定。
var ErrWriteTimeout = errors.New("ssex: write timeout")
ErrWriteTimeout 表示单帧写入超过了写截止时间(见 WithWriteTimeout), 通常意味着客户端读取过慢。它不等于客户端断开,因此不会被判定为 ErrClientGone。
Functions ¶
func Decode ¶
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 ¶
LastEventID 返回 EventSource 自动重连时携带的 `Last-Event-ID` 请求头 (即客户端最后收到的 EventWithID 的 id),无则返回空串。
r 必须非 nil。 服务端据此决定断线续推的起点,由业务选择从哪一条开始重放。
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 (*Hub) Push ¶
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 ¶
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 ¶
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 ¶
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 额外:
- 首个事件自动提交 SSE 响应头,且帧构造失败时不提交(调用方仍可回普通 JSON);
- 跟踪响应是否已开始;
- 提供统一的 ping / error / heartbeat 辅助方法;
- 支持显式终止流(Close),终止后拒绝后续写入。
并发安全:Stream 用互斥锁串行化所有写方法,可从不同 goroutine (如心跳 goroutine + 业务 goroutine)并发调用。
生命周期:底层 http.ResponseWriter 只在 handler 执行期间有效,框架可能池化 复用它——gin 的 c.Writer 就是 Context 的内部字段,handler 返回后归还对象池, 下个请求会把它重置到另一个响应上。因此在 handler goroutine 之外写入的 goroutine 必须在 handler 返回前退出,否则会写到别人的响应上。 心跳这类后台写入的正确收尾写法见 Heartbeat。
func NewStream ¶
NewStream 创建一个 SSE Stream;可用 Option 调整写入行为。
gin 里这样调用:ssex.NewStream(c.Writer, c.Request)。
w 与 r 都必须非 nil,理由见 New。
func (*Stream) Close ¶
Close 终止本流:发送一条名为 close 的事件,并拒绝后续所有写入(返回 ErrStreamClosed)。
为什么需要它:EventSource 在服务端正常结束流后会按重连间隔自动重连, 订单进入终态、LLM 输出结束后直接 return 会让前端反复重连。约定前端监听 close 事件后调用 EventSource.close(),这轮推送才真正结束。
Close 只终结本流的写入许可,不关闭 HTTP 连接——连接在 handler 返回时结束。 即使终止事件写入失败(客户端已断开),流同样标记为已终止。重复调用返回 ErrStreamClosed。
func (*Stream) EventWithID ¶
EventWithID 发送一条带 `id:` 字段的命名 SSE 事件(断线续传,见 Writer.EventWithID); 响应尚未开始时自动先提交 SSE 响应头。
func (*Stream) Heartbeat ¶
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) Send ¶
Send 写出一条 Event,语义与逐字段调用一致:Name 为空时写 data-only 帧, ID 为空时省略 `id:` 行。供 Hub 的消费循环直接写出投递来的事件。
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 ¶
New 创建一个 SSE Writer;可用 Option 调整写入行为。
gin 里这样调用:ssex.New(c.Writer, c.Request)。
w 与 r 都必须非 nil。传 nil 是调用方的编程错误,后续写入会 panic—— 库不为此做防御性降级:一个写不出任何字节的 Writer 比直接 panic 更难定位。
func (*Writer) Comment ¶
Comment 写入一条 SSE 注释帧。 注释帧不会触发前端的业务事件回调,常用于链路保活、调试标记或代理层防空闲断开。 多行文本按行拆成多条注释行,任何换行形式都无法逃出注释语义。
func (*Writer) Data ¶
Data 写入一条 data-only 帧(仅 `data:` 行,无事件名),即 OpenAI 风格的 流式块格式;payload 自动 JSON 序列化,Raw(...) 原样透传—— 终止哨兵可写作 Data(ssex.Raw("[DONE]")),输出字面 `data: [DONE]`。 前端经 EventSource 的 onmessage(默认事件)接收。
func (*Writer) Event ¶
Event 写入一条命名 SSE 事件,payload 自动 JSON 序列化;写入带 per-write deadline。 name 含 \r / \n / NUL 时返回 ErrInvalidArgument,不写出任何字节。
func (*Writer) EventWithID ¶
EventWithID 写入一条带 `id:` 字段的命名 SSE 事件,用于断线续传: EventSource 自动重连时会把最后收到的 id 放进 `Last-Event-ID` 头回传 (服务端用 LastEventID 读取,自行决定从哪续推)。 id 为空串时不输出 `id:` 行,行为等同 Event。 id / name 含 \r / \n / NUL 时返回 ErrInvalidArgument,不写出任何字节。
func (*Writer) Retry ¶
Retry 写入 SSE 的 retry 指令,提示客户端后续重连间隔(毫秒)。 这是 SSE 协议的一部分,浏览器/EventSource 客户端会把它作为建议重连时间使用。 milliseconds 为负时返回 ErrInvalidArgument 且不写出任何字节:SSE 规范只接受 ASCII 数字,客户端会静默忽略非法值,静默失败比显式报错更难排查。
func (*Writer) WriteHeaders ¶
WriteHeaders 提交 SSE 响应头并立即刷给客户端,同时解除 http.Server.WriteTimeout 对本长连接的写截止时间。
返回错误表示响应头没能送达:客户端已断开时可判定为 ErrClientGone。 底层不支持刷新时返回 nil(静默降级)。 响应头一旦提交就不应重复调用,否则标准库会报 superfluous response.WriteHeader。