From 1b9b4c5ed93b0b45e4095a873ce0d9cd09289305 Mon Sep 17 00:00:00 2001 From: liqi Date: Wed, 9 Jul 2025 09:35:06 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D=EF=BC=9A=201.=20=E6=8E=A5?= =?UTF-8?q?=E5=90=AC=E7=94=B5=E8=AF=9D=E4=BB=A5=E5=AE=9E=E7=8E=B0=202.=20?= =?UTF-8?q?=E5=90=8E=E7=AB=AF=E6=8E=A5=E5=8F=A3=E9=80=82=E9=85=8D=E5=AE=9E?= =?UTF-8?q?=E7=8E=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 现有问题: 1. rtc通话问题待修复 2. 点击接听会有挂断的报错,需要检查逻辑 3. 修复lastMessage的更新问题 --- controller/websocket_controller.go | 585 ++++++++++++++++++++--------- models/messages.go | 71 +++- 2 files changed, 466 insertions(+), 190 deletions(-) diff --git a/controller/websocket_controller.go b/controller/websocket_controller.go index fb6ef42..a64612e 100644 --- a/controller/websocket_controller.go +++ b/controller/websocket_controller.go @@ -269,7 +269,7 @@ func (c *WebSocketController) handleClientMessage(senderID string, message []byt return } - // 修复:正确解析通话信令并分发处理 + // 处理通话信令 var callSignal models.CallSignal if err := json.Unmarshal(message, &callSignal); err == nil && callSignal.CallStatus != "" { // 获取发送者用户ID @@ -282,14 +282,30 @@ func (c *WebSocketController) handleClientMessage(senderID string, message []byt switch callSignal.CallStatus { case "invite": c.handleCallInvite(senderID, senderUserID, callSignal) - case "accept": + case "accepted": c.handleCallAccept(senderID, senderUserID, callSignal) - case "reject": + case "rejected": c.handleCallReject(senderID, senderUserID, callSignal) - case "end": + case "ended": c.handleCallEnd(senderID, senderUserID, callSignal) + case "offer": + c.handleCallOffer(senderID, senderUserID, callSignal) + case "answer": + c.handleCallAnswer(senderID, senderUserID, callSignal) case "candidate": c.handleCallCandidate(senderID, senderUserID, callSignal) + case "hangup": + c.handleCallHangup(senderID, senderUserID, callSignal) + case "disconnected": + c.handleCallDisconnected(senderID, senderUserID, callSignal) + case "terminated": + c.handleCallTerminated(senderID, senderUserID, callSignal) + case "no-answer": + c.handleCallNoAnswer(senderID, senderUserID, callSignal) + case "busy": + c.handleCallBusy(senderID, senderUserID, callSignal) + case "failed": + c.handleCallFailed(senderID, senderUserID, callSignal) default: log.Printf("⚠️ 未知的通话动作: %s", callSignal.CallStatus) } @@ -307,11 +323,11 @@ func (c *WebSocketController) handleCallInvite(senderID, senderUserID string, si // 创建通话消息 callMsg := models.ClientReceivedMessage{ SenderID: senderID, - ReceiverID: "", // 通过用户ID发送 + ReceiverID: "", SenderUserID: senderUserID, ReceiverUserID: signal.CalleeID, - MessageType: signal.CallType, // 6=视频, 7=语音 - Content: signal.Data, // 包含SDP offer + MessageType: signal.CallType, + Content: signal.Data, SendTime: time.Now().Format(time.RFC3339), CallID: signal.CallID, CallStatus: "invite", @@ -330,11 +346,11 @@ func (c *WebSocketController) handleCallAccept(senderID, senderUserID string, si // 创建通话接受消息 acceptMsg := models.ClientReceivedMessage{ SenderID: senderID, - ReceiverID: "", // 通过用户ID发送 + ReceiverID: "", SenderUserID: senderUserID, ReceiverUserID: signal.CallerID, - MessageType: signal.CallType, // 6=视频, 7=语音 - Content: signal.Data, // 包含SDP answer + MessageType: signal.CallType, + Content: signal.Data, SendTime: time.Now().Format(time.RFC3339), CallID: signal.CallID, CallStatus: "accepted", @@ -353,11 +369,11 @@ func (c *WebSocketController) handleCallReject(senderID, senderUserID string, si // 创建通话拒绝消息 rejectMsg := models.ClientReceivedMessage{ SenderID: senderID, - ReceiverID: "", // 通过用户ID发送 + ReceiverID: "", SenderUserID: senderUserID, ReceiverUserID: signal.CallerID, - MessageType: signal.CallType, // 6=视频, 7=语音 - Content: signal.Data, // 可包含拒绝原因 + MessageType: signal.CallType, + Content: signal.Data, SendTime: time.Now().Format(time.RFC3339), CallID: signal.CallID, CallStatus: "rejected", @@ -376,11 +392,11 @@ func (c *WebSocketController) handleCallEnd(senderID, senderUserID string, signa // 创建通话结束消息 endMsg := models.ClientReceivedMessage{ SenderID: senderID, - ReceiverID: "", // 通过用户ID发送 + ReceiverID: "", SenderUserID: senderUserID, ReceiverUserID: signal.CalleeID, - MessageType: signal.CallType, // 6=视频, 7=语音 - Content: signal.Data, // 可包含结束原因 + MessageType: signal.CallType, + Content: signal.Data, SendTime: time.Now().Format(time.RFC3339), CallID: signal.CallID, CallStatus: "ended", @@ -392,6 +408,52 @@ func (c *WebSocketController) handleCallEnd(senderID, senderUserID string, signa } } +// 处理Offer信令 +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) + + // 创建Offer消息 + offerMsg := models.ClientReceivedMessage{ + SenderID: senderID, + ReceiverID: "", + SenderUserID: senderUserID, + ReceiverUserID: signal.CalleeID, + MessageType: signal.CallType, + Content: signal.Data, + SendTime: time.Now().Format(time.RFC3339), + CallID: signal.CallID, + CallStatus: "offer", + } + + // 发送Offer给对端 + if err := c.sendMessageToUser(signal.CalleeID, offerMsg); err != nil { + log.Printf("❌ 发送Offer失败: %v | CalleeID=%s", err, signal.CalleeID) + } +} + +// 处理Answer信令 +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) + + // 创建Answer消息 + answerMsg := models.ClientReceivedMessage{ + SenderID: senderID, + ReceiverID: "", + SenderUserID: senderUserID, + ReceiverUserID: signal.CallerID, + MessageType: signal.CallType, + Content: signal.Data, + SendTime: time.Now().Format(time.RFC3339), + CallID: signal.CallID, + CallStatus: "answer", + } + + // 发送Answer给对端 + if err := c.sendMessageToUser(signal.CallerID, answerMsg); err != nil { + log.Printf("❌ 发送Answer失败: %v | CallerID=%s", err, signal.CallerID) + } +} + // 处理通话候选 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) @@ -399,11 +461,11 @@ func (c *WebSocketController) handleCallCandidate(senderID, senderUserID string, // 创建候选消息 candidateMsg := models.ClientReceivedMessage{ SenderID: senderID, - ReceiverID: "", // 通过用户ID发送 + ReceiverID: "", SenderUserID: senderUserID, ReceiverUserID: signal.CalleeID, - MessageType: signal.CallType, // 6=视频, 7=语音 - Content: signal.Data, // 包含ICE候选 + MessageType: signal.CallType, + Content: signal.Data, SendTime: time.Now().Format(time.RFC3339), CallID: signal.CallID, CallStatus: "candidate", @@ -415,13 +477,151 @@ func (c *WebSocketController) handleCallCandidate(senderID, senderUserID string, } } +// 处理通话挂断 +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) + + // 创建挂断消息 + hangupMsg := models.ClientReceivedMessage{ + SenderID: senderID, + ReceiverID: "", + SenderUserID: senderUserID, + ReceiverUserID: signal.CalleeID, + MessageType: signal.CallType, + Content: signal.Data, + SendTime: time.Now().Format(time.RFC3339), + CallID: signal.CallID, + CallStatus: "hangup", + } + + // 发送挂断消息给对端 + if err := c.sendMessageToUser(signal.CalleeID, hangupMsg); err != nil { + log.Printf("❌ 发送通话挂断失败: %v | CalleeID=%s", err, signal.CalleeID) + } +} + +// 处理通话掉线 +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) + + // 创建掉线消息 + disconnectedMsg := models.ClientReceivedMessage{ + SenderID: senderID, + ReceiverID: "", + SenderUserID: senderUserID, + ReceiverUserID: signal.CalleeID, + MessageType: signal.CallType, + Content: signal.Data, + SendTime: time.Now().Format(time.RFC3339), + CallID: signal.CallID, + CallStatus: "disconnected", + } + + // 发送掉线消息给对端 + if err := c.sendMessageToUser(signal.CalleeID, disconnectedMsg); err != nil { + log.Printf("❌ 发送通话掉线失败: %v | CalleeID=%s", err, signal.CalleeID) + } +} + +// 处理通话终止 +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) + + // 创建终止消息 + terminatedMsg := models.ClientReceivedMessage{ + SenderID: senderID, + ReceiverID: "", + SenderUserID: senderUserID, + ReceiverUserID: signal.CalleeID, + MessageType: signal.CallType, + Content: signal.Data, + SendTime: time.Now().Format(time.RFC3339), + CallID: signal.CallID, + CallStatus: "terminated", + } + + // 发送终止消息给对端 + if err := c.sendMessageToUser(signal.CalleeID, terminatedMsg); err != nil { + log.Printf("❌ 发送通话终止失败: %v | CalleeID=%s", err, signal.CalleeID) + } +} + +// 处理无人接听 +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) + + // 创建无人接听消息 + noAnswerMsg := models.ClientReceivedMessage{ + SenderID: senderID, + ReceiverID: "", + SenderUserID: senderUserID, + ReceiverUserID: signal.CallerID, + MessageType: signal.CallType, + Content: signal.Data, + SendTime: time.Now().Format(time.RFC3339), + CallID: signal.CallID, + CallStatus: "no-answer", + } + + // 发送无人接听消息 + if err := c.sendMessageToUser(signal.CallerID, noAnswerMsg); err != nil { + log.Printf("❌ 发送无人接听失败: %v | CallerID=%s", err, signal.CallerID) + } +} + +// 处理忙线状态 +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) + + // 创建忙线消息 + busyMsg := models.ClientReceivedMessage{ + SenderID: senderID, + ReceiverID: "", + SenderUserID: senderUserID, + ReceiverUserID: signal.CallerID, + MessageType: signal.CallType, + Content: signal.Data, + SendTime: time.Now().Format(time.RFC3339), + CallID: signal.CallID, + CallStatus: "busy", + } + + // 发送忙线消息 + if err := c.sendMessageToUser(signal.CallerID, busyMsg); err != nil { + log.Printf("❌ 发送忙线状态失败: %v | CallerID=%s", err, signal.CallerID) + } +} + +// 处理通话失败 +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) + + // 创建失败消息 + failedMsg := models.ClientReceivedMessage{ + SenderID: senderID, + ReceiverID: "", + SenderUserID: senderUserID, + ReceiverUserID: signal.CalleeID, + MessageType: signal.CallType, + Content: signal.Data, + SendTime: time.Now().Format(time.RFC3339), + CallID: signal.CallID, + CallStatus: "failed", + } + + // 发送失败消息 + if err := c.sendMessageToUser(signal.CalleeID, failedMsg); err != nil { + log.Printf("❌ 发送通话失败消息失败: %v | CalleeID=%s", err, signal.CalleeID) + } +} + // SendMessageHandler API消息发送处理器 func (c *WebSocketController) SendMessageHandler(ctx *gin.Context) { start := time.Now() var payload models.SendMessagePayload if err := ctx.ShouldBindJSON(&payload); err != nil { - log.Printf("⚠️ 无效的JSON请求格式: %极简", err) + log.Printf("⚠️ 无效的JSON请求格式: %v", err) ctx.JSON(http.StatusBadRequest, gin.H{"error": "无效的JSON格式"}) return } @@ -545,232 +745,241 @@ func (c *WebSocketController) SendToUserHandler(ctx *gin.Context) { // 如果是通话信令消息,添加通话相关字段 if payload.MessageType == 6 || payload.MessageType == 7 { - // 直接使用内容作为通话信令数据 - clientMsg.CallStatus = "invite" // 默认设置为invite,API调用通常是发起通话 - clientMsg.CallID = utils.GenerateCallID() // 生成唯一的通话ID + clientMsg.CallStatus = payload.CallStatus + clientMsg.CallID = payload.CallID - log.Printf("📞📞 处理API通话信令: \n CallID=%s \n CallType=%d", - clientMsg.CallID, payload.MessageType) + log.Printf("📞📞 处理通话信令消息: \n CallStatus=%s \n CallID=%s", + payload.CallStatus, payload.CallID) } + // 发送消息给目标用户 if err := c.sendMessageToUser(payload.ReceiverUserID, clientMsg); err != nil { + log.Printf("❌ 发送消息失败: %v | ReceiverUserID=%s", err, payload.ReceiverUserID) ctx.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } ctx.JSON(http.StatusOK, gin.H{ - "code": 0, "status": "success", - "message": fmt.Sprintf("消息已发送至用户 %s", payload.ReceiverUserID), - "call_id": clientMsg.CallID, // 返回生成的call_id + "message": fmt.Sprintf("消息已发送到用户: %s", payload.ReceiverUserID), }) - log.Printf("✅ API请求完成: Duration=%s", time.Since(start)) + log.Printf("✅ API用户发送请求完成: Duration=%s", time.Since(start)) } -// 绑定客户端与用户 +// 添加客户端连接 +func (c *WebSocketController) addClient(clientID string, conn *websocket.Conn) { + c.ClientsMux.Lock() + defer c.ClientsMux.Unlock() + c.Clients[clientID] = conn + log.Printf("📊📊📊📊📊📊📊📊 当前连接数: %d", len(c.Clients)) +} + +// 移除客户端连接 +func (c *WebSocketController) removeClient(clientID string) { + c.ClientsMux.Lock() + defer c.ClientsMux.Unlock() + + if _, exists := c.Clients[clientID]; exists { + delete(c.Clients, clientID) + log.Printf("🔌🔌🔌🔌🔌🔌🔌🔌 客户端已断开: ClientID=%s | 剩余连接数: %d", clientID, len(c.Clients)) + + // 清理用户绑定关系 + c.cleanupUserBinding(clientID) + } +} + +// 清理用户绑定关系 +func (c *WebSocketController) cleanupUserBinding(clientID string) { + // 获取用户ID + userID, err := c.getUserIDByClientID(clientID) + if err != nil { + return + } + + // 从用户-客户端映射中移除 + userClientsKey := fmt.Sprintf("%s:%s", models.UserClientKey, userID) + c.RedisCli.SRem(c.RedisCtx, userClientsKey, clientID) + + // 从客户端-用户映射中移除 + clientUserKey := fmt.Sprintf("%s:%s", models.ClientUserKey, clientID) + c.RedisCli.Del(c.RedisCtx, clientUserKey) + + log.Printf("🧹🧹🧹🧹🧹🧹🧹🧹 清理绑定关系: ClientID=%s | UserID=%s", clientID, userID) +} + +// 绑定客户端到用户 func (c *WebSocketController) bindClientToUser(clientID, userID string) error { - if err := c.RedisCli.HSet(c.RedisCtx, models.ClientUserKey, clientID, userID).Err(); err != nil { - return fmt.Errorf("保存client-user映射失败: %w", err) + // 设置客户端->用户映射 + clientUserKey := fmt.Sprintf("%s:%s", models.ClientUserKey, clientID) + if err := c.RedisCli.Set(c.RedisCtx, clientUserKey, userID, 0).Err(); err != nil { + return fmt.Errorf("设置客户端用户映射失败: %v", err) } - userClientKey := fmt.Sprintf("%s:%s", models.UserClientKey, userID) - if err := c.RedisCli.SAdd(c.RedisCtx, userClientKey, clientID).Err(); err != nil { - return fmt.Errorf("保存user-client映射失败: %w", err) + // 添加用户->客户端集合映射 + userClientsKey := fmt.Sprintf("%s:%s", models.UserClientKey, userID) + if err := c.RedisCli.SAdd(c.RedisCtx, userClientsKey, clientID).Err(); err != nil { + return fmt.Errorf("添加用户客户端集合失败: %v", err) } - expiration := 24 * time.Hour - if err := c.RedisCli.Expire(c.RedisCtx, userClientKey, expiration).Err(); err != nil { - log.Printf("⚠️ 设置过期时间失败: %v", err) - } - - log.Printf("🔗🔗🔗🔗🔗🔗🔗🔗 用户绑定成功: \n ClientID=%s → UserID=%s", clientID, userID) + log.Printf("🔗🔗🔗🔗🔗🔗🔗🔗 绑定成功: \n ClientID=%s \n UserID=%s", clientID, userID) return nil } -// 通过用户ID发送消息 -func (c *WebSocketController) sendMessageToUser(userID string, msg models.ClientReceivedMessage) error { - start := time.Now() - log.Printf("👤👤👤👤👤👤👤👤 通过用户ID发送消息: UserID=%s", userID) - - userClientKey := fmt.Sprintf("%s:%s", models.UserClientKey, userID) - clientIDs, err := c.RedisCli.SMembers(c.RedisCtx, userClientKey).Result() +// 根据客户端ID获取用户ID +func (c *WebSocketController) getUserIDByClientID(clientID string) (string, error) { + clientUserKey := fmt.Sprintf("%s:%s", models.ClientUserKey, clientID) + userID, err := c.RedisCli.Get(c.RedisCtx, clientUserKey).Result() if err != nil { - return fmt.Errorf("获取用户clientID失败: %w", err) + return "", fmt.Errorf("获取用户ID失败: %v", err) + } + return userID, nil +} + +// 根据用户ID获取客户端ID列表 +func (c *WebSocketController) getClientIDsByUserID(userID string) ([]string, error) { + userClientsKey := fmt.Sprintf("%s:%s", models.UserClientKey, userID) + clientIDs, err := c.RedisCli.SMembers(c.RedisCtx, userClientsKey).Result() + if err != nil { + return nil, fmt.Errorf("获取客户端ID列表失败: %v", err) + } + return clientIDs, nil +} + +// 发送消息给指定客户端 +func (c *WebSocketController) sendMessageToClient(clientID string, message models.ClientReceivedMessage) error { + c.ClientsMux.RLock() + conn, exists := c.Clients[clientID] + c.ClientsMux.RUnlock() + + if !exists { + // 尝试通过Redis转发到其他节点 + return c.forwardMessageToOtherNodes(clientID, message) + } + + c.WriteMutex.Lock() + defer c.WriteMutex.Unlock() + + if err := conn.WriteJSON(message); err != nil { + log.Printf("❌ 发送消息失败: %v | ClientID=%s", err, clientID) + c.removeClient(clientID) + return err + } + + log.Printf("✅ 消息已发送: \n ClientID=%s \n MsgType=%d \n Content=%s", + clientID, message.MessageType, message.Content) + return nil +} + +// 发送消息给指定用户 +func (c *WebSocketController) sendMessageToUser(userID string, message models.ClientReceivedMessage) error { + clientIDs, err := c.getClientIDsByUserID(userID) + if err != nil { + log.Printf("⚠️ 获取用户客户端列表失败: %v | UserID=%s", err, userID) + return err } if len(clientIDs) == 0 { - return fmt.Errorf("用户未绑定任何客户端") + log.Printf("⚠️ 用户无在线客户端: UserID=%s", userID) + return fmt.Errorf("用户 %s 无在线客户端", userID) } - log.Printf("📡📡📡📡📡📡📡 找到 %d 个关联的客户端: UserID=%s", len(clientIDs), userID) + log.Printf("📤📤📤📤📤📤📤📤 发送消息给用户: \n UserID=%s \n ClientCount=%d \n MsgType=%d", + userID, len(clientIDs), message.MessageType) - var successCount, failCount int - var lastError error + var lastErr error + successCount := 0 for _, clientID := range clientIDs { - targetMsg := msg - targetMsg.ReceiverID = clientID - - if err := c.sendMessageToClient(clientID, targetMsg); err != nil { - log.Printf("❌❌❌❌❌❌❌❌ 发送消息到客户端失败: \n ClientID=%s \n %v", clientID, err) - failCount++ - lastError = err + if err := c.sendMessageToClient(clientID, message); err != nil { + log.Printf("⚠️ 发送到客户端失败: %v | ClientID=%s", err, clientID) + lastErr = err } else { successCount++ } } - log.Printf("📬📬📬📬📬📬📬📬 消息发送完成: \n 成功 %d \n 失败 %d \n UserID=%s \n Duration=%s", - successCount, failCount, userID, time.Since(start)) - - if failCount > 0 { - return fmt.Errorf("部分消息发送失败,最后错误: %w", lastError) + if successCount == 0 { + return fmt.Errorf("所有客户端发送失败,最后错误: %v", lastErr) } + + log.Printf("✅ 用户消息发送完成: \n UserID=%s \n 成功=%d/%d", userID, successCount, len(clientIDs)) return nil } -// 发送消息给客户端(添加写锁保护) -func (c *WebSocketController) sendMessageToClient(clientID string, msg models.ClientReceivedMessage) error { - start := time.Now() - - msgJSON, err := json.Marshal(msg) - if err != nil { - log.Printf("❌❌❌❌❌❌❌❌ 消息序列化失败: %v", err) - return fmt.Errorf("内部错误") - } - - if conn := c.getClient(clientID); conn != nil { - log.Printf("📤📤📤📤📤📤📤📤 向本地客户端发送消息: \n ClientID=%s \n MsgType=%d \n Size=%d bytes", - clientID, msg.MessageType, len(msgJSON)) - - // 加写锁保护 - c.WriteMutex.Lock() - defer c.WriteMutex.Unlock() - - if err := conn.WriteMessage(websocket.TextMessage, msgJSON); err != nil { - log.Printf("❌❌❌极简❌❌❌❌ 发送消息失败: %v | ClientID=%s", err, clientID) - c.removeClient(clientID) - return fmt.Errorf("发送消息失败") - } - - log.Printf("✅ 消息成功发送到本地客户端: \n ClientID=%s \n Duration=%s", - clientID, time.Since(start)) - return nil - } - - log.Printf("🌐🌐🌐🌐🌐🌐🌐🌐 向远程节点转发消息: ClientID=%s", clientID) - +// 转发消息到其他节点 +func (c *WebSocketController) forwardMessageToOtherNodes(clientID string, message models.ClientReceivedMessage) error { redisMsg := models.RedisMessage{ SenderNodeID: c.NodeID, ClientID: clientID, - Message: string(msgJSON), + Message: "", } - redisMsgJSON, err := json.Marshal(redisMsg) + msgBytes, err := json.Marshal(message) if err != nil { - log.Printf("❌❌❌❌❌❌❌❌ Redis消息序列化失败: %v", err) - return fmt.Errorf("内部错误") + return fmt.Errorf("序列化消息失败: %v", err) + } + redisMsg.Message = string(msgBytes) + + redisMsgBytes, err := json.Marshal(redisMsg) + if err != nil { + return fmt.Errorf("序列化Redis消息失败: %v", err) } - if err := c.RedisCli.Publish(c.RedisCtx, "ws_messages", redisMsgJSON).Err(); err != nil { - log.Printf("❌❌❌❌❌❌❌❌ 发布消息到Redis失败: %v", err) - return fmt.Errorf("无法转发消息") + channel := fmt.Sprintf("websocket:forward:%s", clientID) + if err := c.RedisCli.Publish(c.RedisCtx, channel, redisMsgBytes).Err(); err != nil { + return fmt.Errorf("发布Redis消息失败: %v", err) } - log.Printf("📡📡📡📡📡📡📡📡 消息已转发到Redis: \n ClientID=%s \n Size=%d bytes \n Duration=%s", - clientID, len(redisMsgJSON), time.Since(start)) + log.Printf("📡📡📡📡📡📡📡📡 消息已转发到其他节点: \n ClientID=%s \n Channel=%s", clientID, channel) return nil } -// SubscribeToRedis 订阅Redis消息 +// 启动Redis消息监听 func (c *WebSocketController) SubscribeToRedis() { - pubsub := c.RedisCli.Subscribe(c.RedisCtx, "ws_messages") - defer pubsub.Close() + go func() { + pattern := "websocket:forward:*" + pubsub := c.RedisCli.PSubscribe(c.RedisCtx, pattern) + defer pubsub.Close() - ch := pubsub.Channel() - log.Println("🔔🔔🔔🔔🔔🔔🔔🔔 开始监听Redis消息...") + log.Printf("📡📡📡📡📡📡📡📡 Redis消息监听已启动: Pattern=%s", pattern) - for msg := range ch { - start := time.Now() - var redisMsg models.RedisMessage - - if err := json.Unmarshal([]byte(msg.Payload), &redisMsg); err != nil { - log.Printf("❌❌❌❌❌❌❌❌ 解析Redis消息失败: %v", err) - continue - } - - log.Printf("📥📥📥📥📥📥📥📥 收到Redis消息: \n Sender=%s \n ClientID=%s \n Size=%d bytes", - redisMsg.SenderNodeID, redisMsg.ClientID, len(msg.Payload)) - - if redisMsg.SenderNodeID == c.NodeID { - log.Println("ℹℹℹℹℹℹℹℹ️ 忽略本节点转发的消息") - continue - } - - if conn := c.getClient(redisMsg.ClientID); conn != nil { - // 加写锁保护 - c.WriteMutex.Lock() - err := conn.WriteMessage(websocket.TextMessage, []byte(redisMsg.Message)) - c.WriteMutex.Unlock() - - if err != nil { - log.Printf("❌❌❌❌❌❌❌❌ 处理Redis消息失败: %v | ClientID=%s", err, redisMsg.ClientID) - c.removeClient(redisMsg.ClientID) + for msg := range pubsub.Channel() { + var redisMsg models.RedisMessage + if err := json.Unmarshal([]byte(msg.Payload), &redisMsg); err != nil { + log.Printf("⚠️ 解析Redis消息失败: %v", err) continue } - log.Printf("✅ Redis消息已处理: \n ClientID=%s \n Duration=%s", - redisMsg.ClientID, time.Since(start)) - } else { - log.Printf("⚠️ 目标客户端不在本节点: ClientID=%s", redisMsg.ClientID) + // 忽略自己发送的消息 + if redisMsg.SenderNodeID == c.NodeID { + continue + } + + // 检查目标客户端是否在本节点 + c.ClientsMux.RLock() + conn, exists := c.Clients[redisMsg.ClientID] + c.ClientsMux.RUnlock() + + if !exists { + continue + } + + var clientMsg models.ClientReceivedMessage + if err := json.Unmarshal([]byte(redisMsg.Message), &clientMsg); err != nil { + log.Printf("⚠️ 解析客户端消息失败: %v", err) + continue + } + + c.WriteMutex.Lock() + if err := conn.WriteJSON(clientMsg); err != nil { + log.Printf("❌ 转发消息失败: %v | ClientID=%s", err, redisMsg.ClientID) + c.removeClient(redisMsg.ClientID) + } else { + log.Printf("✅ 转发消息成功: \n ClientID=%s \n FromNode=%s", + redisMsg.ClientID, redisMsg.SenderNodeID) + } + c.WriteMutex.Unlock() } - } -} - -// ========================== 辅助方法 ========================== // - -// 添加客户端到映射 -func (c *WebSocketController) addClient(clientID string, conn *websocket.Conn) { - c.ClientsMux.Lock() - defer c.ClientsMux.Unlock() - c.Clients[clientID] = conn - log.Printf("📌📌📌📌 添加客户端到连接池: \n ClientID=%s \n 当前连接数=%d", - clientID, len(c.Clients)) -} - -// 从映射中移除客户端 -func (c *WebSocketController) removeClient(clientID string) { - c.ClientsMux.Lock() - defer c.ClientsMux.Unlock() - - if _, exists := c.Clients[clientID]; !exists { - return - } - - delete(c.Clients, clientID) - log.Printf("🗑🗑🗑🗑️ 从连接池移除客户端: \n ClientID=%s \n 当前连接数=%d", - clientID, len(c.Clients)) - - if err := c.RedisCli.HDel(c.RedisCtx, models.ClientUserKey, clientID).Err(); err != nil { - log.Printf("⚠️ 删除client_user_mapping失败: %v", err) - } -} - -// 获取客户端连接 -func (c *WebSocketController) getClient(clientID string) *websocket.Conn { - c.ClientsMux.RLock() - defer c.ClientsMux.RUnlock() - return c.Clients[clientID] -} - -// 通过客户端ID获取用户ID -func (c *WebSocketController) getUserIDByClientID(clientID string) (string, error) { - userID, err := c.RedisCli.HGet(c.RedisCtx, models.ClientUserKey, clientID).Result() - if err == redis.Nil { - return "", fmt.Errorf("客户端未绑定用户") - } else if err != nil { - return "", fmt.Errorf("获取用户ID失败: %w", err) - } - return userID, nil + }() } diff --git a/models/messages.go b/models/messages.go index bc1b54f..249e297 100644 --- a/models/messages.go +++ b/models/messages.go @@ -16,6 +16,9 @@ type SendToUserPayload struct { ReceiverUserID string `json:"receiver_user_id"` // 接收者用户ID MessageType int `json:"message_type"` // 消息类型(0-8) MessageContent string `json:"message_content"` // 消息内容 + // 通话专用字段 + CallID string `json:"call_id,omitempty"` // 通话唯一ID + CallStatus string `json:"call_status,omitempty"` // 通话状态:invite/accepted/rejected/ended/candidate/offer/answer } // 客户端接收消息结构(最终格式) @@ -30,7 +33,7 @@ type ClientReceivedMessage struct { // 通话专用字段 CallID string `json:"call_id,omitempty"` // 通话唯一ID - CallStatus string `json:"call_status,omitempty"` // 通话状态:invite/accepted/rejected/ended/candidate + CallStatus string `json:"call_status,omitempty"` // 通话状态:invite/accepted/rejected/ended/candidate/offer/answer } // Redis 消息结构 @@ -48,7 +51,7 @@ type BindRequest struct { // 通话信令结构 type CallSignal struct { - CallStatus string `json:"call_status"` // 动作: invite/accept/reject/end/candidate + CallStatus string `json:"call_status"` // 动作: invite/accept/reject/end/candidate/offer/answer CallType int `json:"call_type"` // 通话类型:6=视频,7=语音 CallID string `json:"call_id"` // 通话唯一ID CallerID string `json:"caller_id"` // 主叫用户ID @@ -56,8 +59,72 @@ type CallSignal struct { Data string `json:"data"` // 数据(SDP/ICE候选) } +// WebRTC信令数据结构 +type WebRTCSignalData struct { + Type string `json:"type"` // offer/answer/candidate + SDP string `json:"sdp,omitempty"` // SDP数据 + Candidate interface{} `json:"candidate,omitempty"` // ICE候选数据 + SDPMid string `json:"sdpMid,omitempty"` // SDP媒体标识 + SDPMLineIndex int `json:"sdpMLineIndex,omitempty"` // SDP行索引 +} + +// 通话状态常量 +const ( + CallStatusInvite = "invite" // 发起邀请 + CallStatusAccepted = "accepted" // 接受通话 + CallStatusRejected = "rejected" // 拒绝通话 + CallStatusEnded = "ended" // 结束通话 + CallStatusOffer = "offer" // WebRTC Offer + CallStatusAnswer = "answer" // WebRTC Answer + CallStatusCandidate = "candidate" // ICE候选 + CallStatusHangup = "hangup" // 挂断 + CallStatusDisconnected = "disconnected" // 断开连接 + CallStatusTerminated = "terminated" // 终止 + CallStatusNoAnswer = "no-answer" // 无人接听 + CallStatusBusy = "busy" // 忙线 + CallStatusFailed = "failed" // 失败 + CallStatusConnecting = "connecting" // 连接中 +) + +// 消息类型常量 +const ( + MessageTypeText = 0 // 文本消息 + MessageTypeImage = 1 // 图片消息 + MessageTypeAudio = 2 // 音频消息 + MessageTypeVideo = 3 // 视频消息 + MessageTypePrescription = 4 // 处方消息 + MessageTypeMedicalRecord = 5 // 病例消息 + MessageTypeVideoCall = 6 // 视频通话 + MessageTypeAudioCall = 7 // 语音通话 + MessageTypeFile = 8 // 文件消息 +) + // Redis键常量 const ( ClientUserKey = "client_user_mapping" // clientid -> userid映射 UserClientKey = "user_client_mapping" // userid -> clientid列表 ) + +// 响应结构 +type APIResponse struct { + Status string `json:"status"` + Message string `json:"message"` + Data interface{} `json:"data,omitempty"` + Error string `json:"error,omitempty"` +} + +// 健康检查响应 +type HealthResponse struct { + Status string `json:"status"` + Node string `json:"node"` + Port string `json:"port"` + Time string `json:"time"` +} + +// 绑定响应 +type BindResponse struct { + Status string `json:"status"` + Message string `json:"message"` + UserID string `json:"user_id"` + ClientID string `json:"client_id"` +}