---
url: https://node-ts.github.io/bus/transports/custom.md
description: >-
  Adapt another message broker by implementing the Transport interface, and
  check it with the @node-ts/bus-test conformance suite.
---

# Custom transports

To use a message broker that doesn't have a transport yet, write an adapter that implements the `Transport` interface from `@node-ts/bus-core`, and run the `@node-ts/bus-test` conformance suite against it. This page walks through both.

## Implementing Transport

A transport sends commands, publishes events, and reads, deletes, returns and fails messages on the service queue. These methods are optional, and called during the bus' lifecycle:

* `connect(options)` at `initialize()`, to connect to the broker. `options.concurrency` is the number of messages the bus handles at once.
* `initialize({ handlerRegistry, sendOnly })` next, to create the queues and subscribe the service queue to every message the bus handles.
* `start()` and `stop()`, when the bus starts and stops reading.
* `disconnect()` and `dispose()`, at `dispose()`.

`prepare(coreDependencies)` is called by `build()`. Keep the dependencies it's given: the bus' `messageSerializer`, to write and read message bodies so their Dates and classes are restored, its `loggerFactory` and its `retryStrategy`.

This skeleton adapts an imaginary broker client:

```ts
import {
  CoreDependencies,
  DEFAULT_DEAD_LETTER_QUEUE_NAME,
  Logger,
  Transport,
  TransportConfiguration,
  TransportInitializationOptions,
  TransportMessage
} from '@node-ts/bus-core'
import { Command, Event, MessageAttributes } from '@node-ts/bus-messages'
import { BrokerClient, BrokerMessage } from './broker-client'

const MAX_ATTEMPTS = 10

export interface MyTransportConfiguration extends TransportConfiguration {
  connectionString: string
}

export class MyTransport implements Transport<BrokerMessage> {
  private coreDependencies: CoreDependencies
  private logger: Logger

  constructor(
    private readonly configuration: MyTransportConfiguration,
    private readonly client: BrokerClient
  ) {}

  // Called by Bus.configure().build() with the bus' serializer, logger and retry strategy
  prepare(coreDependencies: CoreDependencies): void {
    this.coreDependencies = coreDependencies
    this.logger = coreDependencies.loggerFactory('my-org:my-transport')
  }

  async connect(): Promise<void> {
    await this.client.connect()
  }

  async disconnect(): Promise<void> {
    await this.client.close()
  }

  // Create the service queue, and subscribe it to every message the bus handles
  async initialize({
    handlerRegistry,
    sendOnly
  }: TransportInitializationOptions): Promise<void> {
    if (sendOnly) {
      return
    }
    const { queueName } = this.configuration
    await this.client.createQueue(queueName)
    await this.client.createQueue(this.deadLetterQueueName)
    const topics = [
      ...handlerRegistry.getMessageNames(),
      // Topics of the messages handled by withCustomHandler
      ...handlerRegistry.getExternallyManagedTopicIdentifiers()
    ]
    for (const topic of topics) {
      await this.client.subscribe(queueName, topic)
    }
  }

  async publish<TEvent extends Event>(
    event: TEvent,
    attributes?: MessageAttributes
  ): Promise<void> {
    await this.dispatch(event, attributes)
  }

  async send<TCommand extends Command>(
    command: TCommand,
    attributes?: MessageAttributes
  ): Promise<void> {
    await this.dispatch(command, attributes)
  }

  async readNextMessage(): Promise<
    TransportMessage<BrokerMessage> | undefined
  > {
    const raw = await this.client.receive(this.configuration.queueName)
    if (!raw) {
      return undefined
    }
    return {
      id: raw.id,
      raw,
      // Restores the message's Dates and classes from the bus' message types
      domainMessage: this.coreDependencies.messageSerializer.deserialize(
        raw.body
      ),
      attributes: {
        correlationId: raw.headers.correlationId,
        attributes: JSON.parse(raw.headers.attributes ?? '{}'),
        stickyAttributes: JSON.parse(raw.headers.stickyAttributes ?? '{}')
      }
    }
  }

  async deleteMessage(message: TransportMessage<BrokerMessage>): Promise<void> {
    await this.client.ack(this.configuration.queueName, message.raw.id)
  }

  async returnMessage(message: TransportMessage<unknown>): Promise<void> {
    const raw = message.raw as BrokerMessage
    if (raw.deliveryCount >= MAX_ATTEMPTS) {
      this.logger.warn('Message is out of attempts', { message })
      await this.fail(message)
      await this.deleteMessage(message as TransportMessage<BrokerMessage>)
      return
    }
    const delay = this.coreDependencies.retryStrategy.calculateRetryDelay(
      raw.deliveryCount - 1
    )
    await this.client.retry(this.configuration.queueName, raw.id, delay)
  }

  async fail(message: TransportMessage<unknown>): Promise<void> {
    await this.client.moveTo(
      this.deadLetterQueueName,
      message.raw as BrokerMessage
    )
  }

  private get deadLetterQueueName(): string {
    return (
      this.configuration.deadLetterQueueName ?? DEFAULT_DEAD_LETTER_QUEUE_NAME
    )
  }

  private async dispatch(
    message: Command | Event,
    attributes: MessageAttributes = { attributes: {}, stickyAttributes: {} }
  ): Promise<void> {
    await this.client.publish(
      message.$name,
      this.coreDependencies.messageSerializer.serialize(message),
      {
        ...(attributes.correlationId && {
          correlationId: attributes.correlationId
        }),
        attributes: JSON.stringify(attributes.attributes),
        stickyAttributes: JSON.stringify(attributes.stickyAttributes)
      }
    )
  }
}

```

