Files
chat-demo-server/websocket/index.js
T

254 lines
5.3 KiB
JavaScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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, ws);
// 通知所有在线用户当前在线用户列表
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;