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

@wnlx/redisq

v0.3.0

Published

Event-driven Redis Streams message broker for TypeScript

Readme


RedisQ is a lightweight, type-safe event queue for TypeScript built on Redis Streams and Consumer Groups.

If you've ever used Node.js's EventEmitter, RedisQ will feel instantly familiar — except your events are processed reliably across multiple workers and even multiple machines.

q.on("send-email", async (job) => {
  await sendEmail(job.data);
});

await q.add("send-email", {
  to: "[email protected]",
  subject: "Welcome!",
  body: "Thanks for joining!",
});

Why RedisQ?

Most Redis queue libraries are job-oriented.

RedisQ is event-oriented.

Instead of defining processors and queues, you publish typed events and subscribe to them using a familiar API inspired by Node.js.

Features

  • ✅ Event-driven API
  • ✅ Type-safe payloads with full inference
  • ✅ Redis Streams + Consumer Groups
  • ✅ Multi-worker processing
  • ✅ Unlimited numeric priority system (0 = most critical)
  • ✅ Delayed events
  • ✅ Automatic retries with exponential backoff
  • ✅ Dead-letter queue (with retry method)
  • ✅ Pending recovery (XAUTOCLAIM)
  • ✅ Graceful shutdown
  • ✅ Runtime metrics
  • ✅ Job cancellation

Installation

bun add @wnlx/redisq

or

npm install @wnlx/redisq

Quick Start

RedisQ separates producing and consuming. A producer just constructs the queue and calls add(). A consumer additionally registers handlers and calls startWorker().

Producer

import Redis from "ioredis";
import { RedisQ } from "@wnlx/redisq";

interface Events {
  "send-email": {
    to: string;
    subject: string;
    body: string;
  };
}

const redis = new Redis("redis://localhost:6379");
const q = new RedisQ<Events>({ redis });

// No startWorker() needed — producers just add jobs
await q.add(
  "send-email",
  { to: "[email protected]", subject: "Welcome!", body: "Thanks for joining!" },
  { priority: 0, delay: 5000 },
);

Consumer

import Redis from "ioredis";
import { RedisQ } from "@wnlx/redisq";

interface Events {
  "send-email": {
    to: string;
    subject: string;
    body: string;
  };
}

const redis = new Redis("redis://localhost:6379");

const q = new RedisQ<Events>({
  redis,
  onError(error, context) {
    console.error(context, error);
  },
});

q.on("send-email", async (job) => {
  console.log(job.data.to);
});

await q.startWorker();

// On shutdown:
await q.stopWorker();

Programming Model

RedisQ follows the same mental model as Node.js events.

| Node.js | RedisQ | | ---------------- | ----------------- | | emitter.on() | q.on() | | emitter.emit() | q.add() | | Same process | Multiple workers | | In-memory | Redis Streams | | Fire-and-forget | Reliable delivery |

You write event-driven code exactly as you would in Node.js, while RedisQ handles distribution, retries, delayed execution, acknowledgements, and recovery.


Architecture

                 Producer
                    │
            q.add("send-email", data, { priority })
                    │
      ┌─────────────┼─────────────┐
      ▼             ▼             ▼
  Stream :0     Stream :1     Stream :P (dynamic)
 (Critical)      (High)      (Lower Priority)
      │             │             │
      └─────────────┬─────────────┘
                    │
            Consumer Group Polling
           (Highest Priority First)
                    │
            ┌───────┴────────┐
            ▼                ▼
         Worker A        Worker B
            │                │
            └────── ACK ─────┘

Failure
  │
  ▼
Retry Queue (delay sorted set)
  │
  ▼
Dead Letter Queue

Public API

new RedisQ<Events>(options)

Creates a new queue instance. RedisQ duplicates two internal Redis connections from the one you provide and manages their lifecycle. Your original connection is never closed by RedisQ unless specified in startWorker(closeOriginalConnection?).


add(event, data, options?)

Publish a typed event. Returns the generated job ID. Works on both producer and consumer instances, with or without startWorker().

