queue#
MongoDB-backed background jobs with persistent status and concurrent workers.
Requirements#
MongoDB configured as a replica set for change streams. Pass a native database object. The runtime must provide crypto.randomUUIDv7().
Install#
npm i @nodedk/queue
Usage#
Start a worker with your application's database:
// worker.js
var createQueue = require('@nodedk/queue')
module.exports = async function startWorker(db) {
var queue = createQueue({ db })
await queue.listen(async function (job) {
if (job.action === 'greet') {
console.log('Hello ' + job.data.name)
await queue.update(job._id, { progress: 100 })
}
})
return queue
}
Add work from another process or request handler using the same database:
var createQueue = require('@nodedk/queue')
module.exports = async function addGreeting(db, name) {
var queue = createQueue({ db })
var job = await queue.add('greet', { name })
return job._id
}
API#
createQueue({ db })#
Returns { add, cancel, get, listen, sweep, update } synchronously. db is required. Instances use the job collection; all methods on the returned queue return promises.
queue.add(action, data = {})#
Stores and resolves to a new queued job. action is the application action name; data is the application's payload.
queue.listen(handler)#
Claims existing queued jobs and watches for new inserts. handler(job) may be asynchronous. Different jobs run concurrently; a resolved handler completes the job, and a rejected or throwing handler marks it failed.
listen() resolves after scheduling the existing queue, without waiting for handlers to finish. It keeps its change stream open and exposes no stop method. Updates beneath data do not trigger handlers.
queue.get(id)#
Resolves to the job with this UUIDv7 string ID, or null when missing.
queue.update(id, values)#
Shallowly sets the supplied object fields beneath job.data and refreshes updated_at. Resolves to the MongoDB updateOne result. Root job metadata is unchanged.
async function reportProgress(queue, id) {
await queue.update(id, { progress: 50, current_step: 'rendering' })
return queue.get(id)
}
queue.cancel(id)#
Cancels and resolves to a queued job. Returns null when missing or already claimed. Running handlers continue.
queue.sweep(before = new Date())#
Deletes completed, failed, and cancelled jobs whose finished_at is at or before the Date cutoff. Resolves to the MongoDB deleteMany result. Queued and running jobs remain stored.
async function clean(queue) {
var before = new Date(Date.now() - 7 * 86400000)
var result = await queue.sweep(before)
console.log(result.deletedCount)
}
Job data#
var job = {
_id: '0198e15c-8200-7000-8000-000000000001',
action: 'greet',
status: 'running',
error: null,
created_at: new Date('2026-08-25T12:00:00Z'),
updated_at: new Date('2026-08-25T12:00:01Z'),
started_at: new Date('2026-08-25T12:00:01Z'),
finished_at: null,
data: { name: 'Ada', progress: 50 }
}
_id: UUIDv7 string.action: application-defined dispatch name.status:queued,running,completed,failed, orcancelled.error:null, or{ name, message, stack }after a handler failure.created_at,updated_at: creation and most recent update timestamps.started_at,finished_at: lifecycle timestamps, initiallynull.data: application-owned fields.
queued -> running -> completed
-> failed
queued -> cancelled
Processing behavior#
Each job is inserted once and claimed with an atomic queued to running update. Only one listener successfully claims it. Jobs are concurrent and unordered.
A claimed job is never retried or reclaimed automatically. If its process stops, it stays running; inspect it before deciding how to recover its external work. Terminal jobs stay stored until swept.
Created by Vidar Eldøy