rhythmjs

Search documentation

Search guides, the tutorial and every package.

On this page

createQueueService and queueModule#

TypeScript
function createQueueService<TJobs extends JobMap>(
  options?: QueueModuleOptions,
): QueueService<TJobs>;

const queueModule: {
  forRoot<TJobs extends JobMap>(options?: QueueModuleOptions): Rhythm<...>;
};

createQueueService picks the engine: options.engine when given, else redisEngine(options.redis, { prefix }) when options.redis is set (key prefix `${prefix ?? "rhythm"}:${name ?? "queue"}`), else memoryEngine(). defaultJobOptions merge under per-call options. queueModule.forRoot provides { queueService } and disposes it on teardown, workers first, then the engine.

QueueModuleOptions#
{ name?: string; prefix?: string; defaultJobOptions?: JobOptions; redis?: string | Bun.RedisClient; engine?: QueueEngine }.

QueueService<TJobs>#

TypeScript
interface QueueService<TJobs extends JobMap> {
  add<K extends keyof TJobs & string>(name: K, payload: TJobs[K], options?: JobOptions): Promise<string>;
  addBulk(jobs: BulkJobInput<TJobs>[]): Promise<string[]>;
  schedule<K extends keyof TJobs & string>(name: K, repeat: RepeatOptions, payload: TJobs[K], options?: JobOptions): Promise<void>;
  unschedule(name: keyof TJobs & string): Promise<void>;
  process(processors: JobProcessors<TJobs>, options?: ProcessOptions): QueueWorkerHandle;
  counts(): Promise<Record<string, number>>;
  close(): Promise<void>;
}
add(name, payload, options?)#
Enqueues one job and resolves to its id: options.jobId when given, else crypto.randomUUID(). A delay parks the job as delayed until ready.
addBulk(jobs)#
Enqueues many { name, payload, options? } entries and resolves to their ids in order.
schedule(name, repeat, payload, options?)#
Registers or replaces the repeat named name. Requires repeat.every or repeat.pattern; throws otherwise. The first occurrence is computed immediately; each firing enqueues a normal job (with options) and re-arms until limit runs out.
unschedule(name)#
Removes the repeat. Unknown names are a no-op.
process(processors, options?)#
Starts a worker: concurrency parallel pull loops (default 1) that also claim due repeats each tick. Returns a handle whose close() stops the loops and awaits in-flight jobs.
counts()#
Engine counts (waiting, delayed, completed, failed, repeats) plus this service's in-flight active.
close()#
Closes every worker this service started, then the engine (an owned redis connection closes; a passed client stays open).

Job and process options#

JobOptions#
{ delay?: number; attempts?: number; priority?: number; backoff?: number | BackoffOptions; jobId?: string }. attempts is the total including the first (minimum 1). priority: higher runs sooner. A numeric backoff is a fixed delay; exponential backoff waits delay × 2^(attempt−1). Pending jobs sharing a jobId are deduplicated.
BackoffOptions#
{ type: "fixed" | "exponential"; delay: number }.
RepeatOptions#
{ pattern?: string; every?: number; timezone?: string; limit?: number }. pattern is evaluated by @rhythmjs/schedule/cron; timezone applies to it; limit caps total runs.
ProcessOptions#
{ concurrency?: number; pollInterval?: number; onCompleted?: (name, jobId) => void; onError?: (error) => void; onFailed?: (name, jobId, error) => void }. pollInterval (default 20 ms) is the idle sleep between empty polls. onFailed fires once per job, after the final attempt. onError receives engine failures (backend errors, not job errors); the loop sleeps one pollInterval and resumes instead of dying.
JobInfo#
{ id: string; name: string; attemptsMade: number }, the second argument to every processor; attemptsMade counts the current attempt, starting at 1.
JobProcessors<TJobs>#
{ [K in keyof TJobs]?: (payload: TJobs[K], job: JobInfo) => unknown }.
BulkJobInput<TJobs> / QueueWorkerHandle / QueueContext<TJobs>#
The addBulk entry union; { close(): Promise<void> }; and { queueService: QueueService<TJobs> }.

memoryEngine and redisEngine#

TypeScript
function memoryEngine(): QueueEngine;

interface RedisEngineOptions { prefix?: string }
interface RedisLike {
  send(command: string, args: string[]): Promise<unknown>;
  close(): void;
}
function redisEngine(
  connection?: string | Bun.RedisClient | RedisLike,
  options?: RedisEngineOptions,
): QueueEngine;

memoryEngine: arrays and counters, priority-aware take, delayed promotion by timestamp, jobId dedupe via a pending-id set. redisEngine: with a URL string it constructs and owns a Bun.RedisClient (closed by close()); with a passed client or Bun.redis default it never closes what it does not own. Keys under prefix (default rhythm:queue): :waiting (list), :delayed and :repeat (sorted sets scored by readiness), :repeat:data (hash), :completed/:failed (counters). Promotion and repeat claims use ZREM as the atomic claim, so concurrent workers never double-take.

The QueueEngine contract#

TypeScript
interface QueueEngine {
  add(job: StoredJob): Promise<void>;
  addBulk(jobs: StoredJob[]): Promise<void>;
  take(now: number): Promise<StoredJob | null>;
  requeue(job: StoredJob, readyAt: number): Promise<void>;
  record(kind: "completed" | "failed"): Promise<void>;
  setRepeat(spec: RepeatSpec): Promise<void>;
  clearRepeat(name: string): Promise<void>;
  claimDueRepeats(now: number): Promise<RepeatSpec[]>;
  counts(): Promise<Record<string, number>>;
  close(): Promise<void>;
}
StoredJob#
The engine-level job record: { id, name, payload, priority, attemptsMade, maxAttempts, backoffType, backoffDelay, readyAt }.
RepeatSpec#
A stored repeat: { name, payload, options?, every?, pattern?, timezone?, remaining?, nextAt }. claimDueRepeats removes and returns every spec due at now; the service enqueues the job and re-arms survivors via setRepeat.
take(now)#
Promotes due delayed jobs itself, then returns the next ready job (highest priority first) or null.
requeue(job, readyAt)#
Parks a failed job for its next attempt at readyAt.