npm package discovery and stats viewer.

Discover Tips

  • General search

    [free text search, go nuts!]

  • Package details

    pkg:[package-name]

  • User packages

    @[username]

Sponsor

Optimize Toolset

I’ve always been into building performant and accessible sites, but lately I’ve been taking it extremely seriously. So much so that I’ve been building a tool to help me optimize and monitor the sites that I build to make sure that I’m making an attempt to offer the best experience to those who visit them. If you’re into performant, accessible and SEO friendly sites, you might like it too! You can check it out at Optimize Toolset.

About

Hi, 👋, I’m Ryan Hefner  and I built this site for me, and you! The goal of this site was to provide an easy way for me to check the stats on my npm packages, both for prioritizing issues and updates, and to give me a little kick in the pants to keep up on stuff.

As I was building it, I realized that I was actually using the tool to build the tool, and figured I might as well put this out there and hopefully others will find it to be a fast and useful way to search and browse npm packages as I have.

If you’re interested in other things I’m working on, follow me on Twitter or check out the open source projects I’ve been publishing on GitHub.

I am also working on a Twitter bot for this site to tweet the most popular, newest, random packages from npm. Please follow that account now and it will start sending out packages soon–ish.

Open Software & Tools

This site wouldn’t be possible without the immense generosity and tireless efforts from the people who make contributions to the world and share their work via open source initiatives. Thank you 🙏

© 2026 – Pkg Stats / Ryan Hefner

izi-queue

v0.8.0

Published

A minimal, reliable, database-backed job queue for Node.js inspired by Oban

Readme

izi-queue

CI npm version License: MIT Node.js Version TypeScript

A minimal, reliable, database-backed job queue for Node.js inspired by Oban.

Why izi-queue?

  • No extra infrastructure - Use your existing PostgreSQL, SQLite, or MySQL database
  • Transactional - Enqueue inside your own transaction, so a job never outlives the write it belongs to
  • Durable - Jobs live in your database, survive restarts, and are recovered when a node dies
  • Simple API - Define workers, insert jobs, done
  • TypeScript-first - Full type safety and excellent DX
  • Battle-tested patterns - Inspired by Oban's proven design

Table of Contents

Installation

npm install izi-queue

Install the database driver you need:

# PostgreSQL
npm install pg

# SQLite
npm install better-sqlite3

# MySQL
npm install mysql2

Quick Start

import { IziQueue, defineWorker, WorkerResults, createSQLiteAdapter } from 'izi-queue';
import Database from 'better-sqlite3';

// 1. Define a worker
const sendEmailWorker = defineWorker('send_email', async (job) => {
  const { to, subject } = job.args as { to: string; subject: string };
  console.log(`Sending email to ${to}: ${subject}`);
  return WorkerResults.ok();
});

// 2. Create the queue
const db = new Database('jobs.db');
const queue = new IziQueue({
  database: createSQLiteAdapter(db),
  queues: { default: 10 }, // queue name: concurrency limit
});

// 3. Register worker, run migrations, and start
queue.register(sendEmailWorker);
await queue.migrate();
await queue.start();

// 4. Insert jobs. Job arguments go under `args`.
await queue.insert('send_email', {
  args: { to: '[email protected]', subject: 'Welcome!' },
});

insert takes the worker name (or definition) and a single options object. Everything other than args is optional:

await queue.insert('send_email', {
  args: { to: '[email protected]' },
  queue: 'emails',
  priority: 0,
  maxAttempts: 5,
  scheduledAt: new Date(Date.now() + 3600000),
  tags: ['welcome'],
});

Features

Job Scheduling

// Run immediately
await queue.insert('send_email', { args });

// Schedule for later
await queue.insert('send_email', {
  args,
  scheduledAt: new Date(Date.now() + 3600000), // 1 hour from now
});

Retries with Backoff

A failed job is retried up to maxAttempts times (default 20), waiting between attempts according to a backoff curve. The default curve is polynomial, matching Oban's default:

const myWorker = defineWorker('my_worker', async (job) => {
  // Automatic backoff on failure.
  // Formula: attempt^4 + 15 + rand(0..10) * attempt seconds
  return WorkerResults.error('Something went wrong');
}, {
  maxAttempts: 5, // Retry up to 5 times
});

