rhythmjs

Search documentation

Search guides, the tutorial and every package.

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.

Shell
bun add @rhythmjs/queue @rhythmjs/rhythm

Queues 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.

TypeScript
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.

TypeScript
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#

TypeScript
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#

TypeScript
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 jobs

Jobs 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#

TypeScript
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.