const WebSocket = require("ws"); const url = require("url"); const db = require("../db/mysql"); const { verifyToken } = require("../utils/jwt"); const { addConnection, getConnection, removeConnection, getOnlineUsers, broadcast } = require("./connectionManager"); // 创建WebSocket服务 function initWebSocket(server) { const wss = new WebSocket.Server({ server, }); wss.on("connection", (ws, request) => { console.log("有新的WebSocket连接"); // 获取URL中的token const query = url.parse(request.url, true).query; const token = query.token; if (!token) { console.log("没有提供token"); ws.close(); return; } // 验证JWT let user; try { user = verifyToken(token); } catch (error) { console.log("JWT验证失败:", error.message); ws.close(); return; } console.log("WebSocket用户:", user); const username = user.username; // 保存连接 addConnection(username, ws); // 通知所有在线用户当前在线用户列表 broadcast({ type: "userOnline", userList: getOnlineUsers() }); // 接收客户端消息 ws.on("message", (message) => { try { const data = JSON.parse(message.toString()); console.log("收到WebSocket消息:", data); handleMessage(username, data); } catch (error) { console.error("WebSocket消息处理失败:", error); } }); // 连接关闭 ws.on("close", () => { removeConnection(username); // 通知所有在线用户当前在线用户列表 broadcast({ type: "userOnline", userList: getOnlineUsers() }); }); // WebSocket错误 ws.on("error", (error) => { console.error(`用户 ${username} WebSocket错误:`, error); }); }); console.log("WebSocket服务已启动"); return wss; } // 处理WebSocket消息 function handleMessage(username, data) { switch (data.type) { case "sendMessage": handleSendMessage(username, data); break; case "ack": handleAck(data); break; default: console.log("未知消息类型:", data.type); } } // 处理发送消息 function handleSendMessage(senderName, data) { const messageId = data.messageId; const receiverName = data.receiverName; const content = data.content; console.log( `${senderName} → ${receiverName}:${content}` ); const sql = ` INSERT INTO messages (message_id, sender_name, receiver_name, content) VALUES (?, ?, ?, ?) `; db.query(sql, [ messageId, senderName, receiverName, content ], (error, result) => { if (error) { console.error(error); const senderWs = getConnection(senderName); if (senderWs) { senderWs.send( JSON.stringify({ type: "sendError", messageId: messageId, message: "消息保存失败" }) ); } return; } // 数据库确认写入后,才确认发送并推送给接收者 const senderWs = getConnection(senderName); if (senderWs) { senderWs.send( JSON.stringify({ type: "sendSuccess", messageId: messageId }) ); } const receiverWs = getConnection(receiverName); if (!receiverWs) { console.log(`${receiverName} 当前不在线`); return; } receiverWs.send( JSON.stringify({ type: "newMessage", message: { messageId, senderName, receiverName, content }, }) ); }); } function handleAck(receiverName, data) { const messageId = data.messageId; const senderName = data.senderName; console.log( `${receiverName} 收到消息 ${messageId},发送 ACK` ); // 找到原发送者 const senderWs = getConnection(senderName); if (!senderWs) { console.log( `${senderName} 当前不在线` ); return; } // 把 ACK 转给原发送者 senderWs.send( JSON.stringify({ type: "ack", messageId: messageId }) ); console.log( `ACK ${messageId} 已发送给 ${senderName}` ); } module.exports = initWebSocket;