diff --git a/.gitignore b/.gitignore index 6d9594a..9b5834d 100644 --- a/.gitignore +++ b/.gitignore @@ -1 +1,2 @@ /logs/ +/dist/ diff --git a/README.md b/README.md new file mode 100644 index 0000000..3fbd5d6 --- /dev/null +++ b/README.md @@ -0,0 +1,64 @@ +# 角色 +你是一个资深的前端开发工程师、UI设计师、UX设计师、IM系统设计师,你本次的任务帮助一个用户复刻{index.html}中的代码并且优化 + +# 技术选型 +1. vue3+js+tailwindcss+antd-vue4.x+axios + +# vue配置 +```json +{ + "name": "spa-view", + "private": true, + "version": "0.0.0", + "type": "module", + "scripts": { + "dev": "vite", + "build": "vite build", + "preview": "vite preview" + }, + "dependencies": { + "@ant-design-vue/pro-layout": "^3.2.5", + "@ant-design/icons-vue": "^7.0.1", + "@fortawesome/fontawesome-free": "^6.7.2", + "@tailwindcss/vite": "^4.1.4", + "@vueuse/head": "^2.0.0", + "ant-design-vue": "^4.1.1", + "axios": "^1.6.2", + "lucide-vue-next": "^0.507.0", + "pinia": "^2.1.7", + "swiper": "^11.0.5", + "tailwindcss": "^4.1.4", + "vue": "^3.5.13", + "vue-easy-lightbox": "^1.19.0", + "vue-router": "^4.2.5" + }, + "devDependencies": { + "@vitejs/plugin-vue": "^5.2.2", + "sass": "^1.71.0", + "vite": "^6.3.1" + } +} + +``` + +# 技术要求 +1. UI风格不变,但是请求封装成了request.js +2. 将组件拆分细化 +3. 功能不能减少、尤其是输入框的UI、交互 +4. 登录独立出一个页面,在index.js中配置,未登录用户自动重定向到登录页 +5. 要保证所有功能都可以使用 +6. 要能够全局切换亮色暗色模式 +7. 聊天气泡的UI可以变一下,时间和已读状态要在气泡外面,显示的名称、头像要根据当前聊天用户来变化(我的肯定是用自己当前登录的,对方的话就是当前和谁聊天就显示谁) +8. 在未选中聊天人的时候,不展示输入框,并且右侧用svg插画来展示左侧点击用户开始聊天 + +# 已配置好的内容 +1. tailwindcss和antd-vue已经配置好,直接使用即可 +2. 接口代理已经配置好,直接使用即可 + +# 需要优化的点 +1. 图片、视频预览的可以使用第三方开源库来实现 +2. 音频样式需要优化(类似微信的音波浮动) +3. 聊天记录、用户列表信息考虑使用前端的本地数据库来记录(indexedDB) + +# 任务 +你本次的任务是根据我上传的{index.html}来使用技术栈:{技术选型}完成,请结合{vue配置}、{技术要求}和{已配置好的内容}、{需要优化的点}来作为提示,完成本项目(本项目只需要适配各个PC端即可,无需考虑平板、手机,对于笔记本来说就页面宽高都沾满就好了),在开始写代码钱先总结一份我的需求说明,写完之前要检查所有功能是否可用 \ No newline at end of file diff --git a/README2.md b/README2.md new file mode 100644 index 0000000..cfb6987 --- /dev/null +++ b/README2.md @@ -0,0 +1,204 @@ +# 现在需要修复的问题 +1. 媒体预览的组件有问题:Uncaught TypeError: $setup.previewMedia is not a function + at _createElementBlock.onClick._cache.._cache. (MessageBubble.vue:27:21) +2. 点击视频聊天或者语音聊天后应该请求弹框(等待对方接听,对方传输接收则开始聊天,对方拒绝则给出消息提示然后存在聊天记录中),对方也要弹出视频/语音通话邀请,可以接收或者拒绝。 +3. 根据{前端获取信令示例}和{后端交互代码}来完成整个聊天功能,并且要考虑用户掉线、断网、网络异常等情况,以及用户退出登录、用户被删除等情况。 +4. 在好友管理点击聊天后要吧侧边导航变成聊天 +5. 聊天记录列表和好友管理中中,点击对方的头像要有对方的信息 +6. 完整的基于WebRTC实现语音/视频通话功能 + + + + +# 前端获取信令示例 +```javascript +// 发起通话邀请 +function startCall(targetUserId, isVideo) { + const callId = generateUniqueId(); + const callType = isVideo ? 6 : 7; // 6=视频,7=语音 + + const callSignal = { + action: "invite", + call_type: callType, + call_id: callId, + caller_id: currentUserId, + callee_id: targetUserId, + data: JSON.stringify(offer) // SDP offer + }; + + const payload = { + sender_user_id: currentUserId, + receiver_user_id: targetUserId, + message_type: callType, + message_content: JSON.stringify(callSignal) + }; + + // 调用/send-to-user接口 + fetch('/send-to-user', { + method: 'POST', + headers: { + 'Content-Type': 'application/json' + }, + body: JSON.stringify(payload) + }).then(response => { + if (!response.ok) { + throw new Error('Failed to send call invite'); + } + return response.json(); + }).then(data => { + console.log('Call invite sent:', data); + }).catch(error => { + console.error('Error sending call invite:', error); + }); +} + +// 接受通话 +function acceptCall(callId, callType) { + const callSignal = { + action: "accept", + call_id: callId, + data: JSON.stringify(answer) // SDP answer + }; + + const payload = { + sender_user_id: currentUserId, + receiver_user_id: callerUserId, + message_type: callType, + message_content: JSON.stringify(callSignal) + }; + + fetch('/send-to-user', { + method: 'POST', + headers: { + 'Content-Type': 'application/json' + }, + body: JSON.stringify(payload) + }); +} + +// 拒绝通话 +function rejectCall(callId, callType, reason) { + const callSignal = { + action: "reject", + call_id: callId, + data: reason || "用户忙" + }; + + const payload = { + sender_user_id: currentUserId, + receiver_user_id: callerUserId, + message_type: callType, + message_content: JSON.stringify(callSignal) + }; + + fetch('/send-to-user', { + method: 'POST', + headers: { + 'Content-Type': 'application/json' + }, + body: JSON.stringify(payload) + }); +} + +// 结束通话 +function endCall(callId, callType) { + const callSignal = { + action: "end", + call_id: callId + }; + + const payload = { + sender_user_id: currentUserId, + receiver_user_id: otherUserId, + message_type: callType, + message_content: JSON.stringify(callSignal) + }; + + fetch('/send-to-user', { + method: 'POST', + headers: { + 'Content-Type': 'application/json' + }, + body: JSON.stringify(payload) + }); +} + +// 发送ICE候选 +function sendCandidate(callId, callType, candidate, targetUserId) { + const callSignal = { + action: "candidate", + call_id: callId, + data: JSON.stringify(candidate) + }; + + const payload = { + sender_user_id: currentUserId, + receiver_user_id: targetUserId, + message_type: callType, + message_content: JSON.stringify(callSignal) + }; + + fetch('/send-to-user', { + method: 'POST', + headers: { + 'Content-Type': 'application/json' + }, + body: JSON.stringify(payload) + }); +} +``` + +## 后端交互代码 +1. {websocket_controller.go} + +## 消息结构 +1. {messages.go} + +## 通话流程图 +```mermaid +graph TD + A[用户A点击通话按钮] --> B{选择通话类型} + B --> |语音通话| C1[设置 message_type=7] + B --> |视频通话| C2[设置 message_type=6] + + C1 --> D[生成唯一CallID] + D --> E[创建PeerConnection对象] + E --> F[获取本地媒体流] + F --> G[创建SDP offer] + G --> H[构造信令对象] + + H --> I[发送邀请信令] + I --> J[后端转发给用户B] + + J --> K[用户B收到来电通知] + K --> L{用户B选择} + L --> |接听| M[创建SDP answer] + L --> |拒绝| N[发送拒绝信令] + + M --> O[发送接受信令] + O --> P[后端转发给用户A] + + P --> Q[用户A设置远程描述] + Q --> R[交换ICE候选] + + R --> S[建立P2P连接] + S --> T[开始语音通话] + + N --> U[用户A收到拒绝通知] + U --> V[结束通话流程] + + subgraph 信令交换 + I[发送邀请信令] + O[发送接受信令] + R[交换ICE候选] + end + + subgraph WebRTC连接 + E[创建PeerConnection] + F[获取本地媒体流] + G[创建SDP offer] + M[创建SDP answer] + Q[设置远程描述] + S[建立P2P连接] + end +``` \ No newline at end of file diff --git a/controller/websocket_controller.go b/controller/websocket_controller.go index 068c242..0b7333a 100644 --- a/controller/websocket_controller.go +++ b/controller/websocket_controller.go @@ -29,6 +29,7 @@ var upgrader = websocket.Upgrader{ type WebSocketController struct { Clients map[string]*websocket.Conn ClientsMux sync.RWMutex + WriteMutex sync.Mutex // 添加写锁 RedisCli *redis.Client NodeID string Port string @@ -50,7 +51,7 @@ func (c *WebSocketController) ConfigureSystem() { // 只有当端口未设置时才从环境变量获取 if c.Port == "" { c.Port = utils.GetEnv("PORT", "12080") - log.Printf("📡📡📡📡📡📡📡📡 使用环境变量设置端口: %s", c.Port) + log.Printf("📡📡📡📡📡📡📡📡📡📡📡📡📡📡📡📡 使用环境变量设置端口: %s", c.Port) } log.SetPrefix(fmt.Sprintf("[Node:%s] ", c.NodeID)) @@ -79,7 +80,7 @@ func (c *WebSocketController) registerNode() { if err := c.RedisCli.Set(c.RedisCtx, key, value, 30*time.Second).Err(); err != nil { log.Printf("⚠️ 节点注册失败: %v", err) } else { - log.Printf("📌📌📌📌📌📌📌📌 节点已注册: %s = %s", key, value) + log.Printf("📌📌📌📌📌📌📌📌📌📌📌📌📌📌📌📌 节点已注册: %s = %s", key, value) } time.Sleep(20 * time.Second) @@ -90,7 +91,7 @@ func (c *WebSocketController) registerNode() { func (c *WebSocketController) configureLogger() { logDir := utils.GetEnv("LOG_DIR", "./logs") if err := os.MkdirAll(logDir, 0755); err != nil { - log.Fatalf("❌❌❌❌ 创建日志目录失败: %v", err) + log.Fatalf("❌❌❌❌❌❌❌❌ 创建日志目录失败: %v", err) } logFreq := utils.GetEnv("LOG_ROTATE_FREQ", "daily") @@ -118,7 +119,7 @@ func (c *WebSocketController) InitRedisClient() { }) if err := c.checkRedisConnection(); err != nil { - log.Fatalf("❌❌❌❌ Redis连接失败: %v", err) + log.Fatalf("❌❌❌❌❌❌❌❌ Redis连接失败: %v", err) } } @@ -131,16 +132,16 @@ func (c *WebSocketController) checkRedisConnection() error { // PrintStartupInfo 打印启动信息 func (c *WebSocketController) PrintStartupInfo() { hostname, _ := os.Hostname() - log.Printf("🚀🚀🚀🚀 WebSocket服务启动: NodeID=%s", c.NodeID) - log.Printf("🌐🌐🌐🌐 监听端口: %s", c.Port) - log.Printf("📡📡📡📡 Redis地址: %s", utils.GetEnv("REDIS_ADDR", "localhost:6379")) - log.Printf("💻💻💻💻 主机: %s", hostname) - log.Printf("🕒🕒🕒🕒🕒🕒🕒🕒🕒 启动时间: %s", time.Now().Format("2006-01-02 15:04:05")) - log.Printf("🔔🔔🔔🔔 支持消息类型: \n 0:文本\n 1:图片\n 2:音频\n 3:视频\n 4:处方\n 5:病例\n 6:视频通话") - log.Printf("🔑🔑🔑🔑 用户绑定键: \n %s\n %s", models.ClientUserKey, models.UserClientKey) - log.Printf("📬📬📬📬 新增接口: POST /send-to-user (直接通过用户ID发送消息)") - log.Printf("🔗🔗🔗🔗 集群节点注册: %s:%s", utils.GetOutboundIP(), c.Port) - log.Println("🔗🔗🔗🔗 等待客户端连接...") + log.Printf("🚀🚀🚀🚀🚀🚀🚀🚀 WebSocket服务启动: NodeID=%s", c.NodeID) + log.Printf("🌐🌐🌐🌐🌐🌐🌐🌐 监听端口: %s", c.Port) + log.Printf("📡📡📡📡📡📡📡📡 Redis地址: %s", utils.GetEnv("REDIS_ADDR", "localhost:6379")) + log.Printf("💻💻💻💻💻💻💻💻 主机: %s", hostname) + log.Printf("🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒🕒 启动时间: %s", time.Now().Format("2006-01-02 15:04:05")) + log.Printf("🔔🔔🔔🔔🔔🔔🔔🔔 支持消息类型: \n 0:文本\n 1:图片\n 2:音频\n 3:视频\n 4:处方\n 5:病例\n 6:视频通话\n 7:语音通话\n 8:文件消息") + log.Printf("🔑🔑🔑🔑🔑🔑🔑🔑 用户绑定键: \n %s\n %s", models.ClientUserKey, models.UserClientKey) + log.Printf("📬📬📬📬📬📬📬📬 新增接口: POST /send-to-user (直接通过用户ID发送消息)") + log.Printf("🔗🔗🔗🔗🔗🔗🔗🔗 集群节点注册: %s:%s", utils.GetOutboundIP(), c.Port) + log.Println("🔗🔗🔗🔗🔗🔗🔗🔗 等待客户端连接...") } // HealthHandler 健康检查处理器 @@ -157,7 +158,7 @@ func (c *WebSocketController) HealthHandler(ctx *gin.Context) { func (c *WebSocketController) HandleWebSocket(ctx *gin.Context) { start := time.Now() clientIP := ctx.ClientIP() - log.Printf("👤👤👤👤 客户端连接中: IP=%s", clientIP) + log.Printf("👤👤👤👤👤👤👤👤 客户端连接中: IP=%s", clientIP) conn, err := upgrader.Upgrade(ctx.Writer, ctx.Request, nil) if err != nil { @@ -188,7 +189,7 @@ func (c *WebSocketController) HandleWebSocket(ctx *gin.Context) { messageType, message, err := conn.ReadMessage() if err != nil { if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway) { - log.Printf("❌❌❌❌ 连接意外断开: %v | ClientID=%s", err, clientID) + log.Printf("❌❌❌❌❌❌❌❌ 连接意外断开: %v | ClientID=%s", err, clientID) } else { log.Printf("⚠️ 连接正常关闭: ClientID=%s", clientID) } @@ -203,16 +204,16 @@ func (c *WebSocketController) HandleWebSocket(ctx *gin.Context) { // 处理客户端消息 func (c *WebSocketController) handleClientMessage(senderID string, message []byte) { - log.Printf("📥📥📥📥 收到客户端消息: \n SenderID=%s \n Size=%d bytes", senderID, len(message)) + log.Printf("📥📥📥📥📥📥📥📥 收到客户端消息: \n SenderID=%s \n Size=%d bytes", senderID, len(message)) var payload models.SendMessagePayload if err := json.Unmarshal(message, &payload); err == nil && payload.RequestType != "" { - log.Printf("📦📦📦📦 解析JSON消息成功: Type=%s", payload.RequestType) + log.Printf("📦📦📦📦📦📦📦📦 解析JSON消息成功: Type=%s", payload.RequestType) if payload.RequestType == "bind" && payload.SenderUserID != "" { - log.Printf("🔗🔗🔗🔗 处理绑定请求: \n ClientID=%s \n UserID=%s", senderID, payload.SenderUserID) + log.Printf("🔗🔗🔗🔗🔗🔗🔗🔗 处理绑定请求: \n ClientID=%s \n UserID=%s", senderID, payload.SenderUserID) if err := c.bindClientToUser(senderID, payload.SenderUserID); err != nil { - log.Printf("❌❌❌❌ 绑定失败: %v", err) + log.Printf("❌❌❌❌❌❌❌❌ 绑定失败: %v", err) } else { log.Printf("✅ 绑定成功: \n ClientID=%s \n UserID=%s", senderID, payload.SenderUserID) } @@ -221,7 +222,7 @@ func (c *WebSocketController) handleClientMessage(senderID string, message []byt if payload.RequestType == "send_message" { if payload.TargetClientID != "" { - log.Printf("📨📨📨📨 处理客户端发起的发送请求: \n TargetID=%s \n MsgType=%d", + log.Printf("📨📨📨📨📨📨📨📨 处理客户端发起的发送请求: \n TargetID=%s \n MsgType=%d", payload.TargetClientID, payload.MessageType) clientMsg := models.ClientReceivedMessage{ @@ -239,7 +240,7 @@ func (c *WebSocketController) handleClientMessage(senderID string, message []byt } if payload.ReceiverUserID != "" { - log.Printf("📨📨📨📨 处理客户端发起的用户发送请求: \n ReceiverUserID=%s \n MsgType=%d", + log.Printf("📨📨📨📨📨📨📨📨 处理客户端发起的用户发送请求: \n ReceiverUserID=%s \n MsgType=%d", payload.ReceiverUserID, payload.MessageType) clientMsg := models.ClientReceivedMessage{ @@ -260,7 +261,7 @@ func (c *WebSocketController) handleClientMessage(senderID string, message []byt var clientMsg models.ClientReceivedMessage if err := json.Unmarshal(message, &clientMsg); err == nil && clientMsg.ReceiverID != "" { - log.Printf("📨📨📨📨 处理客户端封装消息: \n TargetID=%s \n MsgType=%d", + log.Printf("📨📨📨📨📨📨📨📨 处理客户端封装消息: \n TargetID=%s \n MsgType=%d", clientMsg.ReceiverID, clientMsg.MessageType) clientMsg.SenderID = senderID @@ -268,16 +269,159 @@ func (c *WebSocketController) handleClientMessage(senderID string, message []byt return } + // 修复:正确解析通话信令并分发处理 + var callSignal models.CallSignal + if err := json.Unmarshal(message, &callSignal); err == nil && callSignal.Action != "" { + // 获取发送者用户ID + senderUserID, _ := c.getUserIDByClientID(senderID) + + log.Printf("📞📞 处理通话信令: \n Action=%s \n CallID=%s \n CallType=%d \n SenderUserID=%s \n CallerID=%s \n CalleeID=%s", + callSignal.Action, callSignal.CallID, callSignal.CallType, senderUserID, callSignal.CallerID, callSignal.CalleeID) + + // 处理不同类型的通话动作 + switch callSignal.Action { + case "invite": + c.handleCallInvite(senderID, senderUserID, callSignal) + case "accept": + c.handleCallAccept(senderID, senderUserID, callSignal) + case "reject": + c.handleCallReject(senderID, senderUserID, callSignal) + case "end": + c.handleCallEnd(senderID, senderUserID, callSignal) + case "candidate": + c.handleCallCandidate(senderID, senderUserID, callSignal) + default: + log.Printf("⚠️ 未知的通话动作: %s", callSignal.Action) + } + return + } + log.Printf("⚠️ 无法识别的消息格式: \n Size=%d bytes \n Message=%s", len(message), string(message)) } +// 处理通话邀请 +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) + + // 创建通话消息 + callMsg := models.ClientReceivedMessage{ + SenderID: senderID, + ReceiverID: "", // 通过用户ID发送 + SenderUserID: senderUserID, + ReceiverUserID: signal.CalleeID, + MessageType: signal.CallType, // 6=视频, 7=语音 + Content: signal.Data, // 包含SDP offer + SendTime: time.Now().Format(time.RFC3339), + CallID: signal.CallID, + CallStatus: "invite", + } + + // 发送邀请给被叫方 + if err := c.sendMessageToUser(signal.CalleeID, callMsg); err != nil { + log.Printf("❌ 发送通话邀请失败: %v | CalleeID=%s", err, signal.CalleeID) + } +} + +// 处理通话接受 +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) + + // 创建通话接受消息 + acceptMsg := models.ClientReceivedMessage{ + SenderID: senderID, + ReceiverID: "", // 通过用户ID发送 + SenderUserID: senderUserID, + ReceiverUserID: signal.CallerID, + MessageType: signal.CallType, // 6=视频, 7=语音 + Content: signal.Data, // 包含SDP answer + SendTime: time.Now().Format(time.RFC3339), + CallID: signal.CallID, + CallStatus: "accepted", + } + + // 发送接受消息给主叫方 + if err := c.sendMessageToUser(signal.CallerID, acceptMsg); err != nil { + log.Printf("❌ 发送通话接受失败: %v | CallerID=%s", err, signal.CallerID) + } +} + +// 处理通话拒绝 +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) + + // 创建通话拒绝消息 + rejectMsg := models.ClientReceivedMessage{ + SenderID: senderID, + ReceiverID: "", // 通过用户ID发送 + SenderUserID: senderUserID, + ReceiverUserID: signal.CallerID, + MessageType: signal.CallType, // 6=视频, 7=语音 + Content: signal.Data, // 可包含拒绝原因 + SendTime: time.Now().Format(time.RFC3339), + CallID: signal.CallID, + CallStatus: "rejected", + } + + // 发送拒绝消息给主叫方 + if err := c.sendMessageToUser(signal.CallerID, rejectMsg); err != nil { + log.Printf("❌ 发送通话拒绝失败: %v | CallerID=%s", err, signal.CallerID) + } +} + +// 处理通话结束 +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) + + // 创建通话结束消息 + endMsg := models.ClientReceivedMessage{ + SenderID: senderID, + ReceiverID: "", // 通过用户ID发送 + SenderUserID: senderUserID, + ReceiverUserID: signal.CalleeID, + MessageType: signal.CallType, // 6=视频, 7=语音 + Content: signal.Data, // 可包含结束原因 + SendTime: time.Now().Format(time.RFC3339), + CallID: signal.CallID, + CallStatus: "ended", + } + + // 发送结束消息给对端 + if err := c.sendMessageToUser(signal.CalleeID, endMsg); err != nil { + log.Printf("❌ 发送通话结束失败: %v | CalleeID=%s", err, signal.CalleeID) + } +} + +// 处理通话候选 +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) + + // 创建候选消息 + candidateMsg := models.ClientReceivedMessage{ + SenderID: senderID, + ReceiverID: "", // 通过用户ID发送 + SenderUserID: senderUserID, + ReceiverUserID: signal.CalleeID, + MessageType: signal.CallType, // 6=视频, 7=语音 + Content: signal.Data, // 包含ICE候选 + SendTime: time.Now().Format(time.RFC3339), + CallID: signal.CallID, + CallStatus: "candidate", + } + + // 发送候选消息给对端 + if err := c.sendMessageToUser(signal.CalleeID, candidateMsg); 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请求格式: %v", err) + log.Printf("⚠️ 无效的JSON请求格式: %极简", err) ctx.JSON(http.StatusBadRequest, gin.H{"error": "无效的JSON格式"}) return } @@ -291,7 +435,7 @@ func (c *WebSocketController) SendMessageHandler(ctx *gin.Context) { } if payload.RequestType == "send_message" { - log.Printf("📤📤📤📤 处理API发送请求: \n SenderUser=%s \n MsgType=%d", + log.Printf("📤📤📤📤📤📤📤📤 处理API发送请求: \n SenderUser=%s \n MsgType=%d", payload.SenderUserID, payload.MessageType) clientMsg := models.ClientReceivedMessage{ @@ -308,11 +452,11 @@ func (c *WebSocketController) SendMessageHandler(ctx *gin.Context) { var sendErr error if payload.TargetClientID != "" { - log.Printf("🎯🎯🎯🎯 目标类型: ClientID | Target=%s", payload.TargetClientID) + log.Printf("🎯🎯🎯🎯🎯🎯🎯🎯 目标类型: ClientID | Target=%s", payload.TargetClientID) sendErr = c.sendMessageToClient(payload.TargetClientID, clientMsg) result = fmt.Sprintf("消息已发送到ClientID: %s", payload.TargetClientID) } else if payload.ReceiverUserID != "" { - log.Printf("🎯🎯🎯🎯 目标类型: UserID | ReceiverUser=%s", payload.ReceiverUserID) + log.Printf("🎯🎯🎯🎯🎯🎯🎯🎯 目标类型: UserID | ReceiverUser=%s", payload.ReceiverUserID) sendErr = c.sendMessageToUser(payload.ReceiverUserID, clientMsg) result = fmt.Sprintf("消息已发送到UserID: %s", payload.ReceiverUserID) } else { @@ -333,7 +477,7 @@ func (c *WebSocketController) SendMessageHandler(ctx *gin.Context) { ctx.JSON(http.StatusBadRequest, gin.H{"error": "不支持的request_type"}) } -// BindHandler 用户绑定处理器 +// 用户绑定处理器 func (c *WebSocketController) BindHandler(ctx *gin.Context) { start := time.Now() @@ -350,10 +494,10 @@ func (c *WebSocketController) BindHandler(ctx *gin.Context) { return } - log.Printf("🔗🔗🔗🔗 处理用户绑定请求: \n UserID=%s \n ClientID=%s", req.UserID, req.ClientID) + log.Printf("🔗🔗🔗🔗🔗🔗🔗🔗 处理用户绑定请求: \n UserID=%s \n ClientID=%s", req.UserID, req.ClientID) if err := c.bindClientToUser(req.ClientID, req.UserID); err != nil { - log.Printf("❌❌❌❌ 绑定失败: %v", err) + log.Printf("❌❌❌❌❌❌❌❌ 绑定失败: %v", err) ctx.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return } @@ -367,7 +511,7 @@ func (c *WebSocketController) BindHandler(ctx *gin.Context) { req.UserID, req.ClientID, time.Since(start)) } -// SendToUserHandler 通过用户ID发送消息处理器 +// SendToUserHandler 通过用户ID发送消息处理器(支持通话信令) func (c *WebSocketController) SendToUserHandler(ctx *gin.Context) { start := time.Now() @@ -385,9 +529,10 @@ func (c *WebSocketController) SendToUserHandler(ctx *gin.Context) { return } - log.Printf("📤📤📤📤 处理API用户发送请求: \n SenderUser=%s → ReceiverUser=%s \n MsgType=%d", + log.Printf("📤📤📤📤📤📤📤📤 处理API用户发送请求: \n SenderUser=%s → ReceiverUser=%s \n MsgType=%d", payload.SenderUserID, payload.ReceiverUserID, payload.MessageType) + // 创建基础消息结构 clientMsg := models.ClientReceivedMessage{ SenderID: "system", ReceiverID: "", @@ -398,6 +543,16 @@ func (c *WebSocketController) SendToUserHandler(ctx *gin.Context) { SendTime: time.Now().Format(time.RFC3339), } + // 如果是通话信令消息,添加通话相关字段 + if payload.MessageType == 6 || payload.MessageType == 7 { + // 直接使用内容作为通话信令数据 + clientMsg.CallStatus = "invite" // 默认设置为invite,API调用通常是发起通话 + clientMsg.CallID = utils.GenerateCallID() // 生成唯一的通话ID + + log.Printf("📞📞 处理API通话信令: \n CallID=%s \n CallType=%d", + clientMsg.CallID, payload.MessageType) + } + if err := c.sendMessageToUser(payload.ReceiverUserID, clientMsg); err != nil { ctx.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()}) return @@ -407,6 +562,7 @@ func (c *WebSocketController) SendToUserHandler(ctx *gin.Context) { "code": 0, "status": "success", "message": fmt.Sprintf("消息已发送至用户 %s", payload.ReceiverUserID), + "call_id": clientMsg.CallID, // 返回生成的call_id }) log.Printf("✅ API请求完成: Duration=%s", time.Since(start)) @@ -428,14 +584,14 @@ func (c *WebSocketController) bindClientToUser(clientID, userID string) error { log.Printf("⚠️ 设置过期时间失败: %v", err) } - log.Printf("🔗🔗🔗🔗 用户绑定成功: \n ClientID=%s → UserID=%s", clientID, userID) + log.Printf("🔗🔗🔗🔗🔗🔗🔗🔗 用户绑定成功: \n ClientID=%s → 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) + log.Printf("👤👤👤👤👤👤👤👤 通过用户ID发送消息: UserID=%s", userID) userClientKey := fmt.Sprintf("%s:%s", models.UserClientKey, userID) clientIDs, err := c.RedisCli.SMembers(c.RedisCtx, userClientKey).Result() @@ -447,7 +603,7 @@ func (c *WebSocketController) sendMessageToUser(userID string, msg models.Client return fmt.Errorf("用户未绑定任何客户端") } - log.Printf("📡📡📡📡 找到 %d 个关联的客户端: UserID=%s", len(clientIDs), userID) + log.Printf("📡📡📡📡📡📡📡 找到 %d 个关联的客户端: UserID=%s", len(clientIDs), userID) var successCount, failCount int var lastError error @@ -457,7 +613,7 @@ func (c *WebSocketController) sendMessageToUser(userID string, msg models.Client targetMsg.ReceiverID = clientID if err := c.sendMessageToClient(clientID, targetMsg); err != nil { - log.Printf("❌❌❌❌ 发送消息到客户端失败: \n ClientID=%s \n %v", clientID, err) + log.Printf("❌❌❌❌❌❌❌❌ 发送消息到客户端失败: \n ClientID=%s \n %v", clientID, err) failCount++ lastError = err } else { @@ -465,7 +621,7 @@ func (c *WebSocketController) sendMessageToUser(userID string, msg models.Client } } - log.Printf("📬📬📬📬 消息发送完成: \n 成功 %d \n 失败 %d \n UserID=%s \n Duration=%s", + log.Printf("📬📬📬📬📬📬📬📬 消息发送完成: \n 成功 %d \n 失败 %d \n UserID=%s \n Duration=%s", successCount, failCount, userID, time.Since(start)) if failCount > 0 { @@ -474,22 +630,26 @@ func (c *WebSocketController) sendMessageToUser(userID string, msg models.Client 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) + 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", + 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) + log.Printf("❌❌❌极简❌❌❌❌ 发送消息失败: %v | ClientID=%s", err, clientID) c.removeClient(clientID) return fmt.Errorf("发送消息失败") } @@ -499,7 +659,7 @@ func (c *WebSocketController) sendMessageToClient(clientID string, msg models.Cl return nil } - log.Printf("🌐🌐🌐🌐 向远程节点转发消息: ClientID=%s", clientID) + log.Printf("🌐🌐🌐🌐🌐🌐🌐🌐 向远程节点转发消息: ClientID=%s", clientID) redisMsg := models.RedisMessage{ SenderNodeID: c.NodeID, @@ -509,16 +669,16 @@ func (c *WebSocketController) sendMessageToClient(clientID string, msg models.Cl redisMsgJSON, err := json.Marshal(redisMsg) if err != nil { - log.Printf("❌❌❌❌ Redis消息序列化失败: %v", err) + log.Printf("❌❌❌❌❌❌❌❌ Redis消息序列化失败: %v", err) return fmt.Errorf("内部错误") } if err := c.RedisCli.Publish(c.RedisCtx, "ws_messages", redisMsgJSON).Err(); err != nil { - log.Printf("❌❌❌❌ 发布消息到Redis失败: %v", err) + log.Printf("❌❌❌❌❌❌❌❌ 发布消息到Redis失败: %v", err) return fmt.Errorf("无法转发消息") } - log.Printf("📡📡📡📡 消息已转发到Redis: \n ClientID=%s \n Size=%d bytes \n Duration=%s", + log.Printf("📡📡📡📡📡📡📡📡 消息已转发到Redis: \n ClientID=%s \n Size=%d bytes \n Duration=%s", clientID, len(redisMsgJSON), time.Since(start)) return nil } @@ -529,28 +689,33 @@ func (c *WebSocketController) SubscribeToRedis() { defer pubsub.Close() ch := pubsub.Channel() - log.Println("🔔🔔🔔🔔 开始监听Redis消息...") + log.Println("🔔🔔🔔🔔🔔🔔🔔🔔 开始监听Redis消息...") 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) + log.Printf("❌❌❌❌❌❌❌❌ 解析Redis消息失败: %v", err) continue } - log.Printf("📥📥📥📥 收到Redis消息: \n Sender=%s \n ClientID=%s \n Size=%d bytes", + 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("ℹℹℹℹ️ 忽略本节点转发的消息") + log.Println("ℹℹℹℹℹℹℹℹ️ 忽略本节点转发的消息") continue } if conn := c.getClient(redisMsg.ClientID); conn != nil { - if err := conn.WriteMessage(websocket.TextMessage, []byte(redisMsg.Message)); err != nil { - log.Printf("❌❌❌❌ 处理Redis消息失败: %v | ClientID=%s", err, redisMsg.ClientID) + // 加写锁保护 + 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) continue } @@ -570,7 +735,7 @@ 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", + log.Printf("📌📌📌📌 添加客户端到连接池: \n ClientID=%s \n 当前连接数=%d", clientID, len(c.Clients)) } @@ -584,7 +749,7 @@ func (c *WebSocketController) removeClient(clientID string) { } delete(c.Clients, clientID) - log.Printf("🗑🗑️ 从连接池移除客户端: \n ClientID=%s \n 当前连接数=%d", + log.Printf("🗑🗑🗑🗑️ 从连接池移除客户端: \n ClientID=%s \n 当前连接数=%d", clientID, len(c.Clients)) if err := c.RedisCli.HDel(c.RedisCtx, models.ClientUserKey, clientID).Err(); err != nil { @@ -598,3 +763,14 @@ func (c *WebSocketController) getClient(clientID string) *websocket.Conn { 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 aa6bdbc..23e94c2 100644 --- a/models/messages.go +++ b/models/messages.go @@ -6,7 +6,7 @@ type SendMessagePayload struct { TargetClientID string `json:"target_client_id"` // 目标客户端ID SenderUserID string `json:"sender_user_id"` // 发送者用户ID ReceiverUserID string `json:"receiver_user_id"` // 接收者用户ID - MessageType int `json:"message_type"` // 消息类型(0-6) + MessageType int `json:"message_type"` // 消息类型(0-8) MessageContent string `json:"message_content"` // 消息内容 } @@ -14,7 +14,7 @@ type SendMessagePayload struct { type SendToUserPayload struct { SenderUserID string `json:"sender_user_id"` // 发送者用户ID ReceiverUserID string `json:"receiver_user_id"` // 接收者用户ID - MessageType int `json:"message_type"` // 消息类型(0-6) + MessageType int `json:"message_type"` // 消息类型(0-8) MessageContent string `json:"message_content"` // 消息内容 } @@ -24,9 +24,13 @@ type ClientReceivedMessage struct { ReceiverID string `json:"receiver_id"` // 接收者连接ID SenderUserID string `json:"sender_user_id"` // 发送者用户ID ReceiverUserID string `json:"receiver_user_id"` // 接收者用户ID - MessageType int `json:"message_type"` // 消息类型(0-6) + MessageType int `json:"message_type"` // 消息类型(0-8) Content string `json:"content"` // 消息内容 SendTime string `json:"send_time"` // 发送时间 + + // 通话专用字段 + CallID string `json:"call_id,omitempty"` // 通话唯一ID + CallStatus string `json:"call_status,omitempty"` // 通话状态:invite/accepted/rejected/ended/candidate } // Redis 消息结构 @@ -42,6 +46,16 @@ type BindRequest struct { ClientID string `json:"client_id"` // 客户端ID } +// 通话信令结构 +type CallSignal struct { + Action string `json:"action"` // 动作: invite/accept/reject/end/candidate + CallType int `json:"call_type"` // 通话类型:6=视频,7=语音 + CallID string `json:"call_id"` // 通话唯一ID + CallerID string `json:"caller_id"` // 主叫用户ID + CalleeID string `json:"callee_id"` // 被叫用户ID + Data string `json:"data"` // 数据(SDP/ICE候选) +} + // Redis键常量 const ( ClientUserKey = "client_user_mapping" // clientid -> userid映射 diff --git a/utils/gen_call_id.go b/utils/gen_call_id.go new file mode 100644 index 0000000..f7d0063 --- /dev/null +++ b/utils/gen_call_id.go @@ -0,0 +1,18 @@ +package utils + +import ( + "crypto/rand" + "encoding/hex" + "fmt" + "time" +) + +// GenerateCallID 生成16字节的唯一通话ID +func GenerateCallID() string { + b := make([]byte, 8) + if _, err := rand.Read(b); err != nil { + // 如果随机数生成失败,使用时间戳作为后备 + return fmt.Sprintf("%d", time.Now().UnixNano()) + } + return hex.EncodeToString(b) +} diff --git a/xk-websocket.exe b/xk-websocket.exe new file mode 100644 index 0000000..3296082 Binary files /dev/null and b/xk-websocket.exe differ