# SSE **Repository Path**: vipkwds/sse ## Basic Information - **Project Name**: SSE - **Description**: 高性能、易集成的 Server-Sent Events 解决方案,支持纯 Go 服务端与客户端闭环 - **Primary Language**: Go - **License**: MIT - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-05-30 - **Last Updated**: 2026-05-30 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # SSE (Server-Sent Events) for Go 高性能、易集成的 Server-Sent Events 解决方案,支持纯 Go 服务端与客户端闭环。 [![Go Version](https://img.shields.io/badge/Go-1.18%2B-blue)](https://golang.org/) [![License](https://img.shields.io/badge/License-MIT-green)](LICENSE) [![Beta](https://img.shields.io/badge/beta版-red)](https://gitee.com/vipkwds/sse/tree/v1.0.0) ## 特性 - **开箱即用** - 一行代码启动 SSE 服务 - **纯 Go 闭环** - 服务端和客户端均为 Go 实现,无需 JavaScript - **消息历史** - 支持断线重连后自动补发消息 - **房间隔离** - 支持消息房间分组 - **可插拔存储** - 内存存储/Redis/自定义 Storage 接口 - **高性能** - 基于 Channel 的异步消息分发 - **自动重连** - 客户端自动重连机制 - **双向心跳保活** - Server→Client 心跳 + Client→Server ack 超时检测 - **低依赖** - 最小化外部依赖 ## 安装 ```bash go get gitee.com/vipkwds/sse ``` ## 快速开始 ### Server (服务端) ```go package main import ( "fmt" "net/http" "gitee.com/vipkwds/sse" ) func main() { // 创建 SSE Handler handler := sse.NewHandler(nil) // 路由挂载 http.Handle("/sse/", handler) http.Handle("/api/", handler) fmt.Println("SSE Server started on :8080") http.ListenAndServe(":8080", nil) } ``` ### Client (客户端) ```go package main import ( "context" "fmt" "gitee.com/vipkwds/sse" ) func main() { ctx := context.Background() // 创建 SSE 客户端 client := sse.NewSSEClient("http://localhost:8080/sse/user123", sse.WithSSEClientOnMessage(func(data []byte) { fmt.Printf("收到消息: %s\n", data) }), sse.WithSSEClientOnConnect(func() { fmt.Println("已连接服务器") }), ) // 连接 if err := client.Connect(ctx); err != nil { panic(err) } // 保持运行 select {} } ``` ## API 参考 ### Handler (推荐方式) ```go // 创建 Handler handler := sse.NewHandler(nil) // 自定义配置 handler := sse.NewHandler(nil, sse.WithStorage(sse.NewMemoryStorage(1024)), sse.WithConfig(&sse.HandlerConfig{ HistoryDuration: 5 * time.Minute, CleanupInterval: 1 * time.Minute, }), ) // 挂载到路由 http.Handle("/sse/", handler) // GET /sse/*clientId http.Handle("/api/", handler) // JSON API // 直接操作 handler.SendToClient("client1", "event", data) handler.Broadcast("notice", map[string]string{"msg": "hello"}) handler.JoinRoom("client1", "room1") handler.Shutdown() ``` ### Hub (底层接口) ```go hub := sse.NewHub(nil) // 注册客户端 client := hub.NewClient("client-id", metadata) hub.Register(client) // 发送消息 hub.SendToClient("client-id", "event", data) hub.SendToRoom("room-name", "event", data) hub.Broadcast("event", data, exclude...) // 房间管理 hub.JoinRoom("client-id", "room-name") hub.LeaveRoom("client-id", "room-name") // 统计 stats := hub.GetStats() fmt.Printf("Clients: %d, Rooms: %d\n", stats.TotalClients, stats.TotalRooms) // 优雅关闭 hub.Shutdown() ``` ### SSEClient (Go 客户端) ```go client := sse.NewSSEClient("http://localhost:8080/sse/user123", // 自定义请求头 sse.WithSSEClientHeader("Authorization", "Bearer token"), // 最大重试次数 (默认 10) sse.WithSSEClientMaxRetry(5), // 重连间隔 (默认 3秒) sse.WithSSEClientReconnect(5 * time.Second), // 开启客户端回复 ack(留空则自动从 SSE URL 推导 ack 地址) sse.WithClientAck(""), // 消息回调 sse.WithSSEClientOnMessage(func(event string, data []byte) { fmt.Printf("收到事件 %s: %s\n", event, data) }), // 连接成功回调 sse.WithSSEClientOnConnect(func() { ... }), // 断开连接回调 sse.WithSSEClientOnDisconnect(func() { ... }), // 错误回调 sse.WithSSEClientOnError(func(err error) { ... }), ) // 连接 client.Connect(ctx) // 发送消息 (通过 HTTP POST) client.SendMessage(ctx, "/api/send", map[string]any{ "clientId": "user123", "event": "message", "data": map[string]string{"text": "hello"}, }) // 获取最后的 Event ID (用于重连) lastID := client.LastEventID() // 关闭 client.Close() ``` ### Storage 接口 (自定义存储) ```go type Storage interface { // 保存消息 Save(msg *Message) // 遍历指定时间戳之后的消息 RangeSince(clientID string, since int64, fn func(msg *Message) bool) // 清理过期消息 Cleanup(maxAge time.Duration) } // 使用自定义存储 type RedisStorage struct { ... } func (s *RedisStorage) Save(msg *sse.Message) { ... } func (s *RedisStorage) RangeSince(clientID string, since int64, fn func(msg *sse.Message) bool) bool { ... } func (s *RedisStorage) Cleanup(maxAge time.Duration) { ... } handler := sse.NewHandler(nil, sse.WithStorage(&RedisStorage{...})) ``` ## JSON API (内置) Handler 挂载到 `/api/` 时自动提供以下接口: | 端点 | 方法 | 说明 | |------|------|------| | `/api/stats` | GET | 获取连接统计 | | `/api/send` | POST | 发送消息给客户端 | | `/api/broadcast` | POST | 广播消息 | | `/api/room/send` | POST | 发送消息到房间 | | `/api/room/join` | POST | 加入房间 | | `/api/room/leave` | POST | 离开房间 | | `/api/ack` | POST | 处理客户端心跳应答 | ### 示例 ```bash # 发送消息 curl -X POST http://localhost:8080/api/send \ -H "Content-Type: application/json" \ -d '{"clientId":"user123","event":"message","data":{"text":"Hello!"}}' # 广播 curl -X POST http://localhost:8080/api/broadcast \ -H "Content-Type: application/json" \ -d '{"event":"notice","data":{"msg":"System announcement"}}' # 加入房间 curl -X POST http://localhost:8080/api/room/join \ -H "Content-Type: application/json" \ -d '{"clientId":"user123","room":"room1"}' ``` ## 重连机制 ### 服务端 服务端自动维护消息历史(新消息补发): - 客户端断开时,服务端保留消息历史(默认 5 分钟) - 客户端重连时,自动发送 `Last-Event-ID` 请求头 - 服务端根据时间戳自动补发断线期间的消息 ### 客户端 客户端自动重连: - 连接失败时,自动根据 `reconnect` 间隔重试 - `maxRetry` 次重试后放弃 - 重连时自动携带上次的 `Last-Event-ID` ```go client := sse.NewSSEClient("http://localhost:8080/sse/user123", sse.WithSSEClientReconnect(3 * time.Second), sse.WithSSEClientMaxRetry(10), ) ``` ## 双向心跳保活 双向心跳保活需要 Server 和 Client 两端同时开启。 ### 服务端配置 ```go hub := sse.NewHub(&sse.Config{ HeartbeatInterval: 30, // Server→Client 心跳间隔(秒,0=关闭) ClientAckTimeout: 10, // Client ack 超时阈值(秒,0=关闭双向保活) ClientAckMaxTimes: 3, // 连续超时次数超限则强制断开 }) ``` ### 客户端配置 ```go client := sse.NewSSEClient("http://localhost:8080/sse/user123", // 开启 Client ack 回写(留空则自动从 SSE URL 推导) sse.WithClientAck("http://localhost:8080/api/ack"), ) ``` ### 交互流程 1. Server 每 `HeartbeatInterval` 秒向所有 Client 发送 `event: _heartbeat` 心跳 2. Client 收到心跳后解析 JSON 中的 `time` 字段,自动 POST ack 到 `/api/ack` 3. Server 每 5 秒检查一次 Client 最后 ack 时间 4. 若距上次 ack 超时 `ClientAckTimeout` 秒,计数 +1 5. 累计超 `ClientAckMaxTimes` 次,强制断开该 Client ### Storage 要求 双向心跳保活要求 Storage 实现以下三个方法(未实现则自动关闭): ```go type Storage interface { // ... 基础方法 ... SetClientAck(clientID string, ts int64) error GetClientAck(clientID string) (int64, error) DelClientAck(clientID string) error } ``` `MemoryStorage` 默认实现了这三个方法。自定义存储只需实现对应方法即可无缝接入。 ## 消息 ID 格式 消息 ID 格式: `{uuid}-{timestamp}` 示例: `550e8400-e29b-41d4-a716-446655440000-1609459200000` - **UUID** - 唯一标识 - **Timestamp** - 毫秒时间戳,可用于判断消息时效 ## SSE 协议格式 ``` id: event: data: ``` ## 完整示例 ### 双向通信示例 ```go package main import ( "context" "fmt" "log" "net/http" "os" "os/signal" "syscall" "time" "gitee.com/vipkwds/sse" ) func main() { ctx := context.Background() // 1. 启动 SSE Server handler := sse.NewHandler(nil) http.Handle("/sse/", handler) http.Handle("/api/", handler) go func() { fmt.Println("Server started on :8080") log.Fatal(http.ListenAndServe(":8080", nil)) }() // 2. 启动 SSE Client client := sse.NewSSEClient("http://localhost:8080/sse/user123", sse.WithSSEClientOnMessage(func(data []byte) { fmt.Printf("客户端收到: %s\n", data) }), sse.WithSSEClientOnConnect(func() { fmt.Println("客户端已连接") }), ) if err := client.Connect(ctx); err != nil { panic(err) } // 3. Server 广播消息 time.Sleep(100 * time.Millisecond) handler.Broadcast("notice", map[string]string{ "msg": "Hello from Server!", }) // 4. Client 发送消息 time.Sleep(100 * time.Millisecond) client.SendMessage(ctx, "/api/send", map[string]any{ "clientId": "user123", "event": "message", "data": map[string]string{"text": "Hello from Client!"}, }) // 5. 等待中断信号 sig := make(chan os.Signal, 1) signal.Notify(sig, syscall.SIGINT, syscall.SIGTERM) <-sig // 6. 清理 client.Close() handler.Shutdown() fmt.Println("已退出") } ``` ### 房间聊天示例 ```go // 加入房间 handler.JoinRoom("user1", "room-general") // 向房间发送消息 handler.SendToRoom("room-general", "message", map[string]string{ "user": "user1", "content": "Hello everyone!", }) // 离开房间 handler.LeaveRoom("user1", "room-general") ``` ## 配置项 ### HandlerConfig ```go type HandlerConfig struct { // 消息历史保留时间 (默认 5 分钟) HistoryDuration time.Duration // 历史清理间隔 (默认 1 分钟) CleanupInterval time.Duration // 连接超时 (默认 30 秒) ConnectTimeout time.Duration } ``` ### Config (Hub) ```go type Config struct { // 心跳间隔秒数 (默认 30,0=关闭) HeartbeatInterval int // Client ack 超时秒数 (默认 0=关闭双向保活) ClientAckTimeout int // 连续 ack 超时多少次后主动断开 (默认 3) ClientAckMaxTimes int // 消息缓冲区大小 (默认 256) MessageBufferSize int // 是否启用日志 (默认 true) EnableLog bool } ``` ## 项目结构 ``` sse/ ├── sse.go # Hub + Client (Server端核心) ├── history.go # Storage 接口 + MemoryStorage ├── handler.go # HTTP Handler (开箱即用) ├── client.go # SSE Client (纯 Go 客户端) └── *_test.go # 测试文件 ``` ## 依赖 - github.com/google/uuid (生成消息 ID) ## 许可证 MIT