Skip to content

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.

sh
npm i -D @node-ts/bus-test jest
sh
pnpm add -D @node-ts/bus-test jest
sh
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 and SQS transports' tests.

Contributing a transport

To contribute your transport to @node-ts/bus, add it as packages/bus-<broker> in the repository and open a pull request. CONTRIBUTING.md covers the conventions it follows.

See also ​

Released under the MIT License.