Skip to content

NATS ​

NATS is the most natural fit of the four: core subjects are a backplane and a request/reply fabric by design, and JetStream adds the durable half — feeds and work queues — without a second system.

js
const { connect, headers, createInbox } = require('@nats-io/transport-node');
const { jetstream, jetstreamManager } = require('@nats-io/jetstream');
const { createNatsBroker } = require('@alexify/wrpc/broker/nats');

const nc = await connect({ servers: process.env.NATS_URL });
const broker = createNatsBroker({ nc, headers, createInbox, jetstream, jetstreamManager });

const server = new Server({ router, backplane: broker.backplane, port: 8000 });

headers is a package export, not a method on the connection — that is why it is injected separately. Leave jetstream/jetstreamManager out and you get a backplane + RPC broker with no log and no queue; everything else works unchanged.

OptionDefault
nc—The connection; close() never drains it — its lifetime is yours
headers—The headers() factory
jetstream, jetstreamManager—Both, or neither
createInboxa uuid subjectThe createInbox() export
prefix'wrpc'Subject and stream-name namespace
ackWait30000JetStream ack window for queue consumers (ms)
maxAckPendingJetStream's defaultmax_ack_pending of a queue's durable consumer: the group's cap on unacked messages, shared by every instance. prefetch is each instance's own
stream{}Extra stream config, e.g. { queue: { storage: 'memory' } }

What maps to what ​

CapabilityNATS
backplanecore subjects: PUB/SUB, at-most-once, no persistence
logone JetStream stream per topic; the message sequence is the feed's resume token
queuea work-queue stream with one durable pull consumer per group; prefetch is counted per instance, maxAckPending caps the group
directcore subjects: plain subscriptions for inboxes, queue groups for a service address, native reply subjects

Names become one subject token ​

A room, a topic, a queue and an address are application strings; a NATS subject is a dot-separated hierarchy where * and > are wildcards. So every name is encoded into exactly one token: a room called room:* subscribes to that room and to nothing else, and orders.eu cannot collide with a subject someone else meant. Stream names get the same treatment (. is illegal in one).

Sharp edges ​

  • A JetStream sequence carries no epoch. An id from another stream looks like a future position in this one, so a feed answers 400 (past the tip) where a log with an epoch would say 410. Sign the ids (brokerFeed({ secret })) when clients may send ids from elsewhere.
  • A catch-up page stops at the stream's last sequence, read before the pull, instead of parking for the fetch's expiry once the stream has nothing more — and a page the expiry cut short is asked for again, never taken for the tip.
  • Retention is yours to choose. The adapter creates a log stream with the server's defaults; pass stream.log (max_msgs, max_age, storage) for something else. A feed resuming past a purge answers 410, which onGap turns into a snapshot.
  • release() costs a republish. JetStream's nak always counts a delivery, and a release must not, so the adapter re-publishes the message with its attempt carried in a header and terminates the original. retry() is the cheap path — a plain nak(delay).
  • A slow handler keeps its lease. While a delivery is in flight the adapter calls working() every third of ackWait, and it stops the moment the consumer stops — so a consumer that stops hands its messages back after ackWait rather than holding them forever. A paused consumer (a draining node) keeps the leases of what it still holds until those settle.
  • A durable consumer keeps the configuration it was created with. The group's max_ack_pending and ack_wait are set when the first instance creates the durable; an instance that binds later with other values is served with the old ones and logs broker.nats.consumer.config once — update the consumer (nats consumer edit) to change them.
  • A reader's consumer is deleted when the read is done — a catch-up page's after the page, a live tail's when the tail stops — and reaped by the server 30 seconds after a reader that died without saying so.
  • Core NATS never queues. A direct.send to an address nobody listens on is dropped by the server, as core NATS always does: the caller learns from its own timeout rather than a 503.
  • Payloads are capped (1 MiB by default). A wRPC call or a cluster fetchClients answer above it is refused by the server — raise max_payload or keep the answers small.

Running the tests ​

pnpm test runs the NATS suites against an in-repo fake NATS + JetStream. Against a real server:

bash
docker compose up -d nats
NATS_URL=nats://127.0.0.1:4222 node --test tests/broker/nats.integration.test.js

Released under the MIT License.