Skip to content

Commit 55d2aa0

Browse files
committed
fix(mongodb): 修复并发连接、输入验证和内存泄漏问题
- 添加并发连接锁(_connecting),防止高并发下创建多个连接 - lib/index.js: 在上层添加连接锁 - lib/mongodb/index.js: 在 MongoDB 适配器层添加连接锁 - 为 collection() 方法添加参数校验 - 集合名必须是非空字符串,否则抛出 INVALID_COLLECTION_NAME - 数据库名(如提供)必须是非空字符串,否则抛出 INVALID_DATABASE_NAME - 修复 close() 方法未清理缓存导致的内存泄漏 - 清理 _iidCache(实例ID缓存) - 清理 _connecting(连接锁) - 新增完整的连接管理测试套件 - test/connection-simple.test.js: 5个核心测试(并发、验证、泄漏) - test/connection.test.js: 完整测试套件(28个测试用例) - 更新文档 - CHANGELOG.md: 添加修复说明和测试说明 - README.md: 新增连接管理章节(并发保护、参数验证、资源清理) 测试结果: ✅ 5/5 通过,耗时 0.24秒
1 parent 55e7640 commit 55d2aa0

7 files changed

Lines changed: 876 additions & 33 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 23 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,21 @@
33
所有显著变更将记录在此文件,遵循 Keep a Changelog 与语义化版本(SemVer)。
44

55
## [未发布]
6+
### 修复
7+
- **并发连接问题**:修复高并发调用 `connect()` 时创建多个连接的问题
8+
- 在 `lib/index.js` 和 `lib/mongodb/index.js` 两层添加连接锁(`_connecting`)
9+
- 并发请求现在会等待同一个连接 Promise,确保只建立一个连接
10+
- 添加完整的测试用例验证 10 个并发请求共享同一连接
11+
- **输入验证缺失**:为 `collection()` 方法添加参数校验
12+
- 集合名必须是非空字符串,否则抛出 `INVALID_COLLECTION_NAME` 错误
13+
- 数据库名(如果提供)必须是非空字符串,否则抛出 `INVALID_DATABASE_NAME` 错误
14+
- 添加友好的错误消息,包含参数要求说明
15+
- 添加完整的测试用例覆盖空字符串、null、纯空格、非字符串等边界情况
16+
- **内存泄漏**:修复 `close()` 方法未清理缓存的问题
17+
- 清理 `_iidCache`(实例 ID 缓存),防止多次连接-关闭循环累积内存
18+
- 清理 `_connecting` 锁,避免连接状态残留
19+
- 添加完整的测试用例验证多次连接-关闭循环无内存泄漏
20+
621
### 新增
722
- **文档增强**:为 `distinct` 方法添加完整的文档、示例和测试用例
823
- 新增 `docs/distinct.md`:详细的 distinct 方法使用文档,包含参数说明、使用模式、性能优化建议、常见问题等
@@ -33,7 +48,14 @@
3348
### 性能
3449
- totals:新增 5s 窗口的 inflight 去重,避免同一形状的并发计数击穿。
3550

51+
### 测试
52+
- **新增连接管理测试套件**:`test/connection-simple.test.js`
53+
- 验证并发连接只建立一个实例(10 并发测试)
54+
- 验证集合名和数据库名的输入校验(空、null、空格、非字符串)
55+
- 验证多次连接-关闭循环的资源清理(3 次循环测试)
56+
- 验证连接锁正确清理
57+
3658
### 说明
37-
- 推荐发布类型:`minor`(x.y.z -> x.(y+1).0),因新增能力向后兼容且为按需启用。
59+
- 推荐发布类型:`patch`(x.y.z -> x.y.(z+1)),因为是 bug 修复
3860

3961
[未发布]: ./

‎README.md‎

