Compare commits

...
2 Commits
Author SHA1 Message Date
Seven 74ed1535ed 增加ack确认消息机制 2026-09-23 15:45:27 +08:00
Seven aedd5d1f3f 优化websocket监听消息 2026-09-23 10:20:56 +08:00
+81 -8
View File
@@ -1,5 +1,6 @@
const WebSocket = require("ws"); const WebSocket = require("ws");
const url = require("url"); const url = require("url");
const db = require("../db/mysql");
const { verifyToken } = require("../utils/jwt"); const { verifyToken } = require("../utils/jwt");
@@ -24,27 +25,22 @@ function initWebSocket(server) {
if (!token) { if (!token) {
console.log("没有提供token"); console.log("没有提供token");
ws.close(); ws.close();
return; return;
} }
// 验证JWT // 验证JWT
let user; let user;
try { try {
user = verifyToken(token); user = verifyToken(token);
} catch (error) { } catch (error) {
console.log("JWT验证失败:", error.message); console.log("JWT验证失败:", error.message);
ws.close(); ws.close();
return; return;
} }
console.log("WebSocket用户:", user); console.log("WebSocket用户:", user);
const username = user.username; const username = user.username;
// 保存连接 // 保存连接
addConnection(username, ws); addConnection(username, ws);
@@ -52,7 +48,6 @@ function initWebSocket(server) {
ws.on("message", (message) => { ws.on("message", (message) => {
try { try {
const data = JSON.parse(message.toString()); const data = JSON.parse(message.toString());
console.log("收到WebSocket消息:", data); console.log("收到WebSocket消息:", data);
handleMessage(username, data); handleMessage(username, data);
@@ -83,6 +78,9 @@ function handleMessage(username, data) {
case "sendMessage": case "sendMessage":
handleSendMessage(username, data); handleSendMessage(username, data);
break; break;
case "ack":
handleAck(data);
break;
default: default:
console.log("未知消息类型:", data.type); console.log("未知消息类型:", data.type);
@@ -91,6 +89,7 @@ function handleMessage(username, data) {
// 处理发送消息 // 处理发送消息
function handleSendMessage(senderName, data) { function handleSendMessage(senderName, data) {
const messageId = data.messageId;
const receiverName = data.receiverName; const receiverName = data.receiverName;
const content = data.content; const content = data.content;
@@ -98,9 +97,49 @@ function handleSendMessage(senderName, data) {
`${senderName} → ${receiverName}:${content}` `${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
})
);
}
// 找到接收者的WebSocket // 找到接收者的WebSocket
const receiverWs = getConnection(receiverName); const receiverWs = getConnection(receiverName);
if (!receiverWs) { if (!receiverWs) {
console.log(`${receiverName} 当前不在线`); console.log(`${receiverName} 当前不在线`);
return; return;
@@ -111,12 +150,46 @@ function handleSendMessage(senderName, data) {
JSON.stringify({ JSON.stringify({
type: "newMessage", type: "newMessage",
message: { message: {
messageId,
senderName, senderName,
receiverName, receiverName,
content, 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; module.exports = initWebSocket;