A Node.js/TypeScript service bus for RabbitMQ, wire-compatible with go-rabbitbus and python-rabbitbus.
It implements the same messaging patterns and AMQP conventions as the Go and Python libraries so Node services can exchange commands, events, and RPC calls with existing Go and Python services without changing the transport layer.
- Command-Reply — point-to-point messaging between named services.
- Publish/Subscribe — events via durable topic exchanges.
- RPC — synchronous request/response over RabbitMQ.
- Node.js 18+
- RabbitMQ 3.8+ or RabbitMQ 4.x
npm install logiqbits-rabbitbusThe package ships dual ESM (dist/esm) and CJS (dist/cjs) builds with TypeScript declarations, so both import and require work.
import { builder, BusMessage, Durable, type Invocation, type Message } from "logiqbits-rabbitbus";
class Command1 implements Message {
data = "";
schemaName(): string {
return "example.Command1";
}
}
const bus = builder()
.bus("amqp://guest:guest@localhost")
.withPolicies(new Durable())
.build("node.svc");
bus.handleMessage(Command1, async (invocation: Invocation, message: BusMessage) => {
const cmd = message.payload as Command1;
console.log(`received command: ${cmd.data}`);
});
await bus.start();
// send a command to another service
await bus.send("go.svc", new BusMessage(Object.assign(new Command1(), { data: "hello" })));Every payload must implement the Message interface by providing a schemaName() method. The default JsonSerializer wraps the payload in a small envelope:
{"schema_name":"example.Command1","payload":{"data":"hello"}}The body is compact UTF-8 JSON (no spaces). The schema_name is also sent as the AMQP header x-msg-name; the receiver looks up the message type by that header (not by the body) and instantiates it with Object.assign(new PayloadClass(), payload), so payload classes must have the same property names as the sender's fields.
Wire-format caveat (same as the Go README): payload key naming is your responsibility. TypeScript property names must match the Go struct JSON tags / Python attribute names of the services you interop with, e.g.
datain TS matchesData string \json:"data"`in Go andself.data` in Python.
If your service receives a message type only as a reply or RPC response (not as a direct handler), register it with the serializer:
bus.registerMessageType(Reply1);import { builder, BusMessage, Durable, type Invocation, type Message } from "logiqbits-rabbitbus";
class Command1 implements Message {
data = "";
schemaName(): string {
return "example.Command1";
}
}
class Reply1 implements Message {
data = "";
schemaName(): string {
return "example.Reply1";
}
}
const bus = builder()
.bus("amqp://guest:guest@localhost")
.withPolicies(new Durable())
.build("node.svc");
bus.handleMessage(Command1, async (invocation: Invocation, message: BusMessage) => {
const cmd = message.payload as Command1;
console.log(`received: ${cmd.data}`);
await invocation.reply(new BusMessage(Object.assign(new Reply1(), { data: "ok" })));
});
await bus.start();
// send a command and (optionally) handle the reply elsewhere
await bus.send("another.svc", new BusMessage(Object.assign(new Command1(), { data: "hello" })));class OrderCreated implements Message {
orderId = "";
schemaName(): string {
return "example.OrderCreated";
}
}
const bus = builder()
.bus("amqp://guest:guest@localhost")
.withPolicies(new Durable())
.build("notifications.svc");
bus.handleEvent("events", "order.*", OrderCreated, async (invocation: Invocation, message: BusMessage) => {
const evt = message.payload as OrderCreated;
console.log(`order created: ${evt.orderId}`);
});
await bus.start();
// publish an event
await bus.publish("events", "order.created", new BusMessage(Object.assign(new OrderCreated(), { orderId: "123" })));class RpcRequest implements Message {
data = "";
schemaName(): string {
return "example.RpcRequest";
}
}
class RpcResponse implements Message {
data = "";
schemaName(): string {
return "example.RpcResponse";
}
}
const bus = builder()
.bus("amqp://guest:guest@localhost")
.withPolicies(new Durable())
.build("node.rpc.svc");
bus.handleMessage(RpcRequest, async (invocation: Invocation, message: BusMessage) => {
const req = message.payload as RpcRequest;
await invocation.reply(new BusMessage(Object.assign(new RpcResponse(), { data: `hello ${req.data}` })));
});
bus.registerMessageType(RpcResponse);
await bus.start();
// call a remote service and wait (up to timeoutMs) for a reply
const reply = await bus.rpc("go.rpc.svc", new BusMessage(Object.assign(new RpcRequest(), { data: "world" })), 5000);
console.log((reply.payload as RpcResponse).data);bus.rpc uses a fresh short-lived connection per call, publishes to the target service queue with an exclusive auto-delete reply queue, and throws RpcTimeoutError when no reply arrives within the timeout (default 5000 ms).
The library defaults to JSON serialization and follows the exact wire conventions of the Python port (which is wire-compatible with the Go JSON serializer path): compact {"schema_name":...,"payload":...} bodies, the four x-msg-* headers (always present, empty string when unset), messageId as 32-char hex, contentEncoding: "json", replies correlated via replyTo on the default exchange, and case-insensitive handler matching with */? topic wildcards. The service queue is named exactly after the service name; when deadlettering is configured the queue gets the x-dead-letter-exchange argument and is also bound back to the DLX, matching the Python/Go topology.
class Command1 implements Message {
data = "";
schemaName(): string {
return "example.Command1";
}
}type Command1 struct {
Data string `json:"data"`
}
func (Command1) SchemaName() string { return "example.Command1" }import (
"github.com/logiqbits/go-rabbitbus/gbus"
"github.com/logiqbits/go-rabbitbus/gbus/builder"
"github.com/logiqbits/go-rabbitbus/gbus/policy"
"github.com/logiqbits/go-rabbitbus/gbus/serialization"
)
bus := builder.New().Bus("amqp://guest:guest@localhost").
WithSerializer(serialization.NewJsonSerializer()).
WithPolicies(&policy.Durable{}).
Build("go.svc")After that, bus.Send(...), bus.Publish(...), and bus.RPC(...) from Go will interoperate with this Node implementation, and so will the Python port.
const bus = builder()
.bus("amqp://guest:guest@localhost")
.withPolicies(new Durable())
.workerNum(4, 10)
.purgeOnStartup()
.withDeadlettering("dead-letter-exchange")
.withSerializer(new JsonSerializer())
.heartbeat(30)
.build("node.svc");| Builder method | Description |
|---|---|
bus(url) |
AMQP connection URL. |
withPolicies(...policies) |
Default policies applied to every outgoing message (Durable → persistent, NonDurable → transient). Without a policy the delivery mode is left at the broker default. |
workerNum(workers, prefetchCount) |
Number of consumer worker channels and AMQP prefetch count. |
purgeOnStartup() |
Delete the service queue on startup. |
withDeadlettering(exchange) |
Route rejected/poison messages to the given exchange (fanout, durable); the service queue is also bound back to it for redelivery. |
withSerializer(serializer) |
Use a custom serializer instead of the default JSON serializer. |
heartbeat(seconds) |
Heartbeat interval for every connection the bus opens; null keeps the URL value. |
- Commands are published on the default exchange with the routing key set to the target service name,
mandatory: true, on a confirm-mode channel. A broker nack or an unroutable return raisesPublishError; a transport failure triggers one publisher reconnect + retry before raisingPublishError. - Events are published to the named topic exchange with the topic as routing key.
- Handlers run per delivery: success acks, a handler/decode exception nacks with
requeue=false(→ DLX when configured), and an unmatched message is acked with a warning.
bus.health(), bus.queueDepth(name), bus.purge(name), bus.deadletterCount() (depth of <dlx>.parking), and await bus.shutdown(drain?) are available while the bus is started. The bus keeps two connections (consumer + publisher, heartbeat 30 by default).
Build (dual ESM/CJS output with declarations):
npm install
npm run buildRun the bidirectional Go ↔ Node e2e test against a local RabbitMQ broker:
# terminal 1: start the Go service
cd tests/e2e_go
go run .
# terminal 2 (within ~5 seconds): run the Node service/client
npm run test:e2eThe test exercises Command-Reply, Pub/Sub, and RPC in both directions. Both services use purgeOnStartup, so start the Node side shortly after the Go side (the Node script waits for the Go queue before sending; the Go service waits 8 s before sending).
Run the hardening checks (publisher confirms, unroutable → PublishError, reconnect/retry, queue depth/purge, heartbeat survival, deadletter count):
npm run test:hardeningRun the Python ↔ Node interop check with the unmodified Python e2e service:
# terminal 1: Node service impersonating "go.e2e.service"
npx tsx tests/interop_node4python.ts
# terminal 2: the Python service/client from python-rabbitbus
cd ../python-rabbitbus
.venv/bin/python tests/e2e_python.pyThe tests require a local RabbitMQ at amqp://guest:guest@localhost:5672.
Apache License 2.0