`Message.flags` was parsed and stored but never read. An OP_MSG request with
moreToCome set is fire-and-forget: the client will not read a reply. Sending one
anyway leaves it unread in the socket, so the next command on that connection
reads the previous command's reply and waits forever for its own.
This is not a corner case. Every unacknowledged write uses it, and the Node
driver sends `endSessions` with `writeConcern: {w: 0}` whenever a client closes
-- so an ordinary application that never asks for w:0 still hits it. Before:
insertOne({w: 0}) -> ok, acknowledged=false
countDocuments() (same conn) -> BSON element "cursor" is missing
The command still runs; only the reply is suppressed.
The e2e case pins maxPoolSize to 1, because with a larger pool the driver may
hand the next operation a different connection and hide the bug. It asserts the
connection still works afterwards, which is the part that matters -- not that
the unacknowledged write itself returned.
192 lines
8.6 KiB
JavaScript
192 lines
8.6 KiB
JavaScript
// End-to-end test: official MongoDB Node.js driver against multiforadb.
|
|
const { MongoClient, ObjectId } = require('mongodb');
|
|
|
|
const URL = 'mongodb://127.0.0.1:27020';
|
|
const results = [];
|
|
function check(name, cond, detail = '') {
|
|
results.push({ name, ok: !!cond, detail: String(detail) });
|
|
if (!cond) console.error(` ✗ ${name} ${detail}`);
|
|
}
|
|
async function main() {
|
|
const client = new MongoClient(URL, { serverSelectionTimeoutMS: 5000 });
|
|
await client.connect();
|
|
const db = client.db('e2e');
|
|
const users = db.collection('users');
|
|
await users.drop().catch(() => {});
|
|
|
|
// --- insert ---
|
|
await users.insertOne({ name: 'alice', age: 30, tags: ['a', 'b'] });
|
|
const many = await users.insertMany([
|
|
{ name: 'bob', age: 25, tags: ['b'] },
|
|
{ name: 'carol', age: 35, tags: ['c', 'a'] },
|
|
{ name: 'dave', age: 40, tags: [] },
|
|
]);
|
|
check('insertMany acknowledged', many.acknowledged === true, many);
|
|
check('auto _id assigned', ObjectId.isValid(many.insertedIds[0]));
|
|
|
|
// --- find: filters, operators, sort, skip, limit, projection ---
|
|
const gt = await users.find({ age: { $gt: 28 } }).sort({ age: -1 }).toArray();
|
|
check('find $gt + sort desc', gt.map((d) => d.name).join(',') === 'dave,carol,alice', gt.map((d) => d.name));
|
|
|
|
const inq = await users.find({ name: { $in: ['alice', 'bob'] } }).count();
|
|
check('find $in count', inq === 2, inq);
|
|
|
|
const rgx = await users.find({ name: /^[bc]/ }).toArray();
|
|
check('find $regex', rgx.length === 2, rgx.map((d) => d.name));
|
|
|
|
const exists = await users.find({ tags: { $exists: true } }).count();
|
|
check('find $exists', exists === 4, exists);
|
|
|
|
const lim = await users.find({}).sort({ age: 1 }).skip(1).limit(2).toArray();
|
|
check('find skip+limit+sort', lim.map((d) => d.name).join(',') === 'alice,carol', lim.map((d) => d.name));
|
|
|
|
const proj = await users.findOne({ name: 'alice' }, { projection: { _id: 0, name: 1 } });
|
|
check('projection', proj.name === 'alice' && proj.age === undefined, JSON.stringify(proj));
|
|
|
|
const dot = await users.findOne({ 'tags.0': 'c' });
|
|
check('dot path + array', dot?.name === 'carol');
|
|
|
|
// --- count ---
|
|
check('countDocuments', (await users.countDocuments({})) === 4);
|
|
check('countDocuments with filter', (await users.countDocuments({ age: { $gte: 30 } })) === 3);
|
|
check('estimatedDocumentCount', (await users.estimatedDocumentCount()) === 4);
|
|
|
|
// --- update ---
|
|
const u1 = await users.updateOne({ name: 'alice' }, { $set: { vip: true }, $inc: { age: 1 } });
|
|
check('updateOne nModified', u1.modifiedCount === 1, u1);
|
|
const alice = await users.findOne({ name: 'alice' });
|
|
check('updateOne $set+$inc applied', alice.vip === true && alice.age === 31, JSON.stringify(alice));
|
|
|
|
const um = await users.updateMany({}, { $set: { seen: true } });
|
|
check('updateMany', um.modifiedCount === 4, um);
|
|
|
|
const push = await users.updateOne({ name: 'dave' }, { $push: { tags: 'x' } });
|
|
check('$push', push.modifiedCount === 1);
|
|
check('$push visible', (await users.findOne({ name: 'dave' })).tags.length === 1);
|
|
|
|
const ups = await users.updateOne({ name: 'erin' }, { $set: { age: 28 } }, { upsert: true });
|
|
check('upsert', ups.upsertedCount === 1 && ups.matchedCount === 0, ups);
|
|
check('upsert doc exists', (await users.findOne({ name: 'erin' }))?.age === 28);
|
|
|
|
// --- findOneAndUpdate (findAndModify) ---
|
|
const fam = await users.findOneAndUpdate(
|
|
{ name: 'bob' },
|
|
{ $set: { lucky: true } },
|
|
{ returnDocument: 'after' },
|
|
);
|
|
const famDoc = fam.value ?? fam;
|
|
check('findOneAndUpdate returns new', famDoc.lucky === true, JSON.stringify(fam));
|
|
const famRemove = await users.findOneAndDelete({ name: 'erin' });
|
|
const famDelDoc = famRemove.value ?? famRemove;
|
|
check('findOneAndDelete', famDelDoc?.name === 'erin', JSON.stringify(famRemove));
|
|
|
|
// --- aggregate ---
|
|
const grp = await users
|
|
.aggregate([
|
|
{ $match: { age: { $gte: 25 } } },
|
|
{ $group: { _id: '$tags.length', total: { $sum: '$age' } } },
|
|
{ $sort: { _id: 1 } },
|
|
])
|
|
.toArray();
|
|
check('aggregate $match+$group+$sum+$sort', grp.length >= 1 && grp.some((g) => g.total > 0), JSON.stringify(grp));
|
|
|
|
const cnt = await users.aggregate([{ $match: {} }, { $count: 'n' }]).toArray();
|
|
check('aggregate $count', cnt[0]?.n === 4, JSON.stringify(cnt));
|
|
|
|
// $sort with NO preceding $group. This shape used to kill the server: the
|
|
// stage materialized its document list from the reply arena and it was then
|
|
// freed with the general allocator. Every aggregate case above happens to
|
|
// sort after a $group, which leaves the stream already materialized and the
|
|
// guilty branch unreached -- so the bug survived the whole suite. The second
|
|
// aggregate is the part that actually proves recovery: if the server died,
|
|
// this connection is gone.
|
|
const sorted = await users.aggregate([{ $sort: { age: -1 } }]).toArray();
|
|
const ages = sorted.map((d) => d.age);
|
|
check(
|
|
'aggregate bare $sort (no $group) returns sorted docs',
|
|
ages.length === 4 && ages.every((a, i) => i === 0 || ages[i - 1] >= a),
|
|
JSON.stringify(ages),
|
|
);
|
|
const stillAlive = await users.aggregate([{ $match: {} }, { $count: 'n' }]).toArray();
|
|
check('server survives a bare $sort pipeline', stillAlive[0]?.n === 4, JSON.stringify(stillAlive));
|
|
|
|
// --- a database-level command must not leak the catalog lock ---
|
|
// db.aggregate() sends {aggregate: 1}, which names no collection. dispatch
|
|
// used to resolve the namespace after taking the catalog lock and bail with a
|
|
// plain return, leaking it shared forever. Reads kept working, so the damage
|
|
// only showed on the next write that had to create a collection -- which is
|
|
// the second half of this check, and would hang rather than fail.
|
|
let dbLevelErr = null;
|
|
try {
|
|
await db.aggregate([{ $listLocalSessions: {} }]).toArray();
|
|
} catch (e) {
|
|
dbLevelErr = e;
|
|
}
|
|
check(
|
|
'database-level aggregate gives a real error, not an empty reply',
|
|
dbLevelErr !== null && typeof dbLevelErr.message === 'string' && dbLevelErr.message !== 'n/a',
|
|
String(dbLevelErr && dbLevelErr.message).slice(0, 60),
|
|
);
|
|
const afterDbLevel = await db.collection('lock_probe').insertOne({ _id: 1 });
|
|
check('a write creating a collection still completes afterwards', afterDbLevel.insertedId === 1);
|
|
|
|
// --- unacknowledged writes must not desync the connection ---
|
|
// An OP_MSG request with moreToCome set gets no reply. Sending one anyway
|
|
// left it unread in the socket, so the *next* command on that connection
|
|
// read the wrong reply. maxPoolSize 1 pins both operations to one socket,
|
|
// which is what makes the bug visible; with a larger pool the driver may
|
|
// hand out a different connection and hide it. The driver also does this to
|
|
// itself on close, via endSessions with {w: 0}.
|
|
const w0client = new MongoClient(URL, { maxPoolSize: 1 });
|
|
try {
|
|
await w0client.connect();
|
|
const w0 = w0client.db('e2e').collection('unack');
|
|
await w0.deleteMany({});
|
|
await w0.insertOne({ _id: 1, v: 'acknowledged' });
|
|
const unack = await w0.insertOne({ _id: 2, v: 'unacknowledged' }, { writeConcern: { w: 0 } });
|
|
check('unacknowledged insert is not acknowledged', unack.acknowledged === false, JSON.stringify(unack));
|
|
// The assertion that matters: the same socket still works afterwards.
|
|
const after = await w0.countDocuments({});
|
|
check('connection survives an unacknowledged write', after === 2, `count=${after}`);
|
|
} finally {
|
|
await w0client.close();
|
|
}
|
|
|
|
// --- duplicate key ---
|
|
let dupErr = null;
|
|
try {
|
|
await users.insertOne({ _id: many.insertedIds[0], name: 'clobber' });
|
|
} catch (e) {
|
|
dupErr = e;
|
|
}
|
|
check('duplicate key rejected', dupErr?.code === 11000, dupErr?.message);
|
|
|
|
// --- listCollections / listDatabases ---
|
|
const colls = await db.listCollections({}, { nameOnly: true }).toArray();
|
|
check('listCollections', colls.some((c) => c.name === 'users'), JSON.stringify(colls));
|
|
const dbs = await client.db('admin').admin().listDatabases();
|
|
check('listDatabases', dbs.databases.some((d) => d.name === 'e2e'), JSON.stringify(dbs.databases.map((d) => d.name)));
|
|
|
|
// --- delete ---
|
|
const del1 = await users.deleteOne({ name: 'dave' });
|
|
check('deleteOne', del1.deletedCount === 1, del1);
|
|
const delMany = await users.deleteMany({});
|
|
check('deleteMany', delMany.deletedCount === 3, delMany);
|
|
check('empty after delete', (await users.countDocuments({})) === 0);
|
|
|
|
await client.close();
|
|
|
|
const failed = results.filter((r) => !r.ok);
|
|
console.log(`\n${results.length - failed.length}/${results.length} checks passed`);
|
|
if (failed.length) {
|
|
console.log('FAILED:', failed.map((f) => f.name).join(', '));
|
|
process.exit(1);
|
|
}
|
|
console.log('E2E_OK');
|
|
}
|
|
|
|
main().catch((e) => {
|
|
console.error('E2E_FAIL', e);
|
|
process.exit(1);
|
|
});
|