Retry horizon

The polynomial curve is deliberately steep: it spreads a 20-attempt budget across days rather than exhausting it in a few hours, which is what makes maxAttempts: 20 a meaningful safety margin instead of a formality. With the default options (jitter shown at its midpoint):

| attempt | delay | cumulative time | |---------|----------|---------------------| | 1 | ~21s | ~21s | | 5 | ~11min | ~19min | | 10 | ~2.8hr | ~7.2hr | | 15 | ~14.1hr | ~49.8hr (~2.1 days) | | 20 | ~44.5hr | ~201hr (~8.4 days) |

Behavior change: prior to backoff curve support, the default curve was exponential and capped at 15 + 2^10 ≈ 17 minutes per attempt, so all 20 attempts of a default worker completed within ~3 hours. If you rely on the old, faster-exhausting timeline, opt back in explicitly (see below) or reduce maxAttempts.

Choosing a strategy

backoff on defineWorker accepts either a custom function, or a config object selecting and tuning one of the two built-in curves:

// Opt back into the original exponential curve (basePad + multiplier *
// 2^min(attempt, maxPower) seconds, with ±jitterPercent jitter). Plateaus at
// ~17 minutes once attempt reaches maxPower (default 10).
const legacyBackoffWorker = defineWorker('legacy_worker', async (job) => {
  return WorkerResults.error('Something went wrong');
}, {
  backoff: { strategy: 'exponential' },
});

// Cap either curve's delay, in seconds
const cappedBackoffWorker = defineWorker('capped_worker', async (job) => {
  return WorkerResults.error('Something went wrong');
}, {
  backoff: { maxDelay: 3600 }, // never wait longer than 1 hour
});

// Or bypass curves entirely with a custom function
const customBackoffWorker = defineWorker('custom_worker', async (job) => {
  return WorkerResults.ok();
}, {
  backoff: (job) => job.attempt * 60, // Linear: 60s, 120s, 180s...
});

calculateBackoff(attempt, options?) is also exported directly if you want to compute a delay outside of a worker definition.

Priority Queues

await queue.insert('urgent_task', { args, priority: 0 });      // High priority (lower = higher)
await queue.insert('background_task', { args, priority: 10 }); // Low priority

Unique Jobs

Prevent duplicate jobs from being enqueued:

await queue.insert('send_digest', {
  args,
  unique: {
    fields: ['worker', 'args'], // Uniqueness based on these fields
    keys: ['userId'],           // Or compare only these keys within args
    period: 3600,               // Only one per hour (in seconds), or 'infinity'
    states: ['scheduled', 'available', 'executing'],
  },
});

Use insertWithResult when you need to know whether a duplicate was found:

const { job, conflict } = await queue.insertWithResult('send_digest', {
  args,
  unique: { period: 3600 },
});

Insertion is atomic, so concurrent callers across multiple nodes produce exactly one job. Two notes on the semantics:

  • Only jobs inserted with unique participate. A job enqueued without unique options carries no uniqueness key and will not block a later unique insert. This matches Oban, where uniqueness is a property of how a job was enqueued.
  • Argument order does not matter. { a: 1, b: 2 } and { b: 2, a: 1 } are the same job on every adapter.

Worker Isolation

Run workers in isolated threads with resource limits:

const isolatedWorker = defineWorker('heavy_computation', async (job) => {
  // Runs in a separate worker thread
  return WorkerResults.ok();
}, {
  isolation: {
    isolated: true,
    workerPath: './workers/heavy-computation.js',
    resourceLimits: {
      maxOldGenerationSizeMb: 128,   // caps the worker's V8 heap
      maxYoungGenerationSizeMb: 32,
    },
  },
  timeout: 30000, // 30 seconds
});

Configure the thread pool on the queue:

const queue = new IziQueue({
  database: createSQLiteAdapter(db),
  queues: { default: 4 },
  isolation: { minThreads: 0, maxThreads: 4, idleTimeoutMs: 30000 },
});

A job that arrives while every thread is busy waits for one rather than failing, so a queue's concurrency may exceed maxThreads without jobs burning retry attempts on the pool being full. The job's timeout covers that wait as well as its execution, so a saturated pool cannot leave work outstanding indefinitely.

