Files
zogo/wsc/session_client.go
T

268 lines
9.5 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// session_client.go 提供会话层客户端(SessionClient):在 Client(传输层)之上封装
// 「token 延迟连接、连上即注册、应用层心跳、重连回调」等生产编排,业务零样板直接用。
package wsc
import (
"net/http"
"strings"
"sync"
"time"
"github.com/sirupsen/logrus"
)
// SessionConfig 会话层配置(在 ClientConfig 之上做封装)。
type SessionConfig struct {
URL string // WebSocket 地址(必填),如 wss://host/api/ws
Token string // 鉴权 token,拼接在 URL ?token= 后;为空时不主动连接,待 UpdateToken 注入后再连
Header http.Header // 自定义请求头(如 Authorization),与 Token 二选一或并存
AutoReconnect bool // 是否自动重连,默认 true
HeartbeatInterval time.Duration // 应用层心跳(action=ping)间隔,默认 25s;<=0 关闭
}
// SessionClient 与服务器的会话客户端,基于 Client(传输层)封装:
// 负责连接、自动重连、发送注册/心跳并接收响应。
//
// 关键行为(生产验证过的编排):
// - token 为空时绝不空转重连(避免无凭证时的指数退避风暴),
// 等 UpdateToken 注入 token 后自动发起首次连接;
// - 连接成功(首次/重连)自动发送 register 注册请求;
// - 注册完成后再触发 OnReconnected——保证服务端先认识这个连接,再接收业务数据。
// 业务在 OnReconnected 中做全量对账重放,兜底断连期间丢失的增量上报。
//
// 用法:
//
// c := wsc.NewSessionClient(wsc.SessionConfig{URL: cfg.WsUrl, Token: cfg.Token})
// c.OnReconnected = func() { /* 全量对账 */ }
// c.OnBind("biz.action", func(resp *BizResp) { ... })
// c.Start(&RegisterReq{...})
// defer c.Close()
type SessionClient struct {
// Cli 底层传输层客户端,可直接设置 OnConnected / OnDisconnected /
// OnReconnecting 回调,或调用 On / Send 等方法。
Cli *Client
cfg SessionConfig
// registerPayload 注册请求载荷(业务自定义结构,如主机名/IP/版本号),
// 连接成功后发送;登录后可通过 UpdateRegisterInfo 更新。
registerPayload any
// OnReconnected 连接建立(首次连接与断线重连成功)后触发的回调。
// 业务可在此发起全量对账重放,兜底断连期间丢失的上报。
OnReconnected func()
heartbeatStop chan struct{}
closeOnce sync.Once
connectOnce sync.Once // 保证首次 Connect 仅执行一次(由 Start 或 UpdateToken 触发)
}
// NewSessionClient 创建会话客户端(尚未连接)。构造时会注册内置消息路由,
// 调用方应在 Start/Connect 前完成自定义配置(On/OnBind/OnReconnected 等)。
func NewSessionClient(cfg SessionConfig) *SessionClient {
url := appendToken(cfg.URL, cfg.Token)
wscCfg := DefaultClientConfig(url)
if !cfg.AutoReconnect {
wscCfg.AutoReconnect = false
}
// 加快重连节奏:默认退避上限 30s 体感太慢,收敛到 5s(与参考实现一致)。
wscCfg.MaxReconDelay = 5 * time.Second
if cfg.Header != nil {
wscCfg.Header = cfg.Header
}
if cfg.HeartbeatInterval == 0 {
cfg.HeartbeatInterval = 25 * time.Second
}
cli := NewClient(wscCfg)
c := &SessionClient{Cli: cli, cfg: cfg, heartbeatStop: make(chan struct{})}
c.registerHandlers()
return c
}
// registerHandlers 注册内置消息路由(在 Connect 之前调用)。
func (c *SessionClient) registerHandlers() {
// pong — 服务端对心跳的响应
c.Cli.On(ActionPong, BindPayload(func(resp *PingResp) {
logrus.Debugf("[wsc-session] pong received, time=%d", resp.Time)
}))
// register — 注册(上线)响应(请求与响应共用同一 action)
c.Cli.On(ActionRegister, BindPayload(func(resp *RegisterResp) {
if resp.Success {
logrus.Infof("[wsc-session] register success: %s", resp.Message)
} else {
logrus.Warnf("[wsc-session] register failed: %s", resp.Message)
}
}))
}
// ============================================================
// 连接控制
// ============================================================
// Start 启动客户端并连接。连接/重连成功后自动发送注册请求并启动心跳,
// 调用方无需在外部自行编排心跳循环(registerPayload 为 nil 时仅连接不注册)。
//
// 若 token 为空则跳过连接,等待 UpdateToken 注入 token 后再发起首次连接。
func (c *SessionClient) Start(registerPayload any) error {
c.registerPayload = registerPayload
prev := c.Cli.OnConnected
c.Cli.OnConnected = func() {
if prev != nil {
prev()
}
logrus.Infof("[wsc-session] connected to %s", c.cfg.URL)
if c.registerPayload != nil {
if err := c.SendRegister(c.registerPayload); err != nil {
logrus.Errorf("[wsc-session] send register failed: %v", err)
}
}
// 注册完成后再触发业务回调(全量对账重放),
// 保证服务端先认识这个连接,再接收业务数据。
if c.OnReconnected != nil {
c.OnReconnected()
}
}
go c.heartbeatLoop()
// 没有 token 肯定连不上,跳过连接,等 UpdateToken 注入 token 后再连。
if c.cfg.Token == "" {
logrus.Warnf("[wsc-session] token empty, skip connect, wait for UpdateToken")
return nil
}
c.doConnect()
return nil
}
// doConnect 发起首次连接(仅执行一次,由 Start 或 UpdateToken 触发)。
// 连接在独立 goroutine 中重试,不阻塞调用方(避免 UpdateToken/Start 的调用者被重连循环卡死)。
func (c *SessionClient) doConnect() {
c.connectOnce.Do(func() {
go func() {
if err := c.Cli.Connect(); err != nil {
logrus.Errorf("[wsc-session] connect loop exited: %v", err)
}
}()
})
}
// heartbeatLoop 内部心跳协程:连接状态下按间隔发送 ping,断开或 Close 时退出。
// 业务编排收口在客户端内,调用方不必在外部再写一层心跳逻辑。
func (c *SessionClient) heartbeatLoop() {
if c.cfg.HeartbeatInterval <= 0 {
return
}
ticker := time.NewTicker(c.cfg.HeartbeatInterval)
defer ticker.Stop()
for {
select {
case <-c.heartbeatStop:
return
case <-ticker.C:
if !c.Connected() {
continue
}
if err := c.SendPing(); err != nil {
logrus.Errorf("[wsc-session] heartbeat send failed: %v", err)
}
}
}
}
// UpdateRegisterInfo 更新注册信息(如登录后拿到业务ID)。若已连接则立即重发注册,
// 确保服务端拿到完整字段;未连接则仅暂存,待连接成功后由 OnConnected 发送。
func (c *SessionClient) UpdateRegisterInfo(payload any) {
c.registerPayload = payload
if c.Connected() {
if err := c.SendRegister(payload); err != nil {
logrus.Errorf("[wsc-session] re-send register failed: %v", err)
}
}
}
// UpdateToken 动态更新鉴权 token(重新拼接到 URL)。
// 若此前因 token 为空未连接,则在此发起首次连接;已连接/正在连接则不受影响
// (connectOnce 保证仅连一次,后续断连由底层 Client 自动重连并使用新 URL)。并发安全。
func (c *SessionClient) UpdateToken(token string) {
c.cfg.Token = token
c.Cli.SetURL(appendToken(c.cfg.URL, token))
// token 就绪:若此前因无 token 未连接,则在此发起首次连接(仅一次,异步不阻塞)。
if token != "" {
logrus.Infof("[wsc-session] token updated, triggering connect")
c.doConnect()
}
}
// Close 关闭连接、停止重连与心跳协程。
func (c *SessionClient) Close() {
c.closeOnce.Do(func() {
close(c.heartbeatStop)
})
c.Cli.Close()
}
// Connected 返回当前是否已连接。
func (c *SessionClient) Connected() bool {
return c.Cli.Connected()
}
// Done 返回一个通道,客户端完全关闭后关闭。
func (c *SessionClient) Done() <-chan struct{} {
return c.Cli.Done()
}
// ============================================================
// 发送消息
// ============================================================
// SendPing 发送心跳(action=ping)。
func (c *SessionClient) SendPing() error {
return c.Cli.Send(ActionPing, &PingReq{})
}
// SendRegister 发送注册(上线)请求(action=register)。
func (c *SessionClient) SendRegister(payload any) error {
return c.Cli.Send(ActionRegister, payload)
}
// Send 发送任意自定义消息(并发安全)。
func (c *SessionClient) Send(action string, payload any) error {
return c.Cli.Send(action, payload)
}
// SendRaw 发送已序列化的原始字节(并发安全)。
// 注意:重连时底层 writeCh 会重建、旧缓冲数据被丢弃,
// 上层务必「未连接不投递」,未落盘的数据等待 OnReconnected 全量对账。
func (c *SessionClient) SendRaw(data []byte) error {
return c.Cli.SendRaw(data)
}
// ============================================================
// 接收消息:注册自定义处理器(直接透传到底层 Cli)
// ============================================================
// On 注册 action 对应的消息处理器,回调拿到原始 payload(json.RawMessage),
// 在 Connect/Start 之前调用。
func (c *SessionClient) On(action string, handler PayloadHandler) {
c.Cli.On(action, handler)
}
// OnBind 注册带类型反序列化的消息处理器(自动将 payload 反序列化到结构体指针)。
func (c *SessionClient) OnBind(action string, fn any) {
c.Cli.On(action, BindPayload(fn))
}
// appendToken 把 token 拼接到 URL 查询串。
func appendToken(url, token string) string {
if token == "" {
return url
}
sep := "?"
if strings.Contains(url, "?") {
sep = "&"
}
return url + sep + "token=" + token
}