@nxgt/janus-webhooks-redis
v0.5.0
Published
The Redis queue for @nxgt/janus-webhooks on @nxgt/redis: deliveries that outlive the process, claimed atomically under a lease
Maintainers
Readme
@nxgt/janus-webhooks-redis
The Redis queue for @nxgt/janus-webhooks:
deliveries wait in Redis, on
@nxgt/redis's connection,
shared by every process of your application. A retry waiting when a process
crashes, restarts or is redeployed is sent by the next one, and a request cut
short is sent again once its lease lapses.
It implements the WebhookQueue port, and passes the
@nxgt/janus-webhooks/conformance suite against a real Redis, outages
included, on every CI run.
0.x. A minor version may still change the surface; the changelog says how.
Install
bun add @nxgt/janus-webhooks-redis @nxgt/janus-webhooks @nxgt/janus @nxgt/redis
bun add zod # @nxgt/redis's peer, zod 4, if your application has none
bun add -d typescript # 6Every peer is required:
@nxgt/janus-webhooks, whose port it implements;@nxgt/janus, whoseUserEventit holds and whoseStoreFailureit throws — a peer, never a dependency, soinstanceof StoreFailureholds in your code;@nxgt/redis,>=0.3.1 <1, whose connection wraps Bun's ownRedisClient, so this runs on Bun;typescript6.
It needs Redis 7.0 or later, or Valkey — see
what Redis must be configured with.
Like @nxgt/janus, it expects "moduleResolution": "bundler".
Usage
import { webhooks } from '@nxgt/janus-webhooks';
import { createRedisWebhookQueue } from '@nxgt/janus-webhooks-redis';
import { connectRedis } from '@nxgt/redis';
const secret = process.env.CRM_WEBHOOK_SECRET; // whsec_…, from mintWebhookSecret()
if (!secret) throw new Error('CRM_WEBHOOK_SECRET is not set');
const redis = await connectRedis(process.env.REDIS_URL ?? 'redis://localhost:6379', {
enableOfflineQueue: false, // an outage fails the insert at once, not after 31 s
});
export const listener = webhooks({
endpoints: [{ id: 'crm', url: 'https://crm.example.com/hooks/janus', secrets: [secret] }],
queue: createRedisWebhookQueue(redis), // every process passes the same Redis and prefix
});
// janus({ …, events: listener })
process.on('SIGTERM', async () => {
await listener.close(); // gives nothing up: what waits is sent by the next process
process.exit(0);
});Give each endpoint an id when you first pass a queue: it is what the queue
knows the endpoint by. Wiring covers the connection,
the prefix, and what Redis must be configured with.
API
| Export | What it is |
| --- | --- |
| createRedisWebhookQueue(redis, options?) | Returns a WebhookQueue, for webhooks({ queue }). redis is what connectRedis returns. It connects to nothing and creates nothing: the keys are made by the first insert. |
| RedisWebhookQueueOptions | { prefix? }: what every key starts with, janus:webhooks: by default. Every process sharing deliveries passes the same one. |
What Redis holds
Each delivery's event and endpoint id — never a URL or a secret — with one Lua script per method, so no two claims answer one delivery. Wiring has each key and how long it stays.
Traps
Hold the new event types back until every process sharing a queue is upgraded. A delivery of
user.secondFactorEnabledoruser.secondFactorDisabled, written by 0.2.0, isSTORE_FAILED(a reply that is not … a user event type) in a 0.1.x process that claims it — and so is one ofuser.recoveryCodesRegeneratedoruser.recoveryCodeUsed, written by 0.3.0, in a 0.2.x process, and one ofuser.passwordChangedoruser.emailChanged, written by 0.4.0, in a 0.3.x process, and one ofuser.newDeviceSignedIn, written by 0.5.0, in a 0.4.x process — and that endpoint's claims fail there until every process is upgraded. Upgrade@nxgt/janus,@nxgt/janus-webhooksand@nxgt/janus-webhooks-redistogether — their peer ranges move as one — with each endpoint limited to thetypesthe older processes know, and drop the limit once every process runs the new versions:webhooks({ queue, endpoints: [{ id: 'crm', url, secrets: [secret], types: ['user.created', 'user.emailVerified', 'user.passwordReset', 'user.deleted'] }], });That list is for 0.1.x processes. When the oldest run 0.2.x, which read the second factor's events but not the recovery codes', add those two:
types: ['user.created', 'user.emailVerified', 'user.passwordReset', 'user.secondFactorEnabled', 'user.secondFactorDisabled', 'user.deleted'],When the oldest run 0.3.x, which read the recovery codes' events but not the change events', add those two:
types: ['user.created', 'user.emailVerified', 'user.passwordReset', 'user.secondFactorEnabled', 'user.secondFactorDisabled', 'user.recoveryCodesRegenerated', 'user.recoveryCodeUsed', 'user.deleted'],When the oldest run 0.4.x, which read the change events but not the new device's, add those two:
types: ['user.created', 'user.emailVerified', 'user.passwordReset', 'user.passwordChanged', 'user.emailChanged', 'user.secondFactorEnabled', 'user.secondFactorDisabled', 'user.recoveryCodesRegenerated', 'user.recoveryCodeUsed', 'user.deleted'],Eviction is data loss. A Redis whose
maxmemory-policyevicts keys drops deliveries without a word — no retry, noonGivingUp. Run this on a Redis withnoeviction; a full one then refuses the insert withOOM, whichjanusreports asJANUS_EVENT_FAILED.Persistence is what makes it durable. Without RDB or AOF, restarting Redis loses every delivery waiting in it. Enable AOF (
appendfsync everysecloses at most a second of inserts on a crash) or RDB snapshots.Redis Cluster is not supported. A script touches the keys of a delivery and of its endpoint, which may live in different slots. Use a single Redis, or a primary with replicas.
By default, an outage makes every flow wait 31 seconds. The listener awaits the insert, and Bun's client queues commands while it reconnects. With
enableOfflineQueue: falsethe insert fails at once, the flow goes on, and the event is reported asJANUS_EVENT_FAILED.One prefix per application. Two applications sharing a prefix share deliveries, and each gives up the other's endpoints as
endpointRemovedafterorphanGrace.
Documentation
- Guides: wiring the queue, the prefix, what Redis must be configured with, what Redis holds and how each method stays atomic
- Troubleshooting: look up the error message you see
- Roadmap: what is next, and what is not planned
@nxgt/janus-webhooks's queues guide: claims, leases, orphans, and what a queue must do
Type safety, counted
Four plausible mistakes are refused by the compiler, each with a
@ts-expect-error case in test/types/queue.ts:
- a URL instead of a connection;
- Bun's client instead of the connection that holds it;
- the promise
connectRedisanswers, not awaited; - a
prefixthat is not a string.
Licence
MIT
