增加注释

This commit is contained in:
2025-12-02 14:09:01 +08:00
parent d77f06ad9a
commit 71cc36c7a5
4 changed files with 250 additions and 38 deletions

View File

@@ -1,3 +1,8 @@
/**
* package controller
* 作用:处理核心业务逻辑,包括 WebSocket 连接管理、消息转发、Redis 集群通信以及 API 接口实现。
* 该包是连接网络层Gin与数据层Gorm/Redis的桥梁。
*/
package controller
import (
@@ -27,7 +32,15 @@ var upgrader = websocket.Upgrader{
},
}
// WebSocketController WebSocket控制器结构体
/**
* WebSocketController
* 结构体定义WebSocket 服务的主控制器。
* 包含:
* - 客户端连接池 (Clients)
* - 读写锁 (ClientsMux) 保护连接池安全
* - Redis 客户端 (RedisCli) 用于集群通信
* - 数据库连接 (DB)
*/
type WebSocketController struct {
Clients map[string]*websocket.Conn
ClientsMux sync.RWMutex
@@ -40,6 +53,11 @@ type WebSocketController struct {
DB *gorm.DB // 数据库连接
}
/**
* NewWebSocketController
* 功能:创建并初始化 WebSocketController 实例。
* @param db *gorm.DB 数据库连接实例
*/
func NewWebSocketController(db *gorm.DB) *WebSocketController {
return &WebSocketController{
Clients: make(map[string]*websocket.Conn),
@@ -48,7 +66,11 @@ func NewWebSocketController(db *gorm.DB) *WebSocketController {
}
}
// ConfigureSystem 配置系统参数
/**
* ConfigureSystem
* 功能:加载系统配置,如 NodeID 和端口,并设置日志格式。
* 也会启动节点注册协程。
*/
func (c *WebSocketController) ConfigureSystem() {
c.NodeID = utils.GetEnv("NODE_ID", "local")
@@ -69,7 +91,11 @@ func (c *WebSocketController) ConfigureSystem() {
go c.registerNode() // 现在 RedisCli 已初始化
}
// 注册节点到Redis
/**
* registerNode
* 功能:定时将当前节点的 IP 和端口注册到 Redis 中,以便其他服务发现。
* 这是一个心跳机制,每 20 秒执行一次。
*/
func (c *WebSocketController) registerNode() {
// 添加空值检查
if c.RedisCli == nil {
@@ -91,7 +117,10 @@ func (c *WebSocketController) registerNode() {
}
}
// 配置日志记录器
/**
* configureLogger
* 功能:初始化日志目录和写入器配置。
*/
func (c *WebSocketController) configureLogger() {
logDir := utils.GetEnv("LOG_DIR", "./logs")
if err := os.MkdirAll(logDir, 0755); err != nil {
@@ -107,12 +136,18 @@ func (c *WebSocketController) configureLogger() {
log.SetOutput(c.getLogWriter(logDir, logFreq))
}
// 获取日志文件写入器
/**
* getLogWriter
* 功能:工厂方法,获取 DailyFileWriter 实例。
*/
func (c *WebSocketController) getLogWriter(logDir, freq string) *utils.DailyFileWriter {
return utils.NewDailyFileWriter(logDir, freq, c.Logger, c.NodeID)
}
// InitRedisClient 初始化Redis客户端
/**
* InitRedisClient
* 功能:建立 Redis 连接并验证连通性。
*/
func (c *WebSocketController) InitRedisClient() {
redisAddr := utils.GetEnv("REDIS_ADDR", "localhost:6379")
redisPassword := utils.GetEnv("REDIS_PASSWORD", "")
@@ -133,7 +168,10 @@ func (c *WebSocketController) checkRedisConnection() error {
return err
}
// PrintStartupInfo 打印启动信息
/**
* PrintStartupInfo
* 功能在控制台打印详细的服务启动信息包括节点ID、端口、Redis状态及支持的消息类型。
*/
func (c *WebSocketController) PrintStartupInfo() {
hostname, _ := os.Hostname()
log.Printf("🚀🚀🚀🚀🚀🚀🚀🚀 WebSocket服务启动: NodeID=%s", c.NodeID)
@@ -148,7 +186,11 @@ func (c *WebSocketController) PrintStartupInfo() {
log.Println("🔗🔗🔗🔗🔗🔗🔗🔗 等待客户端连接...")
}
// HealthHandler 健康检查处理器
/**
* HealthHandler
* 功能API 接口,用于健康检查。
* 路径GET /api/health
*/
func (c *WebSocketController) HealthHandler(ctx *gin.Context) {
ctx.JSON(http.StatusOK, gin.H{
"status": "ok",
@@ -158,7 +200,11 @@ func (c *WebSocketController) HealthHandler(ctx *gin.Context) {
})
}
// IsUserOnline 检查某个用户ID是否在线
/**
* IsUserOnline
* 功能API 接口,检查指定用户当前是否在线(通过 Redis 查询)。
* 路径GET /api/check-user-online
*/
func (c *WebSocketController) IsUserOnline(ctx *gin.Context) {
// 优先尝试查询参数,再尝试路径参数
userID := ctx.Query("user_id")
@@ -200,7 +246,11 @@ func (c *WebSocketController) IsUserOnline(ctx *gin.Context) {
})
}
// HandleWebSocket WebSocket连接处理器
/**
* HandleWebSocket
* 功能:处理 WebSocket 握手升级,并开启消息监听循环。
* 路径GET /ws
*/
func (c *WebSocketController) HandleWebSocket(ctx *gin.Context) {
start := time.Now()
clientIP := ctx.ClientIP()
@@ -248,7 +298,12 @@ func (c *WebSocketController) HandleWebSocket(ctx *gin.Context) {
}
}
// 处理客户端消息
/**
* handleClientMessage
* 功能:处理从客户端接收到的消息,支持绑定、发消息、通话信令等操作。
* @param senderID string 发送者的客户端ID
* @param message []byte 原始消息体
*/
func (c *WebSocketController) handleClientMessage(senderID string, message []byte) {
log.Printf("📥📥📥📥📥📥📥📥 收到客户端消息: \n SenderID=%s \n Size=%d bytes", senderID, len(message))
@@ -361,7 +416,10 @@ func (c *WebSocketController) handleClientMessage(senderID string, message []byt
log.Printf("⚠️ 无法识别的消息格式: \n Size=%d bytes \n Message=%s", len(message), string(message))
}
// 处理通话邀请
/**
* handleCallInvite
* 功能处理WebRTC通话邀请信令。
*/
func (c *WebSocketController) handleCallInvite(senderID, senderUserID string, signal models.CallSignal) {
log.Printf("📞 处理通话邀请: \n CallID=%s \n CallType=%d \n From=%s \n To=%s",
signal.CallID, signal.CallType, senderUserID, signal.CalleeID)
@@ -385,7 +443,10 @@ func (c *WebSocketController) handleCallInvite(senderID, senderUserID string, si
}
}
// 处理通话接受
/**
* handleCallAccept
* 功能:处理被叫方接受通话的信令。
*/
func (c *WebSocketController) handleCallAccept(senderID, senderUserID string, signal models.CallSignal) {
log.Printf("✅ 处理通话接受: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CallerID)
@@ -408,7 +469,10 @@ func (c *WebSocketController) handleCallAccept(senderID, senderUserID string, si
}
}
// 处理通话拒绝
/**
* handleCallReject
* 功能:处理被叫方拒绝通话的信令。
*/
func (c *WebSocketController) handleCallReject(senderID, senderUserID string, signal models.CallSignal) {
log.Printf("❌ 处理通话拒绝: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CallerID)
@@ -431,7 +495,10 @@ func (c *WebSocketController) handleCallReject(senderID, senderUserID string, si
}
}
// 处理通话结束
/**
* handleCallEnd
* 功能:处理通话结束/挂断信令。
*/
func (c *WebSocketController) handleCallEnd(senderID, senderUserID string, signal models.CallSignal) {
log.Printf("📞 处理通话结束: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CalleeID)
@@ -454,7 +521,10 @@ func (c *WebSocketController) handleCallEnd(senderID, senderUserID string, signa
}
}
// 处理Offer信令
/**
* handleCallOffer
* 功能:处理 WebRTC 的 Offer 信令SDP交换
*/
func (c *WebSocketController) handleCallOffer(senderID, senderUserID string, signal models.CallSignal) {
log.Printf("📞 处理Offer信令: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CalleeID)
@@ -477,7 +547,10 @@ func (c *WebSocketController) handleCallOffer(senderID, senderUserID string, sig
}
}
// 处理Answer信令
/**
* handleCallAnswer
* 功能:处理 WebRTC 的 Answer 信令SDP交换
*/
func (c *WebSocketController) handleCallAnswer(senderID, senderUserID string, signal models.CallSignal) {
log.Printf("📞 处理Answer信令: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CallerID)
@@ -500,7 +573,10 @@ func (c *WebSocketController) handleCallAnswer(senderID, senderUserID string, si
}
}
// 处理通话候选
/**
* handleCallCandidate
* 功能:处理 ICE Candidate 信令(网络协商)。
*/
func (c *WebSocketController) handleCallCandidate(senderID, senderUserID string, signal models.CallSignal) {
log.Printf("📶 处理通话候选: \n CallID=%s From=%s \n To=%s", signal.CallID, senderUserID, signal.CalleeID)
@@ -523,7 +599,10 @@ func (c *WebSocketController) handleCallCandidate(senderID, senderUserID string,
}
}
// 处理通话挂断
/**
* handleCallHangup
* 功能:处理主动挂断通话。
*/
func (c *WebSocketController) handleCallHangup(senderID, senderUserID string, signal models.CallSignal) {
log.Printf("📞 处理通话挂断: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CalleeID)
@@ -546,7 +625,10 @@ func (c *WebSocketController) handleCallHangup(senderID, senderUserID string, si
}
}
// 处理通话掉线
/**
* handleCallDisconnected
* 功能:处理连接意外断开。
*/
func (c *WebSocketController) handleCallDisconnected(senderID, senderUserID string, signal models.CallSignal) {
log.Printf("📞 处理通话掉线: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CalleeID)
@@ -569,7 +651,10 @@ func (c *WebSocketController) handleCallDisconnected(senderID, senderUserID stri
}
}
// 处理通话终止
/**
* handleCallTerminated
* 功能:处理通话终止信号。
*/
func (c *WebSocketController) handleCallTerminated(senderID, senderUserID string, signal models.CallSignal) {
log.Printf("📞 处理通话终止: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CalleeID)
@@ -592,7 +677,10 @@ func (c *WebSocketController) handleCallTerminated(senderID, senderUserID string
}
}
// 处理无人接听
/**
* handleCallNoAnswer
* 功能:处理无人接听状态。
*/
func (c *WebSocketController) handleCallNoAnswer(senderID, senderUserID string, signal models.CallSignal) {
log.Printf("📞 处理无人接听: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CallerID)
@@ -615,7 +703,10 @@ func (c *WebSocketController) handleCallNoAnswer(senderID, senderUserID string,
}
}
// 处理忙线状态
/**
* handleCallBusy
* 功能:处理忙线状态。
*/
func (c *WebSocketController) handleCallBusy(senderID, senderUserID string, signal models.CallSignal) {
log.Printf("📞 处理忙线状态: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CallerID)
@@ -638,7 +729,10 @@ func (c *WebSocketController) handleCallBusy(senderID, senderUserID string, sign
}
}
// 处理通话失败
/**
* handleCallFailed
* 功能:处理通话建立失败状态。
*/
func (c *WebSocketController) handleCallFailed(senderID, senderUserID string, signal models.CallSignal) {
log.Printf("📞 处理通话失败: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CalleeID)
@@ -661,7 +755,11 @@ func (c *WebSocketController) handleCallFailed(senderID, senderUserID string, si
}
}
// SendMessageHandler API消息发送处理器
/**
* SendMessageHandler
* 功能API 接口用于服务器端主动发送消息HTTP -> WebSocket
* 路径POST /api/send
*/
func (c *WebSocketController) SendMessageHandler(ctx *gin.Context) {
start := time.Now()
@@ -724,7 +822,11 @@ func (c *WebSocketController) SendMessageHandler(ctx *gin.Context) {
ctx.JSON(http.StatusBadRequest, gin.H{"error": "不支持的request_type"})
}
// 用户绑定处理器
/**
* BindHandler
* 功能API 接口,用于强制绑定 ClientID 和 UserID。
* 路径POST /api/bind
*/
func (c *WebSocketController) BindHandler(ctx *gin.Context) {
start := time.Now()
@@ -758,7 +860,11 @@ func (c *WebSocketController) BindHandler(ctx *gin.Context) {
req.UserID, req.ClientID, time.Since(start))
}
// SendToUserHandler 通过用户ID发送消息处理器支持通话信令
/**
* SendToUserHandler
* 功能API 接口,向指定用户的所有在线设备发送消息,并持久化到数据库。
* 路径POST /api/send-to-user
*/
// SendToUserHandler 通过用户ID发送消息处理器支持通话信令
func (c *WebSocketController) SendToUserHandler(ctx *gin.Context) {
start := time.Now()
@@ -841,7 +947,11 @@ func (c *WebSocketController) SendToUserHandler(ctx *gin.Context) {
log.Printf("✅ API用户发送请求完成: Duration=%s", time.Since(start))
}
// 获取聊天记录(分页)
/**
* GetMessagesHandler
* 功能API 接口,分页查询历史聊天记录。
* 路径GET /api/messages
*/
func (c *WebSocketController) GetMessagesHandler(ctx *gin.Context) {
// 解析查询参数
roomId := ctx.Query("room_id")
@@ -892,7 +1002,11 @@ func (c *WebSocketController) GetMessagesHandler(ctx *gin.Context) {
})
}
// 同步聊天记录(全部)
/**
* SyncMessagesHandler
* 功能API 接口,获取全部聊天记录(用于同步)。
* 路径GET /api/messages/sync
*/
func (c *WebSocketController) SyncMessagesHandler(ctx *gin.Context) {
roomId := ctx.Query("room_id")
senderUserId := ctx.Query("sender_user_id")
@@ -921,7 +1035,11 @@ func (c *WebSocketController) SyncMessagesHandler(ctx *gin.Context) {
})
}
// 添加客户端连接
/**
* addClient
* 功能:将新的 WebSocket 连接添加到本地连接池。
* 注意:使用 Mutex 锁保证线程安全。
*/
func (c *WebSocketController) addClient(clientID string, conn *websocket.Conn) {
c.ClientsMux.Lock()
defer c.ClientsMux.Unlock()
@@ -929,7 +1047,10 @@ func (c *WebSocketController) addClient(clientID string, conn *websocket.Conn) {
log.Printf("📊📊📊📊📊📊📊📊 当前连接数: %d", len(c.Clients))
}
// 移除客户端连接
/**
* removeClient
* 功能:从本地连接池移除客户端,并清理 Redis 中的绑定关系。
*/
func (c *WebSocketController) removeClient(clientID string) {
c.ClientsMux.Lock()
defer c.ClientsMux.Unlock()
@@ -943,7 +1064,10 @@ func (c *WebSocketController) removeClient(clientID string) {
}
}
// 清理用户绑定关系
/**
* cleanupUserBinding
* 功能:从 Redis 中清除用户与客户端ID的映射关系。
*/
func (c *WebSocketController) cleanupUserBinding(clientID string) {
// 获取用户ID
userID, err := c.getUserIDByClientID(clientID)
@@ -962,7 +1086,12 @@ func (c *WebSocketController) cleanupUserBinding(clientID string) {
log.Printf("🧹🧹🧹🧹🧹🧹🧹🧹 清理绑定关系: ClientID=%s | UserID=%s", clientID, userID)
}
// 绑定客户端到用户
/**
* bindClientToUser
* 功能:在 Redis 中建立 UserID 和 ClientID 的双向映射。
* 映射1client_user_mapping:clientID -> userID (string)
* 映射2user_client_mapping:userID -> [clientID1, clientID2] (set)
*/
func (c *WebSocketController) bindClientToUser(clientID, userID string) error {
// 设置客户端->用户映射
clientUserKey := fmt.Sprintf("%s:%s", models.ClientUserKey, clientID)
@@ -1000,7 +1129,11 @@ func (c *WebSocketController) getClientIDsByUserID(userID string) ([]string, err
return clientIDs, nil
}
// 发送消息给指定客户端
/**
* sendMessageToClient
* 功能:向指定 ClientID 发送消息。
* 逻辑:如果 Client 在本节点,直接通过 WebSocket 发送;如果不在,尝试通过 Redis 广播转发。
*/
func (c *WebSocketController) sendMessageToClient(clientID string, message models.ClientReceivedMessage) error {
c.ClientsMux.RLock()
conn, exists := c.Clients[clientID]
@@ -1025,7 +1158,10 @@ func (c *WebSocketController) sendMessageToClient(clientID string, message model
return nil
}
// 发送消息给指定用户
/**
* sendMessageToUser
* 功能:向指定用户的所有设备发送消息。
*/
func (c *WebSocketController) sendMessageToUser(userID string, message models.ClientReceivedMessage) error {
clientIDs, err := c.getClientIDsByUserID(userID)
if err != nil {
@@ -1061,7 +1197,10 @@ func (c *WebSocketController) sendMessageToUser(userID string, message models.Cl
return nil
}
// 转发消息到其他节点
/**
* forwardMessageToOtherNodes
* 功能:当目标客户端不在本机时,通过 Redis Pub/Sub 将消息广播到集群其他节点。
*/
func (c *WebSocketController) forwardMessageToOtherNodes(clientID string, message models.ClientReceivedMessage) error {
redisMsg := models.RedisMessage{
SenderNodeID: c.NodeID,
@@ -1089,7 +1228,10 @@ func (c *WebSocketController) forwardMessageToOtherNodes(clientID string, messag
return nil
}
// 启动Redis消息监听
/**
* SubscribeToRedis
* 功能:订阅 Redis 广播频道,监听来自其他节点的消息转发请求。
*/
func (c *WebSocketController) SubscribeToRedis() {
go func() {
pattern := "websocket:forward:*"