const jobId = await q.add(
  "send-email",
  { to: "[email protected]", subject: "Hi", body: "Hello" },
  {
    priority: 0, // optional number: 0 is most critical, higher numbers are less critical (default: 0)
    delay: 10000, // optional milliseconds
  },
);

Jobs with no delay (or delay: 0) are dispatched to their respective priority stream immediately. Jobs with a delay are held in a sorted set and promoted to their priority stream when their time comes.


on(event, handler)

Register an event handler. The payload is automatically typed from your event interface. Only relevant on consumer instances.

q.on("resize-image", async (job) => {
  console.log(job.data.path); // string
  console.log(job.data.width); // number
});

cancel(jobId)

Best-effort cancellation. Returns true if the job was removed, false otherwise.

const cancelled = await q.cancel(jobId);

Note: Cancellation only works for jobs still in the delay queue. Once a job has been promoted to the stream it will be processed — cancel() returns false in that case.


retryDLQJobs(event?)

Re-queues jobs from the Dead Letter Queue back into the active queue and returns the number of jobs requeued. Optionally filter by event name or an array of event names.

// Retry all DLQ jobs
const count = await q.retryDLQJobs();

// Retry a single event type
await q.retryDLQJobs("send-email");

// Retry multiple event types
await q.retryDLQJobs(["send-email", "resize-image"]);

startWorker(closeOriginalConnection?)

Starts the worker, scheduler, and reclaim loops. Creates the Consumer Group automatically if it does not exist. Calling startWorker() multiple times is safe — concurrent calls wait for the same initialization to complete.

Only call this on consumer instances. Producer instances do not need it.

// Standard usage
await q.startWorker();

// If you no longer need the original Redis connection after starting:
await q.startWorker(true);

Passing true closes the original connection you provided in options, since RedisQ has already duplicated it internally. Useful when the original connection was created solely for constructing RedisQ.


stopWorker()

Gracefully shuts down all loops and closes both internal Redis connections. A stopped instance cannot be restarted.


getStats()

Returns a snapshot of runtime metrics for the current process. In a multi-worker deployment each worker tracks its own counters independently.

type RedisQStats = {
  added: number;
  cancelled: number;
  processed: number;
  retried: number;
  failed: number;
  deadLettered: number;
  reclaimed: number;
};

Type Safety

Payloads are fully inferred from your event interface — no casts needed anywhere.

interface Events {
  "send-email": {
    to: string;
    subject: string;
    body: string;
  };
  "resize-image": {
    path: string;
    width: number;
    height: number;
  };
}

const q = new RedisQ<Events>({ redis });

await q.add("resize-image", {
  path: "/tmp/image.png",
  width: 800,
  height: 600,
});

q.on("resize-image", async (job) => {
  job.data.width; // number ✓
});

The onProcessStart and onProcessEnd hooks are also fully typed — narrowing event automatically narrows job.data:

const q = new RedisQ<Events>({
  redis,
  onProcessStart(event, job) {
    if (event === "send-email") {
      console.log(job.data.to); // string ✓
    }
    if (event === "resize-image") {
      console.log(job.data.width); // number ✓
    }
  },
});

Priority Queues

RedisQ includes a built-in unlimited numeric priority system:

  • Precedence: 0 is the lowest number and represents the highest priority / most critical task. The higher the number (1, 2, 3, 10, 100...), the less critical the task.
  • Unbounded: No arbitrary limits or pre-defined fixed tiers — any non-negative integer is a valid priority.
  • Default: When omitted, priority defaults to 0 (or defaultPriority configured in options).
// Critical emergency job (processed first)
await q.add("send-alert", { text: "Database unreachable!" }, { priority: 0 });

// High priority job
await q.add("send-alert", { text: "High memory usage" }, { priority: 1 });

// Normal standard job
await q.add("send-alert", { text: "Daily report ready" }, { priority: 2 });