Lines changed: 107 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
- [invalidate(op) 用法](#invalidate)
2121
- [跨库访问注意事项](#cross-db)
2222
- [说明](#notes)
23+
- [连接管理](#连接管理)
2324
- [事件(Mongo)](#事件mongo)
2425
- [健康检查与事件(Mongo)](#健康检查与事件mongo)
2526

@@ -824,6 +825,112 @@ console.table(page2.meta.steps);
824825
> - 默认不返回 meta,需显式开启;开销很小,仅一次时间戳与对象组装。
825826
> - includeCache 仅包含去敏维度(如 cacheTtl 等,具体依实现)。
826827
828+
## 连接管理
829+
830+
### 并发连接保护
831+
`connect()` 方法内置并发锁机制,确保高并发场景下只建立一个连接:
832+
833+
```js
834+
const MonSQLize = require('monsqlize');
835+
const msq = new MonSQLize({
836+
type: 'mongodb',
837+
databaseName: 'example',
838+
config: { uri: 'mongodb://localhost:27017' },
839+
});
840+
841+
// 高并发场景:10 个并发请求
842+
const promises = Array(10).fill(null).map(() => msq.connect());
843+
const results = await Promise.all(promises);
844+
845+
// 所有请求返回同一个连接对象
846+
console.log(results[0] === results[1]); // true
847+
```
848+
849+
**特性**:
850+
- ✅ 首次调用建立连接,后续调用直接返回缓存的连接对象
851+
- ✅ 并发请求等待同一个 Promise,避免重复连接
852+
- ✅ 连接失败或成功后自动清理锁状态
853+
854+
### 参数验证
855+
`collection()` 和 `db()` 方法内置参数校验,确保接收合法参数:
856+
857+
```js
858+
const { collection, db } = await msq.connect();
859+
860+
// ✅ 正常使用
861+
const users = collection('users');
862+
const orders = db('shop').collection('orders');
863+
864+
// ❌ 无效参数(会抛出错误)
865+
try {
866+
collection(''); // 错误:INVALID_COLLECTION_NAME
867+
collection(null); // 错误:INVALID_COLLECTION_NAME
868+
collection(123); // 错误:INVALID_COLLECTION_NAME
869+
db('').collection('test'); // 错误:INVALID_DATABASE_NAME
870+
} catch (err) {
871+
console.error(err.code, err.message);
872+
}
873+
```
874+
875+
**验证规则**:
876+
- 集合名必须是**非空字符串**,不允许 null/undefined/空字符串/纯空格
877+
- 数据库名(如果提供)必须是**非空字符串**
878+
- 错误信息明确指出问题和要求
879+
880+
### 资源清理
881+
`close()` 方法会正确清理所有资源,防止内存泄漏:
882+
883+
```js
884+
const MonSQLize = require('monsqlize');
885+
const msq = new MonSQLize({
886+
type: 'mongodb',
887+
databaseName: 'example',
888+
config: { uri: 'mongodb://localhost:27017' },
889+
});
890+
891+
// 多次连接-关闭循环(安全)
892+
for (let i = 0; i < 5; i++) {
893+
await msq.connect();
894+
const { collection } = await msq.connect();
895+
896+
// 使用连接...
897+
await collection('test').find({ query: {} });
898+
899+
// 关闭连接
900+
await msq.close();
901+
// ✅ 内存已正确清理
902+
}
903+
```
904+
905+
**清理内容**:
906+
- ✅ 关闭 MongoDB 客户端连接
907+
- ✅ 清理实例 ID 缓存(`_iidCache`)
908+
- ✅ 清理连接锁(`_connecting`)
909+
- ✅ 释放所有内部引用
910+
911+
**注意事项**:
912+
- 多次调用 `close()` 是安全的,不会抛出错误
913+
- 关闭后再调用 `connect()` 会重新建立连接
914+
- 建议在应用关闭时调用 `close()` 释放资源
915+
916+
### 错误处理
917+
918+
```js
919+
try {
920+
const msq = new MonSQLize({
921+
type: 'mongodb',
922+
databaseName: 'example',
923+
config: { uri: 'mongodb://invalid-host:27017' },
924+
});
925+
926+
await msq.connect();
927+
} catch (err) {
928+
// 连接失败错误
929+
console.error('连接失败:', err.message);
930+
// ✅ 连接锁已自动清理,可以安全重试
931+
}
932+
```
933+
827934
## 事件(Mongo)
828935
- 事件基于 Node.js EventEmitter,进程内有效:
829936
- `connected`: `{ type, db, scope, iid? }`

‎lib/index.js‎

Lines changed: 32 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -69,7 +69,7 @@ module.exports = class {
6969
this.defaults = Object.freeze(this.defaults);
7070
}
7171

72-
/**
72+
/**
7373
* 连接数据库并返回访问集合/表的对象
7474
* @returns {{collection: Function, db: Function}} 返回包含 collection 与 db 方法的对象
7575
* @throws {Error} 当连接失败时抛出错误
@@ -79,22 +79,38 @@ module.exports = class {
7979
if (this.dbInstance) {
8080
return this.dbInstance;
8181
}
82+
83+
// 防止并发连接:使用连接锁
84+
if (this._connecting) {
85+
return this._connecting;
86+
}
8287

83-
// 使用 ConnectionManager 建立连接
84-
const { collection, db, instance } = await (ConnectionManager.connect(
85-
this.type,
86-
this.databaseName,
87-
this.config,
88-
this.cache,
89-
this.logger,
90-
this.defaults,
91-
));
88+
try {
89+
this._connecting = (async () => {
90+
// 使用 ConnectionManager 建立连接
91+
const { collection, db, instance } = await ConnectionManager.connect(
92+
this.type,
93+
this.databaseName,
94+
this.config,
95+
this.cache,
96+
this.logger,
97+
this.defaults,
98+
);
9299

93-
// 保存连接状态(关键:缓存对象,保证多次调用幂等返回同一形态/引用)
94-
this.dbInstance = { collection, db };
95-
this._adapter = instance;
100+
// 保存连接状态(关键:缓存对象,保证多次调用幂等返回同一形态/引用)
101+
this.dbInstance = { collection, db };
102+
this._adapter = instance;
96103

97-
return this.dbInstance;
104+
return this.dbInstance;
105+
})();
106+
107+
const result = await this._connecting;
108+
this._connecting = null;
109+
return result;
110+
} catch (err) {
111+
this._connecting = null;
112+
throw err;
113+
}
98114
}
99115

100116
/**
@@ -113,7 +129,7 @@ module.exports = class {
113129
return { ...this.defaults };
114130
}
115131

116-
/**
132+
/**
117133
* 关闭底层数据库连接(释放资源)
118134
*/
119135
async close() {
@@ -122,6 +138,7 @@ module.exports = class {
122138
}
123139
this._adapter = null;
124140
this.dbInstance = null;
141+
this._connecting = null;
125142
}
126143

127144
/**

‎lib/mongodb/index.js‎

Lines changed: 51 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -41,7 +41,7 @@ module.exports = class {
4141
this.emit = this._emitter.emit.bind(this._emitter);
4242
}
4343

44-
/**
44+
/**
4545
* 连接到MongoDB数据库
4646
* @param {Object} config - MongoDB连接配置
4747
* @param {string} config.uri - MongoDB连接URI
@@ -54,20 +54,34 @@ module.exports = class {
5454
if (this.client) {
5555
return this.client;
5656
}
57+
58+
// 防止并发连接:使用连接锁
59+
if (this._connecting) {
60+
return this._connecting;
61+
}
62+
5763
this.config = config;
64+
5865
try {
59-
const {client, db} = await connectMongo({
60-
databaseName: this.databaseName,
61-
config: this.config,
62-
logger: this.logger,
63-
defaults: this.defaults,
64-
type: this.type,
65-
});
66-
this.client = client;
67-
this.db = db;
68-
try { this.emit && this.emit('connected', { type: this.type, db: this.databaseName, scope: this.defaults?.namespace?.scope }); } catch(_) {}
69-
return this.client;
66+
this._connecting = (async () => {
67+
const {client, db} = await connectMongo({
68+
databaseName: this.databaseName,
69+
config: this.config,
70+
logger: this.logger,
71+
defaults: this.defaults,
72+
type: this.type,
73+
});
74+
this.client = client;
75+
this.db = db;
76+
try { this.emit && this.emit('connected', { type: this.type, db: this.databaseName, scope: this.defaults?.namespace?.scope }); } catch(_) {}
77+
return this.client;
78+
})();
79+
80+
const result = await this._connecting;
81+
this._connecting = null;
82+
return result;
7083
} catch (err) {
84+
this._connecting = null;
7185
try { this.emit && this.emit('error', { type: this.type, db: this.databaseName, error: String(err && (err.message || err)) }); } catch(_) {}
7286
throw err;
7387
}
@@ -109,12 +123,27 @@ module.exports = class {
109123
);
110124
}
111125

112-
collection(databaseName, collectionName) {
126+
collection(databaseName, collectionName) {
113127
if (!this.client) {
114128
const err = new Error('MongoDB is not connected. Call connect() before accessing collections.');
115129
err.code = 'NOT_CONNECTED';
116130
throw err;
117131
}
132+
133+
// 输入验证:集合名称必须是非空字符串
134+
if (!collectionName || typeof collectionName !== 'string' || collectionName.trim() === '') {
135+
const err = new Error('Collection name must be a non-empty string.');
136+
err.code = 'INVALID_COLLECTION_NAME';
137+
throw err;
138+
}
139+
140+
// 输入验证:数据库名称如果提供,必须是非空字符串
141+
if (databaseName !== undefined && databaseName !== null && (typeof databaseName !== 'string' || databaseName.trim() === '')) {
142+
const err = new Error('Database name must be a non-empty string or null/undefined.');
143+
err.code = 'INVALID_DATABASE_NAME';
144+
throw err;
145+
}
146+
118147
const effectiveDbName = databaseName || this.databaseName;
119148
const db = this.client.db(effectiveDbName);
120149
const collection = db.collection(collectionName);
@@ -564,7 +593,7 @@ module.exports = class {
564593
};
565594
}
566595

567-
/**
596+
/**
568597
* 关闭连接并释放资源
569598
*/
570599
async close() {
@@ -573,6 +602,14 @@ module.exports = class {
573602
}
574603
this.client = null;
575604
this.db = null;
605+
this._connecting = null;
606+
607+
// 清理实例ID缓存,防止内存泄漏
608+
if (this._iidCache) {
609+
this._iidCache.clear();
610+
this._iidCache = null;
611+
}
612+
576613
try { this.emit && this.emit('closed', { type: this.type, db: this.databaseName }); } catch(_) {}
577614
return true;
578615
}

0 commit comments

Comments
 (0)