gows

package
v1.4.47 Latest Latest
Warning

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

Go to latest
Published: Aug 27, 2026 License: MIT Imports: 23 Imported by: 0

README

gows

Go 语言 WebSocket 服务端封装库,提供连接管理、消息读写、心跳保活和全局分发功能。

✨ 特性

  • 🔒 Channel 驱动写入模型:非阻塞 WriteJSONCtx,慢客户端自动丢弃消息,指数退避 + 随机 jitter
  • 🛡️ 安全防护:默认拒绝跨域、全局升级速率限制、单 IP 连接数限制
  • 📡 分布式 Dispatcher:连接注册、广播、按 UID 发送、条件过滤、跨实例消息分发
  • 🔗 多种后端支持:单机模式(无中间件)或 RabbitMQ 分布式模式(Producer 缓存池)
  • 💓 可配置心跳:自定义间隔、Pong 超时、Ping 消息内容
  • 🔗 链路追踪:内置 OpenTelemetry,所有读写方法提供 Ctx 变体(WriteJSONCtx/ReadMessageCtx
  • 📊 连接统计:收发计数、队列状态、错误计数、远程地址、健康检测(IsAlive
  • 🪝 生命周期钩子:升级前鉴权(JWT)、升级后注入、关闭回调
  • ⏱️ 主动读超时:独立于心跳,防止 ReadMessage 永久阻塞
  • 📏 消息大小限制:独立的读写大小限制(readLimit/writeLimit),防止恶意大消息
  • 🗂️ Client 代理分发:客户端自带 SendToUIDCtx/BroadcastCtx 方法,无需依赖全局 DefaultDispatcher

📦 安装

go get github.com/18721889353/sunshine/pkg/gows
依赖安装
go get github.com/gorilla/websocket
go get github.com/gin-gonic/gin

🚀 快速开始

1️⃣ 基础 WebSocket 升级
package main

import (
    "net/http"
    "github.com/gin-gonic/gin"
    "github.com/18721889353/sunshine/pkg/gows"
)

func main() {
    r := gin.Default()

    r.GET("/ws", func(c *gin.Context) {
        // 升级 HTTP 为 WebSocket 连接
        client, err := gows.Upgrade(c,
            gows.WithCheckOrigin(func(r *http.Request) bool {
                return true // 生产环境请限制可信域名
            }),
        )
        if err != nil {
            c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
            return
        }
        defer client.Close()

        // 读写循环(使用 Ctx 变体进行链路追踪)
        for {
            data, err := client.ReadMessageCtx(c.Request.Context())
            if err != nil {
                break
            }
            _ = client.WriteJSONCtx(c.Request.Context(), gows.Message{
                Type: "reply",
                Msg:  string(data),
            })
        }
    })

    _ = r.Run(":8080")
}
2️⃣ 启用心跳保活
client, err := gows.Upgrade(c,
    gows.WithHeartbeat(), // 30s 自动心跳
)
3️⃣ 注册全局分发
client, err := gows.Upgrade(c,
    gows.WithDispatcher(gows.DefaultDispatcher),
)
// 关闭时自动从 Dispatcher 注销
defer client.Close()

// 通过 Client 代理发送消息,无需直接引用 Dispatcher
client.BroadcastCtx(ctx, gows.Message{Type: "notice", Msg: "系统公告"})
client.SendToUIDCtx(ctx, "user-123", gows.Message{Type: "private", Msg: "你好"})
4️⃣ 完整示例
package main

import (
    "log"
    "net/http"
    "github.com/gin-gonic/gin"
    "github.com/18721889353/sunshine/pkg/gows"
)

func main() {
    r := gin.Default()

    r.GET("/ws", func(c *gin.Context) {
        client, err := gows.Upgrade(c,
            gows.WithHeartbeat(),
            gows.WithDispatcher(gows.DefaultDispatcher),
        )
        if err != nil {
            log.Printf("upgrade failed: %v", err)
            return
        }
        defer client.Close()

        for {
            data, err := client.ReadMessageCtx(c.Request.Context())
            if err != nil {
                break
            }
            log.Printf("收到消息: %s, 来自: %s", string(data), client.RemoteAddr())
            _ = client.WriteJSONCtx(c.Request.Context(), gows.Message{
                Type: "ack",
                Msg:  "已收到: " + string(data),
            })
        }
    })

    // 在其他地方广播消息(通过 Client 代理方法)
    // client.BroadcastCtx(c.Request.Context(), gows.Message{
    //     Type: "notice",
    //     Msg:  "系统公告",
    // })

    _ = r.Run(":8080")
}

📡 Dispatcher 用法

创建与注册
// 使用默认全局实例
dispatcher := gows.DefaultDispatcher

// 或创建独立实例(单机模式)
dispatcher := gows.NewDispatcher(nil)

// 或创建分布式实例(需 RabbitMQ)
backend := gows.NewRabbitMQBackend("amqp://...", "ws:messages")
dispatcher := gows.NewDispatcher(backend)
dispatcher.Start(ctx)
发送消息
// 广播给所有连接
dispatcher.BroadcastCtx(ctx, gows.Message{Type: "notice", Msg: "全员通知"})

// 按条件过滤广播(仅本地)
dispatcher.BroadcastFilterCtx(ctx, msg, func(c *gows.Client) bool {
    return c.UID() == "vip-user"
})

// 按 UID 发送(跨实例)
dispatcher.SendToUIDCtx(ctx, "user-123", gows.Message{Type: "personal", Msg: "你好"})

// 按多 UID 发送(跨实例)
dispatcher.SendToMultiUIDCtx(ctx, []string{"uid1", "uid2"}, gows.Message{Type: "batch"})
Client 代理方法

当 Client 已通过 WithDispatcher 注册到分发中心后,可直接在 Client 上调用发送方法:

// 向单个用户发送
client.SendToUIDCtx(ctx, "target_uid", gows.Message{Type: "notify", Msg: "你有新消息"})

// 向多个用户发送
client.SendToMultiUIDCtx(ctx, []string{"uid1", "uid2"}, gows.Message{Type: "batch", Data: scores})

// 向所有在线用户广播
client.BroadcastCtx(ctx, gows.Message{Type: "announcement", Msg: "系统维护通知"})
遍历与管理
// 获取在线数量(本地 + 远端)
count := dispatcher.Len()

// 本地在线数
local := dispatcher.LenLocal()

// 遍历所有连接
dispatcher.Range(func(c *gows.Client) bool {
    log.Printf("在线: %s (%s)", c.UID(), c.RemoteAddr())
    return true // false 停止遍历
})

// 获取连接快照
clients := dispatcher.Clients()
for _, c := range clients {
    stats := c.Stats()
    log.Printf("UID=%s 已发=%d 已收=%d", stats.UID, stats.NumSent, stats.NumReceived)
}

// 统计信息
stats := dispatcher.Stats()
log.Printf("当前在线: %d (本地: %d, 远端: %d)",
    stats.TotalConnections, stats.LocalConnections, stats.RemoteUIDs)

// 清理已关闭的僵尸连接
cleaned := dispatcher.CleanupDeadConns()

// 获取全局唯一 UID 列表
uids := dispatcher.ConnectedUIDs()

// 最大连接限制
log.Printf("最大连接: %d, 已拒绝: %d", dispatcher.MaxConnections(), dispatcher.TotalRejected())

// 分布式模式:启动后台接收
// dispatcher.Start(ctx)
// defer dispatcher.Stop()

⚙️ 高级配置

JWT Token 鉴权
// 解析 JWT token(支持 Bearer 前缀),提取用户标识
uid, err := gows.ParseTokenCtx(ctx, tokenString)
if err != nil {
    // token 过期: errors.Is(err, jwt.ErrTokenExpired)
    // token 格式正确但缺 uid: errors.Is(err, gows.ErrTokenInvalid)
    return
}

在 Upgrade 前钩子中结合使用:

r.GET("/ws", func(c *gin.Context) {
    client, err := gows.Upgrade(c,
        gows.WithBeforeUpgrade(func(c *gin.Context) error {
            token := c.Query("token")
            if token == "" {
                return fmt.Errorf("missing token")
            }
            uid, err := gows.ParseTokenCtx(c.Request.Context(), token)
            if err != nil {
                return err
            }
            c.Set("uid", uid) // 后续可通过 WithClientUID 获取
            return nil
        }),
        gows.WithClientUIDFromHeader("X-UID"),
    )
})
分布式部署
// 使用 RabbitMQ 后端,跨实例消息分发
backend := gows.NewRabbitMQBackend(
    "amqp://user:pass@host:5672/",
    "ws:messages",
)
d := gows.NewDispatcher(backend, gows.WithMaxConnections(10000))
d.Start(ctx)
defer d.Stop()

// 注册连接后,即可跨实例发送
d.RegisterCtx(ctx, client)
d.SendToUIDCtx(ctx, "user-123", msg)
d.BroadcastCtx(ctx, msg)

// 查看在线统计
log.Printf("在线: %d, 最大限制: %d", d.Len(), d.MaxConnections())
自定义 Backend
// Backend 接口允许接入任何消息中间件
type Backend interface {
    Publish(ctx context.Context, msg *PubSubMessage) error
    Receive(ctx context.Context) (<-chan *PubSubMessage, error)
    Close() error
}

// 内置实现:RabbitMQ Fanout 交换机
Upgrader 选项
client, err := gows.Upgrade(c,
    // 自定义跨域检查
    gows.WithCheckOrigin(func(r *http.Request) bool {
        return r.Header.Get("Origin") == "https://example.com"
    }),
    // 读写缓冲区大小
    gows.WithBufferSize(8192, 8192),
    // 子协议协商
    gows.WithSubprotocols("v2", "v1"),
    // 启用压缩(默认开启)
    gows.WithEnableCompression(true),
    // 设置用户标识
    gows.WithClientUID("user-123"),
    // 心跳
    gows.WithHeartbeat(),
    // 心跳高级配置
    gows.WithHeartbeatOptions(
        gows.WithHeartbeatInterval(15*time.Second),
    ),
    // 全局分发
    gows.WithDispatcher(gows.DefaultDispatcher),
    // 启用分布式注册
    gows.WithEnableDistributed(true),
    // 读写队列容量(默认 64)
    gows.WithQueueSize(128, 128),
    // 升级前钩子(鉴权)
    gows.WithBeforeUpgrade(func(c *gin.Context) error {
        token := c.Query("token")
        if token == "" {
            return fmt.Errorf("missing token")
        }
        return nil
    }),
    // 升级后钩子(链路追踪)
    gows.WithAfterUpgrade(func(c *gin.Context, client *gows.Client) {
        log.Printf("客户端 %s 已连接", client.RemoteAddr())
    }),
    // 升级失败回调
    gows.WithErrorHandler(func(c *gin.Context, err error) {
        c.JSON(http.StatusUnauthorized, gin.H{"error": err.Error()})
    }),
)
Client 选项
// Client 通过 NewClient 构造函数创建,ctx 用于链路追踪和生命周期管理
client := gows.NewClient(ctx, rawConn, "uid",
    // 写入队列大小(默认 64)
    gows.WithWriteQueueSize(128),
    // 读取队列大小(默认 64)
    gows.WithReadQueueSize(128),
)

// 单条消息大小限制、读写超时等通过 Upgrade 选项自动传播(Upgrade → Client)
Dispatcher 选项
// 分布式模式:设置 WorkerPool 大小(默认 4)
d := gows.NewDispatcher(backend,
    gows.WithMaxConnections(10000),
    gows.WithWorkerPool(8), // 消息投递 worker 数
)
RabbitMQ Backend 选项
// 直接通过 URL 创建
backend := gows.NewRabbitMQBackend("amqp://...", "ws:messages")

// 或复用已有连接
conn, _ := gorabbitmq.NewConnection(ctx, "amqp://...")
backend := gows.NewRabbitMQBackendFromConn(conn, "ws:messages")
Client 统计信息
stats := client.Stats()
fmt.Printf(`客户端状态:
  UID:            %s
  远程地址:       %s
  是否存活:       %v
  已发送:         %d
  已接收:         %d
  写入队列容量:   %d
  写入队列长度:   %d
  读取队列容量:   %d
  读取队列长度:   %d
  是否已关闭:     %v
  写入错误次数:   %d
  读取错误次数:   %d
  上次写入错误:   %s
  上次读取错误:   %s
  上次写入时间:   %s
  上次读取时间:   %s
`, stats.UID, stats.RemoteAddr, stats.IsAlive, stats.NumSent, stats.NumReceived,
    stats.WriteQueueSize, stats.WriteQueueLen, stats.ReadQueueSize, stats.ReadQueueLen,
    stats.IsClosed, stats.WriteErrCount, stats.ReadErrCount, stats.LastWriteErr,
    stats.LastReadErr, stats.LastWriteTime, stats.LastReadTime)

// Dispatcher 统计
dStats := dispatcher.Stats()
fmt.Printf(`分发器状态:
  总连接数:       %d
  本地连接数:     %d
  远端 UID 数:    %d
  最大连接限制:   %d
  被拒绝连接数:   %d
`, dStats.TotalConnections, dStats.LocalConnections,
    dStats.RemoteUIDs, dStats.MaxConnections, dStats.TotalRejected)

📋 API 参考

核心类型
// 通用消息结构
type Message struct {
    Type string      `json:"type"`
    Msg  string      `json:"msg,omitempty"`
    Data any         `json:"data,omitempty"`
}

// 分布式消息结构(Backend 传输用)
type PubSubMessage struct {
    InstanceID string          `json:"instance_id"`
    Type       string          `json:"type"`       // "broadcast" | "send_to_uid"
    UIDs       []string        `json:"uids,omitempty"`
    Payload    json.RawMessage `json:"payload"`
}
Client 方法
方法 说明
WriteJSONCtx(ctx, v) error 异步非阻塞发送 JSON 消息(带链路追踪)
WriteRawCtx(ctx, data) error 直接写入预序列化的原始字节
ReadMessageCtx(ctx) ([]byte, error) 带链路追踪的消息读取(阻塞)
Close() error 优雅关闭(幂等)
UID() string 获取用户标识
RemoteAddr() string 获取远程地址
Context() context.Context 获取关联上下文(Close 后自动取消)
IsAlive() bool 检测连接是否健康存活
Done() <-chan struct{} 连接关闭信号
SetCloseHook(fn func()) 设置关闭钩子
Stats() ClientStats 获取统计信息
SendToUIDCtx(ctx, uid, v) 通过关联 Dispatcher 向指定 UID 发送(代理方法)
SendToMultiUIDCtx(ctx, uids, v) 通过关联 Dispatcher 向多 UID 发送(代理方法)
BroadcastCtx(ctx, v) 通过关联 Dispatcher 广播(代理方法)
ClientStats 字段
字段 类型 说明
UID string 用户标识
RemoteAddr string 远程地址
NumSent int64 已发送消息数
NumReceived int64 已接收消息数
WriteQueueSize int 写入队列容量
WriteQueueLen int 写入队列当前长度
ReadQueueSize int 读取队列容量
ReadQueueLen int 读取队列当前长度
IsClosed bool 是否已关闭
IsAlive bool 是否健康存活
WriteErrCount int64 写入失败累计次数
ReadErrCount int64 读取失败累计次数
LastWriteErr string 最近一次写入错误信息
LastReadErr string 最近一次读取错误信息
LastWriteTime string 最后一次成功写入时间(ISO8601)
LastReadTime string 最后一次成功读取时间(ISO8601)
Dispatcher 方法
方法 说明
RegisterCtx(ctx, client) error 注册连接(带链路追踪)
UnregisterCtx(ctx, client) error 注销连接(带链路追踪)
BroadcastCtx(ctx, v) 广播消息(跨实例)
BroadcastFilterCtx(ctx, v, filter) 条件广播(仅本地)
SendToUIDCtx(ctx, uid, v) 按 UID 发送(跨实例)
SendToMultiUIDCtx(ctx, uids, v) 按多 UID 发送(跨实例)
Range(fn) 遍历本地连接
Len() int 全局在线连接数(本地 + 远端)
LenLocal() int 本地在线连接数
Clients() []*Client 本地连接快照
ConnectedUIDs() []string 全局唯一 UID 列表(去重)
MaxConnections() int 最大连接数限制
TotalRejected() int 因达到上限被拒绝的累计连接数
CleanupDeadConns() int 清理已关闭的僵尸连接
Stats() DispatcherStats 统计信息
Start(ctx) 启动分布式后端(单机无需调用)
Stop() 停止分布式后端
DispatcherStats 字段
字段 类型 说明
TotalConnections int 当前在线连接总数(本地 + 远端)
LocalConnections int 本地在线连接数
RemoteUIDs int 远端实例上的唯一 UID 数
MaxConnections int 最大连接数限制
TotalRejected int 因达到上限被拒绝的累计连接数
Upgrade 选项
选项 说明
WithCheckOrigin(fn) 跨域检查(默认拒绝所有)
WithBufferSize(read, write) 读写缓冲区大小
WithSubprotocols(protocols...) 子协议协商
WithEnableCompression(enable) 压缩开关
WithHeartbeat() 启用心跳(30s)
WithHeartbeatOptions(opts...) 心跳高级配置
WithDispatcher(d) 注册到分发中心
WithEnableDistributed(enabled) 启用分布式分发注册
WithClientUID(uid) 设置用户标识
WithErrorHandler(fn) 升级失败回调
WithBeforeUpgrade(fn) 升级前钩子
WithAfterUpgrade(fn) 升级后钩子
WithRateLimit(rps, burst) 全局升级速率限制
WithMaxConnPerIP(n) 单 IP 最大连接数
WithReadLimit(n) 单条消息读取大小限制
WithWriteLimit(n) 单条消息写入大小限制
WithReadTimeout(d) 读取超时(Upgrade → Client 传播)
WithWriteTimeout(d) 写入超时(Upgrade → Client 传播)
WithQueueSize(writeSize, readSize) 读写队列容量
Heartbeat 选项
选项 说明
WithHeartbeatInterval(d) 心跳间隔
WithPongTimeout(d) Pong 等待超时(默认 10s)
WithPingWriteWait(d) Ping 控制帧写入超时(默认 5s)
工具函数
函数 说明
ParseTokenCtx(ctx, token) (string, error) 解析 JWT token 提取用户标识
NewClient(ctx, conn, uid, opts...) *Client 直接创建客户端(非 Upgrade 场景)
NewDispatcher(backend, opts...) *DistributedDispatcher 创建分发中心
NewRabbitMQBackend(url, exchange) *RabbitMQBackend 创建 RabbitMQ 分布式后端
NewRabbitMQBackendFromConn(conn, exchange) *RabbitMQBackend 复用已有连接创建后端
SetupPongHandler(conn, interval, timeout) 设置 Pong 处理器和 ReadDeadline
StartHeartbeat(client, opts...) 启动 Ping/Pong 心跳循环
错误标识
错误 说明
ErrWriteQueueFull 写入队列满,消息被丢弃
ErrWriteLimitExceeded 单条消息超出写入大小限制
ErrMaxConnections Dispatcher 达到最大连接数上限
ErrTokenInvalid JWT token 缺少 uid 字段

💡 最佳实践

1. 优雅关闭
client, err := gows.Upgrade(c, gows.WithDispatcher(dispatcher))
if err != nil {
    return
}
defer client.Close()

// client.Context() 会在 Close 时自动取消
client.Context().Done() // 可用于 select 监听
2. 并发写入安全

WriteJSONCtx 是线程安全的,可以在多个 goroutine 中同时调用:

// 业务协程
go func() {
    for {
        client.WriteJSONCtx(ctx, Message{Type: "heartbeat"})
        time.Sleep(30 * time.Second)
    }
}()

// 推送协程
go func() {
    for msg := range msgCh {
        client.WriteJSONCtx(ctx, msg)
    }
}()
3. Dispatcher 与业务整合
type ChatRoom struct {
    dispatcher *gows.DistributedDispatcher
}

func (r *ChatRoom) Join(ctx context.Context, client *gows.Client) {
    if err := r.dispatcher.RegisterCtx(ctx, client); err != nil {
        log.Printf("注册失败: %v", err)
        return
    }
    client.SetCloseHook(func() {
        r.dispatcher.UnregisterCtx(client.Context(), client)
        client.BroadcastCtx(context.Background(), gows.Message{
            Type: "system",
            Msg:  client.UID() + " 离开了房间",
        })
    })
}
4. 错误处理
err := client.WriteJSONCtx(ctx, msg)
if err == gows.ErrWriteQueueFull {
    // 队列满,消息被丢弃,可降级处理
    log.Warn("消息丢弃: 客户端消费过慢")
} else if err == websocket.ErrCloseSent {
    // 连接已关闭
    return
}
5. 使用 Client 代理发送
// 推荐: 通过 Client 代理方法发送,无需引用全局 Dispatcher
client.SendToUIDCtx(ctx, "target_uid", msg)
client.BroadcastCtx(ctx, msg)

// 仅当需要遍历或管理连接时才直接操作 Dispatcher
dispatcher.Range(func(c *gows.Client) bool {
    log.Printf("在线: %s", c.UID())
    return true
})

⚠️ 注意事项

  • 资源释放Upgrade 返回的 *Client 必须调用 Close(),建议使用 defer
  • 跨域:默认拒绝所有来源,必须使用 WithCheckOrigin 显式设置允许的来源
  • 限流:生产环境建议设置 WithRateLimitWithMaxConnPerIP 防止连接风暴
  • 链路追踪:优先使用 Ctx 变体(WriteJSONCtx/ReadMessageCtx),传递请求上下文以保留 request_id
  • 写队列满WriteJSONCtx 在队列满时返回 ErrWriteQueueFull 而非阻塞
  • 消息大小限制:生产环境建议设置 WithReadLimit()WithWriteLimit() 防止恶意大消息
  • ReadTimeout:未启用心跳时建议设置 WithReadTimeout() 防止 ReadMessage 永久阻塞
  • JWT 鉴权:使用前需确保 jwt.Init() 已被调用,通常在应用启动时初始化
  • Client 代理方法:仅在 Client 已通过 WithDispatcher 注册后生效,未注册时静默跳过

Documentation

Overview

Package gows 提供 WebSocket 服务端封装,包含连接管理、消息读写、心跳保活和全局分发。

Index

Constants

View Source
const (
	MsgTypeBroadcast     = "broadcast"
	MsgTypeSendToUID     = "send_to_uid"
	MsgTypeClientOnline  = "client_online"
	MsgTypeClientOffline = "client_offline"
)

PubSubMessage 的消息类型常量。

Variables

View Source
var DefaultDispatcher = NewDispatcher(nil)

DefaultDispatcher 默认全局 Dispatcher 单例(单机模式)。

View Source
var ErrClientNotFound = errors.New("dispatcher: client not found")

ErrClientNotFound 未找到指定 UID 的客户端连接

View Source
var ErrClientNotRegistered = errors.New("dispatcher: client not registered")

ErrClientNotRegistered 客户端未注册到 Dispatcher

View Source
var ErrEmptyUID = errors.New("dispatcher: empty uid")

ErrEmptyUID UID 为空时禁止注册

View Source
var ErrMaxConnections = errors.New("dispatcher: max connections reached")

ErrMaxConnections 达到最大连接数限制的错误

View Source
var ErrNilContext = errors.New("ctx must not be nil")

ErrNilContext 传入 nil context 的错误标识。 当 ReadMsgFromClientReadCh 接收到 nil context 时返回此错误,由调用方自行修复。

View Source
var ErrNoDispatcher = errors.New("client has no dispatcher configured")

ErrNoDispatcher Client 未关联 Dispatcher,无法执行跨实例发送/广播。

View Source
var ErrTokenInvalid = errors.New("token is invalid: missing uid")

ErrTokenInvalid token 无效(格式正确但缺少 uid 字段)

View Source
var ErrWriteLimitExceeded = errors.New("write message exceeds size limit")

ErrWriteLimitExceeded 单条消息大小超过写入限制的错误标识。 当消息超过 WithWriteLimit 设置的字节数时 WriteJSON/WriteRaw 返回此错误。

View Source
var ErrWriteQueueFull = errors.New("write queue is full, message dropped")

ErrWriteQueueFull 写入队列已满,消息被丢弃的错误标识。 当客户端写入缓冲区满载时 WriteJSON 返回此错误, 发送方可根据此错误判断是否为短暂拥塞,决定是否降级处理。

Functions

func ParseTokenCtx added in v1.4.36

func ParseTokenCtx(ctx context.Context, tokenString string) (string, error)

ParseTokenCtx 解析并验证 JWT token,返回用户标识(带 context 的版本)。

参数:

  • ctx: 上下文,用于链路追踪传播(其中应包含 request_id)
  • tokenString: JWT token 字符串,支持带 "Bearer " 前缀或不带前缀两种格式

返回:

  • string: 用户标识(UID),token 有效时返回
  • error: 解析失败时返回具体错误信息

注意:

  • 使用前需确保 jwt.Init() 已被调用,通常在应用启动时初始化
  • token 过期返回 jwt.ErrTokenExpired,调用方可用 errors.Is 判断
  • token 格式正确但缺少 uid 字段返回 ErrTokenInvalid
  • 自动去除 "Bearer " 前缀,兼容 Authorization header 传入的 token

func SetupPongHandler added in v1.4.36

func SetupPongHandler(conn *websocket.Conn, interval, pongTimeout time.Duration)

SetupPongHandler 在底层 WebSocket 连接上设置 Pong 处理器和读取超时。

核心机制:

  • 设置 SetPongHandler,每次收到 Pong 响应帧时刷新 ReadDeadline
  • ReadDeadline = interval + 3 × pongTimeout(保证大于心跳间隔,避免竞态)
  • 配合 StartHeartbeat 发送的 Ping 控制帧,构成协议级连接活性检测
  • 当网络断开时,ReadMessage 会在 ReadDeadline 过后返回 timeout 错误

参数:

  • conn: 底层 WebSocket 连接
  • interval: 心跳间隔,用于计算 ReadDeadline 的下限
  • pongTimeout: Pong 超时基准值

func StartHeartbeat

func StartHeartbeat(client *Client, opts ...HeartbeatOption)

StartHeartbeat 启动 WebSocket 协议级 Ping/Pong 心跳保活循环。

设计说明:

  • 通过 conn.WriteControl 发送 PingMessage 控制帧
  • 对端 WebSocket 协议栈自动回复 Pong 响应帧
  • PongHandler 每次收到 Pong 时刷新 ReadDeadline
  • 网络断开时 Pong 超时 → ReadDeadline 过期 → ReadMessage 返回超时错误
  • 心跳协程退出时主动关闭 TCP 连接,触发 read 方感知断开

参数:

  • client: 需要保活的客户端连接
  • opts: 可选心跳配置(间隔、Pong 超时等)

Types

type Backend added in v1.4.37

type Backend interface {
	// Publish 发布消息到分布式后端。
	// 根据消息类型决定路由:
	//   - broadcast: 广播到所有实例
	//   - send_to_uid: 按 UIDs 路由到对应持久化队列
	//   - client_online/client_offline: 广播到所有实例
	Publish(ctx context.Context, msg *PubSubMessage) error

	// SubscribeBroadcast 订阅广播消息,返回消息通道。
	// 所有 broadcast/client_online/client_offline 类型消息通过此通道接收。
	// ctx 取消时关闭通道并释放资源。
	SubscribeBroadcast(ctx context.Context) (<-chan *PubSubMessage, error)

	// Subscribe 订阅指定 UID 的消息队列。
	// 为每个在线用户创建独立消费者,绑定到该 UID 的持久化队列。
	// 用户断线后队列保留(消息不丢失),重连后继续消费。
	Subscribe(ctx context.Context, uid string) (<-chan *PubSubMessage, error)

	// Unsubscribe 取消订阅指定 UID 的消息队列。
	// 关闭消费者但保留队列及其中的消息,支持离线消息积压。
	Unsubscribe(ctx context.Context, uid string) error

	// Close 释放后端资源(如关闭 RabbitMQ 连接和所有消费者)。
	Close() error
}

Backend 分布式后端接口,抽象跨实例消息传递。 所有实现必须 goroutine 安全。

内置实现:

  • NewRabbitMQBackend(url, exchange) — 基于 RabbitMQ Fanout+Direct 交换机

可自行实现接入 Kafka / Redis Pub/Sub / NATS 等中间件。

type Client

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

Client 代表一个 WebSocket 客户端连接。 采用 Channel 驱动读写模型替代传统 Mutex 锁,核心设计原则:

  • writeCh 缓冲通道解耦业务协程和网络写入协程,发送方永不阻塞
  • msgFromChToWs 后台 goroutine 串行化消费通道数据,规避锁竞争
  • readCh 缓冲通道解耦底层连接和业务读取协程
  • msgFromWsToCh 后台 goroutine 持续从底层连接读取消息并推入 readCh
  • 队列满载时自动丢弃消息,防止慢客户端拖慢整体吞吐
  • 原子 CAS 保证关闭幂等,关闭信号通知所有监听协程优雅退出
  • 网络波动保护:写入失败时指数退避重试(含 jitter),Ping/Pong 协议级保活
  • 健康指标:记录读写时间与错误计数,支持 IsAlive 探测(定义在 health.go)
  • 读写超时:msgFromWsToCh 和 msgFromChToWs 均支持独立超时,不依赖心跳机制
  • dispatcher 字段用于代理 SendToUIDCtx / BroadcastCtx 等分发方法

clientConfig 以内嵌方式提供 dispatcher/readTimeout/writeTimeout/readLimit/writeLimit 配置, 消除构造函数中逐字段手动拷贝的样板代码。writeChSize/readChSize 作为内嵌字段保留(由 clientConfig 提供), 仅在 make(chan) 初始化时使用,运行时读取队列容量通过 cap(c.writeCh) 获取。

func Upgrade

func Upgrade(c *gin.Context, opts ...UpgradeOption) (*Client, error)

Upgrade 将 HTTP 请求升级为 WebSocket 长连接,并返回封装后的 Client。

参数:

  • c: Gin 上下文,提供 ResponseWriter 和 Request
  • opts: 可选配置参数

返回:

  • *Client: 封装后的 WebSocket 客户端,具备非阻塞写入能力
  • error: 升级失败时返回错误

使用示例:

client, err := gows.Upgrade(c,
    gows.WithHeartbeat(),
    gows.WithDispatcher(gows.DefaultDispatcher),
)

Upgrade 负责:

  1. 执行升级前钩子(可选)
  2. 创建 websocket.Upgrader 并应用配置
  3. 执行 HTTP→WebSocket 升级
  4. 执行升级后钩子(可选)
  5. 用 *websocket.Conn 创建 *Client(含 msgFromChToWs)
  6. 可选启动心跳、注册 Dispatcher

安全保护:

  • CORS: 默认拒绝所有来源,需显式调用 WithCheckOrigin
  • 限流: WithWsRateLimit 设置全局升级速率
  • 单IP限制: WithMaxConnPerIP 设置单IP最大连接数

调用方需负责 defer client.Close() 确保资源释放。

func (*Client) BroadcastCtx added in v1.4.37

func (c *Client) BroadcastCtx(ctx context.Context, v any) error

BroadcastCtx 通过关联的 Dispatcher 向所有在线客户端广播消息。 仅在 Client 已注册到 Dispatcher 时有效。

func (*Client) BroadcastReliableCtx added in v1.4.37

func (c *Client) BroadcastReliableCtx(ctx context.Context, v any) error

BroadcastReliableCtx 通过关联的 Dispatcher 进行可靠广播(按 UID 下发)。 仅在 Client 已注册到 Dispatcher 时有效。

func (*Client) Close

func (c *Client) Close() error

Close 优雅关闭 WebSocket 连接,触发关闭信号并释放所有资源。 关闭顺序: 执行关闭钩子 → 取消上下文 → 等待 msgFromChToWs 完全退出 (drain writeCh 剩余消息后退出,确保所有排队消息已写入) → 关闭底层连接(触发 msgFromWsToCh 的 ReadMessage 返回错误) → 等待 msgFromWsToCh 完全退出。 通过 atomic.CompareAndSwap 保证幂等性,首次调用执行完整关闭流程, 后续调用直接返回 nil。 返回:

  • error: 首次关闭底层连接失败时返回 error,重复关闭返回 nil

func (*Client) Context

func (c *Client) Context() context.Context

Context 返回客户端的关联上下文。 Close 后该上下文被取消,可通过 ctx.Done() 感知连接关闭事件。

func (*Client) DisconnectByUID added in v1.4.39

func (c *Client) DisconnectByUID(ctx context.Context, uid string) error

DisconnectByUID 通过关联的 Dispatcher 断开指定 UID 的连接。 仅在 Client 已注册到 Dispatcher 时有效。

func (*Client) DisconnectByUIDs added in v1.4.39

func (c *Client) DisconnectByUIDs(ctx context.Context, uids ...string) int

DisconnectByUIDs 通过关联的 Dispatcher 断开多个 UID 的连接。 仅在 Client 已注册到 Dispatcher 时有效。

func (*Client) Done

func (c *Client) Done() <-chan struct{}

Done 返回一个只读 channel,当连接关闭时该 channel 被 close。 调用方可通过 select 或 <-client.Done() 感知连接断开事件, 常用于心跳协程和消息读取协程的退出通知。 返回:

  • <-chan struct{}: 关闭信号接收通道,关闭时立即返回零值

func (*Client) IsAlive added in v1.4.36

func (c *Client) IsAlive() bool

IsAlive 判断客户端连接是否处于健康状态。 返回 false 的场景:

  • Close() 已调用
  • msgFromChToWs 已因写入重试全部失败退出(closeWsConn 已触发)
  • 超过 3 个心跳周期无成功写入(疑似僵尸连接)

注意:

  • 此方法返回 true 不代表底层网络一定可达,仅表示组件内部状态正常
  • 精确的活性检测依赖 Ping/Pong 协议级心跳 + ReadDeadline 联动

func (*Client) ReadMsgFromClientReadCh added in v1.4.37

func (c *Client) ReadMsgFromClientReadCh(ctx context.Context) ([]byte, error)

ReadMsgFromClientReadCh 同步阻塞读取客户端发送的一条消息。

func (*Client) RemoteAddr

func (c *Client) RemoteAddr() string

RemoteAddr 返回客户端的远程网络地址(IP:Port)。

func (*Client) SendToMultiUIDCtx added in v1.4.37

func (c *Client) SendToMultiUIDCtx(ctx context.Context, uids []string, v any) error

SendToMultiUIDCtx 通过关联的 Dispatcher 向多个 UID 的用户发送消息。 发送者自身不会收到消息(自动过滤)。 仅在 Client 已注册到 Dispatcher 时有效。

func (*Client) SendToUIDCtx added in v1.4.37

func (c *Client) SendToUIDCtx(ctx context.Context, uid string, v any) error

SendToUIDCtx 通过关联的 Dispatcher 向指定 UID 的用户发送消息。 发送者自身不会收到消息(自动过滤)。 仅在 Client 已注册到 Dispatcher 时有效(通过 Upgrade 传入 WithDispatcher)。

func (*Client) SetCloseHook

func (c *Client) SetCloseHook(fn func())

SetCloseHook 设置关闭回调钩子,在 Close 时自动调用。 可用于执行 Dispatcher 注销等清理操作。 参数:

  • fn: 回调函数,在底层连接关闭前执行

func (*Client) Stats

func (c *Client) Stats() ClientStats

Stats 返回客户端连接的实时统计信息。 各字段均为 goroutine 安全读取。

func (*Client) UID

func (c *Client) UID() string

UID 返回当前连接关联的用户标识。

func (*Client) WriteJSONToClientWriteCh added in v1.4.37

func (c *Client) WriteJSONToClientWriteCh(ctx context.Context, v any) error

WriteJSONToClientWriteCh 向客户端写入 JSON 消息,队列满时阻塞等待。 阻塞期间只响应 clientCtx 关闭,不受上游 ctx 取消影响。

func (*Client) WriteRawToClientWriteCh added in v1.4.37

func (c *Client) WriteRawToClientWriteCh(ctx context.Context, data []byte) error

WriteRawToClientWriteCh 直接写入预序列化的原始字节数据,跳过 JSON 序列化步骤。 队列满时阻塞等待,阻塞期间只响应 clientCtx 关闭,不受上游 ctx 取消影响。

type ClientStats

type ClientStats struct {
	UID            string // 用户标识
	RemoteAddr     string // 远程地址
	NumSent        int64  // 已发送消息数
	NumReceived    int64  // 已接收消息数
	WriteQueueSize int    // 写入队列容量
	WriteQueueLen  int    // 写入队列当前长度
	ReadQueueSize  int    // 读取队列容量
	ReadQueueLen   int    // 读取队列当前长度
	IsClosed       bool   // 是否已关闭
	IsAlive        bool   // 是否健康(基于 IsAlive() 判断)
	WriteErrCount  int64  // 写入失败累计次数
	ReadErrCount   int64  // 读取失败累计次数
	LastWriteErr   string // 最近一次写入错误信息
	LastReadErr    string // 最近一次读取错误信息
	LastWriteTime  string // 最后一次成功写入时间(ISO8601)
	LastReadTime   string // 最后一次成功读取时间(ISO8601)
}

ClientStats 客户端连接统计信息。

type DispatcherOption added in v1.4.36

type DispatcherOption func(*dispatcherOptions)

DispatcherOption 分发中心配置选项函数类型

func WithMaxConnections added in v1.4.36

func WithMaxConnections(n int) DispatcherOption

WithMaxConnections 设置最大连接数限制。 参数:

  • n: 最大连接数。达到上限时 Register 返回 ErrMaxConnections。 0 表示不限制(默认)。

func WithWorkerPool added in v1.4.37

func WithWorkerPool(n int) DispatcherOption

WithWorkerPool 设置后台 Worker 协程数,用于消息投递的背压保护。 参数:

  • n: Worker 协程数(默认 4,建议 2~16)。Worker 池满时消息降级为串行处理。 0 或 1 表示关闭 WorkerPool,receiveLoop 串行投递。

type DispatcherStats

type DispatcherStats struct {
	TotalConnections int // 当前在线连接总数(本地 + 远端)
	LocalConnections int // 本地在线连接数
	RemoteUIDs       int // 远端实例上的唯一 UID 数
	MaxConnections   int // 最大连接数限制(0=不限制)
}

DispatcherStats 分发中心统计信息。

type DistributedDispatcher added in v1.4.37

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

DistributedDispatcher 消息分发中心。

负责管理所有在线 WebSocket 客户端连接,支持跨实例消息转发。 所有公开方法均为 goroutine 安全,适用于高并发业务场景。

分布式架构

实例A (alice)          Backend (RabbitMQ)         实例B (bob)
  │                                                     │
  ├─ SendToUIDCtx("bob")                                │
  │   └─ Publish ─────────────────────────────────────►├─ receiveLoop
  │                                                      └─ client.WriteJSON ✓

工作原理

  • Send 系列方法将消息序列化为 JSON → 发布到 Backend
  • receiveLoop 协程接收所有实例的消息 → 通过 WorkerPool 并发投递到本地连接
  • 发送方也通过 Backend 接收自己发布的消息(统一处理路径)
  • 通过 InstanceID 跳过自发布消息,避免回环干扰

背压保护

  • receiveLoop 使用 WorkerPool 异步处理消息,避免单条慢消息阻塞后续投递
  • WorkerPool 满时阻塞发送,天然反压到消息源

网络波动保护

  • Broadcast 系列方法自动清理已关闭的僵尸连接
  • 连接管理使用 sync.Map,无锁化高并发安全

func NewDispatcher

func NewDispatcher(backend Backend, opts ...DispatcherOption) *DistributedDispatcher

NewDispatcher 创建并初始化一个新的消息分发中心。 参数:

  • backend: 分布式后端,传 nil 表示单机模式(仅供本地消息投递)
  • opts: 可选配置(如 WithMaxConnections、WithWorkerPool)

示例(单机):

d := gows.NewDispatcher(nil)
d.RegisterCtx(ctx, client)

示例(分布式):

backend := gows.NewRabbitMQBackend("amqp://...", "ws:messages")
d := gows.NewDispatcher(backend, gows.WithMaxConnections(10000))
d.Start(ctx)

func (*DistributedDispatcher) BroadcastCtx added in v1.4.37

func (dd *DistributedDispatcher) BroadcastCtx(ctx context.Context, v any) error

BroadcastCtx 向所有在线客户端广播消息(跨实例)。 广播使用临时队列(exclusive+auto-delete),不做持久化,重启后离线消息不保留。

func (*DistributedDispatcher) BroadcastFilterCtx added in v1.4.37

func (dd *DistributedDispatcher) BroadcastFilterCtx(ctx context.Context, v any, filter func(*Client) bool) error

BroadcastFilterCtx 向满足 filter 条件的本地客户端广播消息(仅本地)。 使用 deliverPool 并发投递,与 deliverBroadcast / WriteRawToLocalUIDs 保持一致的并发模式。

func (*DistributedDispatcher) BroadcastReliableCtx added in v1.4.37

func (dd *DistributedDispatcher) BroadcastReliableCtx(ctx context.Context, v any) error

BroadcastReliableCtx 可靠广播:转为按所有在线 UID 逐个下发。 每条消息走 UID 持久化队列,自动支持离线积压和重连补推。 适用于需要确保所有用户(含离线用户)最终能收到的广播场景。

func (*DistributedDispatcher) CleanupDeadConns added in v1.4.37

func (dd *DistributedDispatcher) CleanupDeadConns() int

CleanupDeadConns 清理所有已关闭的僵尸连接,释放 Dispatcher 内存。 同时递减 localUIDCounts 引用计数,确保 hasLocalUID 查询不漂移。 返回清理的连接数。

func (*DistributedDispatcher) Clients added in v1.4.37

func (dd *DistributedDispatcher) Clients() []*Client

Clients 返回当前所有已注册客户端的快照切片。

func (*DistributedDispatcher) ConnectedUIDs added in v1.4.37

func (dd *DistributedDispatcher) ConnectedUIDs() []string

ConnectedUIDs 返回所有实例上的在线 UID 列表(去重)。

func (*DistributedDispatcher) DisconnectByUID added in v1.4.39

func (dd *DistributedDispatcher) DisconnectByUID(ctx context.Context, uid string) error

DisconnectByUID 断开指定 UID 的客户端连接。 如果该 UID 有多个连接(同账号多设备),全部断开。 返回:

  • error: 未找到该 UID 时返回 ErrClientNotFound

func (*DistributedDispatcher) DisconnectByUIDs added in v1.4.39

func (dd *DistributedDispatcher) DisconnectByUIDs(ctx context.Context, uids ...string) int

DisconnectByUIDs 断开多个指定 UID 的客户端连接。 参数:

  • uids: 要断开的 UID 列表

返回:

  • int: 实际断开的连接数

func (*DistributedDispatcher) Len added in v1.4.37

func (dd *DistributedDispatcher) Len() int

Len 返回全局在线连接数(本地 + 远端)。

func (*DistributedDispatcher) LenLocal added in v1.4.37

func (dd *DistributedDispatcher) LenLocal() int

LenLocal 返回本地在线连接数。

func (*DistributedDispatcher) MaxConnections added in v1.4.37

func (dd *DistributedDispatcher) MaxConnections() int

MaxConnections 返回最大连接数限制(0=不限制)。

func (*DistributedDispatcher) Range added in v1.4.37

func (dd *DistributedDispatcher) Range(fn func(*Client) bool)

Range 遍历所有已注册客户端。 如果 fn 返回 false 则停止遍历。 注意:fn 中调用 Delete 等修改操作是安全的,但可能导致遍历结果不一致。 如需一致性快照,使用 Clients() 替代。

func (*DistributedDispatcher) RegisterCtx added in v1.4.37

func (dd *DistributedDispatcher) RegisterCtx(ctx context.Context, client *Client) error

RegisterCtx 将客户端连接注册到 Dispatcher 全局列表(带自定义上下文)。

注册流程依次执行:最大连接数检查 → 空 UID 拒绝 → SSO 踢旧连接 → 入表计数 → 事件广播 → 消费队列。分布式模式下自动向所有实例广播上线通知。

参数:

  • ctx:用于跨实例链路追踪的上下文,tracing 信息随事件消息传播
  • client:已建立的 WebSocket 客户端连接,需包含 uid、remoteAddr 等有效信息

返回值:

  • nil:注册成功,client 已加入全局列表并开始接收广播/点对点消息
  • ErrMaxConnections:达到最大连接数上限,调用方应执行 client.Close() 释放资源
  • ErrEmptyUID:客户端 UID 为空,禁止注册

func (*DistributedDispatcher) SendToMultiUIDCtx added in v1.4.37

func (dd *DistributedDispatcher) SendToMultiUIDCtx(ctx context.Context, uids []string, v any) error

SendToMultiUIDCtx 向多个指定 UID 的客户端发送消息(跨实例)。

func (*DistributedDispatcher) SendToUIDCtx added in v1.4.37

func (dd *DistributedDispatcher) SendToUIDCtx(ctx context.Context, uid string, v any) error

SendToUIDCtx 向指定 UID 的客户端发送消息(跨实例)。 所有实例通过 Backend 接收后在本地查找并投递。 自动清理已关闭的僵尸连接。

func (*DistributedDispatcher) Start added in v1.4.37

func (dd *DistributedDispatcher) Start(ctx context.Context)

Start 启动后台接收协程和 WorkerPool,开始监听其他实例的消息。 仅分布式模式下调用(backend != nil),单机模式无需调用。 ctx 用于链路追踪。 Start 是幂等的,多次调用安全。

func (*DistributedDispatcher) Stats added in v1.4.37

Stats 返回分发中心的实时统计信息。

func (*DistributedDispatcher) Stop added in v1.4.37

func (dd *DistributedDispatcher) Stop()

Stop 停止后台接收协程,释放后端资源。 关闭顺序: 关闭 UID 订阅 → 关闭 Backend → 释放 WorkerPool → 关闭 receiveLoop。

func (*DistributedDispatcher) UnregisterCtx added in v1.4.37

func (dd *DistributedDispatcher) UnregisterCtx(ctx context.Context, client *Client) error

UnregisterCtx 从 Dispatcher 全局列表中移除客户端连接(带自定义上下文)。

注销流程依次执行:从 clients 删除 → 事件广播 → 递减 UID 计数 → 清理 SSO → 停止消费。分布式模式下自动向所有实例广播下线通知。

参数:

  • ctx:携带 tracing 信息的上下文
  • client:要移除的客户端连接

返回值:

  • nil:注销成功

func (*DistributedDispatcher) WriteRawToLocalUIDs added in v1.4.37

func (dd *DistributedDispatcher) WriteRawToLocalUIDs(span trace.Span, uids []string, payload json.RawMessage) error

WriteRawToLocalUIDs 向本地指定 UID 的客户端投递预序列化的原始消息。

type HeartbeatOption

type HeartbeatOption func(*heartbeatOptions)

HeartbeatOption 心跳配置选项函数类型。 采用 Functional Options 模式,支持可扩展的配置传递。

func WithHeartbeatInterval

func WithHeartbeatInterval(interval time.Duration) HeartbeatOption

WithHeartbeatInterval 设置心跳发送间隔。 参数:

  • interval: 间隔时间,建议 10~60 秒。<=0 时使用默认值 30 秒。

func WithPingWriteWait added in v1.4.36

func WithPingWriteWait(timeout time.Duration) HeartbeatOption

WithPingWriteWait 设置 Ping 控制帧的写入超时时间。 参数:

  • timeout: 超过此时间 Ping 帧写入失败,认为连接异常。默认 5 秒。

func WithPongTimeout added in v1.4.36

func WithPongTimeout(timeout time.Duration) HeartbeatOption

WithPongTimeout 设置 Pong 超时时间。 参数:

  • timeout: 超过此时间未收到 Pong 响应帧,触发连接关闭。 建议为心跳间隔的 1/3 ~ 1/2,默认 10 秒。

type Message

type Message struct {
	Type string `json:"type"`           // 消息类型标识(如 "ping", "notify", "greeting")
	Msg  string `json:"msg,omitempty"`  // 消息内容文本(可选)
	Data any    `json:"data,omitempty"` // 附加业务数据(可选,任意类型)
}

Message WebSocket 服务端通用消息结构。 所有服务端推送消息统一使用此结构序列化, 客户端通过 Type 字段区分消息类型进行分发处理。

type PubSubMessage added in v1.4.37

type PubSubMessage struct {
	InstanceID string          `json:"instance_id"`    // 发送方实例 ID
	Type       string          `json:"type"`           // MsgTypeBroadcast | MsgTypeSendToUID | MsgTypeClientOnline | MsgTypeClientOffline
	UIDs       []string        `json:"uids,omitempty"` // send_to_uid 的目标 UID 列表
	Payload    json.RawMessage `json:"payload"`        // JSON 序列化的业务消息
}

PubSubMessage 跨实例分发的消息结构。 InstanceID 用于防止消息回环(自己发的消息自己不再重复处理)。

type RabbitMQBackend added in v1.4.37

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

RabbitMQBackend 基于 RabbitMQ 的分布式后端实现。

使用两个交换机分离广播和单播流量:

  • Fanout 交换机 (exchange+":broadcast"): 广播消息,所有实例订阅
  • Direct 交换机 (exchange): UID 消息,按 routing_key=uid 投递

工作原理:

  • 广播消息: PublishFanout → Fanout 交换机 → 每个实例的广播队列
  • 按 UID 消息: PublishDirect → Direct 交换机 → routing_key=uid
  • 用户断线后队列保留,重连后继续消费积压消息

适用场景:

  • 已有 RabbitMQ 基础设施的团队
  • 需要离线消息补推(消息不丢失)
  • WebSocket 分布式架构

func NewRabbitMQBackend added in v1.4.37

func NewRabbitMQBackend(url string, exchange string, connOpts ...gorabbitmq.ConnectionOption) *RabbitMQBackend

NewRabbitMQBackend 创建基于 RabbitMQ 的分布式后端。 自动创建两个交换机:

  • exchange (Direct): 用于 UID 消息路由
  • exchange+":broadcast" (Fanout): 用于广播消息

func NewRabbitMQBackendFromConn added in v1.4.37

func NewRabbitMQBackendFromConn(conn *gorabbitmq.Connection, exchange string) *RabbitMQBackend

NewRabbitMQBackendFromConn 使用已有连接创建后端。 注意:此方式无法创建 ProducerPool(无 URL),所有路径使用 per-publish Producer。

func (*RabbitMQBackend) Close added in v1.4.37

func (b *RabbitMQBackend) Close() error

Close 关闭所有消费者、ProducerPool 和 AMQP 连接。

func (*RabbitMQBackend) Publish added in v1.4.37

func (b *RabbitMQBackend) Publish(ctx context.Context, msg *PubSubMessage) error

Publish 发布消息到 RabbitMQ。 根据消息类型路由:

  • broadcast/client_online/client_offline: 使用 ProducerPool + PublishFanout
  • send_to_uid: 使用 PublishDirect,routing_key=uid

func (*RabbitMQBackend) SetPoolCap added in v1.4.44

func (b *RabbitMQBackend) SetPoolCap(initialCap, maxCap int)

SetPoolCap 设置 RabbitMQ 连接池容量参数。 应在首次发布消息前调用,否则将使用默认值。

参数:

  • initialCap: 初始连接数,<=0 则使用默认值。
  • maxCap: 最大连接数,<=0 则使用默认值。

func (*RabbitMQBackend) SetReceiveChanSize added in v1.4.44

func (b *RabbitMQBackend) SetReceiveChanSize(size int)

SetReceiveChanSize 设置接收消息通道缓冲区大小。 应在首次订阅广播前调用,否则将使用默认值。

func (*RabbitMQBackend) Subscribe added in v1.4.37

func (b *RabbitMQBackend) Subscribe(ctx context.Context, uid string) (<-chan *PubSubMessage, error)

Subscribe 订阅指定 UID 的持久化消息队列。 使用 Direct 交换机,routing_key=uid。

func (*RabbitMQBackend) SubscribeBroadcast added in v1.4.37

func (b *RabbitMQBackend) SubscribeBroadcast(ctx context.Context) (<-chan *PubSubMessage, error)

SubscribeBroadcast 订阅广播消息,返回消息通道。 创建独立队列绑定到 Fanout 交换机,接收所有广播消息。

func (*RabbitMQBackend) Unsubscribe added in v1.4.37

func (b *RabbitMQBackend) Unsubscribe(_ context.Context, uid string) error

Unsubscribe 取消订阅指定 UID 的消息队列。 关闭消费者但保留队列及消息,支持离线消息积压。

type UpgradeOption

type UpgradeOption func(*upgradeOptions)

UpgradeOption 升级配置选项函数类型。 采用 Functional Options 模式,支持可扩展的配置传递, 新增配置项无需修改 Upgrade 的函数签名。

func WithBufferSize

func WithBufferSize(read, write int) UpgradeOption

WithBufferSize 设置 WebSocket 读写缓冲区大小(字节)。 参数:

  • read: 读取缓冲区大小,<=0 时使用默认值 4096
  • write: 写入缓冲区大小,<=0 时使用默认值 4096

对于传输大消息的业务(如文件流、大 JSON),建议增大到 16384 以上。

func WithCheckOrigin

func WithCheckOrigin(fn func(r *http.Request) bool) UpgradeOption

WithCheckOrigin 设置 WebSocket 跨域检查函数。 参数:

  • fn: 接收 *http.Request,返回 true 表示允许该来源连接

默认拒绝所有来源。生产环境必须根据实际域名配置。

func WithClientUID

func WithClientUID(uid string) UpgradeOption

WithClientUID 设置客户端用户标识,传递给 Client 用于业务追踪。 参数:

  • uid: 用户唯一标识(如从 JWT 解析的 sub)

func WithDispatcher

func WithDispatcher(d *DistributedDispatcher) UpgradeOption

WithDispatcher 设置升级后自动将客户端注册到指定 Dispatcher。 参数:

  • d: 分发中心实例,通常传入 ws.DefaultDispatcher

注册后可通过 Dispatcher 全局管理在线连接(如广播消息)。 连接关闭时自动从 Dispatcher 注销。

func WithEnableCompression

func WithEnableCompression(enable bool) UpgradeOption

WithEnableCompression 设置是否启用 WebSocket 压缩。 参数:

  • enable: true 启用(默认),false 禁用。

压缩可减少带宽消耗,但会增加 CPU 开销。 内网环境或传输已压缩数据(如视频流)时建议关闭。

func WithEnableDistributed added in v1.4.37

func WithEnableDistributed(enabled bool) UpgradeOption

WithEnableDistributed 设置是否启用分布式分发注册。 参数:

  • enabled: true 启用分布式模式,连接自动注册到 Dispatcher

需同时调用 WithDispatcher 设置目标 Dispatcher 实例。

func WithHeartbeat

func WithHeartbeat() UpgradeOption

WithHeartbeat 启用心跳保活机制,升级后自动在后台协程运行 StartHeartbeat。 心跳间隔固定为 30 秒,发送 {"type":"ping"} 消息。

func WithHeartbeatOptions

func WithHeartbeatOptions(opts ...HeartbeatOption) UpgradeOption

WithHeartbeatOptions 启用心跳保活并传入高级配置参数。 参数:

  • opts: 心跳配置选项,如 WithHeartbeatInterval、WithPingMessage 等

示例:

ws.Upgrade(c,
    ws.WithHeartbeatOptions(
        ws.WithHeartbeatInterval(15*time.Second),
        ws.WithPingMessage(func() any { return Message{Type: "ping", Data: time.Now().Unix()} }),
    ),
)

func WithMaxConnPerIP added in v1.4.37

func WithMaxConnPerIP(n int) UpgradeOption

WithMaxConnPerIP 设置单个 IP 的最大 WebSocket 连接数。 参数:

  • n: 单 IP 最大连接数(如 10 表示每个 IP 最多建立 10 个连接)

超出限制的 Upgrade 调用返回 max connections per IP 错误。 默认不限制。连接关闭时自动从计数器中移除。

func WithQueueSize added in v1.4.37

func WithQueueSize(writeSize, readSize int) UpgradeOption

WithQueueSize 设置客户端读写队列缓冲区容量。 参数:

  • writeSize: 写入队列容量,<=0 时使用默认值 1024。增大可减少高并发时的丢包,但占用更多内存。
  • readSize: 读取队列容量,<=0 时使用默认值 1024。增大可应对消费速度跟不上生产速度的场景。

队列满载时新消息会被丢弃(写入)或丢失(读取),这是背压保护机制。 建议根据业务峰值 QPS 和消息大小估算,通常在 1024~4096 之间。

func WithReadLimit

func WithReadLimit(limit int64) UpgradeOption

WithReadLimit 设置单条消息的最大读取字节数(Upgrade 入口)。 超过此大小的消息将被拒绝,连接关闭。 默认 0 表示不限制(采用 gorilla/websocket 默认值 32768)。 此配置会在 Upgrade 内部自动传播到创建的 Client。

func WithReadTimeout added in v1.4.37

func WithReadTimeout(timeout time.Duration) UpgradeOption

WithReadTimeout 设置单次 ReadMessage 的超时时间。 参数:

  • timeout: 读取超时时间(如 60*time.Second)。超过此时间未收到消息, ReadMessage 返回超时错误,触发连接断开和重连。

此超时会在 Upgrade 内部自动传播到创建的 Client,无需额外传入 NewClient。 默认 0 表示不设置主动超时,由心跳机制间接管理 ReadDeadline。

func WithSSO added in v1.4.37

func WithSSO() UpgradeOption

WithSSO 启用单点登录模式,同一 UID 仅保留一个有效连接。 当已在线用户再次建立连接时,旧连接会被自动关闭(后登录踢前登录)。

func WithSubprotocols

func WithSubprotocols(protocols ...string) UpgradeOption

WithSubprotocols 设置 WebSocket 子协议协商列表。 参数:

  • protocols: 服务端支持的子协议列表,按优先级排列

客户端在 Sec-WebSocket-Protocol 头中声明支持的协议列表, 服务端从中选择一个匹配的返回,用于协议版本协商或多协议支持。

func WithWriteLimit added in v1.4.37

func WithWriteLimit(limit int64) UpgradeOption

WithWriteLimit 设置单条消息的最大写入字节数(Upgrade 入口)。 超过此大小的消息将被 WriteJSON/WriteRaw 拒绝,返回 ErrWriteLimitExceeded。 默认 0 表示不限制。 此配置会在 Upgrade 内部自动传播到创建的 Client。

func WithWriteMessageType added in v1.4.38

func WithWriteMessageType(msgType int) UpgradeOption

WithWriteMessageType 设置 WebSocket 消息帧类型。 参数:

  • msgType: 消息帧类型,可选 websocket.TextMessage(1) 或 websocket.BinaryMessage(2)

默认 TextMessage。二进制协议(如 protobuf)请传 websocket.BinaryMessage。

func WithWriteTimeout added in v1.4.37

func WithWriteTimeout(timeout time.Duration) UpgradeOption

WithWriteTimeout 设置单次写入的超时时间。 参数:

  • timeout: 写入超时时间(如 30*time.Second)。超过此时间写入操作未完成, 将触发重试机制。

此超时会在 Upgrade 内部自动传播到创建的 Client 的 writeWithRetry。 默认 0 表示使用默认值 10s。

func WithWsRateLimit added in v1.4.37

func WithWsRateLimit(rps int, burst int) UpgradeOption

WithWsRateLimit 设置全局 WebSocket 升级速率限制。 参数:

  • rps: 每秒最大升级次数(如 100 表示每秒最多处理 100 次握手)
  • burst: 最大突发量(如 20 表示短时间内允许最多 20 个突发连接)

超出限制的 Upgrade 调用返回 rate limited 错误。 默认不限制。

Directories

Path Synopsis
Package service 业务逻辑层 - 提供 WebSocket 消息推送能力
Package service 业务逻辑层 - 提供 WebSocket 消息推送能力

Jump to

Keyboard shortcuts

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