Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
58 changes: 32 additions & 26 deletions mongodb-queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -185,19 +185,11 @@ export class MongoDBQueue<T = any> {

public async ping(ack: string, opts: PingOptions = {}): Promise<string> {
const visibility = opts.visibility || this.visibility;
const query: Filter<Partial<Message<T>>> = {
ack: ack,
visible: {$gt: now()},
};
const update: UpdateFilter<Message<T>> = {
$set: {
visible: nowPlusSecs(visibility),
},
};
const options = {
returnDocument: 'after',
includeResultMetadata: true,
} satisfies FindOneAndUpdateOptions;

if (opts.resetTries) {
update.$set = {
Expand All @@ -210,18 +202,10 @@ export class MongoDBQueue<T = any> {
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<string> {
const query: Filter<Partial<Message<T>>> = {
ack: ack,
visible: {$gt: now()},
};
const update: UpdateFilter<Message<T>> = {
$set: {
deleted: new Date(),
Expand All @@ -230,15 +214,20 @@ export class MongoDBQueue<T = any> {
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<string> {
const update: UpdateFilter<Message<T>> = {
$set: {
visible: now(),
},
$unset: {
ack: 1,
tries: 1,
},
};
return this.updateByAck(ack, update);
}

public async clean(): Promise<void> {
Expand Down Expand Up @@ -274,4 +263,21 @@ export class MongoDBQueue<T = any> {
deleted: {$exists: true},
});
}

private async updateByAck(ack: string, update: UpdateFilter<Message<T>>): Promise<string> {
const query: Filter<Partial<Message<T>>> = {
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;
}
}
2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
@@ -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": {
Expand Down
48 changes: 48 additions & 0 deletions test/re-enqueue.js
Original file line number Diff line number Diff line change
@@ -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();
Comment thread
dodoarg marked this conversation as resolved.
});

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();
});
});
1 change: 1 addition & 0 deletions test/setup.js
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ const collections = [
'dead-queue',
'queue-2',
'dead-queue-2',
're-enqueue',
];

module.exports = async function() {
Expand Down
Loading