Documentation
¶
Overview ¶
Package gows 提供 WebSocket 服务端封装,包含连接管理、消息读写、心跳保活和全局分发。
Index ¶
- Constants
- Variables
- func ParseTokenCtx(ctx context.Context, tokenString string) (string, error)
- func SetupPongHandler(conn *websocket.Conn, interval, pongTimeout time.Duration)
- func StartHeartbeat(client *Client, opts ...HeartbeatOption)
- type Backend
- type Client
- func (c *Client) BroadcastCtx(ctx context.Context, v any) error
- func (c *Client) BroadcastReliableCtx(ctx context.Context, v any) error
- func (c *Client) Close() error
- func (c *Client) Context() context.Context
- func (c *Client) DisconnectByUID(ctx context.Context, uid string) error
- func (c *Client) DisconnectByUIDs(ctx context.Context, uids ...string) int
- func (c *Client) Done() <-chan struct{}
- func (c *Client) IsAlive() bool
- func (c *Client) ReadMsgFromClientReadCh(ctx context.Context) ([]byte, error)
- func (c *Client) RemoteAddr() string
- func (c *Client) SendToMultiUIDCtx(ctx context.Context, uids []string, v any) error
- func (c *Client) SendToUIDCtx(ctx context.Context, uid string, v any) error
- func (c *Client) SetCloseHook(fn func())
- func (c *Client) Stats() ClientStats
- func (c *Client) UID() string
- func (c *Client) WriteJSONToClientWriteCh(ctx context.Context, v any) error
- func (c *Client) WriteRawToClientWriteCh(ctx context.Context, data []byte) error
- type ClientStats
- type DispatcherOption
- type DispatcherStats
- type DistributedDispatcher
- func (dd *DistributedDispatcher) BroadcastCtx(ctx context.Context, v any) error
- func (dd *DistributedDispatcher) BroadcastFilterCtx(ctx context.Context, v any, filter func(*Client) bool) error
- func (dd *DistributedDispatcher) BroadcastReliableCtx(ctx context.Context, v any) error
- func (dd *DistributedDispatcher) CleanupDeadConns() int
- func (dd *DistributedDispatcher) Clients() []*Client
- func (dd *DistributedDispatcher) ConnectedUIDs() []string
- func (dd *DistributedDispatcher) DisconnectByUID(ctx context.Context, uid string) error
- func (dd *DistributedDispatcher) DisconnectByUIDs(ctx context.Context, uids ...string) int
- func (dd *DistributedDispatcher) Len() int
- func (dd *DistributedDispatcher) LenLocal() int
- func (dd *DistributedDispatcher) MaxConnections() int
- func (dd *DistributedDispatcher) Range(fn func(*Client) bool)
- func (dd *DistributedDispatcher) RegisterCtx(ctx context.Context, client *Client) error
- func (dd *DistributedDispatcher) SendToMultiUIDCtx(ctx context.Context, uids []string, v any) error
- func (dd *DistributedDispatcher) SendToUIDCtx(ctx context.Context, uid string, v any) error
- func (dd *DistributedDispatcher) Start(ctx context.Context)
- func (dd *DistributedDispatcher) Stats() DispatcherStats
- func (dd *DistributedDispatcher) Stop()
- func (dd *DistributedDispatcher) UnregisterCtx(ctx context.Context, client *Client) error
- func (dd *DistributedDispatcher) WriteRawToLocalUIDs(span trace.Span, uids []string, payload json.RawMessage) error
- type HeartbeatOption
- type Message
- type PubSubMessage
- type RabbitMQBackend
- func (b *RabbitMQBackend) Close() error
- func (b *RabbitMQBackend) Publish(ctx context.Context, msg *PubSubMessage) error
- func (b *RabbitMQBackend) SetPoolCap(initialCap, maxCap int)
- func (b *RabbitMQBackend) SetReceiveChanSize(size int)
- func (b *RabbitMQBackend) Subscribe(ctx context.Context, uid string) (<-chan *PubSubMessage, error)
- func (b *RabbitMQBackend) SubscribeBroadcast(ctx context.Context) (<-chan *PubSubMessage, error)
- func (b *RabbitMQBackend) Unsubscribe(_ context.Context, uid string) error
- type UpgradeOption
- func WithBufferSize(read, write int) UpgradeOption
- func WithCheckOrigin(fn func(r *http.Request) bool) UpgradeOption
- func WithClientUID(uid string) UpgradeOption
- func WithDispatcher(d *DistributedDispatcher) UpgradeOption
- func WithEnableCompression(enable bool) UpgradeOption
- func WithEnableDistributed(enabled bool) UpgradeOption
- func WithHeartbeat() UpgradeOption
- func WithHeartbeatOptions(opts ...HeartbeatOption) UpgradeOption
- func WithMaxConnPerIP(n int) UpgradeOption
- func WithQueueSize(writeSize, readSize int) UpgradeOption
- func WithReadLimit(limit int64) UpgradeOption
- func WithReadTimeout(timeout time.Duration) UpgradeOption
- func WithSSO() UpgradeOption
- func WithSubprotocols(protocols ...string) UpgradeOption
- func WithWriteLimit(limit int64) UpgradeOption
- func WithWriteMessageType(msgType int) UpgradeOption
- func WithWriteTimeout(timeout time.Duration) UpgradeOption
- func WithWsRateLimit(rps int, burst int) UpgradeOption
Constants ¶
const ( MsgTypeBroadcast = "broadcast" MsgTypeSendToUID = "send_to_uid" MsgTypeClientOnline = "client_online" MsgTypeClientOffline = "client_offline" )
PubSubMessage 的消息类型常量。
Variables ¶
var DefaultDispatcher = NewDispatcher(nil)
DefaultDispatcher 默认全局 Dispatcher 单例(单机模式)。
var ErrClientNotFound = errors.New("dispatcher: client not found")
ErrClientNotFound 未找到指定 UID 的客户端连接
var ErrClientNotRegistered = errors.New("dispatcher: client not registered")
ErrClientNotRegistered 客户端未注册到 Dispatcher
var ErrEmptyUID = errors.New("dispatcher: empty uid")
ErrEmptyUID UID 为空时禁止注册
var ErrMaxConnections = errors.New("dispatcher: max connections reached")
ErrMaxConnections 达到最大连接数限制的错误
var ErrNilContext = errors.New("ctx must not be nil")
ErrNilContext 传入 nil context 的错误标识。 当 ReadMsgFromClientReadCh 接收到 nil context 时返回此错误,由调用方自行修复。
var ErrNoDispatcher = errors.New("client has no dispatcher configured")
ErrNoDispatcher Client 未关联 Dispatcher,无法执行跨实例发送/广播。
var ErrTokenInvalid = errors.New("token is invalid: missing uid")
ErrTokenInvalid token 无效(格式正确但缺少 uid 字段)
var ErrWriteLimitExceeded = errors.New("write message exceeds size limit")
ErrWriteLimitExceeded 单条消息大小超过写入限制的错误标识。 当消息超过 WithWriteLimit 设置的字节数时 WriteJSON/WriteRaw 返回此错误。
var ErrWriteQueueFull = errors.New("write queue is full, message dropped")
ErrWriteQueueFull 写入队列已满,消息被丢弃的错误标识。 当客户端写入缓冲区满载时 WriteJSON 返回此错误, 发送方可根据此错误判断是否为短暂拥塞,决定是否降级处理。
Functions ¶
func ParseTokenCtx ¶ added in v1.4.36
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
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 负责:
- 执行升级前钩子(可选)
- 创建 websocket.Upgrader 并应用配置
- 执行 HTTP→WebSocket 升级
- 执行升级后钩子(可选)
- 用 *websocket.Conn 创建 *Client(含 msgFromChToWs)
- 可选启动心跳、注册 Dispatcher
安全保护:
- CORS: 默认拒绝所有来源,需显式调用 WithCheckOrigin
- 限流: WithWsRateLimit 设置全局升级速率
- 单IP限制: WithMaxConnPerIP 设置单IP最大连接数
调用方需负责 defer client.Close() 确保资源释放。
func (*Client) BroadcastCtx ¶ added in v1.4.37
BroadcastCtx 通过关联的 Dispatcher 向所有在线客户端广播消息。 仅在 Client 已注册到 Dispatcher 时有效。
func (*Client) BroadcastReliableCtx ¶ added in v1.4.37
BroadcastReliableCtx 通过关联的 Dispatcher 进行可靠广播(按 UID 下发)。 仅在 Client 已注册到 Dispatcher 时有效。
func (*Client) Close ¶
Close 优雅关闭 WebSocket 连接,触发关闭信号并释放所有资源。 关闭顺序: 执行关闭钩子 → 取消上下文 → 等待 msgFromChToWs 完全退出 (drain writeCh 剩余消息后退出,确保所有排队消息已写入) → 关闭底层连接(触发 msgFromWsToCh 的 ReadMessage 返回错误) → 等待 msgFromWsToCh 完全退出。 通过 atomic.CompareAndSwap 保证幂等性,首次调用执行完整关闭流程, 后续调用直接返回 nil。 返回:
- error: 首次关闭底层连接失败时返回 error,重复关闭返回 nil
func (*Client) DisconnectByUID ¶ added in v1.4.39
DisconnectByUID 通过关联的 Dispatcher 断开指定 UID 的连接。 仅在 Client 已注册到 Dispatcher 时有效。
func (*Client) DisconnectByUIDs ¶ added in v1.4.39
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
IsAlive 判断客户端连接是否处于健康状态。 返回 false 的场景:
- Close() 已调用
- msgFromChToWs 已因写入重试全部失败退出(closeWsConn 已触发)
- 超过 3 个心跳周期无成功写入(疑似僵尸连接)
注意:
- 此方法返回 true 不代表底层网络一定可达,仅表示组件内部状态正常
- 精确的活性检测依赖 Ping/Pong 协议级心跳 + ReadDeadline 联动
func (*Client) ReadMsgFromClientReadCh ¶ added in v1.4.37
ReadMsgFromClientReadCh 同步阻塞读取客户端发送的一条消息。
func (*Client) SendToMultiUIDCtx ¶ added in v1.4.37
SendToMultiUIDCtx 通过关联的 Dispatcher 向多个 UID 的用户发送消息。 发送者自身不会收到消息(自动过滤)。 仅在 Client 已注册到 Dispatcher 时有效。
func (*Client) SendToUIDCtx ¶ added in v1.4.37
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) WriteJSONToClientWriteCh ¶ added in v1.4.37
WriteJSONToClientWriteCh 向客户端写入 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
SendToMultiUIDCtx 向多个指定 UID 的客户端发送消息(跨实例)。
func (*DistributedDispatcher) SendToUIDCtx ¶ added in v1.4.37
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
func (dd *DistributedDispatcher) Stats() DispatcherStats
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 错误。 默认不限制。
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
Package service 业务逻辑层 - 提供 WebSocket 消息推送能力
|
Package service 业务逻辑层 - 提供 WebSocket 消息推送能力 |