feat(pool): 添加基于BadgerDB的连接池实现
- 新增BadgerPool结构体,支持WebSocket和TCP连接类型 - 实现连接的增删改查功能,包括内存缓存机制提升性能 - 添加按类型查询连接、统计连接数量等辅助方法 - 实现清理非活跃连接的功能,支持定期维护 - 更新示例代码以处理初始化错误并改进错误处理 - 添加BadgerDB依赖及其相关间接依赖包
This commit is contained in:
@@ -19,7 +19,10 @@ func NewWs() *Manager {
|
||||
}
|
||||
|
||||
// 2. 创建管理器
|
||||
m := NewManager(customConfig)
|
||||
m, err := NewManager(customConfig)
|
||||
if err != nil {
|
||||
log.Fatalf("Failed to create manager: %v", err)
|
||||
}
|
||||
|
||||
// 3. 覆盖业务回调(核心:自定义消息处理逻辑)
|
||||
// 连接建立回调
|
||||
|
||||
+78
-3
@@ -9,6 +9,7 @@ import (
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"git.magicany.cc/black1552/gf-common/pool"
|
||||
"github.com/gogf/gf/v2/encoding/gjson"
|
||||
"github.com/gogf/gf/v2/os/gctx"
|
||||
"github.com/gogf/gf/v2/os/gtime"
|
||||
@@ -92,7 +93,8 @@ type Connection struct {
|
||||
type Manager struct {
|
||||
config *Config // 配置
|
||||
upgrader *websocket.Upgrader // HTTP升级器
|
||||
connections map[string]*Connection // 所有在线连接(connID -> Connection)
|
||||
connections map[string]*Connection // 内存中的连接(connID -> Connection)
|
||||
badgerPool *pool.BadgerPool // BadgerDB连接池
|
||||
mutex sync.RWMutex // 读写锁(保护connections)
|
||||
// 业务回调:收到消息时触发(用户自定义处理逻辑)
|
||||
OnMessage func(connID string, msgType int, data any)
|
||||
@@ -148,7 +150,7 @@ func (c *Config) Merge(other *Config) *Config {
|
||||
}
|
||||
|
||||
// NewManager 创建连接管理器
|
||||
func NewManager(config *Config) *Manager {
|
||||
func NewManager(config *Config) (*Manager, error) {
|
||||
defaultConfig := DefaultConfig()
|
||||
finalConfig := defaultConfig.Merge(config)
|
||||
// 初始化升级器
|
||||
@@ -170,10 +172,17 @@ func NewManager(config *Config) *Manager {
|
||||
},
|
||||
}
|
||||
|
||||
// 初始化BadgerDB连接池
|
||||
badgerPool, err := pool.NewBadgerPool()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to create badger pool: %w", err)
|
||||
}
|
||||
|
||||
return &Manager{
|
||||
config: finalConfig,
|
||||
upgrader: upgrader,
|
||||
connections: make(map[string]*Connection),
|
||||
badgerPool: badgerPool,
|
||||
mutex: sync.RWMutex{},
|
||||
// 默认回调(用户可覆盖)
|
||||
OnMessage: func(connID string, msgType int, data any) {
|
||||
@@ -185,7 +194,7 @@ func NewManager(config *Config) *Manager {
|
||||
OnDisconnect: func(connID string, err error) {
|
||||
log.Printf("[默认回调] 连接[%s]已关闭:%v", connID, err)
|
||||
},
|
||||
}
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Upgrade HTTP升级为WebSocket连接
|
||||
@@ -237,6 +246,24 @@ func (m *Manager) Upgrade(w http.ResponseWriter, r *http.Request, connID string)
|
||||
m.connections[connID] = wsConn
|
||||
m.mutex.Unlock()
|
||||
|
||||
// 存储到BadgerDB
|
||||
connInfo := &pool.ConnectionInfo{
|
||||
ID: connID,
|
||||
Type: pool.ConnTypeWebSocket,
|
||||
Address: r.RemoteAddr,
|
||||
IsActive: true,
|
||||
LastUsed: time.Now(),
|
||||
CreatedAt: time.Now(),
|
||||
Data: map[string]interface{}{
|
||||
"origin": r.Header.Get("Origin"),
|
||||
"userAgent": r.Header.Get("User-Agent"),
|
||||
},
|
||||
}
|
||||
if err := m.badgerPool.Add(connInfo); err != nil {
|
||||
log.Printf("[错误] 存储连接到BadgerDB失败:%v", err)
|
||||
// 不影响连接建立,仅记录错误
|
||||
}
|
||||
|
||||
// 触发连接建立回调
|
||||
m.OnConnect(connID)
|
||||
|
||||
@@ -282,6 +309,18 @@ func (c *Connection) ReadPump() {
|
||||
return
|
||||
}
|
||||
|
||||
// 更新最后使用时间
|
||||
now := time.Now()
|
||||
// 从BadgerDB获取连接信息并更新
|
||||
connInfo, err := c.manager.badgerPool.Get(c.connID)
|
||||
if err == nil && connInfo != nil {
|
||||
connInfo.LastUsed = now
|
||||
if err := c.manager.badgerPool.Update(connInfo); err != nil {
|
||||
log.Printf("[错误] 更新BadgerDB连接信息失败:%v", err)
|
||||
// 不影响消息处理,仅记录错误
|
||||
}
|
||||
}
|
||||
|
||||
// 尝试解析JSON格式的心跳消息(精准判断,替代包含判断)
|
||||
isHeartbeat := false
|
||||
// 先尝试解析为JSON对象
|
||||
@@ -369,6 +408,19 @@ func (c *Connection) Send(data []byte) error {
|
||||
if err != nil {
|
||||
return fmt.Errorf("发送消息失败:%w", err)
|
||||
}
|
||||
|
||||
// 更新最后使用时间
|
||||
now := time.Now()
|
||||
// 从BadgerDB获取连接信息并更新
|
||||
connInfo, err := c.manager.badgerPool.Get(c.connID)
|
||||
if err == nil && connInfo != nil {
|
||||
connInfo.LastUsed = now
|
||||
if err := c.manager.badgerPool.Update(connInfo); err != nil {
|
||||
log.Printf("[错误] 更新BadgerDB连接信息失败:%v", err)
|
||||
// 不影响消息发送,仅记录错误
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
}
|
||||
@@ -394,6 +446,12 @@ func (c *Connection) Close(err error) {
|
||||
delete(c.manager.connections, c.connID)
|
||||
c.manager.mutex.Unlock()
|
||||
|
||||
// 从BadgerDB移除
|
||||
if err := c.manager.badgerPool.Remove(c.connID); err != nil {
|
||||
log.Printf("[错误] 从BadgerDB移除连接失败:%v", err)
|
||||
// 不影响连接关闭,仅记录错误
|
||||
}
|
||||
|
||||
// 触发断开回调
|
||||
c.manager.OnDisconnect(c.connID, err)
|
||||
|
||||
@@ -462,12 +520,18 @@ func (m *Manager) GetAllConn() map[string]*Connection {
|
||||
return connCopy
|
||||
}
|
||||
|
||||
// GetConn 获取指定连接
|
||||
func (m *Manager) GetConn(connID string) *Connection {
|
||||
m.mutex.RLock()
|
||||
defer m.mutex.RUnlock()
|
||||
return m.connections[connID]
|
||||
}
|
||||
|
||||
// GetAllConnIDs 获取所有在线连接的ID列表
|
||||
func (m *Manager) GetAllConnIDs() ([]string, error) {
|
||||
return m.badgerPool.GetAllConnIDs()
|
||||
}
|
||||
|
||||
// CloseAll 关闭所有连接
|
||||
func (m *Manager) CloseAll() {
|
||||
m.mutex.RLock()
|
||||
@@ -486,3 +550,14 @@ func (m *Manager) CloseAll() {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Close 关闭管理器,清理资源
|
||||
func (m *Manager) Close() error {
|
||||
// 关闭所有连接
|
||||
m.CloseAll()
|
||||
// 关闭BadgerDB连接池
|
||||
if m.badgerPool != nil {
|
||||
return m.badgerPool.Close()
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user