Files

11 KiB

@wrnexus/queue

A background job queue with delays, retries + exponential backoff, recurring jobs, and concurrent workers.

Part of the WrNexus framework — an SSR-first, Bun-native full-stack web framework.

Overview

@wrnexus/queue is a server-side durable job queue. You register named workers, enqueue jobs (optionally delayed or recurring), and let the queue poll and run them on a timer — with per-job retry limits and doubling backoff between attempts. defineQueue persists to .wrnexus/queue.sqlite by default, while low-level createQueue remains an intentionally in-memory primitive. Reach for it when you need to defer work (emails, webhooks, cleanup) off the request path without a heavyweight external broker. Tests can drive it deterministically via drain().

Installation

bun add @wrnexus/queue

Private package — the machine must be authenticated to the wrnexus npm org (a read token in ~/.npmrc). Requires Bun (Node is not supported).

API

The package exports a single factory plus its supporting types.

createQueue(options?): Queue

Creates a new queue instance.

function createQueue(options?: QueueOptions): Queue;

QueueOptions

Option Type Default Description
maxAttempts number 3 Default max attempts per job before it is dead-lettered.
backoffMs number 1000 Base retry backoff in ms; doubles per attempt.
pollMs number 250 Poll interval used once start() is called (ms).
onFailed (job: Job, error: unknown) => void Called when a job exhausts its attempts.
concurrency number unlimited Maximum jobs claimed by one drain() call.
capacity number unlimited Optional maximum queued jobs; omitted queues never scan merely to enforce a hidden cap.
now () => number Date.now Clock injection for deterministic tests.

Queue

The object returned by createQueue.

Method Signature Description
add add<T>(name, data: T, options?: AddOptions): Promise<Job<T>> Enqueue a job under a worker name. Returns the created job.
process process<T>(name, handler: JobHandler<T>): void Register the worker that runs jobs of the given name.
drain drain(now?: number): Promise<number> Run every job whose runAt ≤ now, once. Returns how many ran.
start start(): void Begin polling every pollMs. No-op if already started.
stop stop(): void Stop the poll timer.
shutdown shutdown({ force? }): Promise<void> Stop accepting jobs and await active work; force aborts it.
size size(): number Number of jobs currently queued.
get/list get(id) / list(name?) Inspect defensive copies of pending jobs.
cancel cancel(id): boolean Remove queued work or abort an active handler.
failed failed(): Job[] Inspect exhausted jobs in the dead-letter collection.
retry retry(id): Promise<boolean> Reset and requeue a dead-lettered job.

AddOptions

Option Type Description
delayMs number Delay before the job becomes runnable (ms).
maxAttempts number Max attempts before dead-lettering. Defaults to the queue's maxAttempts.
repeat number Re-enqueue this job this many ms after each successful run (recurring).
priority number Higher values are selected first among due jobs.
idempotencyKey string Return the matching pending job instead of enqueueing a duplicate.

JobHandler<T>

type JobHandler<T = unknown> = (
  job: Job<T>,
  context: { signal: AbortSignal },
) => void | Promise<void>;

Job<T>

interface Job<T = unknown> {
  id: string; // e.g. "job_1"
  name: string;
  data: T;
  attempts: number;
  maxAttempts: number;
  runAt: number; // epoch ms; job runs when now ≥ runAt
  repeat?: number; // if set, re-enqueue this many ms after each success
}

Usage

