Queue@rhythmjs/queue
Two engines, one contract
In-memory by default, Bun's native redis client for distribution: key layout, atomic claims, repeatable schedules through the cron engine, and the structural contract both engines implement.
On this page
Choosing an engine#
The service is engine-agnostic. With no options it runs on memoryEngine(); setting redis selects redisEngine on Bun's native client; passing engine injects your own. Producers and processors never change.
queueModule.forRoot<AppJobs>(); // in-memory
queueModule.forRoot<AppJobs>({ redis: process.env.REDIS_URL }); // distributed
queueModule.forRoot<AppJobs>({ engine: myEngine }); // customThe memory engine#
memoryEngine() is in-process and zero-dependency: waiting and delayed jobs live in arrays, due delayed jobs are promoted on every take, the highest priority wins (FIFO within), and a pending-id set implements jobId dedupe. Perfect for tests, CLIs, and single-instance apps; nothing survives the process.
The redis engine#
redisEngine(connection?, { prefix? }) distributes on Bun's native redis client. Pass a redis:// URL (the engine constructs, owns, and closes the client), an existing Bun.RedisClient (left open on close), or nothing to use Bun.redis and its REDIS_URL. It actually accepts any RedisLike (send/close), which is how the package's own tests run against an in-memory fake, no server needed.
import { createQueueService, redisEngine } from "@rhythmjs/queue";
const service = createQueueService<AppJobs>({
engine: redisEngine("redis://localhost:6379", { prefix: "app:jobs" }),
});The key layout is plain Redis: ready work in the {prefix}:waiting list; delayed jobs in the {prefix}:delayed sorted set scored by readiness time; repeats in {prefix}:repeat (a sorted set of names) plus a {prefix}:repeat:data hash; counters at {prefix}:completed and {prefix}:failed. The default prefix is rhythm:queue; through the module it composes as {options.prefix ?? "rhythm"}:{options.name ?? "queue"}.
Every claim goes through an atomic ZREM: whichever worker removes the member owns it, so concurrent workers never double-take a delayed job or a due repeat.
Repeatable schedules#
await queueService.schedule("order.process", { pattern: "0 3 * * *", timezone: "UTC" }, { orderId: "nightly" });
await queueService.schedule("email.send", { every: 60_000, limit: 10 }, { to: "digest@x", subject: "digest" });
await queueService.unschedule("order.process");every repeats on a fixed interval; pattern runs through @rhythmjs/schedule's cron engine with full timezone and DST semantics; one of the two is required (otherwise schedule throws). limit counts down and the schedule retires at zero. Each firing enqueues an ordinary job that inherits defaultJobOptions plus the options you scheduled with.
On the redis engine, due repeats are claimed with the same atomic ZREM before being re-armed, so a schedule fires once across any number of app replicas.
The engine contract#
Both engines implement the structural QueueEngine interface (add, addBulk, take(now) (the engine promotes its own due delayed jobs), requeue(job, readyAt) for retries, record("completed" | "failed"), setRepeat/clearRepeat/claimDueRepeats(now), counts, and close) over the StoredJob and RepeatSpec shapes. Tests can substitute a recording engine wholesale.
What the queue does not have#
This package has zero queue dependencies. What it does not have: stalled-job recovery (a job popped by a worker that crashes mid-run is lost), Lua-scripted atomic state transitions, per-job history, pause/resume, and rate limiting. See the API reference for every export.