Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9ef71f51d4 |
@@ -1,24 +1,27 @@
|
|||||||
# zogo
|
# zogo
|
||||||
|
|
||||||
zeroone Go 通用库(单仓多包)。被各后端服务以 Go module 依赖引用:
|
zeroone Go 通用库(单仓多包,顶层即包)。被各后端服务以 Go module 依赖引用:
|
||||||
`go get git.zeroonesoft.cn/golib/zogo/xxx`
|
`go get git.zeroonesoft.cn/golib/zogo/xxx`
|
||||||
|
|
||||||
## 包清单
|
## 包清单
|
||||||
|
|
||||||
| 包 | 用途 | 主要依赖 |
|
| 包 | 用途 | 文档 |
|
||||||
|---|---|---|
|
|---|---|---|
|
||||||
| `httpx` | Gin 统一出入参:OkJson/ErrorJson(CodeMsg 错误码透传)、Handle/HandleResult、Parse(uri→query→header→json 顺序绑定 + binding 校验收口) | gin, validator |
|
| `httpx` | Gin 统一出入参:OkJson/ErrorJson(CodeMsg 错误码透传)、Handle/HandleResult、Parse(uri→query→header→json 顺序绑定 + binding 校验收口) | [httpx/README.md](httpx/README.md) |
|
||||||
| `logger` | logrus 日志初始化(文件轮转) | logrus, rotatelogs |
|
| `logger` | logrus 日志初始化(控制台彩色 + 按天轮转文件 + 历史压缩) | [logger/README.md](logger/README.md) |
|
||||||
| `wsc` | WebSocket 服务框架(会话/房间/路由/反射绑定) | gin, gorilla/websocket, uuid, logrus |
|
| `wsc` | WebSocket 服务端框架(会话/房间/路由)+ 配套客户端(自动重连),同一 Message 信封 | [wsc/README.md](wsc/README.md) |
|
||||||
| `wscclient` | WebSocket 客户端(自动重连) | gorilla/websocket |
|
| `jwtx` | 轻量 HMAC-SHA256 令牌签发与校验 | [jwtx/README.md](jwtx/README.md) |
|
||||||
| `jwtx` | 轻量 HMAC-SHA256 JWT 签发与校验 | 无 |
|
|
||||||
|
> 历史说明:v0.1.0 曾存在独立包 `wscclient`,v0.2.0 起并入 `wsc`(`wsc.Client`)——
|
||||||
|
> 客户端与服务端共用信封定义,协议单一来源防两端漂移。
|
||||||
|
|
||||||
## 约定
|
## 约定
|
||||||
|
|
||||||
- 通用性准入:被 ≥2 个服务复制使用过的包才进入本仓库;仅单服务使用的留在服务内。
|
- 通用性准入:被 ≥2 个服务复制使用过的包才进入本仓库;仅单服务使用的留在服务内。
|
||||||
- 改动流程:在本仓库修复/演进 → 打 tag(如 v0.1.1)→ 各服务 `go get -u git.zeroonesoft.cn/golib/zogo@v0.1.1`。
|
- 改动流程:在本仓库修复/演进 → 打 tag(如 v0.2.0)→ 各服务 `go get git.zeroonesoft.cn/golib/zogo@v0.2.0`。
|
||||||
- 本地联动开发:在 `D:\zomaintain\backend\go.work` 中 `use D:/golib/zogo`(go.work 不提交)。
|
- 本地联动开发:在 `D:\zomaintain\backend\go.work` 中 `use D:/golib/zogo`(go.work 不提交)。
|
||||||
- go.mod 声明 go 1.24(不高于最低版本消费方),不要随意上调。
|
- go.mod 声明 go 1.24(不高于最低版本消费方),不要随意上调。
|
||||||
|
- 每个包自带 README.md(定位/用法/注意事项),改 API 先改文档。
|
||||||
|
|
||||||
## 私有模块拉取配置(每台开发机/CI 一次)
|
## 私有模块拉取配置(每台开发机/CI 一次)
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,31 @@
|
|||||||
|
# httpx
|
||||||
|
|
||||||
|
Gin 服务的统一出入参收口。
|
||||||
|
|
||||||
|
## 响应
|
||||||
|
|
||||||
|
统一信封 `{"code":0,"msg":"OK","data":...}`:
|
||||||
|
|
||||||
|
- `OkJson(c, data)` —— 成功出参
|
||||||
|
- `ErrorJson(c, err)` —— 错误出参;`*CodeMsg` 携带的错误码**透出**(如 401),其余错误折叠为 code 1
|
||||||
|
- `wsc.New(401, "未登录")` 构造带码错误
|
||||||
|
|
||||||
|
## Handler 收口
|
||||||
|
|
||||||
|
```go
|
||||||
|
// 泛型:Parse 取参 → 调 logic → 统一出参
|
||||||
|
httpx.Handle(c, func(req *CreateReq) (*CreateResp, error) { return svc.Create(req) })
|
||||||
|
httpx.HandleResult(c, data, err) // 非泛型场景
|
||||||
|
```
|
||||||
|
|
||||||
|
## Parse 参数绑定
|
||||||
|
|
||||||
|
`Parse(c, &req)` 按固定顺序绑定:**uri → query → header → json body**,全部绑定后统一跑
|
||||||
|
binding tag 校验(`bindingValidator`,与 gin 同源 tag 名)。
|
||||||
|
|
||||||
|
注意事项:
|
||||||
|
|
||||||
|
- 包 `init()` 会设置 `binding.Validator = nil`(关闭 gin 内置校验,避免 query 阶段误报
|
||||||
|
json 字段的 required)。**引入本包即全局生效**,服务内不要绕过 Parse 自行 ShouldBind。
|
||||||
|
- 请求结构体可实现 `Validate() error` 接口做自定义校验,在 binding 校验之前调用。
|
||||||
|
- body 解析用 `json.NewDecoder` 直解,不触发 gin 内部校验。
|
||||||
@@ -0,0 +1,23 @@
|
|||||||
|
# jwtx
|
||||||
|
|
||||||
|
轻量 HMAC-SHA256 令牌(JWT 风格,无第三方依赖),用于登录态/服务间鉴权。
|
||||||
|
|
||||||
|
## 用法
|
||||||
|
|
||||||
|
```go
|
||||||
|
import "git.zeroonesoft.cn/golib/zogo/jwtx"
|
||||||
|
|
||||||
|
token, err := jwtx.Sign(uid, 3600) // 签发,TTL 秒
|
||||||
|
claims, err := jwtx.Parse(token) // 校验签名+过期,返回 *Claims{Uid, Exp, Msg}
|
||||||
|
```
|
||||||
|
|
||||||
|
## 格式
|
||||||
|
|
||||||
|
`base64url(payload).base64url(hmac-sha256(payload))`,两段式(无 header 段),payload 为
|
||||||
|
`Claims` JSON。
|
||||||
|
|
||||||
|
## 注意
|
||||||
|
|
||||||
|
- **签名密钥当前硬编码在包内**(`zo-maintain-jwt-secret-2026`)。仅适用于内网/测试场景;
|
||||||
|
用于生产鉴权前应改为注入式(加 `SetSecret` 或从环境读取)。
|
||||||
|
- 过期校验在 `Parse` 内完成(`Exp` 与当前时间比较)。
|
||||||
@@ -0,0 +1,27 @@
|
|||||||
|
# logger
|
||||||
|
|
||||||
|
logrus 封装:控制台彩色输出(Windows colorable)+ 按天轮转文件 + 历史日志压缩。
|
||||||
|
|
||||||
|
## 用法
|
||||||
|
|
||||||
|
```go
|
||||||
|
import "git.zeroonesoft.cn/golib/zogo/logger"
|
||||||
|
|
||||||
|
logger.SetProjectPrefix("D:/gopath/your-service/") // 可选:日志里 file 字段去掉公共前缀
|
||||||
|
logger.InitLog(logrus.InfoLevel, true, "D:/logs/app") // 可变参数传日志目录则同时落文件
|
||||||
|
```
|
||||||
|
|
||||||
|
## 行为
|
||||||
|
|
||||||
|
- 控制台:强制彩色、完整时间戳、caller 显示为 `相对路径:行号`
|
||||||
|
- 落文件(传了 logDir 时):
|
||||||
|
- `app-YYYY-MM-DD.log` 全级别
|
||||||
|
- `error-YYYY-MM-DD.log` 仅 Error 及以上
|
||||||
|
- 按天轮转(file-rotatelogs)
|
||||||
|
- `logger.StartLogCompressor(logDir, interval)`:后台协程周期压缩历史日志(每日 0:00-0:30 跳过窗口)
|
||||||
|
- `logger.CompressHistoryLogs(logDir)`:手动触发一次压缩
|
||||||
|
|
||||||
|
## 注意
|
||||||
|
|
||||||
|
初始化后请直接用 `logrus.Info` 等包级函数打日志,**不要**保存 `logrus.Logger` 实例调用,
|
||||||
|
否则 caller 文件位置信息会不准。
|
||||||
@@ -0,0 +1,44 @@
|
|||||||
|
# wsc
|
||||||
|
|
||||||
|
WebSocket 服务端框架 + 配套客户端,**同一 Message 信封、同一处协议定义**。
|
||||||
|
|
||||||
|
信封格式:`{"action":"...","payload":...}`(`Message` 结构体,payload 为 `json.RawMessage`)。
|
||||||
|
心跳约定:双方默认 60s 读超时 + 周期 ping(客户端 `PongTimeout` 需 ≥ 服务端 ping 周期 54s,默认值即满足)。
|
||||||
|
|
||||||
|
## 服务端
|
||||||
|
|
||||||
|
```go
|
||||||
|
import "git.zeroonesoft.cn/golib/zogo/wsc"
|
||||||
|
|
||||||
|
wsc.HandleWS(r, "/ws/chat", wsc.QuickServer(onConnect, onDisconnect), func(router *wsc.Router) {
|
||||||
|
// 带响应:fn 返回 (resp, error),resp 序列化后以同名 action+".resp" 语义回包
|
||||||
|
router.On("ping", wsc.Bind(func(ctx *wsc.Context, req *PingReq) (*PingResp, error) {
|
||||||
|
return &PingResp{Pong: true}, nil
|
||||||
|
}))
|
||||||
|
// 无响应 / 客户端约定固定回包类型时
|
||||||
|
router.On("register", wsc.BindNoResp(func(ctx *wsc.Context, req *RegisterReq) error { ... }))
|
||||||
|
})
|
||||||
|
```
|
||||||
|
|
||||||
|
- `Server`:连接池与房间管理(`GetOrCreateRoom`/`ActiveCount`…),`OnConnect/OnDisconnect/OnUpgrade` 钩子
|
||||||
|
- `Router`:action 分发,反射校验 handler 签名;`Use` 挂中间件
|
||||||
|
- `Context`:单会话上下文,实现 `context.Context`;`WriteMessage`/`WriteError`/`JoinRoom`/`Set`…,`WriteError` 使用 `WsActionError` action
|
||||||
|
- gin 耦合仅在 `HandleWS`/`newContext`/`OnUpgrade` 三处入口,核心零 gin
|
||||||
|
|
||||||
|
## 客户端
|
||||||
|
|
||||||
|
```go
|
||||||
|
cli := wsc.NewClient(wsc.DefaultClientConfig("wss://cloud/ws/gateway?token=xxx"))
|
||||||
|
cli.On("kick", wsc.BindPayload(func(req *KickReq) { ... }))
|
||||||
|
cli.OnConnected = func() { cli.Send("register", &RegisterReq{...}) }
|
||||||
|
cli.SetURL("wss://cloud/ws/gateway?token=新token") // 重连时生效
|
||||||
|
_ = cli.Connect() // AutoReconnect=true 时阻塞到连上(指数退避,稳定 10s 后重置)
|
||||||
|
defer cli.Close()
|
||||||
|
```
|
||||||
|
|
||||||
|
- 自动重连:指数退避(1s→30s 封顶),连接稳定满 10s 才重置退避,防抖动时"永远 1s 一连"
|
||||||
|
- 断线期间 `Send` 返回错误;重连成功重建读写泵,`OnConnected` 重新触发(在此做重注册/重订阅)
|
||||||
|
|
||||||
|
## 线程安全
|
||||||
|
|
||||||
|
`Client` 与 `Session` 的公开方法均可并发调用;发送走带缓冲通道 + 单写泵。
|
||||||
@@ -1,6 +1,7 @@
|
|||||||
// Package wscclient 提供 wsc 服务端的客户端封装。
|
// client.go 提供 wsc 服务端的客户端封装(Client)。
|
||||||
// 对称设计:cli.On == router.On,cli.Send == ctx.WriteMessage。
|
// 对称设计:cli.On == router.On,cli.Send == ctx.WriteMessage。
|
||||||
package wscclient
|
// 客户端与服务端共用同一 Message 信封与协议约定,改协议两边一起改。
|
||||||
|
package wsc
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
@@ -20,8 +21,8 @@ import (
|
|||||||
// 配置
|
// 配置
|
||||||
// ============================================================
|
// ============================================================
|
||||||
|
|
||||||
// Config 客户端配置。
|
// ClientConfig 客户端配置。
|
||||||
type Config struct {
|
type ClientConfig struct {
|
||||||
URL string // WebSocket 地址(必填)
|
URL string // WebSocket 地址(必填)
|
||||||
Header http.Header // 自定义请求头(如 token)
|
Header http.Header // 自定义请求头(如 token)
|
||||||
DialTimeout time.Duration // 连接超时,默认 10s
|
DialTimeout time.Duration // 连接超时,默认 10s
|
||||||
@@ -35,9 +36,9 @@ type Config struct {
|
|||||||
SendBufferSize int // 发送缓冲区大小,默认 256
|
SendBufferSize int // 发送缓冲区大小,默认 256
|
||||||
}
|
}
|
||||||
|
|
||||||
// DefaultConfig 返回推荐默认配置。
|
// DefaultClientConfig 返回推荐默认配置。
|
||||||
func DefaultConfig(url string) Config {
|
func DefaultClientConfig(url string) ClientConfig {
|
||||||
return Config{
|
return ClientConfig{
|
||||||
URL: url,
|
URL: url,
|
||||||
DialTimeout: 10 * time.Second,
|
DialTimeout: 10 * time.Second,
|
||||||
AutoReconnect: true,
|
AutoReconnect: true,
|
||||||
@@ -52,11 +53,11 @@ func DefaultConfig(url string) Config {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// ============================================================
|
// ============================================================
|
||||||
// Handler — 客户端消息处理器
|
// PayloadHandler — 客户端消息处理器
|
||||||
// ============================================================
|
// ============================================================
|
||||||
|
|
||||||
// Handler 客户端消息处理函数。payload 为服务端发来的 JSON,已去掉 action 包装。
|
// PayloadHandler 客户端消息处理函数。payload 为服务端发来的 JSON,已去掉 action 包装。
|
||||||
type Handler func(payload json.RawMessage)
|
type PayloadHandler func(payload json.RawMessage)
|
||||||
|
|
||||||
// ============================================================
|
// ============================================================
|
||||||
// Client
|
// Client
|
||||||
@@ -67,11 +68,11 @@ type Handler func(payload json.RawMessage)
|
|||||||
type Client struct {
|
type Client struct {
|
||||||
url string
|
url string
|
||||||
urlMu sync.RWMutex
|
urlMu sync.RWMutex
|
||||||
cfg Config
|
cfg ClientConfig
|
||||||
|
|
||||||
// 消息路由
|
// 消息路由
|
||||||
mu sync.RWMutex
|
mu sync.RWMutex
|
||||||
routes map[string]Handler
|
routes map[string]PayloadHandler
|
||||||
|
|
||||||
// 连接
|
// 连接
|
||||||
conn *websocket.Conn
|
conn *websocket.Conn
|
||||||
@@ -102,7 +103,7 @@ type Client struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// New 创建客户端(尚未连接)。
|
// New 创建客户端(尚未连接)。
|
||||||
func New(cfg Config) *Client {
|
func NewClient(cfg ClientConfig) *Client {
|
||||||
if cfg.DialTimeout == 0 {
|
if cfg.DialTimeout == 0 {
|
||||||
cfg.DialTimeout = 10 * time.Second
|
cfg.DialTimeout = 10 * time.Second
|
||||||
}
|
}
|
||||||
@@ -128,7 +129,7 @@ func New(cfg Config) *Client {
|
|||||||
return &Client{
|
return &Client{
|
||||||
url: cfg.URL,
|
url: cfg.URL,
|
||||||
cfg: cfg,
|
cfg: cfg,
|
||||||
routes: make(map[string]Handler),
|
routes: make(map[string]PayloadHandler),
|
||||||
writeCh: make(chan []byte, cfg.SendBufferSize),
|
writeCh: make(chan []byte, cfg.SendBufferSize),
|
||||||
closeCh: make(chan struct{}),
|
closeCh: make(chan struct{}),
|
||||||
doneCh: make(chan struct{}),
|
doneCh: make(chan struct{}),
|
||||||
@@ -201,19 +202,24 @@ func (c *Client) getURL() string {
|
|||||||
// ============================================================
|
// ============================================================
|
||||||
|
|
||||||
// On 注册 action 对应的消息处理器。
|
// On 注册 action 对应的消息处理器。
|
||||||
func (c *Client) On(action string, handler Handler) {
|
func (c *Client) On(action string, handler PayloadHandler) {
|
||||||
c.mu.Lock()
|
c.mu.Lock()
|
||||||
c.routes[action] = handler
|
c.routes[action] = handler
|
||||||
c.mu.Unlock()
|
c.mu.Unlock()
|
||||||
}
|
}
|
||||||
|
|
||||||
// Send 发送结构化消息。并发安全。
|
// Send 发送结构化消息。并发安全。
|
||||||
|
// 与服务端同用 Message 信封:{"action":"...","payload":...},payload 为 nil 时字段省略。
|
||||||
func (c *Client) Send(action string, payload any) error {
|
func (c *Client) Send(action string, payload any) error {
|
||||||
m := map[string]any{"action": action}
|
msg := Message{Action: action}
|
||||||
if payload != nil {
|
if payload != nil {
|
||||||
m["payload"] = payload
|
raw, err := json.Marshal(payload)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
msg.Payload = raw
|
||||||
}
|
}
|
||||||
data, err := json.Marshal(m)
|
data, err := json.Marshal(msg)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
@@ -223,16 +229,16 @@ func (c *Client) Send(action string, payload any) error {
|
|||||||
// SendRaw 发送已序列化的字节。并发安全。
|
// SendRaw 发送已序列化的字节。并发安全。
|
||||||
func (c *Client) SendRaw(data []byte) error {
|
func (c *Client) SendRaw(data []byte) error {
|
||||||
if c.closed.Load() {
|
if c.closed.Load() {
|
||||||
return errors.New("wscclient: client closed")
|
return errors.New("wsc: client closed")
|
||||||
}
|
}
|
||||||
select {
|
select {
|
||||||
case c.writeCh <- data:
|
case c.writeCh <- data:
|
||||||
return nil
|
return nil
|
||||||
case <-c.closeCh:
|
case <-c.closeCh:
|
||||||
return errors.New("wscclient: client closed")
|
return errors.New("wsc: client closed")
|
||||||
default:
|
default:
|
||||||
logrus.Warnf("[wscclient] send buffer full, dropping message")
|
logrus.Warnf("[wsc-client] send buffer full, dropping message")
|
||||||
return errors.New("wscclient: send buffer full")
|
return errors.New("wsc: send buffer full")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -240,29 +246,29 @@ func (c *Client) SendRaw(data []byte) error {
|
|||||||
// Bind — 自动 JSON 反序列化
|
// Bind — 自动 JSON 反序列化
|
||||||
// ============================================================
|
// ============================================================
|
||||||
|
|
||||||
// Bind 将带类型的函数包装为 Handler,自动 JSON 反序列化 payload。
|
// BindPayload 将带类型的函数包装为 PayloadHandler,自动 JSON 反序列化 payload。
|
||||||
//
|
//
|
||||||
// cli.On("ping.resp", wscclient.Bind(func(resp *PingResp) {
|
// cli.On("ping.resp", wsc.BindPayload(func(resp *PingResp) {
|
||||||
// log.Println(resp.Message)
|
// log.Println(resp.Message)
|
||||||
// }))
|
// }))
|
||||||
func Bind(fn any) Handler {
|
func BindPayload(fn any) PayloadHandler {
|
||||||
fnVal := reflect.ValueOf(fn)
|
fnVal := reflect.ValueOf(fn)
|
||||||
fnType := fnVal.Type()
|
fnType := fnVal.Type()
|
||||||
|
|
||||||
if fnType.Kind() != reflect.Func || fnType.NumIn() != 1 {
|
if fnType.Kind() != reflect.Func || fnType.NumIn() != 1 {
|
||||||
panic("wscclient.Bind: function must have 1 parameter")
|
panic("wsc.BindPayload: function must have 1 parameter")
|
||||||
}
|
}
|
||||||
|
|
||||||
reqType := fnType.In(0)
|
reqType := fnType.In(0)
|
||||||
if reqType.Kind() != reflect.Ptr {
|
if reqType.Kind() != reflect.Ptr {
|
||||||
panic("wscclient.Bind: parameter must be a pointer")
|
panic("wsc.BindPayload: parameter must be a pointer")
|
||||||
}
|
}
|
||||||
|
|
||||||
return func(raw json.RawMessage) {
|
return func(raw json.RawMessage) {
|
||||||
req := reflect.New(reqType.Elem()).Interface()
|
req := reflect.New(reqType.Elem()).Interface()
|
||||||
if len(raw) > 0 {
|
if len(raw) > 0 {
|
||||||
if err := json.Unmarshal(raw, req); err != nil {
|
if err := json.Unmarshal(raw, req); err != nil {
|
||||||
logrus.Errorf("[wscclient] Bind unmarshal error: %v, raw=%s", err, string(raw))
|
logrus.Errorf("[wsc-client] Bind unmarshal error: %v, raw=%s", err, string(raw))
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -286,7 +292,7 @@ func (c *Client) connect() error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if c.closed.Load() {
|
if c.closed.Load() {
|
||||||
return errors.New("wscclient: closed")
|
return errors.New("wsc: closed")
|
||||||
}
|
}
|
||||||
|
|
||||||
c.reconCount++
|
c.reconCount++
|
||||||
@@ -295,7 +301,7 @@ func (c *Client) connect() error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
delay := c.nextDelay()
|
delay := c.nextDelay()
|
||||||
logrus.Warnf("[wscclient] connect failed (attempt %d): %v, retry in %v", c.reconCount, err, delay)
|
logrus.Warnf("[wsc-client] connect failed (attempt %d): %v, retry in %v", c.reconCount, err, delay)
|
||||||
|
|
||||||
if c.OnReconnecting != nil {
|
if c.OnReconnecting != nil {
|
||||||
c.OnReconnecting(c.reconCount, delay)
|
c.OnReconnecting(c.reconCount, delay)
|
||||||
@@ -304,7 +310,7 @@ func (c *Client) connect() error {
|
|||||||
select {
|
select {
|
||||||
case <-time.After(delay):
|
case <-time.After(delay):
|
||||||
case <-c.closeCh:
|
case <-c.closeCh:
|
||||||
return errors.New("wscclient: closed during reconnect")
|
return errors.New("wsc: closed during reconnect")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -361,7 +367,7 @@ func (c *Client) readPump() {
|
|||||||
defer c.wg.Done()
|
defer c.wg.Done()
|
||||||
defer func() {
|
defer func() {
|
||||||
if r := recover(); r != nil {
|
if r := recover(); r != nil {
|
||||||
logrus.Errorf("[wscclient] readPump panic: %v", r)
|
logrus.Errorf("[wsc-client] readPump panic: %v", r)
|
||||||
}
|
}
|
||||||
c.onDisconnect()
|
c.onDisconnect()
|
||||||
}()
|
}()
|
||||||
@@ -385,16 +391,13 @@ func (c *Client) readPump() {
|
|||||||
if !c.closed.Load() {
|
if !c.closed.Load() {
|
||||||
if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway, websocket.CloseNormalClosure) &&
|
if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway, websocket.CloseNormalClosure) &&
|
||||||
!errors.Is(err, context.DeadlineExceeded) {
|
!errors.Is(err, context.DeadlineExceeded) {
|
||||||
logrus.Errorf("[wscclient] read error: %v", err)
|
logrus.Errorf("[wsc-client] read error: %v", err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
var msg struct {
|
var msg Message
|
||||||
Action string `json:"action"`
|
|
||||||
Payload json.RawMessage `json:"payload,omitempty"`
|
|
||||||
}
|
|
||||||
if err := json.Unmarshal(data, &msg); err != nil {
|
if err := json.Unmarshal(data, &msg); err != nil {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
@@ -406,7 +409,7 @@ func (c *Client) readPump() {
|
|||||||
if ok {
|
if ok {
|
||||||
handler(msg.Payload)
|
handler(msg.Payload)
|
||||||
} else {
|
} else {
|
||||||
logrus.Warnf("[wscclient] unhandled message action=%q, payload=%s", msg.Action, string(msg.Payload))
|
logrus.Warnf("[wsc-client] unhandled message action=%q, payload=%s", msg.Action, string(msg.Payload))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -433,7 +436,7 @@ func (c *Client) writePump() {
|
|||||||
}
|
}
|
||||||
conn.SetWriteDeadline(time.Now().Add(c.cfg.WriteTimeout))
|
conn.SetWriteDeadline(time.Now().Add(c.cfg.WriteTimeout))
|
||||||
if err := conn.WriteMessage(websocket.TextMessage, data); err != nil {
|
if err := conn.WriteMessage(websocket.TextMessage, data); err != nil {
|
||||||
logrus.Errorf("[wscclient] write error: %v", err)
|
logrus.Errorf("[wsc-client] write error: %v", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -515,13 +518,13 @@ func (c *Client) reconnectLoop() {
|
|||||||
|
|
||||||
c.reconCount++
|
c.reconCount++
|
||||||
if c.cfg.MaxRetry > 0 && c.reconCount > c.cfg.MaxRetry {
|
if c.cfg.MaxRetry > 0 && c.reconCount > c.cfg.MaxRetry {
|
||||||
logrus.Errorf("[wscclient] reconnect max retry exceeded (%d)", c.cfg.MaxRetry)
|
logrus.Errorf("[wsc-client] reconnect max retry exceeded (%d)", c.cfg.MaxRetry)
|
||||||
c.Close()
|
c.Close()
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
delay := c.nextDelay()
|
delay := c.nextDelay()
|
||||||
logrus.Warnf("[wscclient] reconnect attempt %d failed: %v, retry in %v", c.reconCount, err, delay)
|
logrus.Warnf("[wsc-client] reconnect attempt %d failed: %v, retry in %v", c.reconCount, err, delay)
|
||||||
|
|
||||||
if c.OnReconnecting != nil {
|
if c.OnReconnecting != nil {
|
||||||
c.OnReconnecting(c.reconCount, delay)
|
c.OnReconnecting(c.reconCount, delay)
|
||||||
Reference in New Issue
Block a user