Application queues belong in app/queues/*.ts and default-export defineQueue(...). Configure persistence once in wrnexus.config.ts:

export default {
  queue: { storage: "database", databaseName: "default" },
  // Or omit queue entirely for .wrnexus/queue.sqlite.
};

Use storage: "memory" only for disposable work. A selected database must exist; WrNexus fails startup/first use with a clear error instead of silently falling back. Each job may declare success(data, job) and failed(data, error, job) lifecycle hooks.

Register workers, enqueue jobs, then start the poller:

import { createQueue } from "@wrnexus/queue";

const queue = createQueue({ maxAttempts: 3, backoffMs: 1000 });

// Register a worker for the "email" job name.
queue.process<{ to: string }>("email", async (job) => {
  await send(job.data.to);
});

// Enqueue a delayed job with up to 3 attempts.
await queue.add("email", { to: "a@b.com" }, { delayMs: 5000, maxAttempts: 3 });

queue.start(); // begin polling; queue.stop() to halt

Application queues with defineQueue

WrNexus applications can place typed queue modules in app/queues. The dev and production servers discover them, start them after runtime initialization, and shut them down gracefully:

import { defineQueue } from "@wrnexus/queue";

export default defineQueue({
  name: "email",
  jobs: {
    send: {
      maxAttempts: 3,
      idempotency: ({ messageId }: { messageId: number }) => `message:${messageId}`,
      run: async ({ messageId }) => sendMessage(messageId),
    },
  },
});

Callers get a typed producer and inspection helpers:

const job = await email.send.add({ messageId: 42 });
await email.send.status(job.id);
await email.send.cancel(job.id);
await email.send.retry(job.id);

createDurableQueue also exposes start(), stop(), isRunning(), health() and get() so applications do not need to build their own polling middleware.

SQLite

Use sqliteQueueStore(db) from @wrnexus/queue/sqlite with a WrNexus-compatible SQLite client. installSqliteQueueSchema(db) installs the queue table and its due-job and idempotency indexes.

Use context.signal in network/database calls so forced shutdown and active cancellation finish promptly. For process termination, prefer await queue.shutdown(); use { force: true } only after your grace period.

Durable queue

createDurableQueue({ store }) retains jobs until their handler succeeds and supports atomic leases when a driver implements QueueStore.claim. It exposes the same cancellation/shutdown behavior plus list, failed, and retry. The included memoryQueueStore() is useful for tests; production Redis/SQL drivers should make claim() atomic to prevent two workers executing one job.

Recurring jobs

Pass repeat to re-enqueue a job a fixed interval after each successful run:

queue.process("heartbeat", async () => ping());
await queue.add("heartbeat", {}, { repeat: 60_000 }); // runs ~every minute

Handling permanent failures

When a job's attempts reaches maxAttempts, it is dropped and onFailed fires instead of retrying:

const queue = createQueue({
  onFailed: (job, error) => {
    console.error(`job ${job.id} (${job.name}) gave up`, error);
  },
});

Deterministic testing

Instead of start(), inject a clock and drive the queue with drain():

let clock = 0;
const queue = createQueue({ now: () => clock });

queue.process("task", async () => {
  /* ... */
});
await queue.add("task", {}, { delayMs: 5000 });

clock = 5000;
const ran = await queue.drain(); // => 1

Durable workflows and approvals

createWorkflowEngine(store) executes dependency-ordered steps and persists every transition, result, progress update, failure, cancellation, and approval record. Approval steps pause safely and can resume after a process restart because the snapshot lives in the supplied WorkflowStore.

const workflow = defineDurableWorkflow({
  name: "publish-report",
  steps: [
    { name: "build", run: buildReport },
    { name: "approve", dependsOn: ["build"], approval: true, run: (report) => report },
    { name: "publish", dependsOn: ["approve"], run: publishReport },
  ],
});

const run = await engine.start(workflow, input);
await engine.approve(workflow, run.id, "approve", currentUser.id);

Use memoryWorkflowStore() for tests. Production stores implement the small get, put, and list contract using the same transactional database or durable service as the application.

Retry & backoff behavior

  • On a thrown handler error, the job is retried while attempts < maxAttempts.
  • The next runAt is set to now + backoffMs * 2^(attempts - 1) (exponential backoff): with backoffMs: 1000 the delays are 1s, 2s, 4s, …
  • A job whose worker name has no registered handler stays queued until one is registered (it is not counted as runnable by drain).
  • drain is re-entrant-safe: overlapping calls are skipped while one is running.

Requirements / Notes

  • Bun-only runtime (Node is not supported), consistent with the rest of the WrNexus framework. The queue itself relies only on standard timers (setInterval/clearInterval) and has no runtime dependencies.
  • The default store is in-process, so queued jobs do not survive a restart; a pluggable driver is intended for backing it with Redis/SQL for durability.
  • Works alongside @wrnexus/core for offloading work from the request path.