Queue@rhythmjs/queue
API reference
Every export of @rhythmjs/queue and @rhythmjs/queue/types: the service factory and kernel module, job/repeat/process options, both engines, and the structural engine contract.
On this page
createQueueService and queueModule#
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>#
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.jobIdwhen given, elsecrypto.randomUUID(). Adelayparks 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. Requiresrepeat.everyorrepeat.pattern; throws otherwise. The first occurrence is computed immediately; each firing enqueues a normal job (withoptions) and re-arms untillimitruns out. unschedule(name)#- Removes the repeat. Unknown names are a no-op.
process(processors, options?)#- Starts a worker:
concurrencyparallel pull loops (default 1) that also claim due repeats each tick. Returns a handle whoseclose()stops the loops and awaits in-flight jobs. counts()#- Engine counts (
waiting,delayed,completed,failed,repeats) plus this service's in-flightactive. 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 }.attemptsis the total including the first (minimum 1).priority: higher runs sooner. A numericbackoffis a fixed delay; exponential backoff waitsdelay × 2^(attempt−1). Pending jobs sharing ajobIdare deduplicated.BackoffOptions#{ type: "fixed" | "exponential"; delay: number }.RepeatOptions#{ pattern?: string; every?: number; timezone?: string; limit?: number }.patternis evaluated by@rhythmjs/schedule/cron;timezoneapplies to it;limitcaps 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.onFailedfires once per job, after the final attempt.onErrorreceives engine failures (backend errors, not job errors); the loop sleeps onepollIntervaland resumes instead of dying.JobInfo#{ id: string; name: string; attemptsMade: number }, the second argument to every processor;attemptsMadecounts the current attempt, starting at 1.JobProcessors<TJobs>#{ [K in keyof TJobs]?: (payload: TJobs[K], job: JobInfo) => unknown }.BulkJobInput<TJobs> / QueueWorkerHandle / QueueContext<TJobs>#- The
addBulkentry union;{ close(): Promise<void> }; and{ queueService: QueueService<TJobs> }.
memoryEngine and redisEngine#
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#
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 }.claimDueRepeatsremoves and returns every spec due atnow; the service enqueues the job and re-arms survivors viasetRepeat. 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.