From c03c43a6492a55e8fd1311e67b57d63a35f9dbc9 Mon Sep 17 00:00:00 2001 From: Edoardo Date: Wed, 22 Jul 2026 10:19:38 +0200 Subject: [PATCH 1/6] exposed method, tested, bumped minor version --- mongodb-queue.ts | 25 +++++++++++++++++++++++++ package.json | 2 +- test/re-enqueue.js | 46 ++++++++++++++++++++++++++++++++++++++++++++++ test/setup.js | 1 + 4 files changed, 73 insertions(+), 1 deletion(-) create mode 100644 test/re-enqueue.js diff --git a/mongodb-queue.ts b/mongodb-queue.ts index 8cbd6f0..634c85b 100644 --- a/mongodb-queue.ts +++ b/mongodb-queue.ts @@ -241,6 +241,31 @@ export class MongoDBQueue { return '' + msg.value._id; } + public async reEnqueue(ack: string): Promise { + const query: Filter>> = { + ack: ack, + visible: {$gt: now()}, + }; + const update: UpdateFilter> = { + $set: { + visible: now(), + }, + $unset: { + ack: 1, + tries: 1, + }, + }; + const options = { + returnDocument: 'after', + includeResultMetadata: true, + } satisfies FindOneAndUpdateOptions; + const msg = await this.col.findOneAndUpdate(query, update, options); + if (!msg.value) { + throw new Error('Queue.reEnqueue(): Unidentified ack : ' + ack); + } + return '' + msg.value._id; + } + public async clean(): Promise { const query = { deleted: {$exists: true}, diff --git a/package.json b/package.json index 927b234..53becc8 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@reedsy/mongodb-queue", - "version": "8.2.1", + "version": "8.3.0", "description": "Message queues which uses MongoDB.", "main": "mongodb-queue.js", "scripts": { diff --git a/test/re-enqueue.js b/test/re-enqueue.js new file mode 100644 index 0000000..917efc0 --- /dev/null +++ b/test/re-enqueue.js @@ -0,0 +1,46 @@ +const test = require('tape'); + +const setup = require('./setup.js'); +const {MongoDBQueue} = require('../'); + +setup().then(({client, db}) => { + test('re-enqueue: message becomes immediately available again, with tries and ack reset', async function(t) { + const queue = new MongoDBQueue(db, 're-enqueue', {visibility: 30}); + + const id = await queue.add('Hello, World!'); + const msg = await queue.get(); + + const reEnqueuedId = await queue.reEnqueue(msg.ack); + t.equal(reEnqueuedId, id, 'Re-enqueue keeps the same document id'); + + const requeued = await queue.get(); + t.ok(requeued, 'Message is immediately available again'); + t.equal(requeued.id, id, 'Same document id after re-enqueue'); + t.equal(requeued.tries, 1, 'Tries restarted from zero'); + t.notEqual(requeued.ack, msg.ack, 'The old ack no longer applies'); + + t.end(); + }); + + test("re-enqueue: can't re-enqueue with a stale or unknown ack", async function(t) { + const queue = new MongoDBQueue(db, 're-enqueue', {visibility: 30}); + + await queue.add('Hello, World!'); + const msg = await queue.get(); + await queue.ack(msg.ack); + + const error = await queue.reEnqueue(msg.ack).catch((err) => err); + t.ok(error, 'Got an error when re-enqueuing an already-acked message'); + + const unknownError = await queue.reEnqueue('unknown-ack').catch((err) => err); + t.ok(unknownError, 'Got an error when re-enqueuing an unknown ack'); + + t.end(); + }); + + test('client.close()', function(t) { + t.pass('client.close()'); + client.close(); + t.end(); + }); +}); diff --git a/test/setup.js b/test/setup.js index 59dae7e..8300d5c 100644 --- a/test/setup.js +++ b/test/setup.js @@ -16,6 +16,7 @@ const collections = [ 'dead-queue', 'queue-2', 'dead-queue-2', + 're-enqueue', ]; module.exports = async function() { From aa9a563611cbe79209645be88c60d94e0aedc716 Mon Sep 17 00:00:00 2001 From: Edoardo Date: Wed, 22 Jul 2026 10:24:53 +0200 Subject: [PATCH 2/6] ++ readme doc for new method --- README.md | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/README.md b/README.md index 5684d76..e59b249 100644 --- a/README.md +++ b/README.md @@ -368,6 +368,20 @@ queue.get((err, msg) => { }) ``` +### .reEnqueue() ### + +Atomically resets a message's tries and ack so it's immediately available again, +without inserting a new document: + +```js +const msg = await queue.get(); +const id = await queue.reEnqueue(msg.ack); +// this message is immediately available again, with tries reset to 0 +``` + +Unlike `.ping(ack, { resetTries: true, resetAck: true })`, this also resets the +message's visibility immediately, rather than keeping the current window. + ### .total() ### Returns the total number of messages that has ever been in the queue, including From 1f9e35facdccd14668dbf896fecdefe0f69d8be0 Mon Sep 17 00:00:00 2001 From: Edoardo Date: Wed, 22 Jul 2026 12:22:08 +0200 Subject: [PATCH 3/6] internal refactor --> updateByAck --- mongodb-queue.ts | 59 ++++++++++++++++-------------------------------- 1 file changed, 20 insertions(+), 39 deletions(-) diff --git a/mongodb-queue.ts b/mongodb-queue.ts index 634c85b..b5863c9 100644 --- a/mongodb-queue.ts +++ b/mongodb-queue.ts @@ -185,19 +185,11 @@ export class MongoDBQueue { public async ping(ack: string, opts: PingOptions = {}): Promise { const visibility = opts.visibility || this.visibility; - const query: Filter>> = { - ack: ack, - visible: {$gt: now()}, - }; const update: UpdateFilter> = { $set: { visible: nowPlusSecs(visibility), }, }; - const options = { - returnDocument: 'after', - includeResultMetadata: true, - } satisfies FindOneAndUpdateOptions; if (opts.resetTries) { update.$set = { @@ -210,18 +202,10 @@ export class MongoDBQueue { update.$unset = {ack: 1}; } - const msg = await this.col.findOneAndUpdate(query, update, options); - if (!msg.value) { - throw new Error('Queue.ping(): Unidentified ack : ' + ack); - } - return '' + msg.value._id; + return this.updateByAck('ping', ack, update); } public async ack(ack: string): Promise { - const query: Filter>> = { - ack: ack, - visible: {$gt: now()}, - }; const update: UpdateFilter> = { $set: { deleted: new Date(), @@ -230,22 +214,10 @@ export class MongoDBQueue { visible: 1, }, }; - const options = { - returnDocument: 'after', - includeResultMetadata: true, - } satisfies FindOneAndUpdateOptions; - const msg = await this.col.findOneAndUpdate(query, update, options); - if (!msg.value) { - throw new Error('Queue.ack(): Unidentified ack : ' + ack); - } - return '' + msg.value._id; + return this.updateByAck('ack', ack, update); } public async reEnqueue(ack: string): Promise { - const query: Filter>> = { - ack: ack, - visible: {$gt: now()}, - }; const update: UpdateFilter> = { $set: { visible: now(), @@ -255,15 +227,7 @@ export class MongoDBQueue { tries: 1, }, }; - const options = { - returnDocument: 'after', - includeResultMetadata: true, - } satisfies FindOneAndUpdateOptions; - const msg = await this.col.findOneAndUpdate(query, update, options); - if (!msg.value) { - throw new Error('Queue.reEnqueue(): Unidentified ack : ' + ack); - } - return '' + msg.value._id; + return this.updateByAck('reEnqueue', ack, update); } public async clean(): Promise { @@ -299,4 +263,21 @@ export class MongoDBQueue { deleted: {$exists: true}, }); } + + private async updateByAck(methodName: string, ack: string, update: UpdateFilter>): Promise { + const query: Filter>> = { + ack: ack, + visible: {$gt: now()}, + }; + const options = { + returnDocument: 'after', + includeResultMetadata: true, + } satisfies FindOneAndUpdateOptions; + + const msg = await this.col.findOneAndUpdate(query, update, options); + if (!msg.value) { + throw new Error(`Queue.${methodName}(): Unidentified ack : ` + ack); + } + return '' + msg.value._id; + } } From 475c4708288ede705ff43c209ccb6200d6afaec9 Mon Sep 17 00:00:00 2001 From: Edoardo Date: Wed, 22 Jul 2026 12:41:05 +0200 Subject: [PATCH 4/6] copilot was concerned --- test/re-enqueue.js | 2 ++ 1 file changed, 2 insertions(+) diff --git a/test/re-enqueue.js b/test/re-enqueue.js index 917efc0..75712e8 100644 --- a/test/re-enqueue.js +++ b/test/re-enqueue.js @@ -18,6 +18,8 @@ setup().then(({client, db}) => { t.equal(requeued.id, id, 'Same document id after re-enqueue'); t.equal(requeued.tries, 1, 'Tries restarted from zero'); t.notEqual(requeued.ack, msg.ack, 'The old ack no longer applies'); + // ack so it doesn't linger reserved and leak into the next test on this collection + await queue.ack(requeued.ack); t.end(); }); From 7d9476dc04529dbb32ca61565f8f67acbff38d75 Mon Sep 17 00:00:00 2001 From: Edoardo Date: Wed, 22 Jul 2026 12:49:37 +0200 Subject: [PATCH 5/6] didn't like passing a method name just to format some error msg --- mongodb-queue.ts | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/mongodb-queue.ts b/mongodb-queue.ts index b5863c9..5903253 100644 --- a/mongodb-queue.ts +++ b/mongodb-queue.ts @@ -202,7 +202,7 @@ export class MongoDBQueue { update.$unset = {ack: 1}; } - return this.updateByAck('ping', ack, update); + return this.updateByAck(ack, update); } public async ack(ack: string): Promise { @@ -214,7 +214,7 @@ export class MongoDBQueue { visible: 1, }, }; - return this.updateByAck('ack', ack, update); + return this.updateByAck(ack, update); } public async reEnqueue(ack: string): Promise { @@ -227,7 +227,7 @@ export class MongoDBQueue { tries: 1, }, }; - return this.updateByAck('reEnqueue', ack, update); + return this.updateByAck(ack, update); } public async clean(): Promise { @@ -264,7 +264,7 @@ export class MongoDBQueue { }); } - private async updateByAck(methodName: string, ack: string, update: UpdateFilter>): Promise { + private async updateByAck(ack: string, update: UpdateFilter>): Promise { const query: Filter>> = { ack: ack, visible: {$gt: now()}, @@ -276,7 +276,7 @@ export class MongoDBQueue { const msg = await this.col.findOneAndUpdate(query, update, options); if (!msg.value) { - throw new Error(`Queue.${methodName}(): Unidentified ack : ` + ack); + throw new Error('Queue: Unidentified ack : ' + ack); } return '' + msg.value._id; } From 6e5c2dc449ab8116744414b88c48122c9fbb3000 Mon Sep 17 00:00:00 2001 From: Edoardo Argiolas Date: Wed, 22 Jul 2026 12:50:48 +0200 Subject: [PATCH 6/6] minor tweak to readme msg copilot said so Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- README.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/README.md b/README.md index e59b249..6f7106c 100644 --- a/README.md +++ b/README.md @@ -376,7 +376,7 @@ without inserting a new document: ```js const msg = await queue.get(); const id = await queue.reEnqueue(msg.ack); -// this message is immediately available again, with tries reset to 0 +// this message is immediately available again, and its tries counter restarts (next get() will return tries=1) ``` Unlike `.ping(ack, { resetTries: true, resetAck: true })`, this also resets the