- Services, methods and typed clients
- How do I build my first TypeScript RPC service, step by step?
- How do I expose a service method so it can be called remotely?
- How do I return a complex type over RPC with classType and property?
- Why must removeComments be false in a project that uses @imqueue?
- How do I generate a typed client for a running service?
- Validating input
- Caching results
- Background and delayed work
- PostgreSQL notifications
- GraphQL composition
- Tracing and logging
- Scaling out
- Deploys and shutdown
- Hardening an HTTP gateway
- Encrypting the broker connection
- How do I encrypt the connection between my services and the broker?
- How do I turn on TLS across a fleet without changing application code?
- How do I connect to a broker with a private CA or mutual TLS?
- How do I encrypt a broker fleet that services discover over UDP?
- Why does my TLS connection to the broker fail with "Connection is closed"?
- Where to look next
Services, methods and typed clients
How do I build my first TypeScript RPC service, step by step?
Install the CLI, scaffold a service, implement the methods you want to expose, run it, and generate a typed client from the running process. Five commands, and the only prerequisites are Node.js 22.12+ and a reachable Redis 6.2+.
npm i -g @imqueue/cli
mkdir user-service && cd user-service
imq service create
npm run dev # terminal 1: the service
imq client generate UserService ./src/clients # terminal 2: the client
What the scaffold gives you is a class extending IMQService with @expose() on
each remotely callable method. There is no schema file and no IDL: the generated
client is built by asking the running service to describe itself, which is why
step four has to be running before step five.
Reference: IMQService ·
expose() ·
IMQClient.create(). Worked
walkthroughs: Get started and the Tutorial.
How do I expose a service method so it can be called remotely?
Decorate it with @expose() and give it a JSDoc block with typed @param and
@returns tags. A method without the decorator is absent from the service
description and stays callable in-process only — a remote call to it is rejected
with IMQ_RPC_NO_ACCESS.
import { IMQService, expose } from '@imqueue/rpc';
export class UserService extends IMQService {
/**
* Returns how many users are active
*
* @return {Promise<number>} - the number of active users
*/
@expose()
public async countActive(): Promise<number> {
return this.users.filter(user => user.isActive).length;
}
}
Four things decide whether it actually works. The doc-block is not decoration:
it is the only type source the client generator has, and the documented @param
list is also what the service's argument-count check validates — it must match
the method's real arity or calls fail with IMQ_RPC_INVALID_ARGS_COUNT. Combined
with @lock(), @cache or @logged(), @expose() must be the innermost
decorator, listed last: those replace the method with a (...args) wrapper, and
registering the wrapper records its rest parameter as the method's only argument.
It applies to instance methods only — on a static method it silently registers
under the pseudo-class name Function and the method stays unreachable. And
under standard decorators registration is deferred to an initializer that runs on
first construction, so the description is empty until an instance exists.
Reference: expose() ·
IMQService ·
lock() · cache ·
logged()
How do I return a complex type over RPC with classType and property?
Declare the type as a class, put @classType() on the class and @property() on
each field that crosses the wire. @property()'s first argument is the type in
TypeScript notation, and a second argument of true marks the field optional.
That pair is what lets the service and the generated client agree on the shape.
import { classType, property } from '@imqueue/rpc';
@classType()
export class UserObject {
@property('string')
id: string;
@property('boolean')
isActive: boolean;
// optional — pass true as the second argument
@property('string', true)
nickname?: string;
}
Under standard (TC39) decorators — the protocol @imqueue/rpc targets — a field
decorator cannot see its own class, so @property() only collects field
metadata and @classType() is what flushes it under the class name. Omitting it
raises no error: the type is silently missing from the RPC type description and
generated clients then reference a type nobody declared. @indexed() performs
the same flush plus an index signature, so a class carrying it does not also need
@classType(). A field may name another complex type ('UserCarObject') or an
array of one ('UserCarObject[]'); everything crossing the queue is JSON, so
returning a plain object literal instead of a declared class works at runtime and
types the field any on the client.
Reference: classType() ·
property() ·
indexed()
Why must removeComments be false in a project that uses @imqueue?
Because the JSDoc block on an exposed method is the runtime type source. Standard
decorators provide no runtime type reflection, so the client generator reads the
doc-block and nothing else — strip comments and it sees no types at all:
parameters and return values degrade to any, and the @param list the
service's argument-count check validates against disappears with them.
{
"compilerOptions": {
"target": "es2024",
"lib": ["es2024", "esnext.decorators"],
"experimentalDecorators": false,
// doc-blocks are the only type source the generator reads,
// so stripping comments leaves it nothing
"removeComments": false
}
}
This is a property of the consuming project, not of the packages: your service compiles fine with comments stripped, and the damage shows up later as an untyped or empty generated client.
Reference: expose() ·
@imqueue/rpc package reference ·
migration from 2.x to 3.x
How do I generate a typed client for a running service?
Run imq client generate <ServiceName> <outDir> while the service is up with
Redis reachable. Generation works by asking the running service to describe
itself, so there is no schema file to keep in step — and nothing to generate from
if the process is not running.
imq client generate UserService ./src/clients
The generated file exports exactly one thing: a namespace named after the service
with a lower-case first letter, holding a client class whose name is the
service's with a trailing Service replaced by Client. So UserService gives
you userService.UserClient, and there is no top-level UserClient to import.
Treat the file as a build artefact — regenerate it when the interface changes,
-o overwrites without prompting. IMQClient.create() does the same job
in-process when you would rather not shell out.
import { userService } from './src/clients/UserService.js';
const client = new userService.UserClient();
await client.start();
try {
console.log(await client.countActive());
} finally {
await client.destroy(); // or the process will not exit
}
Reference: IMQClient.create() ·
IMQClient ·
Clients & Versioning
Validating input
How do I validate method arguments with decorators before the method runs?
Declare the rules on the input class — @validate() per field, @validatable()
on the class to seal them — and guard the method with @validated(...). The
check runs before the method body, and the body then receives exactly what the
caller passed.
import { z } from 'zod';
import { validatable, validate, validated } from '@imqueue/validation';
@validatable()
class Credentials {
@validate(z.string().email())
email!: string;
@validate(z.string().min(8))
password!: string;
}
class AuthService {
@validated(Credentials)
async signIn(creds: Credentials): Promise<string> {
return `token-for-${creds.email}`;
}
}
@validatable() is not optional bookkeeping. Field validators are buffered until
a class decorator claims them, so a class that uses @validate() without it
hands its fields to the next class that is sealed — which then rejects valid
input over properties it does not declare, while the class with the real mistake
validates nothing. Note also that arguments are checked but never replaced, so
transforming schemas (z.coerce.number(), .trim(), .default(...)) validate as
expected and change nothing the method sees. And a failure does not reach a
remote caller as a ZodError: @imqueue/rpc converts whatever a method throws
into its own error payload, so the caller sees IMQ_RPC_CALL_ERROR with Zod's
issue list as the message string.
Reference: validated() ·
validate() ·
validatable() ·
schemaOf()
Caching results
How do I cache a service method result and invalidate it when a table row changes?
Decorate the service class with @PgCache() and the method with @cacheWith(),
naming the tables the result depends on. PostgreSQL then decides when the entry
dies: @PgCache() installs a change-notify trigger on each declared table and
subscribes to one LISTEN/NOTIFY channel per table, and a row change drops every
entry tagged with that table. So the entry lives exactly as long as the data
behind it is unchanged — no guessed TTL, no manual del().
import { PgCache, cacheWith } from '@imqueue/pg-cache';
@PgCache({
postgres: process.env.DB_URL!,
redis: { host: 'localhost', port: 6379 },
})
class UserService extends IMQService {
@cacheWith({ channels: ['users'] })
public async list(): Promise<User[]> {
return this.db.query('SELECT * FROM users');
}
}
Two things to know before you rely on it. The triggers and the subscription are
established in start(), so a service that never starts is never cached. And a
ChannelFilter given as an array of ChannelOperation is an exclusion list —
the operations named in it are the ones that do not invalidate — which reads the
opposite way round from how it looks.
Since 5.1.0, start() also waits for those triggers and subscriptions before it
resolves, so awaiting it is enough — up to 5.0.6 it was not, and a row changed in
the ~35ms after it could go unnoticed until the next change to that table or the
TTL. For that measured, along with how long invalidation actually takes after a
row changes, see
cache invalidation across services.
Reference: PgCache() ·
cacheWith() ·
CacheWithOptions ·
ChannelFilter ·
ChannelOperation
What does the pg-cache cacheBy decorator do that cacheWith does not?
@cacheWith() names its tables literally; @cacheBy() derives them from a model
and narrows invalidation using the field map the caller actually asked for, so a
change to a column nobody selected does not drop the entry. It is the same
mechanism with a smaller blast radius.
import { cacheBy } from '@imqueue/pg-cache';
class UserService extends IMQService {
// fields is the 2nd argument, which is where cacheBy looks by default
@cacheBy(User, { ttl: 30000 })
public async list(filter: UserFilter, fields?: any): Promise<User[]> {
return this.repo.find(filter, fields);
}
}
fieldsArg is the zero-based position of that field map among the method's
runtime arguments. Omit it and the decorator looks at the second argument — the
position @imqueue/pg-sequelize services pass it in — and pass -1 to disable
the analysis explicitly. The field map itself is the kind of thing
fieldsMap() from graphql-fields-list produces from an incoming GraphQL
request, which is where most callers get one.
Reference: cacheBy() ·
CacheByOptions ·
CacheByOptions.fieldsArg ·
channelsOf() ·
cacheWith()
Background and delayed work
How do I run a job later with a delay and retry it if it fails?
Push it with { delay } and return a positive number from the handler to have it
come back. The retry policy lives in the handler rather than in queue
configuration, because the handler's return value is the re-scheduling
instruction.
import JobQueue from '@imqueue/job';
const queue = new JobQueue<Email>({ name: 'Email' });
queue.onPop(async (email: Email) => {
try {
await send(email);
return -1; // done; no re-schedule
} catch (err) {
return 60000; // try again in a minute
}
});
await queue.start();
queue.push({ to: 'a@b.c', subject: 'Later' }, {
delay: 3600000, // not before an hour from now
ttl: 86400000, // and stop re-scheduling after a day
});
| Handler outcome | Effect |
|---|---|
| returns a positive number | re-scheduled after that many milliseconds |
returns 0 |
re-scheduled immediately — a hot loop, never use it to stop |
| returns a negative number, or nothing | stops; no re-schedule |
| throws | re-scheduled with the delay the job was originally pushed with |
That last row is the trap: a job pushed with no delay and then throwing is
dropped rather than retried. Catch your own errors and return a number. ttl
bounds how long a job stays worth re-scheduling, counted from the push. Delivery
is at-least-once, so handlers must be idempotent, and job data travels as JSON —
class instances, Date and undefined properties do not arrive as they left.
Since 3.1.0 you do not have to reason the trap out from the table: every handler
failure logs the decision it just made, as retry in <ms> or no retry, with
the message id — along with a retry the ttl suppressed, and a re-schedule whose
own write to redis failed. See
where a failed send() or push() shows up.
Reference: JobQueue ·
JobQueuePopHandler ·
PushOptions.delay ·
PushOptions.ttl ·
JobQueueWorker ·
JobQueuePublisher
PostgreSQL notifications
How do I listen for Postgres notifications with only one replica handling each?
Use PgPubSub and leave singleListener at its default of true. LISTEN/NOTIFY
is a broadcast — every listening connection receives every notification, so a
service scaled to N replicas handles each message N times. In single-listener
mode the replicas compete for a per-channel lock held as a row in PostgreSQL,
only the holder listens, and the rest stay connected as hot standbys.
import { type AnyJson, PgPubSub } from '@imqueue/pg-pubsub';
const pubSub = new PgPubSub({ connectionString: process.env.DB_URL });
pubSub.on('connect', async () => {
await pubSub.listen('UserChanged');
});
pubSub.on('message', (channel: string, payload: AnyJson) =>
handle(channel, payload),
);
await pubSub.connect();
The trade is worth stating plainly: this makes delivery at-most-once. NOTIFY has no backlog, so anything published while no process holds the lock is gone. A clean shutdown releases the lock and a standby takes over at once; an unclean exit leaves the channel unhandled until the next acquire retry. PostgreSQL also caps payloads at 8000 bytes. Where losing a message is unacceptable, pair this with a durable queue rather than replacing one.
Reference: PgPubSub ·
PgPubSubOptions.singleListener ·
PgPubSub.listen() ·
PgPubSub.notify()
How does the PgPubSub inter-process lock work, and when should I turn it off?
PgIpLock holds the lock for one channel as a row in PostgreSQL and retries
acquiring it on an interval; it also installs SIGINT, SIGTERM and SIGABRT
handlers so a shutdown releases the lock and another process can pick the channel
up immediately. Turn it off — singleListener: false, which swaps in NoLock —
only when every replica genuinely needs to see every notification, such as
invalidating an in-process cache in each one.
const pubSub = new PgPubSub({
connectionString: process.env.DB_URL,
singleListener: false, // every replica handles every notification
});
acquireInterval is the tuning knob while the lock is on, and it is a real
trade-off in both directions: too short, with many replicas, floods the database
with lock-acquire requests; too long widens the window in which a silent
disconnect leaves a channel unhandled. You do not normally construct PgIpLock
yourself — PgPubSub does it — but it is exported so the mechanism can be reused
elsewhere.
Reference: PgIpLock ·
NoLock ·
PgPubSubOptions.singleListener ·
PgPubSubOptions.acquireInterval
GraphQL composition
How do I avoid N+1 service calls when resolving nested GraphQL fields?
Take the loading out of the field resolvers. Declare a bulk loader per type and
the requirements between types once at start-up, then make a single load() call
in the top-level resolver: it walks the fields the client actually asked for,
merges everything that needs the same type into one filter, and calls each loader
once per level instead of once per parent object.
import { Dependency } from '@imqueue/graphql-dependency';
import { fieldsMap } from 'graphql-fields-list';
// at start-up, next to the type definitions
Dependency(UserType).defineLoader(async (context, filter, fields) =>
(await context.user.listUser(filter, fields)).data,
);
Dependency(CompanyType).require(UserType, () => ({
as: CompanyType.getFields().employees,
filter: {
[UserType.getFields().companyId.name]: CompanyType.getFields().id,
},
}));
// in the top-level company resolver
async function companies(source, args, context, info) {
const data = await context.company.listCompany(args);
return Dependency(CompanyType).load(data, context, fieldsMap(info));
}
Every object taking part must carry an id — loaded rows are matched back onto
their parents by id and by nothing else. Batching is per level rather than
global: siblings of one type run concurrently, and the next level down waits,
because a child's filter is built from values the parent level has just loaded.
The resolution cache lives for the duration of one load() call, so no request
can serve another request's stale data. On a type-graphql schema use
@DependencyFor() on the decorated class instead — and remember to run every
hook in schemaHooks after building the schema, because nothing else does, and
missing that step leaves the dependency fields silently empty.
For what that is worth in measured calls — the same query at 26 and at 3, why the difference does not show up in latency until the system is busy, and the two ways a dependency comes back silently empty — see the N+1 problem when GraphQL resolvers call microservices.
Reference: Dependency ·
defineLoader() ·
require() ·
load() ·
DependencyFor() ·
schemaHooks
Tracing and logging
How do I register OpenTelemetry instrumentation once at startup so every RPC is traced?
Register ImqueueInstrumentation at start-up, before any client or service is
constructed. Every RPC through @imqueue/rpc then produces a CLIENT span on the
calling side and a SERVER span on the handling side, linked into one trace by
context carried in the IMQ request metadata — with no changes to service or
client code.
import { NodeTracerProvider } from '@opentelemetry/sdk-trace-node';
import { registerInstrumentations } from '@opentelemetry/instrumentation';
import { ImqueueInstrumentation } from '@imqueue/opentelemetry';
new NodeTracerProvider().register();
registerInstrumentations({
instrumentations: [new ImqueueInstrumentation()],
});
"Before any client or service is constructed" is the whole reason this belongs at
start-up. The instrumentation works by patching @imqueue/rpc's exported default
option singletons rather than by hooking module loading, so anything built before
enable() has already copied those defaults and is not traced. The same
mechanism has a second failure mode: if the @imqueue/rpc this package resolves
is a different copy from the one your application imported — duplicate installs
at different tree depths — the patch lands on the wrong singletons and no spans
appear at all. Note also that this package only produces spans; exporting them
is the host application's job.
Since 4.1.0 there is a second entry point for exactly that resolution question.
enable() has to reach for @imqueue/rpc synchronously, which means require;
imqueueInstrumentation() resolves it with import() instead, so the module it
patches is the one the loader gave the application. It returns a promise, so it
has to be awaited before registerInstrumentations sees it, and the
instrumentation it hands back is already patched:
import { imqueueInstrumentation } from '@imqueue/opentelemetry';
registerInstrumentations({
instrumentations: [await imqueueInstrumentation()],
});
The package recommends this form for an ESM application. It is insurance rather
than a fix for something always broken: on Node 24.19 under tsx, require and
import hand back different namespace objects but the same default-option
singletons, so new ImqueueInstrumentation() patched what the application held
and traced normally. Prefer the factory in ESM, and reach for it first if spans
are missing under a loader that compiles TypeScript on the fly.
For what that produces — a measured three-process trace, how to read queue wait out of it, and the failure modes that leave a trace starting one hop in — see distributed tracing over a message queue.
Reference: ImqueueInstrumentation ·
imqueueInstrumentation() ·
traced() ·
traceStart() ·
traceEnd()
How do I write structured JSON logs from a service to a file?
Declare a file transport in LOGGER_TRANSPORTS and import the default logger.
There is nothing to construct and nothing to wire: the default export is already
configured from the environment, and its records go to a winston File transport,
which writes one JSON object per line unless you hand it a different format.
export LOGGER_TRANSPORTS='[{
"type": "file",
"options": { "filename": "/var/log/user-service.log" },
"enabled": true
}]'
export LOGGER_METADATA='{"service":"%name %version"}'
import logger from '@imqueue/async-logger';
logger.info('service started on port %s', port);
%name and %version are substituted from the running service's own
package.json, so one configuration can be shared across services. With no
transports configured the logger still works, console only — that is the intended
local-development mode, not a misconfiguration. One caveat that follows from the
package's name: console writes are deferred with setTimeout, so a process that
exits immediately after logging may lose the tail. Log a tick before exiting if
the last lines matter.
Reference: Logger ·
TransportOptions ·
getTransport() ·
defaultMetadata()
How do I configure async-logger Logger transports?
LOGGER_TRANSPORTS is a JSON array of transport declarations, each one
{ type, options, enabled }. type is 'file' or 'http' and any other value
is rejected at construction time; options is handed straight to the winston
transport constructor; and enabled: false skips a transport entirely, so it can
be left in the configuration and switched off per environment rather than
deleted.
export LOGGER_TRANSPORTS='[{
"type": "file",
"options": { "filename": "/var/log/user-service.log" },
"enabled": true
}, {
"type": "http",
"options": {
"ssl": true,
"port": 443,
"host": "http-intake.logs.datadoghq.com",
"path": "/v1/input/<API_KEY>"
},
"enabled": true
}]'
options is typed as winston's wide LoggerOptions for historical reasons, so
the compiler will not catch a mismatch for you: treat it as FileTransportOptions
or HttpTransportOptions according to type — filename for a file,
host/port/path/ssl for HTTP. LOGGER_METADATA is separate: a JSON object
of default fields attached to every record, with the same %name/%version
expansion.
Reference: TransportOptions.type ·
TransportOptions.options ·
TransportOptions.enabled ·
AsyncLoggerOptions.transports ·
Logger
Why does a failed send() or push() not throw, and where do I see it?
Because neither waits for the broker. send() resolves with a locally generated
message id before redis confirms the write, and push() returns synchronously
— so the rejection has nowhere to go by the time the caller has moved on. That is
the design, and it has not changed. What changed is that the failure is no longer
invisible: since @imqueue/core 3.4.0, @imqueue/rpc 3.7.0 and
@imqueue/job 3.1.0 these paths report through the logger you already
configured, with no verbose flag to turn on and no new API to call.
Set logger in the queue's options — anything matching ILogger, so the
@imqueue/async-logger default export drops in — and
watch for:
| Line | What it tells you |
|---|---|
| a write-failure episode on a queue | a send() was rejected by redis. The first rejection is logged with its operation, message id and code, further ones are counted, and the first success logs the recovery with that count |
[JobQueue] push error: at error level |
a job never made it onto the queue — including a write rejected after push() returned. Carries the queue, the requested delay and ttl, and a code |
a handler failure stating retry in <ms> or no retry |
whether that job is coming back. no retry is the throw-without-a-delay drop becoming provable rather than inferred |
response to request … has no pending call |
a reply arrived for a call nobody is waiting for any more — a late answer, or a backlog left by a process that is gone. This is how you tell it from a service that never answered |
| a subscription established, and restored after a reconnect | the absence of the restore line is what makes a lost subscription provable |
| no subscribers on a channel, or no server to publish to at all | a publish() that resolved successfully while reaching nobody |
Two properties are worth relying on. Nothing is quoted from the error. Only
an allow-listed failure code is printed — an IMQ_-prefixed framework code, a
system E… code, a small integer, a known redis reply code (WRONGTYPE,
NOSCRIPT, LOADING, …) or one of a few known redis-client failure messages
mapped to codes of the framework's own. Everything else, the error's message,
stack and class name included, comes out as unknown, because an application
error may carry personal data and an imq error carries the call arguments in its
properties. No line carries a message payload, call arguments or a raw redis key
either. And a broken logger cannot change what the queue does — every one of
these lines goes through a writer that contains its own failures.
The same reasoning changed @logged() in @imqueue/rpc 3.7.0: it now logs
Class.method() failed, code <code> instead of handing the caught error to the
logger. If you were parsing stacks out of those records, that is the one upgrade
note in this wave. Everything else about it is unchanged — which logger is
resolved, doNotThrow, and the value re-thrown.
Repeating conditions are reported on entering the state rather than per occurrence, so an outage costs a bounded number of lines. Nothing here schedules work or adds a timer.
Each of those lines is reproduced and measured in silent failures in a Node.js message queue, which also maps them onto what is worth alerting on.
Reference: IMQOptions.logger ·
ILogger ·
RedisQueue.send() ·
JobQueueOptions.logger ·
JobQueuePublisher.push() ·
logged()
Scaling out
How do I auto-scale @imqueue services?
Scale on queue depth, and let the service report it. Every IMQService can serve
GET /metrics with the number of messages waiting in its queue — one option
away, no exporter to install — which is the signal an autoscaler actually wants:
work waiting for this service, rather than the CPU it happens to be burning.
If you run KEDA, you may not need the endpoint at all: the backlog is a plain
Redis list at imq:<ServiceName>, so the built-in redis scaler reads it with no
exporter and no code change. Autoscale on queue depth, not
CPU covers that route, this
one, and the clustered case where only this one is correct.
import { IMQService, expose } from '@imqueue/rpc';
export class UserService extends IMQService {
// ...exposed methods
}
const service = new UserService({
metricsServer: {
enabled: true, // off by default
port: 9090, // the default
},
});
await service.start(); // the listener comes up with it
The listener answers exactly one route, in Prometheus exposition format, and 404 for anything else:
$ curl -s localhost:9090/metrics
queue_length{} 17
Point Prometheus at it, expose queue_length as an external metric through
prometheus-adapter, and a HorizontalPodAutoscaler reads it like any other:
metrics:
- type: External
external:
metric:
name: queue_length
target:
type: Value
value: "20" # scale out while more than 20 wait
Four properties of the number decide whether that loop behaves. It counts
messages waiting in the queue's main list: delayed messages not yet due, and
messages already leased to a worker under safeDelivery, are not included, so it
is a backlog gauge and not the amount of outstanding work. It is 0 while the
queue's writer is disconnected, which makes a broker outage indistinguishable
from an empty queue — keep minReplicas above zero so a Redis blip cannot scale
the service to nothing. Every replica of the service reads the same queue and so
reports the same figure, which is why an External target on the value is a
better fit than a per-pod average. And every process that starts the service
binds the port, so under multiProcess the primary plus N workers all try:
either leave it off there or expect N of the N+1 binds to fail. On a
ClusteredRedisQueue the figure is summed across every broker in the fleet.
One shutdown detail: the signal handlers close the listener on SIGINT/SIGTERM,
but destroy() does not — close
service.metricsServer yourself there, or the open listener keeps the process
alive.
Reference: IMQServiceOptions.metricsServer ·
IMQMetricsServerOptions ·
IMQMetricsServerOptions.enabled ·
IMQMetricsServerOptions.port ·
IMQMetricsServerOptions.queueLengthFormatter ·
DEFAULT_IMQ_METRICS_SERVER_OPTIONS ·
IMQService.metricsServer ·
IMessageQueue.queueLength() ·
ClusteredRedisQueue.queueLength() ·
IMQServiceOptions.multiProcess
How do I auto-scale the @imqueue broker?
Run each Redis with one of the two announcer modules and give every service and
client a UDPClusterManager. The module makes a broker announce itself over UDP;
the manager folds announced brokers into a ClusteredRedisQueue round-robin as
they appear and drops them when they stop announcing. Starting or stopping a
redis-server then is the scaling operation — nothing is restarted or
reconfigured.
Which module depends on one question: does your network deliver broadcast?
Both modules ship in one image,
ghcr.io/imqueue/redis-broker, and
IMQ_BROKER_MODE picks which one loads — so the answer changes one environment
variable rather than the deployment:
# L2 segment — bare metal, VMs, a docker bridge, your laptop
docker run -p 6379:6379 -e IMQ_BROKER_MODE=promoter \
ghcr.io/imqueue/redis-broker:7.4
# Kubernetes on GCP or any cloud VPC, where broadcast is dropped: the same
# datagram, unicast to every pod the K8s API lists. DEPLOYMENT_ENV is the
# NAMESPACE despite the name — take it from metadata.namespace, never type it.
docker run -e IMQ_BROKER_MODE=unicaster \
-e DEPLOYMENT_ENV="$POD_NAMESPACE" -e SELECTED_INTERFACES=10. \
ghcr.io/imqueue/redis-broker:7.4
The image also fixes the keyspace-event floor (Ex) in its config rather than
leaving it to the first client that connects, denies CONFIG SET at runtime, and
refuses to start on a misconfiguration that would otherwise discover nothing
silently. deploy/ in that repo carries the ServiceAccount, RBAC and
NetworkPolicy. To compile the modules yourself instead, both repositories are a
make away and remain the source of truth.
The client side does not care which one is running, because both emit the same datagram to the same port:
import { IMQService, UDPClusterManager } from '@imqueue/rpc';
const options = +(process.env.DISABLE_CLUSTER_MANAGER || 0)
? { cluster: [{ host: 'localhost', port: 6379 }] } // static
: { clusterManagers: [new UDPClusterManager()] }; // discovery
const service = new UserService(options);
Apply those same options to every service and every client in the fleet. Replies travel back through whichever broker the request round-robined onto, so a client pinned to one Redis silently misses responses that landed elsewhere.
The announce protocol is worth knowing because it sets the timings you observe. A
broker announces up every REDIS_BROADCAST_INTERVAL seconds (default 1) to
port 63000, and down on graceful shutdown, which removes it at once; a broker
that dies silently is dropped when its advertised liveness timeout plus
aliveTimeoutCorrection (5 s by default) passes without an announcement. So a
join takes about one interval, a graceful exit is immediate, and a crash costs a
few seconds during which sends may still be routed at it. Announcements are not
filtered by queue name — every cluster registered with a manager on that address
and port receives every server announced there, so unrelated fleets need distinct
addresses or ports. The datagrams are plain, unauthenticated UDP: keep port 63000
inside the cluster, and give every broker the same credentials, since any service
may connect to any broker it discovers.
The full recipes — building both modules, the RBAC a unicaster pod needs, and how the fleet behaves as brokers come and go — are in Auto-scaling Redis broker: with and without broadcast.
Reference: UDPClusterManager ·
UDPClusterManagerOptions ·
UDPClusterManagerOptions.port ·
UDPClusterManagerOptions.useAliveCheck ·
UDPClusterManagerOptions.aliveTimeoutCorrection ·
IMQOptions.clusterManagers ·
IMQOptions.cluster ·
ClusteredRedisQueue
Deploys and shutdown
How do I stop a deploy from dropping in-flight requests?
Set IMQ_DRAIN_ENABLE=1. SIGTERM and SIGINT then stop the service consuming,
wait for the requests already running to finish and publish their replies, tear
the transport down and exit 0 — instead of the default, which starts destroy()
without awaiting it and force-exits on a fixed timer, abandoning the handler and
leaving its caller on a promise that never settles.
IMQ_DRAIN_ENABLE=1
That is the whole opt-in. There is nothing to wrap and nothing to register:
tracking sits at the single point where a request is dispatched, so every
@expose()d method is covered automatically. Both controls are constructor
options too, for services that would rather configure in code:
const service = new OrderService({ drain: true, drainTimeout: 4000 });
| control | option | environment | default |
|---|---|---|---|
run a drain on SIGTERM/SIGINT |
drain |
IMQ_DRAIN_ENABLE |
off |
| drain budget, milliseconds | drainTimeout |
IMQ_DRAIN_TIMEOUT |
4000 |
Both are read numerically, like the rest of the IMQ_* family — and a
non-numeric value throws at construction rather than being quietly read as off,
which is the failure mode that would otherwise leave you believing a deploy was
draining when it was not. A well-meant IMQ_DRAIN_ENABLE=true is the exact case:
it coerces to NaN, and silence there would be worse than a crash.
A second signal during a drain exits immediately, so an impatient operator or a
supervisor escalating to a second SIGTERM is never blocked by it.
Available from @imqueue/rpc 3.8.0 and @imqueue/job 3.2.0. It is opt-in
precisely so that upgrading changes nothing: left off, shutdown behaves exactly as
it always has.
Delivery is at-least-once either way. A drain narrows the window in which
in-flight work is lost — SIGKILL, an OOM kill or a lost node still take it with
them — so handlers must stay idempotent regardless.
Reference: IMQService ·
IMQServiceOptions ·
expose(). Measured end to end in
graceful shutdown and zero-drop deploys.
How long should the drain budget be?
Longer than your slowest handler, and shorter than whatever is about to SIGKILL
the process — the budget is only useful in the gap between those two. The 4000
default is sized for the tighter of the two common cases, not the friendlier one.
Two windows tend to bound it:
| what stops the process | window | what fits |
|---|---|---|
imq stop |
~5 s SIGTERM → SIGKILL |
the 4000 default, with about a second to spare for destroy() |
| Kubernetes | 30 s terminationGracePeriodSeconds |
raise drainTimeout to match, or the pod waits on a drain that gave up at 4 s |
Overrunning is bounded rather than fatal: the wait always ends, whatever has not
finished is abandoned, the outcome is logged, and the process exits 0. So the
cost of setting it too low is the very thing you turned the drain on to avoid,
while the cost of setting it too high is a slower rollout — an easy trade to get
right in one direction and not the other.
Under multiProcess, each forked worker drains the work it is holding, and the
cluster primary — which runs a consumer of its own — drains only what that
consumer had. The budget is per process, so a rollout does not pay it N times.
Reference: IMQServiceOptions ·
IMQService
Do background jobs drain too?
Yes, under the same IMQ_DRAIN_ENABLE and IMQ_DRAIN_TIMEOUT, with one addition
that matters more for jobs than for RPC: whatever has not finished when the budget
runs out is put back on the queue rather than dropped.
That difference follows from where a job is in its life when the signal lands. An abandoned RPC request leaves a caller waiting, and that caller will time out and can retry; an abandoned job has already been popped, so without re-queueing there is nobody left holding it and nothing to notice it is gone. Re-queueing turns a deploy that overruns its window into a re-run rather than a silent loss — which again asks that handlers be idempotent, exactly as at-least-once delivery already did.
Reference: JobQueue ·
JobQueueOptions
Why does enabling the drain turn off handleSignals?
Because the queue layer registers its own SIGTERM/SIGINT handlers that exit
the process without waiting for anything, and one of those firing mid-drain would
end the process from the side, at an arbitrary point, with handlers still running.
Enabling drain therefore forces
handleSignals to false on
the service's queue.
It is worth knowing what this does not touch. The drain takes over only the handlers this framework registered, matched by exact function reference — not by clearing the signal's listener list — so a handler installed by your application, by a tracing SDK or by any other library keeps working exactly as before. A drain that silently unhooked another library's shutdown logic would trade one class of lost work for another.
Reference: IMQOptions.handleSignals ·
IMQServiceOptions
Hardening an HTTP gateway
How do I protect an HTTP gateway from too many requests per IP?
Mount HttpProtect. It counts requests per client IP in Redis and, past two
configurable thresholds, first answers 429 and then adds the address to a
persistent block list that is answered with 418.
import HttpProtect from '@imqueue/http-protect';
const protect = new HttpProtect({ ttl: 60, maxRequests: 600, banLimit: 5000 });
app.use(protect.jsonMiddleware());
Three things here are load-bearing and none is visible in a signature. A ban is
permanent — addresses go into a Redis set that is never given an expiry, and no
method in the package removes one, so an address stays banned until something
outside it deletes the key. The counter measures a continuous stream, not a
fixed window — its TTL is pushed back to ttl on every request, so the count
only resets after a full ttl of silence; the default maxRequests of 200
therefore stops a steady one-request-per-second client after about 200 seconds,
not just a 200-request burst. And everything is keyed by the client IP, which
by default comes from proxy headers a client can set, so set getClientIp before
exposing this behind a proxy you do not control.
Reference: HttpProtect ·
HttpProtect.maxRequests ·
HttpProtect.banLimit ·
HttpProtect.ttl ·
HttpProtectOptions.getClientIp
How do I mount HttpProtect as express middleware?
One app.use(), before the routes it protects, with the middleware that matches
how your gateway answers errors: jsonMiddleware() for a JSON API,
textMiddleware() for plain text, middleware() for the bare form.
import HttpProtect from '@imqueue/http-protect';
// 429 then 418, as JSON, on default thresholds and a local Redis
app.use(new HttpProtect().jsonMiddleware());
If you would rather decide what happens yourself — a custom body, a redirect, a
metric — skip the middlewares and call verify() on the request, then act on the
VerificationStatus and httpCode it hands back:
import HttpProtect, { VerificationStatus } from '@imqueue/http-protect';
const protect = new HttpProtect();
const { status, httpCode } = await protect.verify(req);
if (status !== VerificationStatus.SAFE) {
res.status(httpCode).end();
}
Construct one instance and share it: each one holds its own Redis connection.
Reference: HttpProtect.jsonMiddleware() ·
HttpProtect.textMiddleware() ·
HttpProtect.middleware() ·
HttpProtect.verify() ·
VerificationStatus
How do I check whether an IP address is inside a CIDR range?
Build a Networks from your CIDR records and ask it. It covers IPv4 and IPv6 in
one object, dispatching on the address it is given, and answers from sorted
binary ranges rather than by comparing against each network in turn.
import { Networks } from '@imqueue/net';
const allowed = new Networks(['10.0.0.0/8', '192.168.0.0/16', '2001:db8::/32']);
allowed.includes('10.1.2.3'); // true
allowed.includes('8.8.8.8'); // false
allowed.includes('2001:db8::1'); // true
Every record needs an explicit prefix length: a bare address is rejected, so a
single host is 203.0.113.7/32 or 2001:db8::1/128. Anything invalid throws
while parsing rather than being skipped, so one bad entry fails the whole list —
validate with isValid() first when the input is untrusted. cidrToRange() and
ipToInt() are exported for building something else on the same primitives; this
is the package @imqueue/http-protect uses for its own allow-list.
Overlapping records are coalesced at construction, so listing a network beside a
subnet of it is safe but not free of consequence: the two collapse into one
record, length counts fewer networks than you passed, and toArray() gives the
supernet back on its own. The addresses covered are identical either way — it is
only the record list that differs.
import { isValid, Networks } from '@imqueue/net';
const records = input.filter(r => isValid(r.split('/')[0]));
const networks = new Networks(records);
Reference: Networks ·
Networks.includes() ·
isValid() ·
cidrToRange() ·
ipToInt()
Encrypting the broker connection
How do I encrypt the connection between my services and the broker?
Set tls on the queue options. A queue opens several connections — reader,
writer, watcher and subscription — and they are all built by the same factory,
so one setting encrypts the whole bus.
import IMQ from '@imqueue/core';
import { readFileSync } from 'node:fs';
const queue = IMQ.create('user-service', {
host: 'redis.internal',
port: 6380,
tls: { ca: readFileSync('/etc/ssl/internal-ca.crt') },
});
true connects with Node's defaults, verifying the broker against the system
trust store — right for a managed Redis with a publicly signed certificate. An
object is handed to tls.connect() as given, so anything Node accepts works.
@imqueue/rpc caches and @imqueue/job queues take the same option and behave
the same way.
The broker has to be listening for TLS, which in Redis means tls-port. Turn
the plaintext port off as well, because encryption you can opt out of is a
suggestion rather than a guarantee:
redis-server --port 0 --tls-port 6380 \
--tls-cert-file /etc/redis/server.crt \
--tls-key-file /etc/redis/server.key \
--tls-ca-cert-file /etc/redis/ca.crt
If the brokers announce themselves for discovery, use the
redis-broker image instead of those
flags — from v1.2.0 it composes them from IMQ_TLS_*, keeps TLS on 6379 so no
port moves, and announces the port it is really listening on:
docker run -v /path/to/tls:/run/tls:ro \
-e IMQ_TLS_CERT_FILE=/run/tls/broker.crt \
-e IMQ_TLS_KEY_FILE=/run/tls/broker.key \
-e IMQ_TLS_CA_FILE=/run/tls/ca.crt \
ghcr.io/imqueue/redis-broker:7.4
There is no negotiation step and therefore no downgrade: a plaintext client cannot reach a TLS-only broker, and an encrypted client cannot be talked into plaintext by a broker that is not configured for TLS. Both simply fail.
Reference: IMessageQueueAuthConnection.tls ·
JobQueueOptions.tls ·
IRedisCacheOptions
How do I turn on TLS across a fleet without changing application code?
Leave tls unset and configure the environment instead. When the option is
absent, @imqueue/core builds one from IMQ_REDIS_TLS and its companions, so
encrypting every service becomes a deployment change in one place rather than a
pull request per repository.
IMQ_REDIS_TLS=1
IMQ_REDIS_TLS_CA_FILE=/etc/ssl/internal-ca.crt # private CA
IMQ_REDIS_TLS_CERT_FILE=/etc/ssl/service.crt # mutual TLS
IMQ_REDIS_TLS_KEY_FILE=/etc/ssl/service.key
IMQ_REDIS_TLS_KEY_PASSPHRASE=… # encrypted key
IMQ_REDIS_TLS_SERVERNAME=redis.internal # expected cert name
IMQ_REDIS_TLS_REJECT_UNAUTHORIZED=0 # local experiments only
Three rules govern how they combine. Supplying key material is enough on its
own — naming a CA file enables TLS, because there is no other reason to have
named one. The off switch beats everything: IMQ_REDIS_TLS=0 disables TLS
even when certificates are configured, and short-circuits before the files are
read, so a rollback works even if the certificates are already gone from the
image. And options that only shape a connection cannot start one — setting
IMQ_REDIS_TLS_REJECT_UNAUTHORIZED=0 or IMQ_REDIS_TLS_SERVERNAME by itself
does not enable anything.
The same variables cover @imqueue/core, @imqueue/rpc caches and
@imqueue/job queues. Passing tls: false explicitly declines the fallback for
a service that must stay in plaintext; leaving it unset means "ask the
environment".
Reference: envTls() ·
IMessageQueueAuthConnection.tls
How do I connect to a broker with a private CA or mutual TLS?
Pass your trust anchor as ca to verify the broker, and add cert and key to
have the broker verify you:
const tls = {
ca: readFileSync('/etc/ssl/internal-ca.crt'),
cert: readFileSync('/etc/ssl/user-service.crt'), // mutual TLS
key: readFileSync('/etc/ssl/user-service.key'),
};
With tls-auth-clients yes on the broker, a client without a certificate signed
by that CA is refused at the handshake, before it can send AUTH.
Two things catch people. The certificate is verified against the host you
connected to, so reaching a broker by bare IP needs that address in the
certificate as an IP SAN — otherwise set servername to the name the
certificate carries, or connect by that name. And rejectUnauthorized: false
keeps encryption while discarding authentication, which protects you from
someone reading the wire and not at all from someone intercepting it; the queue
logs a warning at construction whenever it is set, for that reason.
Connections are pooled per broker within a process, and the pool key includes a
fingerprint of the TLS configuration, so a queue asking for an encrypted
connection is never handed a plaintext socket that another queue opened first.
Configurations equal by value still share a connection. One caveat: an opaque
object such as a prebuilt SecureContext is fingerprinted by its class name
alone, so pass the certificate material rather than a prebuilt context.
Reference: tlsFingerprint() ·
IMQOptions.cluster ·
IMessageQueueAuthConnection.tls
How do I encrypt a broker fleet that services discover over UDP?
Exactly as above, with one difference that discovery forces: the broker's address is not knowable in advance. Brokers announce whatever IP the scheduler gave them, so no certificate can carry it, and there is no name to connect by either — the fleet is found by announcement rather than by lookup.
Issue one certificate for the whole fleet, carrying a name that will never be resolved, and pin that name on the services:
# broker: ghcr.io/imqueue/redis-broker v1.2.0+
IMQ_TLS_CERT_FILE=/run/tls/broker.crt # CN/SAN: imq-broker.internal
IMQ_TLS_KEY_FILE=/run/tls/broker.key
IMQ_TLS_CA_FILE=/run/tls/ca.crt # signs the client certificates
# service
IMQ_REDIS_TLS_CA_FILE=/run/tls/ca.crt
IMQ_REDIS_TLS_SERVERNAME=imq-broker.internal
servername is not a host to connect to — Node compares it against the
certificate while the connection goes to the announced IP — so a broker pod that
dies and comes back on a different address needs nothing reissued.
Two things are the broker's side of it. Redis serves TLS by setting port 0 and
tls-port, and the announcer modules advertise the port that is listening:
before redis-broker v1.2.0 they advertised port verbatim, so a TLS broker
announced <ip>:0 — an address nothing can connect to, and one
UDPClusterManager discards as malformed, leaving the fleet with no broker and
no error. And the announcement carries one transport for everybody, so switching
a running fleet is a cutover: brokers up with both listeners
(IMQ_TLS_PLAINTEXT=on), then the services and REDIS_BROADCAST_TLS=1 together,
then drop both.
Reference: auto-scaling Redis broker ·
TLS to the broker ·
UDPClusterManager ·
envTls()
Why does my TLS connection to the broker fail with "Connection is closed"?
Because TLS verification fails below the queue, and what comes back up is a closed socket rather than the reason. A wrong trust anchor, a name the certificate does not carry, a missing client certificate and a broker that is not listening for TLS all report the same message. Only the timing differs: a plaintext client against a TLS broker is rejected immediately, while an encrypted client against a plaintext broker takes about ten seconds to give up.
Get the real error from the broker directly, with the same trust anchors:
openssl s_client -connect redis.internal:6380 \
-CAfile /etc/ssl/internal-ca.crt -servername redis.internal
The one TLS failure that is diagnostic is unreadable key material. If a
certificate file is named but cannot be read — a mistyped path, a secret that
did not mount — construction throws an error with the code
IMQ_TLS_MATERIAL_UNREADABLE naming the variable at fault, rather than falling
back to an unencrypted connection. A broken TLS configuration stops the service;
it never quietly starts one that talks in the clear.
Reference: envTls() ·
RedisQueue
Where to look next
The full generated API reference covers every exported symbol of every
documented package, and /api/search-index.json is the
same set as one machine-readable feed — {name, kind, package, url, summary} per
symbol — if you would rather look a name up than browse. For anything this page
does not answer, support lists which repository to file in.