1 Commits
6 changed files with 180 additions and 49 deletions
+11 -8
View File
@@ -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 一次)
+31
View File
@@ -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 内部校验。
+23
View File
@@ -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` 与当前时间比较)。
+27
View File
@@ -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 文件位置信息会不准。
+44
View File
@@ -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` 的公开方法均可并发调用;发送走带缓冲通道 + 单写泵。
+44 -41
View File
@@ -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)