Queue@rhythmjs/queue
Typed jobs in, typed jobs out
Declare a typed job map, produce with delays, priorities, and dedupe ids, and process with named handlers, retries with backoff, and lifecycle-managed workers.
On this page
Install the package#
@rhythmjs/queue is typed job queues for Rhythm on Bun: background work with retries, delays, priorities, and repeatable schedules, declared through a typed job map and processed by lifecycle-managed workers. The engines are in-house - no queue library underneath; Engines & repeats covers them.
bun add @rhythmjs/queue @rhythmjs/rhythmQueues and events are complementary, not interchangeable: an event announces (ephemeral broadcast), a queue job obligates (at-least-once, processed by exactly one worker).
The job map is the contract#
Producers and processors share one interface: keys are job names, values are payload types. Misspelled names and wrong payload shapes are compile errors on both sides.
import type { QueueContext } from "@rhythmjs/queue";
export interface AppJobs {
"email.send": { to: string; subject: string };
"order.process": { orderId: string };
}
export type AppQueueContext = QueueContext<AppJobs>;The kernel module#
queueModule.forRoot<TJobs>(options?) provides queueService and owns the lifecycle: the kernel's teardown() closes workers first (draining in-flight jobs), then the engine. defaultJobOptions merge under every job's own options.
import { Rhythm } from "@rhythmjs/rhythm";
import { queueModule } from "@rhythmjs/queue";
const app = new Rhythm().register(
queueModule.forRoot<AppJobs>({
name: "app",
defaultJobOptions: { attempts: 3, backoff: { type: "exponential", delay: 1000 } },
// redis: process.env.REDIS_URL switches engines without touching producers
}),
({ queueService }) => ({ queueService }),
);Produce#
await queueService.add("email.send", { to: "a@b.c", subject: "hi" }, { delay: 5000 });
await queueService.addBulk([
{ name: "email.send", payload: { to: "x@y.z", subject: "s" } },
{ name: "order.process", payload: { orderId: "o1" }, options: { priority: 1 } },
]);add returns the job id, always a string: crypto.randomUUID() unless you pass jobId. Job options: delay in ms; attempts (total including the first, default 1); backoff, a fixed delay in ms, or { type: "fixed" | "exponential", delay }, where exponential doubles per attempt (delay · 2^(attempt−1)); priority (higher runs sooner, default 0, FIFO within a priority); and jobId: pending jobs (waiting or delayed) with the same id are deduplicated on the memory engine.
Process#
const worker = queueService.process(
{
"email.send": async (payload, job) => sendMail(payload),
"order.process": async (payload, job) => fulfil(payload.orderId),
},
{
concurrency: 4,
onCompleted: (name, id) => metrics.increment(`jobs.${name}.ok`),
onFailed: (name, id, error) => log.error(`${name}#${id} exhausted retries`, error),
},
);
await worker.close(); // drains in-flight jobsJobs dispatch to processors by name; each handler receives its typed payload and a JobInfo - { id, name, attemptsMade }, where attemptsMade is 1 on the first run. A job with no registered processor fails loudly and enters the retry/failed flow like any other error.
A failing job is re-queued with its backoff delay until attempts is exhausted; onFailed fires once per job, after the last attempt, never per retry. concurrency (default 1) runs that many pull loops in parallel; idle workers poll the engine every pollInterval ms (default 20).
Worker loops survive their backend: an engine failure (a Redis connection reset, a corrupted repeat spec) is reported to onError, the loop sleeps one pollInterval, and polling resumes — it never kills the worker.
Introspection#
await queueService.counts();
// { waiting, delayed, completed, failed, repeats, active }active counts jobs currently inside this process's handlers; the rest come from the engine. Repeatable schedules, the two engines, and the engine contract live on the next page: Engines & repeats.