# websocket-module **Repository Path**: dreamwithouttrace/websocket-module ## Basic Information - **Project Name**: websocket-module - **Description**: 嘻嘻嘻嘻嘻嘻嘻嘻嘻嘻 - **Primary Language**: Unknown - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2025-08-03 - **Last Updated**: 2026-08-26 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # WebSocket Module [![Go Version](https://img.shields.io/badge/Go-1.24-blue)](https://go.dev/) [![License](https://img.shields.io/badge/License-MIT-green)](LICENSE) 这是一个高性能、可复用的实时通信服务器模块,支持 **WebSocket** 和 **TCP** 双协议接入,提供了完整的连接管理、消息处理、用户认证、gRPC 服务以及出站连接池等功能。 ## 特性 - **双协议**: WebSocket + TCP 统一接入,共享连接池 - **高性能**: 使用分片连接池管理大量并发连接 - **安全认证**: 支持 Token 认证(WebSocket HTTP 鉴权 / TCP 首帧鉴权) - **消息确认**: 支持消息 ACK 机制,确保消息可靠传递(最多4次指数退避重试) - **心跳保活**: TCP 基于 `tick` 消息心跳,WebSocket 原生 ping/pong + tick 双保险 - **出站连接池**: gRPC 和 TCP 客户端连接池,支持借还和直发两种模式 - **限流保护**: 内置请求限流机制 - **实时统计**: 提供在线用户统计和监控(按性别统计) - **完整 API**: 提供 RESTful API 进行服务器管理 - **gRPC 服务**: 内置 gRPC 服务,支持远程调用(Kick/Exists/Push/Update) ## 目录结构 ``` websocket-module/ ├── core/ # 核心接口和类型定义 │ ├── types.go # 核心类型(User、Connection、Config等) │ ├── base_conn.go # 连接公共基础结构体(ACK/心跳/消息处理) │ ├── connection.go # WebSocket 连接实现 │ └── tcp_connection.go # TCP 连接实现(长度前缀分帧) ├── transport/ # 传输层管理(分片连接池) │ └── transport.go # Transport + SharedSlice(统一管理 WS/TCP) ├── handlers/ # 处理器 │ ├── handshake.go # WebSocket 握手处理 │ ├── tcp_handler.go # TCP 监听器 + 认证 │ └── api.go # RESTful API 处理 ├── middleware/ # 中间件 │ └── middleware.go # 认证和限流中间件 ├── protocol/ # 消息协议定义 │ └── types.go # Message 格式定义 ├── clientpool/ # 出站连接池 │ ├── config.go # 统一连接池配置 │ ├── client_connection.go # ClientConnection 统一连接 │ └── channel_pool.go # ClientPool 连接池实现 ├── service/ # gRPC 服务定义 │ ├── handler.proto # Proto 定义 │ └── service.go # gRPC 服务实现 ├── monitoring/ # 监控配置 ├── examples/ # 使用示例 ├── server.go # 主服务器入口 ├── go.mod └── README.md ``` ## 快速开始 ### 1. 基本使用 ```go package main import ( "log" "gitee.com/dreamwithouttrace/websocket-module/core" "gitee.com/dreamwithouttrace/websocket-module/protocol" "gitee.com/dreamwithouttrace/websocket-module" "github.com/go-redis/redis" ) // User 实现 core.User 接口 type User struct { ID int64 UniqueID string Nickname string Avatar string Sex string DisabledEndTime int64 AppID int64 } func (u *User) GetID() int64 { return u.ID } func (u *User) GetUniqueID() string { return u.UniqueID } func (u *User) GetNickname() string { return u.Nickname } func (u *User) GetAvatar() string { return u.Avatar } func (u *User) GetSex() string { return u.Sex } func (u *User) GetIsApp() int { return 0 } func (u *User) GetIsWx() int { return 0 } func (u *User) LifecycleDelay() {} func (u *User) GetDisabledEndTime() int64 { return u.DisabledEndTime } func (u *User) GetAppID() int64 { return u.AppID } func main() { config := core.DefaultConfig() config.Host = "0.0.0.0" config.Port = 8080 config.TCPHost = "0.0.0.0" config.TCPPort = 9000 config.GrpcHost = "0.0.0.0" config.GrpcPort = 8001 redisClient := redis.NewClient(&redis.Options{ Addr: "localhost:6379", }) queryFunc := func(token string) (core.User, error) { return &User{ID: 1, UniqueID: "user_1", Nickname: "测试用户", Sex: "1"}, nil } // 方式1: 使用默认 Redis token 校验 server := websocket.NewServer(config, redisClient, queryFunc) // 方式2: 自定义 token 校验回调 customValidator := func(token string) error { // 自定义校验逻辑,返回 nil 表示通过 return nil } server = websocket.NewServer(config, redisClient, queryFunc, websocket.WithTokenValidator(customValidator), ) server.SetOnConnect(func(conn core.Connection) error { log.Printf("用户连接: %s", conn.GetUserID()) return nil }) server.SetOnMessage(func(conn core.Connection, user core.User, role interface{}, message *protocol.Message) { log.Printf("收到消息: 类型=%s, 用户=%d", message.Type, message.UserId) }) server.SetOnDisconnect(func(conn core.Connection, user core.User, role interface{}) { log.Printf("用户断开: %s", conn.GetUserID()) }) if err := server.Run(); err != nil { log.Fatalf("服务器错误: %v", err) } } ``` ## TCP 接入 ### 协议 TCP 使用 **4 字节 Big-Endian 长度前缀 + JSON 载荷** 的分帧格式: ``` [4字节长度(大端)] [JSON消息体] ``` ### 认证流程 1. 客户端连接 TCP 端口 2. 客户端发送认证帧:`{"type":"auth","data":{"token":"xxx"}}` 3. 服务端验证后回复:`{"type":"auth_ok","channel":"","user_id":}` 4. 认证失败回复 `{"type":"auth_fail","data":{"reason":"..."}}` 并断开 ### 心跳 TCP 客户端需定期发送 `{"type":"tick"}` 消息,服务端 30 秒检测一次,1 分钟超时则踢出。 ## API 接口 ### WebSocket 连接 - `GET /ws?token=xxx` - WebSocket 连接(需认证) ### TCP 连接 - 连接 `TCPHost:TCPPort`,首帧发认证消息 ### 管理 API | 方法 | 路径 | 描述 | |------|------|------| | POST | `/kick?userId=123` | 踢出指定用户 | | GET | `/online/lists` | 获取在线用户列表 | | GET | `/get/online?id=123` | 获取指定用户信息 | | GET | `/check/content?keyword=xxx` | 内容检查(GET) | | POST | `/check/content` | 内容检查(POST) | | GET | `/statistics` | 获取在线统计 | ## gRPC 服务 | 方法 | 描述 | |------|------| | `Kick(KickRequest) returns (KickResponse)` | 踢出用户 | | `Exists(ExistsRequest) returns (ExistsResponse)` | 检查是否在线 | | `Push(MessageRequest) returns (MessageResponse)` | 推送消息 | | `Update(UpdateRequest) returns (UpdateResponse)` | 强制用户重连 | ## 出站连接池 `clientpool` 包提供统一客户端连接池,用于服务端向 TCP 或 gRPC 下游服务发起调用。 ### 快速使用 调用方只需要传目标地址和消息内容,连接池由 `clientpool` 自动维护: ```go // 发送协议消息 err := clientpool.Send("192.168.1.10:9000", msg) // 发送原始字节 err = clientpool.SendRaw("192.168.1.10:9000", data) // 发送普通内容,内部会包装为 protocol.Message err = clientpool.SendMessage("192.168.1.10:9000", map[string]interface{}{"text": "hello"}) // 请求-响应模式 resp, err := clientpool.Request("192.168.1.10:9000", msg, 5*time.Second) ``` 同一个 `ip:port` 会自动复用同一个内部连接池;外部无需创建、保存或关闭连接池。 ## 核心接口 ### User 接口 ```go type User interface { GetID() int64 GetUniqueID() string GetNickname() string GetAvatar() string GetSex() string // "1"=男, "2"=女, 其他=无性别 GetIsApp() int GetIsWx() int LifecycleDelay() GetDisabledEndTime() int64 GetAppID() int64 } ``` ### Role 接口 ```go type Role interface { GetId() int64 } ``` 连接设置角色后,可以通过 `server.GetRole(role.GetId())` 获取连接,或通过 `server.SendRole(role.GetId(), msg)` 给指定角色推送消息。 断线或被 `Kick` 后,如果用户在 `ReconnectTimeout` 内重连,服务端会复用原 `TCPConnection`/`WebSocketConnection` 对象,并保留 `Waits`、`Role` 等临时状态。 ### Connection 接口 ```go type Connection interface { GetUser() User SetUser(user User) GetChannelID() string GetUserID() string GetIP() string WriteMessage(data *protocol.Message) error WriteByteMessage(data []byte) error Close() error IsConnected() bool SetConnected(connected bool) GetWaits() *Waits GetSendNo() int64 IncrementSendNo() Kick() Ping() GetRole() Role SetRole(role Role) Pong() GetUniqueID() string SetReadDeadline(second int) error ConnectionDurationThisTime() *ClientInfo SetCallbacks(onDisconnect, onMessage, onConnect, onDeliveryFailure, onBeforeMessageLoop) Start(duration time.Duration) } ``` ## 消息格式 ```go type Message struct { Type string `json:"type"` ID string `json:"id"` Timestamp int64 `json:"timestamp"` IsGroup bool `json:"is_group"` NeedAck bool `json:"need_ack"` UserId int64 `json:"user_id"` Channel string `json:"channel"` SendInfo map[string]interface{} `json:"send_info,omitempty"` Data interface{} `json:"data"` } ``` ## 配置说明 | 配置项 | 类型 | 默认值 | 说明 | |--------|------|--------|------| | Host | string | "127.0.0.1" | HTTP/WS 监听地址 | | Port | int | 14101 | HTTP/WS 监听端口 | | TCPHost | string | - | TCP 监听地址(为空则不启动 TCP) | | TCPPort | int | - | TCP 监听端口(为0则不启动) | | GrpcHost | string | - | gRPC 监听地址 | | GrpcPort | int | - | gRPC 监听端口 | | TokenName | string | - | Token 参数名 | | ReadTimeout | Duration | 2s | 读取超时 | | WriteTimeout | Duration | 2s | 写入超时 | | SharedCount | int64 | 1000 | 连接分片数 | | EnableCompression | bool | true | WebSocket 压缩 | | ReconnectTimeout | Duration | 30s | 断线或 Kick 后可复用连接对象的重连窗口 | | ReadBufferSize | int | 4096 | WS 读缓冲区 | | WriteBufferSize | int | 4096 | WS 写缓冲区 | ## 依赖项 - `github.com/gorilla/websocket` - WebSocket 实现 - `github.com/gorilla/mux` - HTTP 路由 - `github.com/go-redis/redis` - Redis 客户端 - `google.golang.org/grpc` - gRPC 框架 - `google.golang.org/protobuf` - Protocol Buffers ## 注意事项 1. **TCP 端口**: `TCPPort` 为 0 时不启动 TCP 服务,保持向后兼容 2. **共享连接池**: WebSocket 和 TCP 连接共用 `Transport` 分片池,通过 `Connection` 接口统一管理 3. **连接安全性**: 每个连接下发独立的 `channelID`,消息中 `channel` 不匹配会立即踢出 4. **Redis 依赖**: 认证依赖 Redis,确保 Redis 服务可用 5. **ACK 机制**: `need_ack: true` 的消息,客户端需回复 `ack:` 字节流来确认 ## 许可证 MIT License