Message brokers
Scaling uses a broker as a backplane: a fan-out that makes a room span every instance, at most once. That is the right contract for presence and typing indicators, and the wrong one for everything a broker is usually deployed for — work that must not be lost, history a client resumes from, services talking to each other without HTTP between them.
@alexify/wrpc/broker is the broker-agnostic core for those. It describes a broker by what it can do, and every wRPC feature built on brokers asks only for the capability it needs:
| Capability | Guarantee | What wRPC builds on it |
|---|---|---|
backplane | at-most-once fan-out | rooms and the cluster across instances |
log | ordered, replayable | durable subscription feeds that resume on any instance |
queue | at-least-once, competing consumers | procedures invoked from a queue, with retry and dead-lettering |
direct | addressable inboxes | RPC between services over the broker itself |
Experimental
The whole @alexify/wrpc/broker* family is @experimental: it may change in a minor release until every adapter has shipped. See Stability.
No broker is a dependency
Like the Redis backplane, every adapter takes your client and checks its shape; @alexify/wrpc installs no broker package and never will. Each broker gets its own subpath, so an application that uses Kafka does not load the NATS adapter:
| Broker | Subpath | backplane | log | queue | direct |
|---|---|---|---|---|---|
| in-process | @alexify/wrpc/broker | ✓ | ✓ | ✓ | ✓ |
| Redis (Streams + pub/sub) | @alexify/wrpc/broker/redis | ✓ | ✓ | ✓ | ✓ |
| NATS + JetStream | @alexify/wrpc/broker/nats | ✓ | ✓ | ✓ | ✓ |
| RabbitMQ (AMQP 0-9-1) | @alexify/wrpc/broker/amqp | ✓ | ✓ | ✓ | ✓ |
| Kafka | @alexify/wrpc/broker/kafka | caveats | ✓ | ✓ | — |
Kafka has no direct: consumer-group rebalances and a topic per instance make it a poor carrier for request/response. Its backplane works, with costs the Kafka page spells out. Brokers that speak one of these protocols need no adapter of their own — Valkey, KeyDB and Dragonfly use the Redis one; Redpanda and Azure Event Hubs' Kafka endpoint use the Kafka one.
Choosing one
| If you… | Use |
|---|---|
| already run a Redis, and want feeds and queues without new infrastructure | Redis |
| want the lowest-latency RPC and a backplane that costs nothing | NATS (JetStream for the durable half) |
| need per-message TTLs, dead-letter exchanges and operator tooling | RabbitMQ |
| keep long histories other systems replay, or already stream through Kafka | Kafka — for feeds and queues; put the backplane elsewhere |
| are writing tests, or running one process | the in-process MemoryBroker below |
Nothing stops you from using two: a Redis backplane with Kafka feeds is a perfectly ordinary deployment, and each binding takes the broker it needs.
The in-process broker
const { MemoryBroker } = require('@alexify/wrpc/broker');
const broker = new MemoryBroker();
const a = new RpcServer({ router, backplane: broker.backplane });
const b = new RpcServer({ router, backplane: broker.backplane });MemoryBroker implements every capability in one process. It is the reference implementation the contract tests are written against, and what makes a multi-instance setup testable without infrastructure: two servers sharing one MemoryBroker behave like two processes sharing a real broker. Delivery is always deferred to a microtask, as a network broker's would be, so code that accidentally relied on synchronous delivery fails in the test suite and not after a deploy.
| Option | Default | |
|---|---|---|
retention.maxEntries | 10000 | Entries kept per log topic; older ones are trimmed |
epoch | random | Stamped into every log id, so an id from another broker is recognized as foreign |
prefix | 'wrpc' | Backplane channel namespace |
logger | console | false silences it |
The contracts
A broker is a plain object: { name, backplane?, log?, queue?, direct?, close() }. isBroker(value) and isBrokerLog/isBrokerQueue/isBrokerDirect are the structural checks every entry point runs.
log
const id = await log.append('orders', JSON.stringify(order), { headers, key });
const read = log.read('orders', { after: lastEventId, signal });
await read.ready; // the position is fixed from here on
for await (const { id, value, headers } of read) { /* … */ }afterreads strictly after an id this log minted. An id older than the retained history — or minted by another log — fails the read with 410; a malformed id, or one past the tip, with 400.- Without
after,from: 'latest'(the default) reads only new entries and'earliest'everything retained. parseId(text)checks an id's syntax. AlastEventIdcomes from the peer, so it is refused before it reaches the broker.- The
ida read yields is the resume token after that entry. For a single-sequence log it is the entry's own id; for Kafka it is a vector of partition offsets.
queue
await queue.produce('invoices', body, { headers });
const consumer = await queue.consume('invoices', async (delivery) => {
try {
await handle(delivery.body);
await delivery.ack();
} catch {
await delivery.retry({ delay: 1000 });
}
}, { prefetch: 16, deadLetter: 'invoices.dlq' });- Consumers of one queue compete: each message reaches one of them, at least once.
prefetchbounds the unsettled deliveries a consumer holds, and they run concurrently. - A delivery settles once — the first of
ack(),retry({ delay }),release()ordeadLetter(reason)wins:
| Settlement | Effect | attempt |
|---|---|---|
ack() | done | — |
retry({ delay }) | redelivered after delay ms | + 1 |
release() | back to the queue now (a draining node handing work over) | unchanged |
deadLetter(reason) | moved to the deadLetter queue with x-wrpc-dead-reason (one line, at most 512 characters) and x-wrpc-attempt headers | — |
- A handler that throws (or rejects) has settled nothing: the adapter retries the delivery after a backoff (50 ms doubling to 1 s), attempt
- 1, and logs it — never a release, which would put the message back at the head of the queue and spin it.
release()is for draining and stopping.
- 1, and logs it — never a release, which would put the message back at the head of the queue and spin it.
- A settlement the broker refuses — a retry's copy it would not take, a commit it would not accept — leaves the message available: the adapter tries the settlement again, then hands the message back (a requeue on RabbitMQ, a seek to its offset on Kafka) rather than moving past it. Logged
broker.<name>.settle;healthyisfalseon Kafka until a settlement lands again. - The attempt counter belongs to the adapter, not the broker: RabbitMQ 4 does not count a requeue, so the adapters carry it in an
x-wrpc-attemptheader where the broker has none. stop()takes no new deliveries; whatever was unsettled is redelivered later.
direct
const stop = await direct.listen('svc.billing', onMessage, { group: 'svc.billing' });
await direct.send('svc.billing', packet, { correlationId, replyTo: direct.inbox(), timeout: 5000 });inbox()answers a unique, routable address.- Listeners sharing a
groupcompete; ungrouped listeners each receive every message. Messages from one sender to one address arrive in order. - Delivery is at-most-once. The RPC binding built on it numbers its frames and treats a gap as a lost connection.
sendmay reject with 503 when the broker knows nobody listens, and may discard a message nobody took withintimeout. stop.healthyis optional: an adapter that knows a listener went deaf — a failing read, a cancelled consumer, a closed channel — answersfalsethere, and the RPC binding'shealthyis made of it. An adapter that cannot tell leaves it out.
backplane
The same three methods Scaling documents. The broker family adds one requirement the contract tests enforce: channel names are literal. A room called room:* must never become a wildcard subscription, and an adapter encodes names into whatever its broker allows.
Names the adapter mints
Every adapter puts a few names of its own on the wire — a consumer name, a reply inbox, a consumer group, a message id — and each takes a generateId option to replace the default:
createRedisBroker({ client, generateId: () => myUlid() });Two things to know. The value is used verbatim: wRPC never truncates it, because trimming an id would quietly weaken the uniqueness you chose it for. And the alphabet is the broker's business, not wRPC's — a generator answering characters Redis refuses in a consumer name, NATS in a subject or RabbitMQ in a queue name fails at the driver, on connect, not here. Stick to [A-Za-z0-9_-] unless you know the broker better than that.
The option is strict: a bad generator is a TypeError at construction. See Identifiers for the whole picture.
Writing an adapter
The contracts are executable. tests/broker/{backplane,log,queue,direct}Contract.js in the repository run against the memory broker, every adapter's in-repo fake, and — in the broker CI jobs — a real server. A new adapter that passes them is a working adapter.
Two building blocks carry the parts every adapter would otherwise get wrong:
TopicTailskeeps one live reader per topic per instance, shared by every local read. A feed is read by many subscriptions at once, and a broker consumer per subscription does not survive Kafka (a group join each) or Redis (a blocked connection each). The adapter supplieslive(a positioned tail) andrange(a catch-up page);TopicTailsjoins them without a gap or a duplicate, and sends a reader that falls behind back torangeso memory per slow subscriber stays bounded. Arangeread straight off a store answers an array, and a short page is the end; one that pages by time — a stream reader, a consumer group's fetch, a pull with an expiry — answers{ entries, done }, so a page cut short by a pause is never taken for the tip. An optionalcontiguous(cursor, entry)lets a log with dense ids refuse a page whose head skipped.encodeToken(name, { safe, escape, maxLength })maps an arbitrary name into one token of a broker's alphabet, injectively, shortening pastmaxLengthwith a digest (AMQP routing keys stop at 255 bytes, Kafka topic names at 249 characters).