resourceLimits are applied when the thread is created, and threads are pooled per limit set — jobs asking for different limits do not share a thread. Note that maxOldGenerationSizeMb caps the V8 heap; Buffer memory lives outside it and is not governed by that setting.

Plugins

Extend functionality with plugins:

import { LifelinePlugin, PrunerPlugin } from 'izi-queue';

const queue = new IziQueue({
  database: createSQLiteAdapter(db),
  queues: { default: 10 },
  plugins: [
    new LifelinePlugin({ rescueAfter: 300 }), // Rescue orphaned jobs after 5 min
    new PrunerPlugin({ maxAge: 86400 }),      // Prune finished jobs older than 24h
  ],
});

Lifeline returns jobs abandoned by a node that stopped heartbeating. It will not touch a job that is still running on a live node, so rescueAfter does not have to exceed your worker timeouts.

Pruner deletes finished jobs, at most batchSize rows per statement (default 5000), looping until the whole eligible backlog is cleared, rather than one unbounded DELETE that would lock the table for as long as a large backlog takes to scan. Staging (moving due scheduled/retryable jobs to available, which runs automatically -- no plugin needed) batches the same way; tune it with stageBatchSize on IziQueue's config.

Every node runs every plugin: there is no leader election yet (#26). With several nodes this duplicates maintenance work against the same rows.

Telemetry

Monitor your queue with the telemetry system:

queue.on('job:complete', ({ job, duration }) => {
  console.log(`Job ${job.id} completed in ${duration}ms`);
});

queue.on('job:error', ({ job, error }) => {
  console.error(`Job ${job.id} failed:`, error);
});

// Subscribe to all events
queue.on('*', ({ event, job }) => {
  metrics.increment(`queue.${event}`, { worker: job?.worker });
});

Available events:

  • job:start, job:complete, job:error, job:cancel, job:snooze
  • job:retry, job:rescue, job:unknown_worker, job:transition_refused
  • job:unique_conflict, job:isolated:start, job:isolated:timeout
  • jobs:pruned
  • queue:start, queue:stop, queue:pause, queue:resume
  • thread:spawn, thread:exit
  • plugin:start, plugin:stop, plugin:error

job:transition_refused fires when a result could not be written because the job had already moved on — cancelled by an operator, or rescued onto another node. jobs:pruned and job:rescue carry result, not job.

Managing Jobs

const job = await queue.getJob(id);

// Cancel one job, or a scoped set
await queue.cancelJob(id);
await queue.cancelJobs({ queue: 'emails' });
await queue.cancelJobs({ worker: 'SendEmail', state: ['available', 'scheduled'] });

// Return discarded or cancelled jobs to the queue
await queue.retryJob(id);
await queue.retryJobs({ worker: 'SendEmail' });

Bulk operations require at least one criterion. cancelJobs({}) throws rather than cancelling everything, so a handler that forwards optional filters cannot wipe the queue when called with none of them. To act on every job, be explicit:

await queue.cancelJobs({ all: true });

Cancelling a job that is currently executing marks it cancelled and causes its result to be discarded when the worker finishes. It does not interrupt the worker mid-run (#30).

Transactional Inserts

Pass your open transaction as tx and the job is committed or discarded with the rest of your work. No job for an order that was never placed, and no job that becomes visible before the row it refers to.

const client = await pool.connect();
try {
  await client.query('BEGIN');

  const { rows } = await client.query(
    'INSERT INTO orders (total) VALUES ($1) RETURNING id',
    [total]
  );

  await queue.insert('send_receipt', {
    args: { orderId: rows[0].id },
    tx: client,          // same connection, same transaction
  });

  await client.query('COMMIT');
} catch (error) {
  await client.query('ROLLBACK');   // the job goes with it
  throw error;
} finally {
  client.release();
}

The handle is whatever your driver uses for a transaction:

| Adapter | Pass as tx | Transaction opened with | | --- | --- | --- | | PostgreSQL | PoolClient | client.query('BEGIN') | | MySQL | PoolConnection | connection.beginTransaction() | | SQLite | the Database | db.exec('BEGIN') |

On PostgreSQL the wake-up notification is issued on your connection too, so a worker is never nudged toward a row that has not committed yet, and is not nudged at all if you roll back.

SQLite has a single connection, so any insert while a transaction is open is already part of it; passing tx makes that explicit and catches the mistake of handing over a different database. Use db.exec('BEGIN') rather than db.transaction(), which is synchronous and cannot await.

MySQL cannot combine unique with a caller-managed transaction and will throw if you try. Its advisory lock is connection-scoped and cannot be held across your commit, which would leave a window for a concurrent node to insert a duplicate. Insert unique jobs outside the transaction, or use PostgreSQL, where the lock is transaction-scoped.

Migrations

migrate() creates and upgrades the schema, and is safe to call from every node at boot — it holds an advisory lock, so concurrent starts do not collide.

await queue.migrate();

It applies DDL against izi_jobs, so on a large table review what a release adds before deploying. Each version's migrations are listed in CHANGELOG.md.

Running Multiple Nodes

Jobs are claimed with FOR UPDATE SKIP LOCKED on PostgreSQL and MySQL, so any number of nodes can share a queue without handing out the same job twice.

Each node records a heartbeat. When one stops heartbeating, the Lifeline plugin returns its in-flight jobs to the queue — and only its jobs, so a slow job on a healthy node is never restarted underneath you.

const queue = new IziQueue({
  database: adapter,
  queues: { default: 10 },
  node: 'worker-1',        // defaults to a random id
  heartbeatInterval: 15000,
});

Two caveats when scaling out:

  • Concurrency limits are per node. Five nodes with { default: 10 } run up to 50 jobs at once (#37).
  • Every node runs every plugin (#26).
  • A worker that blocks the event loop for longer than the node TTL stops the heartbeat and may be treated as dead. Use isolated workers for CPU-bound work.

Worker Results

async perform(job) {
  // Success - job completed
  return WorkerResults.ok();
  return WorkerResults.ok({ processed: 100 }); // With metadata

  // Retry later - job moves to `retryable` and is retried with backoff
  return WorkerResults.error('Temporary failure');

  // Cancel - job moves to `cancelled` and is not retried
  return WorkerResults.cancel('Invalid data');

  // Snooze - reschedule without consuming a retry attempt
  return WorkerResults.snooze(60); // Try again in 60 seconds
}

Database Support

| Database | Adapter | Production Ready | | ---------- | ----------------------- | ---------------- | | PostgreSQL | createPostgresAdapter | Yes | | SQLite | createSQLiteAdapter | Yes | | MySQL | createMySQLAdapter | Yes |

PostgreSQL is recommended for production due to FOR UPDATE SKIP LOCKED support for efficient concurrent job fetching.

MySQL requires 8.0.1 or later for FOR UPDATE SKIP LOCKED. MariaDB does not implement it and is not supported.

Driver versions covered by the test suite: better-sqlite3 11-13, pg 8, mysql2 3.

better-sqlite3 and Node have interlocking floors, and the peer range cannot express that: v12 requires Node 20+ and v13 requires Node 22+. On Node 18, stay on better-sqlite3 11 — a newer driver installs but crashes the process with a segmentation fault rather than failing cleanly.

Examples

Check out the examples directory:

  • Fastify Sample - Full REST API with queue management, multiple queues, and graceful shutdown

Contributing

We welcome contributions! Please see our Contributing Guide for details.

Quick Start for Contributors

# Clone the repository
git clone https://github.com/IagoCavalcante/izi-queue.git
cd izi-queue

# Install dependencies
npm install

# Run tests
npm test

# Run tests with coverage
npm run test:coverage

# Run linting
npm run lint

# Build the project
npm run build

Development Guidelines

  • Write tests for new features
  • Follow existing code patterns (see CLAUDE.md for architecture details)
  • Run npm run lint before committing
  • Keep PRs focused and atomic

Architecture

izi-queue follows these key architectural patterns:

  • Registry Pattern - Global worker registry for dynamic worker management
  • Adapter Pattern - Database adapters for multi-database support
  • Plugin Architecture - Extensible plugin system with lifecycle hooks
  • State Machine - Job state transitions with validation
  • Observable Pattern - Telemetry event system for monitoring

For detailed architecture documentation, see CLAUDE.md.

Acknowledgments

izi-queue is heavily inspired by Oban, the excellent background job library for Elixir. We've adapted many of its battle-tested patterns for the Node.js ecosystem.

License

MIT