// Low priority / background batch job
await q.add("send-alert", { text: "Archiving old logs" }, { priority: 10 });

How Priority Works Under the Hood

  1. Dynamic Stream Partitioning: Each priority level maps to a dedicated stream (<streamKey>:<priority>). Active priorities are tracked in <streamKey>:priorities.
  2. Ascending Polling Order: Workers poll streams via XREADGROUP in ascending numeric order (:0, :1, :2...). Redis Streams inspects streams in the specified argument order, guaranteeing critical (0) jobs are fetched and processed before lower-priority jobs.
  3. Priority-Preserving Delays & Retries: Delayed jobs and retrying jobs preserve their exact numeric priority and are routed to their respective priority stream when due.
  4. Priority-Ordered Crash Recovery: XAUTOCLAIM reclaims orphaned pending jobs in ascending priority order.

Reliability

RedisQ provides at-least-once delivery.

Successful handlers are acknowledged and removed from the stream. Failed handlers are retried using exponential backoff before eventually being moved to a Dead Letter Queue.

Handlers should be idempotent — retries and crash recovery may execute a job more than once.

Retries

Failed handlers are retried up to retryLimit times using exponential backoff with optional jitter. Between retries the job is held in the delay sorted set and re-promoted to the stream when its retry time comes.

Dead Letter Queue

Jobs that exceed the retry limit are moved into a dedicated DLQ stream (redisq:dlq by default) with their final error message and attempt count attached. DLQ jobs can be requeued at any time by calling retryDLQJobs().

Pending Recovery

Workers periodically call XAUTOCLAIM to reclaim messages that have been pending longer than reclaimMinIdleMs. This allows another worker to pick up jobs abandoned by a crashed consumer.


Configuration

interface RedisQOptions<Events> {
  // Required
  redis: Redis;

  // Redis key names
  streamKey?: string; // default: "redisq:stream"
  delayKey?: string; // default: "redisq:delay"
  jobKey?: string; // default: "redisq:job"
  dlqKey?: string; // default: "redisq:dlq"

  // Priority tuning
  defaultPriority?: number; // default: 0 (most critical)

  // Consumer identity
  consumerGroup?: string; // default: "redisq"
  consumerName?: string; // default: "consumer-<uuid>"

  // Tuning
  blockMs?: number; // default: 1000
  schedulerIntervalMs?: number; // default: 250
  batchSize?: number; // default: 50
  streamMaxLen?: number; // default: 10000
  dlqMaxLen?: number; // default: 10000

  // Retries
  retryLimit?: number; // default: 3
  retryBaseDelayMs?: number; // default: 1000
  retryMaxDelayMs?: number; // default: 60000
  retryJitterMs?: number; // default: 250

  // Reclaim
  reclaimIntervalMs?: number; // default: 2000
  reclaimMinIdleMs?: number; // default: 30000
  reclaimBatchSize?: number; // default: 50

  // Hooks
  onProcessStart?(
    event: EventName<Events>,
    job: Job<Events, EventName<Events>>,
  ): void;
  onProcessEnd?(
    event: EventName<Events>,
    job: Job<Events, EventName<Events>>,
  ): void;
  onError?(error: Error, context: string): void;
  onMetric?(metric: RedisQMetric): void;
}

Delivery Semantics

RedisQ guarantees:

  • At-least-once delivery
  • Ordered processing within a stream
  • Reliable acknowledgements via XACK + XDEL
  • Automatic crash recovery via XAUTOCLAIM
  • Distributed processing across multiple workers

Production Tips

  • Use one Consumer Group per application.
  • Give every worker a unique consumerName (the default UUID is fine).
  • Producer instances do not need to call startWorker() — construct and add() is sufficient.
  • Monitor onMetric, onProcessStart, onProcessEnd, and onError for observability.
  • Configure streamMaxLen and dlqMaxLen to match your retention requirements.
  • Enable Redis persistence (AOF or RDB) for durability.
  • Make handlers idempotent — at-least-once delivery means a job may run more than once.

License

MIT