diff --git a/README.md b/README.md index 5684d76..6f7106c 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, and its tries counter restarts (next get() will return tries=1) +``` + +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 diff --git a/mongodb-queue.ts b/mongodb-queue.ts index 8cbd6f0..5903253 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(ack, update); } public async ack(ack: string): Promise { - const query: Filter>> = { - ack: ack, - visible: {$gt: now()}, - }; const update: UpdateFilter> = { $set: { deleted: new Date(), @@ -230,15 +214,20 @@ 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, update); + } + + public async reEnqueue(ack: string): Promise { + const update: UpdateFilter> = { + $set: { + visible: now(), + }, + $unset: { + ack: 1, + tries: 1, + }, + }; + return this.updateByAck(ack, update); } public async clean(): Promise { @@ -274,4 +263,21 @@ export class MongoDBQueue { deleted: {$exists: true}, }); } + + private async updateByAck(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: Unidentified ack : ' + ack); + } + return '' + msg.value._id; + } } 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..75712e8 --- /dev/null +++ b/test/re-enqueue.js @@ -0,0 +1,48 @@ +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'); + // ack so it doesn't linger reserved and leak into the next test on this collection + await queue.ack(requeued.ack); + + 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() {