A few rules the bus relies on:

* `readNextMessage()` returns `undefined` when there's nothing to read, rather than throwing.
* `returnMessage()` makes the message available again after the retry strategy's delay, and moves it to the dead letter queue once it's out of attempts. The conformance suite expects at least 10 attempts.
* `fail()` moves a message straight to the dead letter queue. The bus deletes it from the service queue afterwards with `deleteMessage()`.
* Each transport instance is one queue and one connection. `build()` throws `TransportAlreadyInUse` if two buses are given the same instance.

Pass the transport to the bus configuration:

```ts
const bus = Bus.configure()
  .withMessageTypes(messageTypes)
  .withTransport(
    new MyTransport(
      {
        queueName: 'reservations-service',
        connectionString: 'broker://localhost'
      },
      brokerClient
    )
  )
  .build()
```

## Testing with the conformance suite

`@node-ts/bus-test` runs the same tests against every transport, to check that it sends, publishes, retries and dead-letters messages the way the bus expects, and that messages keep their types, attributes and sticky attributes on a round trip.

::: code-group

```sh [npm]
npm i -D @node-ts/bus-test jest
```

```sh [pnpm]
pnpm add -D @node-ts/bus-test jest
```

```sh [yarn]
yarn add -D @node-ts/bus-test jest
```

:::

Call `transportTests()` inside a `describe()` in your transport's integration test, with:

* the transport to test, fully configured.
* `publishSystemMessage`, which publishes a raw `TestSystemMessage` to the system message topic, with a `systemMessage` attribute set to the value it's given.
* the topic identifier of the system message, which the suite subscribes to with `withCustomHandler`.
* `readAllFromDeadLetterQueue`, which reads, deletes and returns every message on the dead letter queue.

```ts
import { TestSystemMessage, transportTests } from '@node-ts/bus-test'
import { brokerClient } from './broker-client'
import { MyTransport } from './my-transport'

jest.setTimeout(30_000)

describe('MyTransport', () => {
  const transport = new MyTransport(
    {
      queueName: 'bus-test',
      deadLetterQueueName: 'bus-test-dead-letter',
      connectionString: 'broker://localhost'
    },
    brokerClient
  )

  // The suite handles TestSystemMessage with withCustomHandler, subscribed to this topic
  const systemMessageTopic = 'bus-test-system'

  const publishSystemMessage = async (systemMessage: string) =>
    brokerClient.publish(
      systemMessageTopic,
      JSON.stringify(new TestSystemMessage()),
      { attributes: JSON.stringify({ systemMessage }) }
    )

  // Reads and removes every message on the dead letter queue
  const readAllFromDeadLetterQueue = async () => {
    const messages = await brokerClient.readAll('bus-test-dead-letter')
    return messages.map(raw => ({
      message: JSON.parse(raw.body),
      attributes: {
        correlationId: raw.headers.correlationId,
        attributes: JSON.parse(raw.headers.attributes ?? '{}'),
        stickyAttributes: JSON.parse(raw.headers.stickyAttributes ?? '{}')
      }
    }))
  }

  transportTests(
    transport,
    publishSystemMessage,
    systemMessageTopic,
    readAllFromDeadLetterQueue
  )
})
```

The suite builds its own bus around the transport and disposes it when it's done. Create the broker resources it needs before it runs, and remove them afterwards. It uses jest's globals, so run it with jest 29 or later, or a runner with jest-compatible globals.

`messageRoundTripTests(transport)` runs only the round trip tests, for a transport that can't run the full suite. For complete examples, see the [RabbitMQ](https://github.com/node-ts/bus/blob/master/packages/bus-rabbitmq/src/rabbitmq-transport.integration.ts) and [SQS](https://github.com/node-ts/bus/blob/master/packages/bus-sqs/src/sqs-transport.integration.ts) transports' tests.

::: tip Contributing a transport
To contribute your transport to **@node-ts/bus**, add it as `packages/bus-<broker>` in [the repository](https://github.com/node-ts/bus) and open a pull request. [CONTRIBUTING.md](https://github.com/node-ts/bus/blob/master/CONTRIBUTING.md) covers the conventions it follows.
:::

## See also

* [Custom persistence](/persistence/custom)
* [`Transport`](/api/bus-core/interfaces/Transport) and [`transportTests`](/api/bus-test/functions/transportTests) in the API reference
