Compare commits

..
5 Commits
7 changed files with 408 additions and 183 deletions
+208
View File
@@ -0,0 +1,208 @@
# Chat Demo Server 设计说明
## 1. 文档范围
本文依据当前仓库中的代码和 SQL 脚本,记录聊天服务端现有结构、模块职责、接口和消息协议、数据模型及设计取舍。文中“当前行为”指代码已经实现的行为;代码中存在但可能导致运行失败或安全风险的部分会在“已知问题与建议”中单独说明,不将计划性改进描述为已有能力。
## 2. 项目概述
本项目是一个基于 Node.js 的单进程聊天服务端,提供 HTTP REST 接口和 WebSocket 实时连接。Express 和 WebSocket 共用同一个 HTTP Server,监听 3000 端口;MySQL 保存用户、联系人和聊天消息。
当前实现的主要能力:
- 用户注册、用户名密码登录,并在登录成功后签发 JWT。
- 使用 JWT 中间件保护用户信息、联系人和聊天记录查询接口。
- 查询联系人、每个联系人的最近一条消息,以及指定联系人的分页消息。
- WebSocket 连接通过 JWT 验证身份,维护进程内在线用户连接表,并广播在线列表。
- WebSocket 消息处理代码包含消息发送、投递和 ACK 转发逻辑。
## 3. 技术栈与依赖
| 技术/依赖 | 用途 |
| --- | --- |
| Node.js | JavaScript 服务端运行环境 |
| Express 5 | HTTP 服务、JSON 请求体解析和路由分发 |
| `http` | 创建与 Express 共用的 HTTP Server |
| `ws` | WebSocket 服务端实现 |
| `mysql2` | MySQL 单连接及参数化 SQL 查询 |
| `jsonwebtoken` | JWT 签发和验证 |
| `dotenv` | 将 `.env` 配置加载到 `process.env` |
| `nodemon` | 开发时监视文件并重启服务 |
项目使用 CommonJS 模块格式。`npm start` 执行 `node index.js`,`npm run dev` 执行 `nodemon index.js`。当前 `npm test` 是占位脚本,没有自动化测试。
## 4. 代码结构
```text
index.js 服务启动入口、HTTP 与 WebSocket 组装
db/
mysql.js MySQL 连接创建与导出
sql.txt users、contacts、messages 的建表及样例数据脚本
middleware/
authMiddleware.js HTTP 请求 JWT 校验
routes/
login.js 注册和登录接口
user.js 用户信息与联系人接口
message.js 最近聊天、历史消息、REST 消息写入接口
utils/
jwt.js JWT 签发和验证封装
websocket/
index.js WebSocket 握手、协议分发和消息处理
connectionManager.js 进程内在线连接管理及广播
```
### 模块职责与依赖
- `index.js`:加载环境变量,创建 Express 应用和 HTTP Server,启用 JSON 请求解析,注册 `/api` 路由并将同一个 Server 交给 WebSocket 初始化函数,最后监听 3000 端口。
- `db/mysql.js`:创建一个 `mysql2` 单连接,并作为共享模块导出,路由和 WebSocket 处理器直接调用 `db.query`。
- `routes/login.js`:处理公开的 `/login`、`/register`。登录成功后通过 JWT 工具生成 token。
- `routes/user.js`:依赖 `authMiddleware`,从 `req.user.username` 识别当前用户并提供用户信息与联系人查询。
- `routes/message.js`:历史记录及最近聊天查询依赖认证中间件;当前 REST `/sendMessage` 未挂载认证中间件。
- `middleware/authMiddleware.js`:从 `Authorization` 请求头拆出 token,调用 JWT 验证,将 payload 放入 `req.user` 后继续处理请求。
- `utils/jwt.js`:使用 `JWT_SECRET` 签发和验证 token,签发有效期为 10 分钟。
- `websocket/index.js`:从连接 URL 查询参数读取 token 验证客户端身份;连接建立或关闭时广播在线用户列表;按消息 `type` 分发 WebSocket 消息。
- `websocket/connectionManager.js`:使用模块级 `Map` 以用户名为键保存 WebSocket 连接。该状态仅存在于当前 Node 进程内。
- `db/sql.txt`:描述用户、联系人、消息表,以及本地演示用数据。
## 5. 总体架构与请求路径
```mermaid
flowchart LR
Client[客户端]
Entry[index.js / HTTP Server]
Express[Express 路由]
Auth[JWT 中间件与工具]
WS[WebSocket 服务]
Manager[在线连接 Map]
DB[(MySQL)]
Client -->|HTTP /api| Entry
Entry --> Express
Express --> Auth
Express --> DB
Client <-->|WebSocket| Entry
Entry --> WS
WS --> Auth
WS --> Manager
WS --> DB
```
HTTP 请求由 Express 根据 `/api` 前缀分发。需要认证的路由先运行 JWT 中间件,随后执行业务 SQL。WebSocket 则在握手后从 URL 查询参数读取 token,再在连接生命周期内以已验证的用户名处理消息。HTTP 与 WebSocket 最终都使用同一 MySQL 模块。
## 6. 数据模型
数据库脚本目标库为 `chat_mysql`,定义以下三张表:
### `users`
| 字段 | 类型 | 约束/用途 |
| --- | --- | --- |
| `username` | `VARCHAR(50)` | 主键,用户唯一标识 |
| `password` | `VARCHAR(255)` | 非空;当前代码保存并直接比较明文密码 |
### `contacts`
| 字段 | 类型 | 约束/用途 |
| --- | --- | --- |
| `user_name` | `VARCHAR(50)` | 联系人所属用户,外键引用 `users.username` |
| `contact_name` | `VARCHAR(50)` | 联系人用户名,外键引用 `users.username` |
`(user_name, contact_name)` 为联合主键。同一用户不能重复添加同一联系人。脚本中的联系人关系按有向关系保存,双向关系需要分别插入两行。
### `messages`
| 字段 | 类型 | 约束/用途 |
| --- | --- | --- |
| `id` | `INT` | 自增主键 |
| `sender_name` | `VARCHAR(50)` | 非空,外键引用发送者 |
| `receiver_name` | `VARCHAR(50)` | 非空,外键引用接收者 |
| `content` | `VARCHAR(1000)` | 非空,消息正文 |
| `created_at` | `DATETIME` | 默认使用数据库当前时间 |
消息以一条记录表示一次单向发送。双方聊天记录通过发送者和接收者字段组合查询;表中没有当前 WebSocket 处理器所写入的 `message_id` 字段。
## 7. HTTP 接口设计
所有路由统一挂载在 `/api` 下,JSON 请求体由 Express 解析。成功/失败响应大多采用 `{ code, message, data }` 结构,但不同接口的 HTTP 状态码与业务 `code` 使用并不完全一致。
| 方法 | 路径 | 认证 | 输入/用途 |
| --- | --- | --- | --- |
| `POST` | `/api/register` | 否 | body:`username`、`password`;插入用户 |
| `POST` | `/api/login` | 否 | body:`username`、`password`;验证凭据并返回 token、用户名 |
| `GET` | `/api/userInfo` | 是 | 从 token 返回当前用户名 |
| `GET` | `/api/contactsInfo` | 是 | 查询当前用户联系人,返回 `contacts` 数组 |
| `GET` | `/api/recentChats` | 是 | 查询当前用户每个聊天对象的最新消息 |
| `POST` | `/api/contactMessagesInfo` | 是 | body:`contactName`、`limit`、`offset`;分页查询双方消息 |
| `POST` | `/api/sendMessage` | **否** | body:`senderName`、`receiverName`、`content`;将消息插入数据库 |
### 认证约定
HTTP 客户端应在受保护接口的请求头中传入 `Authorization: Bearer <token>`。登录接口响应的 token 字符串自身已带 `Bearer ` 前缀。当前中间件通过空格拆分请求头并取第二段作为 token;缺少请求头或验证失败时返回 401。
JWT payload 当前包含 `username`,过期时间为 10 分钟。WebSocket 客户端把 token 放在连接 URL 查询参数 `token` 中。由于登录响应包含 `Bearer ` 前缀,而 WebSocket 端会直接将查询参数交给 `verifyToken`,客户端需要传纯 JWT 部分,不能把 `Bearer ` 前缀一并作为 token 值。
## 8. WebSocket 协议与消息流程
WebSocket 挂载在 HTTP Server 上,客户端连接地址的形式为 `ws://<host>:3000/?token=<JWT>`。连接验证通过后,服务端按用户名登记连接,并广播在线列表。
### 服务端消息类型
| `type` | 方向 | 当前代码行为 |
| --- | --- | --- |
| `userOnline` | 服务端 -> 所有在线客户端 | 连接建立或关闭后,携带 `userList` 在线用户名数组 |
| `newMessage` | 服务端 -> 接收方 | 携带 `message` 对象,包括消息 ID、发送方、接收方和正文 |
| `sendSuccess` | 服务端 -> 发送方 | 携带客户端提供的 `messageId`,表示发送成功 |
| `sendError` | 服务端 -> 发送方 | 数据库写入回调报错时尝试发送,携带 `messageId` 和错误说明 |
| `ack` | 服务端 -> 原发送方 | 将接收方确认的消息 ID 转发给发送方 |
### 客户端消息类型
- `sendMessage`:处理器读取 `messageId`、`receiverName`、`content`;发送者身份来自已验证 WebSocket 连接,不从消息体读取。
- `ack`:处理器预期使用 `messageId` 和 `senderName` 将 ACK 转发给原发送者。
### 预期发送流程
1. 客户端建立携带 JWT 的 WebSocket 连接。
2. 服务端校验 JWT,登记用户名与连接,并广播在线用户列表。
3. 客户端发送 `sendMessage`,服务端使用连接身份确定发送者。
4. 服务端尝试把消息持久化到 MySQL;接收方在线时向其推送 `newMessage`。
5. 接收方返回 `ack` 后,服务端尝试通知原发送方。
步骤 4 当前实现没有等待数据库回调成功再确认和推送;此外建表脚本与插入语句字段不匹配,因此此流程目前不能视为可靠送达流程。具体差异见第 10 节。
## 9. 关键设计思路
- **按职责拆分模块**:入口负责组装,路由负责 HTTP 接口,middleware 负责认证,utils 封装 JWT,websocket 目录处理长连接和连接表,db 统一导出数据库连接。
- **用户名作为用户标识**:用户表以用户名为主键,联系人关系、消息发送接收双方及在线连接表均使用用户名关联。
- **JWT 无服务端会话存储**:服务端通过签名和过期时间验证 token,HTTP 请求和 WebSocket 握手各自执行验证。
- **聊天列表与历史消息分开查询**:`recentChats` 使用窗口函数按聊天对象分组并选出最新消息;历史消息接口按 limit/offset 分页,避免单次返回全部记录。
- **在线状态在内存维护**:当前使用 Map 直接由用户名定位 WebSocket 连接,便于单进程内向在线用户推送消息,但没有跨进程共享能力。
- **SQL 参数化**:业务查询使用 `?` 占位符与参数数组传值,避免把用户输入直接拼入 SQL 文本。
## 10. 已知问题与后续建议
以下内容是基于当前代码与脚本的直接核对结果,应在扩展功能或对外部署前处理:
1. **数据库凭据写在源码中**:`db/mysql.js` 将 MySQL 账号、密码和数据库名直接传给连接创建函数。建议改为读取环境变量,轮换已经放入源码的凭据,并避免将真实密钥写入文档或版本库。
2. **密码以明文存储和比对**:注册直接写入密码,登录 SQL 按用户名与明文密码匹配。应使用成熟密码哈希方案存储和校验,并考虑统一认证失败响应。
3. **REST 消息写入缺少认证**:`/api/sendMessage` 未使用 `authMiddleware`,客户端可自行指定发送者。应要求认证并从 `req.user` 派生发送者身份,同时校验接收方及联系人关系。(已解决)
4. **WebSocket 数据库写入字段不匹配**:`websocket/index.js` 插入 `message_id`,但 `db/sql.txt` 的 `messages` 表没有该列。需统一消息 ID 设计,例如新增唯一字段及迁移,或使用现有自增 `id` 并调整协议和查询。(已解决)
5. **发送确认早于持久化结果**:`sendSuccess` 和接收方推送发生在异步 `db.query` 回调之前。应将成功确认和实时推送放到成功回调中;失败时只返回错误,不推送为已发送消息。(已解决)
6. **连接管理器广播引用未定义的 `WebSocket`**:`connectionManager.js` 的 `broadcast` 检查 `WebSocket.OPEN`,但文件中没有导入 `ws`。当广播函数被调用时可能抛出 `ReferenceError`。应在此模块引入库常量,或由连接对象状态采用明确且可用的检查方式。(已解决)
7. **ACK 参数调用不一致**:消息分发通过 `handleAck(data)` 调用,而函数签名是 `handleAck(receiverName, data)`;函数内部随后读取第二个参数,可能因 `data` 未定义而失败。应统一函数参数并以当前连接用户名作为 ACK 接收者身份。(已解决)
8. **重复登录连接的关闭竞态**:同一用户名的新连接会覆盖 Map 中旧连接;旧连接关闭时无条件按用户名删除,可能把新连接也从 Map 移除。移除连接时应确认 Map 中仍是即将关闭的那个 WebSocket。(已解决)
9. **分页和输入校验不足**:`contactMessagesInfo` 未验证 `contactName`、`limit`、`offset` 的类型与范围,也未约束最大页大小。建议验证请求参数并设置默认值和上限。
10. **在线状态及可靠性受单进程限制**:Map 不支持多实例共享、服务重启恢复、离线消息投递或持久化 ACK 状态。若部署多实例,需要共享在线状态/消息协调机制;若要求可靠消息,应定义消息状态、幂等键、重试及离线投递策略。
11. **配置与错误响应尚未统一**:数据库连接配置没有从 `.env` 读取;HTTP 接口混用 HTTP 状态码和响应体 `code`,空结果有时返回业务码 401。建议集中配置并统一 API 错误语义及日志策略。
12. **SQL 初始化脚本需要整理验证**:脚本包含清空用户表的语句,执行前会删除该表全部用户数据。应将演示数据与建表迁移分离,并在测试数据库验证脚本。
13. **缺少自动化测试**:当前测试脚本是占位命令。建议优先为认证边界、参数校验、消息持久化失败/成功和 WebSocket 连接生命周期补充测试。
## 11. 本地运行与部署边界
1. 安装 Node.js 及 MySQL,并在 MySQL 中准备 `chat_mysql` 数据库。
2. 按需执行经过校验的建表脚本;注意当前 `sql.txt` 含删除数据语句,不应直接用于已有数据环境。
3. 设置 `JWT_SECRET`。当前 MySQL 参数仍需在 `db/mysql.js` 中配置;建议先将连接参数迁移到环境变量,再部署。
4. 执行 `npm install`,开发环境使用 `npm run dev`,普通启动使用 `npm start`。
5. 服务监听 `http://localhost:3000`,根路径 `/` 返回运行提示;WebSocket 使用相同主机和端口。
当前代码适用于本地演示和单进程验证,不具备生产环境所需的凭据管理、密码安全、输入校验、多实例在线状态和消息可靠性保障。
+8 -101
View File
@@ -11,7 +11,7 @@ VALUES
('query', '123'),
('tom', '123'),
('jack', '123'),
('lucy', '123'),
('lucy', '123'),
('lily', '123');
select * from chat_mysql.users;
@@ -37,13 +37,14 @@ VALUES
('query', 'admin'),
('query', 'tom'),
('tom', 'admin'),
('tom', 'query'),
('tom', 'query'),
('admin', 'lily'),
('lily', 'admin');
select * from chat_mysql.contacts;
CREATE TABLE chat_mysql.messages (
id INT PRIMARY KEY AUTO_INCREMENT,
message_id VARCHAR(64) NOT NULL,
sender_name VARCHAR(50) NOT NULL,
receiver_name VARCHAR(50) NOT NULL,
content VARCHAR(1000) NOT NULL,
@@ -55,103 +56,9 @@ CREATE TABLE chat_mysql.messages (
FOREIGN KEY (receiver_name)
REFERENCES chat_mysql.users(username)
);
INSERT INTO chat_mysql.messages
(sender_name, receiver_name, content, created_at)
VALUES
-- admin ↔ query
('admin', 'query', '你好,最近怎么样?', '2026-09-20 10:00:00'),
('query', 'admin', '挺好的,你呢?', '2026-09-20 10:01:00'),
('admin', 'query', '我也不错,最近在做聊天项目。', '2026-09-20 10:02:00'),
('query', 'admin', '听起来不错,是用 Vue 做的吗?', '2026-09-20 10:03:00'),
('admin', 'query', '对,前端使用 Vue2,后端使用 Node.js。', '2026-09-20 10:04:00'),
('query', 'admin', '那数据库使用 MySQL?', '2026-09-20 10:05:00'),
('admin', 'query', '对,目前正在完善聊天功能。', '2026-09-20 10:06:00'),
-- admin ↔ tom
('admin', 'tom', '你今天有空吗?', '2026-09-20 11:00:00'),
('tom', 'admin', '有啊,怎么了?', '2026-09-20 11:01:00'),
('admin', 'tom', '想和你讨论一下项目。', '2026-09-20 11:02:00'),
('tom', 'admin', '可以,什么时间方便?', '2026-09-20 11:03:00'),
('admin', 'tom', '下午三点怎么样?', '2026-09-20 11:04:00'),
('tom', 'admin', '没问题。', '2026-09-20 11:05:00'),
-- admin ↔ jack
('jack', 'admin', '周末一起吃饭吗?', '2026-09-20 12:00:00'),
('admin', 'jack', '可以啊,去哪吃?', '2026-09-20 12:01:00'),
('jack', 'admin', '附近找一家餐厅吧。', '2026-09-20 12:02:00'),
('admin', 'jack', '好的,晚上六点怎么样?', '2026-09-20 12:03:00'),
('jack', 'admin', '可以,到时候联系。', '2026-09-20 12:04:00'),
-- admin ↔ lucy
('admin', 'lucy', '最近项目进展怎么样?', '2026-09-20 13:00:00'),
('lucy', 'admin', '目前进展还不错。', '2026-09-20 13:01:00'),
('admin', 'lucy', '还有哪些功能没有完成?', '2026-09-20 13:02:00'),
('lucy', 'admin', '主要是消息通知和联系人功能。', '2026-09-20 13:03:00'),
('admin', 'lucy', '好的,我这边也在做消息功能。', '2026-09-20 13:04:00'),
-- query ↔ tom
('query', 'tom', '你知道聊天项目怎么做 WebSocket 吗?', '2026-09-20 14:00:00'),
('tom', 'query', '知道一些,可以用来做实时消息。', '2026-09-20 14:01:00'),
('query', 'tom', '我也准备把项目改成 WebSocket。', '2026-09-20 14:02:00'),
('tom', 'query', '这样聊天消息可以实时推送。', '2026-09-20 14:03:00'),
-- tom ↔ query
('tom', 'query', '联系人比较多的时候怎么办?', '2026-09-20 15:00:00'),
('query', 'tom', '不需要每个联系人都创建一个 WebSocket。', '2026-09-20 15:01:00'),
('tom', 'query', '一个用户维护一个 WebSocket 连接就可以了。', '2026-09-20 15:02:00');
select * from chat_mysql.messages;
INSERT INTO chat_mysql.messages
(sender_name, receiver_name, content, created_at)
VALUES
-- admin ↔ lily
('admin', 'lily', '你好,最近怎么样?', '2026-09-20 10:00:00'),
('lily', 'admin', '挺好的,你呢?', '2026-09-20 10:01:00'),
('admin', 'lily', '我也不错,最近在做聊天项目。', '2026-09-20 10:02:00'),
('lily', 'admin', '听起来不错,是用 Vue 做的吗?', '2026-09-20 10:03:00'),
('admin', 'lily', '对,前端使用 Vue2,后端使用 Node.js。', '2026-09-20 10:04:00'),
('lily', 'admin', '那数据库使用 MySQL?', '2026-09-20 10:05:00'),
('admin', 'lily', '对,目前正在完善聊天功能。', '2026-09-20 10:06:00'),
('lily', 'admin', '聊天记录是怎么保存的?', '2026-09-20 10:07:00'),
('admin', 'lily', '目前是保存在 messages 表里面。', '2026-09-20 10:08:00'),
('lily', 'admin', '一个消息对应一条数据吗?', '2026-09-20 10:09:00'),
('admin', 'lily', '对,每条消息都有自己的 id。', '2026-09-20 10:10:00'),
('lily', 'admin', '那发送人和接收人怎么区分?', '2026-09-20 10:11:00'),
('admin', 'lily', '通过 sender_name 和 receiver_name 两个字段区分。', '2026-09-20 10:12:00'),
('lily', 'admin', '这样查询两个人的聊天记录应该比较方便。', '2026-09-20 10:13:00'),
('admin', 'lily', '对,只需要判断双方的用户名即可。', '2026-09-20 10:14:00'),
('lily', 'admin', '比如查询 admin 和 lily 之间的消息?', '2026-09-20 10:15:00'),
('admin', 'lily', '没错,可以使用 OR 和 AND 组合条件。', '2026-09-20 10:16:00'),
('lily', 'admin', '那如果消息很多怎么办?', '2026-09-20 10:17:00'),
('admin', 'lily', '不能一次把所有消息都返回给前端。', '2026-09-20 10:18:00'),
('lily', 'admin', '应该使用分页查询吧?', '2026-09-20 10:19:00'),
('admin', 'lily', '对,比如一次只查询 10 条或者 20 条。', '2026-09-20 10:20:00'),
('lily', 'admin', '这样可以减少接口返回的数据量。', '2026-09-20 10:21:00'),
('admin', 'lily', '也可以提高页面加载速度。', '2026-09-20 10:22:00'),
('lily', 'admin', '最近聊天列表也是单独查询吗?', '2026-09-20 10:23:00'),
('admin', 'lily', '对,最近聊天只需要获取每个联系人的最新消息。', '2026-09-20 10:24:00'),
('lily', 'admin', '这样左边就不用加载完整聊天记录了。', '2026-09-20 10:25:00'),
('admin', 'lily', '对,点击联系人之后再加载具体的聊天记录。', '2026-09-20 10:26:00'),
('lily', 'admin', '这个设计比较合理。', '2026-09-20 10:27:00'),
('admin', 'lily', '我现在就是准备按照这个思路修改。', '2026-09-20 10:28:00'),
('lily', 'admin', '那 Vue 里面是不是也不用保存所有消息?', '2026-09-20 10:29:00'),
('admin', 'lily', '是的,只需要保存当前联系人的消息数组。', '2026-09-20 10:30:00'),
('lily', 'admin', '分页的时候再把历史消息添加进去?', '2026-09-20 10:31:00'),
('admin', 'lily', '对,可以把下一页的数据添加到数组前面。', '2026-09-20 10:32:00'),
('lily', 'admin', '那新发送的消息就添加到数组最后面?', '2026-09-20 10:33:00'),
('admin', 'lily', '没错,因为新消息的时间是最新的。', '2026-09-20 10:34:00'),
('lily', 'admin', '左边的最近聊天还需要移动到最上面。', '2026-09-20 10:35:00'),
('admin', 'lily', '对,可以使用 unshift 把它移动到第一位。', '2026-09-20 10:36:00'),
('lily', 'admin', '同时更新左边显示的最后一条消息。', '2026-09-20 10:37:00'),
('admin', 'lily', '这样聊天列表就会实时更新了。', '2026-09-20 10:38:00'),
('lily', 'admin', '如果切换联系人,消息数组怎么办?', '2026-09-20 10:39:00'),
('admin', 'lily', '可以按照联系人 id 分别保存。', '2026-09-20 10:40:00'),
('lily', 'admin', '比如 contactMessages[18] 保存 lily 的消息?', '2026-09-20 10:41:00'),
('admin', 'lily', '对,contactMessages[20] 就可以保存另一个人的消息。', '2026-09-20 10:42:00'),
('lily', 'admin', 'Vue2 动态增加这个属性需要注意响应式问题吧?', '2026-09-20 10:43:00'),
('admin', 'lily', '对,可以使用 this.$set 来创建新的消息数组。', '2026-09-20 10:44:00'),
('lily', 'admin', '这样切换联系人就比较方便了。', '2026-09-20 10:45:00'),
('admin', 'lily', '对,聊天页面的数据结构基本就确定了。', '2026-09-20 10:46:00'),
('lily', 'admin', '那接下来就可以测试分页加载了。', '2026-09-20 10:47:00'),
('admin', 'lily', '嗯,先测试一次查询 10 条消息。', '2026-09-20 10:48:00'),
('lily', 'admin', '然后再查询下一页历史消息。', '2026-09-20 10:49:00');
#删除某个人的发送聊天数据
SELECT *
FROM chat_mysql.messages
WHERE sender_name = 'test';
DELETE FROM chat_mysql.messages WHERE sender_name = 'test';
+1
View File
@@ -70,6 +70,7 @@ router.post("/register", (req, res) => {
return res.status(500).json({
code: 500,
message: "注册失败",
data: null
});
}
+7 -1
View File
@@ -44,6 +44,7 @@ router.get(
return res.status(500).json({
code: 500,
message: "查询失败",
data: null
});
}
if (results.length > 0) {
@@ -97,6 +98,7 @@ router.post(
return res.status(500).json({
code: 500,
message: "查询失败",
data: null
});
}
if (results.length > 0) {
@@ -128,7 +130,10 @@ router.post(
);
// 写入聊天消息
router.post("/sendMessage", (req, res) => {
router.post(
"/sendMessage",
authMiddleware,
(req, res) => {
const message = req.body;
const sql = `
INSERT INTO messages
@@ -147,6 +152,7 @@ router.post("/sendMessage", (req, res) => {
return res.status(500).json({
code: 500,
message: "发送消息失败",
data: null
});
}
+2
View File
@@ -11,6 +11,7 @@ router.get(
(req, res) => {
res.json({
code: 200,
message: "查询成功",
data: {
username: req.user.username
}
@@ -40,6 +41,7 @@ router.get(
return res.status(500).json({
code: 500,
message: "查询失败",
data: null
});
}
if (results.length > 0) {
+14 -4
View File
@@ -1,3 +1,5 @@
const WebSocket = require("ws");
const onlineUsers = new Map();
// 保存用户连接
@@ -14,9 +16,13 @@ function getConnection(username) {
}
// 删除用户连接
function removeConnection(username) {
onlineUsers.delete(username);
function removeConnection(username, ws) {
const currentWs = onlineUsers.get(username);
// 只有 Map 里面存的还是这个连接,才允许删除
if (currentWs === ws) {
onlineUsers.delete(username);
}
console.log(`用户 ${username} 已断开`);
console.log("当前在线用户:", [...onlineUsers.keys()]);
}
@@ -30,9 +36,13 @@ function getOnlineUsers() {
function broadcast(data) {
const message = JSON.stringify(data);
for (const ws of onlineUsers.values()) {
for (const [username, ws] of onlineUsers) {
if (ws.readyState === WebSocket.OPEN) {
ws.send(message);
ws.send(message, (error) => {
if (error) {
console.error(`向用户 ${username} 广播失败:`, error);
}
});
}
}
}
+168 -77
View File
@@ -43,6 +43,7 @@ function initWebSocket(server) {
console.log("WebSocket用户:", user);
const username = user.username;
ws.username = username;
// 保存连接
addConnection(username, ws);
@@ -54,20 +55,21 @@ function initWebSocket(server) {
// 接收客户端消息
ws.on("message", (message) => {
try {
const data = JSON.parse(message.toString());
console.log("收到WebSocket消息:", data);
try {
const data = JSON.parse(message.toString());
console.log("收到WebSocket消息:", data);
handleMessage(username, data);
} catch (error) {
console.error("WebSocket消息处理失败:", error);
}
});
// websocket收到消息,必须通过解析data内容才知道谁是发送者谁是接收者
handleMessage(data, ws);
} catch (error) {
console.error("WebSocket消息处理失败:", error);
}
});
// 连接关闭
ws.on("close", () => {
removeConnection(username);
removeConnection(username, ws);
// 通知所有在线用户当前在线用户列表
broadcast({
type: "userOnline",
@@ -87,13 +89,13 @@ function initWebSocket(server) {
}
// 处理WebSocket消息
function handleMessage(username, data) {
function handleMessage(data, ws) {
switch (data.type) {
case "sendMessage":
handleSendMessage(username, data);
handleSendMessage(data, ws);
break;
case "ack":
handleAck(data);
handleAck(data, ws);
break;
default:
@@ -102,31 +104,46 @@ function handleMessage(username, data) {
}
// 处理发送消息
function handleSendMessage(senderName, data) {
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, [
const insertSql = `
INSERT INTO messages
(message_id, sender_name, receiver_name, content)
VALUES (?, ?, ?, ?)
`;
db.query(insertSql, [
messageId,
senderName,
receiverName,
content
], (error, result) => {
const senderWs = getConnection(senderName);
if (error) {
console.error(error);
const senderWs = getConnection(senderName);
if (senderWs) {
senderWs.send(
JSON.stringify({
@@ -138,70 +155,144 @@ function handleSendMessage(senderName, data) {
}
return;
}
});
// 告诉发送者:服务器已经保存成功
const senderWs = getConnection(senderName);
if (senderWs) {
// 插入新消息之后,用messageID查询这条新消息的时间点(为了使用数据库插入时创建的时间点,不能用前端自己生成的)
const timeSql = `
SELECT created_at
FROM messages
WHERE message_id = ?
`;
db.query(timeSql, [messageId], (error, results) => {
if (error) {
console.error("查询消息失败:", error);
if (senderWs) {
senderWs.send(
JSON.stringify({
type: "sendError",
messageId: messageId,
message: "获取消息失败"
})
);
}
return;
}
if (results.length === 0) {
if (senderWs) {
senderWs.send(
JSON.stringify({
type: "sendError",
messageId: messageId,
message: "获取消息失败"
})
);
}
return;
}
const createdAt = results[0].created_at;
// 数据库确认写入后,才把成功写入的消息推送给发送者
if (senderWs) {
senderWs.send(
JSON.stringify({
type: "sendSuccess",
message: {
messageId,
senderName,
receiverName,
content,
createdAt
}
})
);
}
// 把写入的新消息推送给接收者,如果接收者不在线,就不用发送websocket消息给他了,但数据库里是有的
const receiverWs = getConnection(receiverName);
if (!receiverWs) {
console.log(`${receiverName} 当前不在线`);
return;
}
receiverWs.send(
JSON.stringify({
type: "newMessage",
message: {
messageId,
senderName,
receiverName,
content,
createdAt
},
})
);
});
});
}
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: "sendSuccess",
type: "ack",
messageId: messageId
})
);
}
// 找到接收者的WebSocket
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} 当前不在线`
`ACK ${messageId} 已发送给 ${originalSender}`
);
return;
}
// 把 ACK 转给原发送者
senderWs.send(
JSON.stringify({
type: "ack",
messageId:
messageId
})
);
console.log(
`ACK ${messageId} 已发送给 ${senderName}`
);
});
}