rhythmjs

Search documentation

Search guides, the tutorial and every package.

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.

TypeScript
queueModule.forRoot<AppJobs>(); // in-memory
queueModule.forRoot<AppJobs>({ redis: process.env.REDIS_URL }); // distributed
queueModule.forRoot<AppJobs>({ engine: myEngine }); // custom

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

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

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