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; ws.username = username; // 保存连接 addConnection(username, ws); // 通知所有在线用户当前在线用户列表 broadcast({ type: "userOnline", userList: getOnlineUsers() }); // 接收客户端消息 ws.on("message", (message) => { try { const data = JSON.parse(message.toString()); console.log("收到WebSocket消息:", data); // websocket收到消息,必须通过解析data内容才知道谁是发送者谁是接收者 handleMessage(data, ws); } 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(data, ws) { switch (data.type) { case "sendMessage": handleSendMessage(data, ws); break; case "ack": handleAck(data, ws); break; default: console.log("未知消息类型:", data.type); } } // 处理发送消息 function handleSendMessage(data, ws) { const senderName = ws.username; const messageId = data.messageId; const receiverName = data.receiverName; const content = data.content; if (!messageId) { console.log("消息缺少 messageId"); return; } if (!receiverName) { console.log("消息缺少接收者"); return; } if (!content) { console.log("消息内容为空"); return; } 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(data, ws) { const messageId = data.messageId; // 当前 WebSocket 对应的真实用户 const ackUser = ws.username; if (!messageId) { console.log("ACK 缺少 messageId"); return; } const sql = ` SELECT sender_name, receiver_name FROM messages WHERE message_id = ? `; db.query(sql, [messageId], (error, results) => { if (error) { console.error("查询消息失败:", error); return; } if (results.length === 0) { console.log(`消息 ${messageId} 不存在`); return; } const message = results[0]; // 当前用户必须是这条消息的接收者 if (message.receiver_name !== ackUser) { console.log( `${ackUser} 无权确认消息 ${messageId}` ); return; } // 原消息发送者 const originalSender = message.sender_name; console.log( `${ackUser} 已收到来自 ${originalSender} 的消息 ${messageId},发送 ACK` ); const senderWs = getConnection(originalSender); if (!senderWs) { console.log( `${originalSender} 当前不在线` ); return; } senderWs.send( JSON.stringify({ type: "ack", messageId: messageId }) ); console.log( `ACK ${messageId} 已发送给 ${originalSender}` ); }); } module.exports = initWebSocket;