Skip to content

Repository files navigation

node-rabbitbus

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.

Supported patterns

  • Command-Reply — point-to-point messaging between named services.
  • Publish/Subscribe — events via durable topic exchanges.
  • RPC — synchronous request/response over RabbitMQ.

Requirements

  • Node.js 18+
  • RabbitMQ 3.8+ or RabbitMQ 4.x

Installation

npm install logiqbits-rabbitbus

The package ships dual ESM (dist/esm) and CJS (dist/cjs) builds with TypeScript declarations, so both import and require work.

Quick start

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" })));

Message contract

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. data in TS matches Data 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);

Pattern examples

Command-Reply

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" })));

Publish/Subscribe

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" })));

RPC

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).

Interoperability with Go and Python

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.

TypeScript message

class Command1 implements Message {
  data = "";
  schemaName(): string {
    return "example.Command1";
  }
}

Equivalent Go message

type Command1 struct {
    Data string `json:"data"`
}

func (Command1) SchemaName() string { return "example.Command1" }

Go bus configured for interop

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.

Configuration

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.

Publishing semantics

  • 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 raises PublishError; a transport failure triggers one publisher reconnect + retry before raising PublishError.
  • 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.

Management

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).

Development

Build (dual ESM/CJS output with declarations):

npm install
npm run build

Run 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:e2e

The 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:hardening

Run 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.py

The tests require a local RabbitMQ at amqp://guest:guest@localhost:5672.

License

Apache License 2.0

About

Node package for RabbitBus

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages