Адаптер для разных брокеров

Причины появления адаптерного слоя

STOMP.js работает поверх WebSocket и реализует клиентскую часть протокола STOMP, однако сам по себе STOMP не гарантирует единообразие поведения брокеров сообщений. Разные серверные реализации STOMP (например, RabbitMQ, ActiveMQ, Apollo, Spring WebSocket STOMP broker, ActiveMQ Artemis) отличаются не только конфигурацией подключения, но и:

  • форматом URL подключения
  • поддержкой heart-beat механизма
  • поведением при переподключении
  • особенностями маршрутизации destination
  • требованиями к префиксам очередей и топиков
  • реализацией ack/nack логики
  • ограничениями на размер сообщений и headers

Без дополнительного слоя абстракции клиентский код оказывается привязан к конкретному брокеру, что усложняет поддержку и миграцию. Адаптер решает задачу изоляции STOMP.js от инфраструктурных различий.


Базовая идея адаптера

Адаптер представляет собой слой, который унифицирует взаимодействие с разными брокерами через единый интерфейс.

Ключевая цель:

STOMP.js → всегда работает одинаково, независимо от брокера

Адаптер выполняет:

  • нормализацию конфигурации подключения
  • преобразование destination-адресов
  • настройку heartbeat
  • управление reconnect-стратегией
  • унификацию подписок и отправки сообщений
  • приведение событий брокера к единому формату

Контракт адаптера

Типовой интерфейс адаптера можно описать следующим образом:

class StompBrokerAdapter {
  connect(config) {}
  disconnect() {}

  subscribe(destination, handler, options) {}
  unsubscribe(subscriptionId) {}

  publish(destination, body, headers) {}

  isConnected() {}
}

Этот интерфейс скрывает специфику STOMP.js и брокера.


Нормализация конфигурации подключения

Разные брокеры требуют разных параметров подключения.

Пример различий:

  • RabbitMQ Web STOMP: ws://host:15674/ws
  • ActiveMQ: ws://host:61614/stomp
  • Spring WebSocket STOMP: ws://host/ws-endpoint

Адаптер вводит единый конфиг:

const config = {
  broker: 'rabbitmq',
  host: 'localhost',
  port: 15674,
  endpoint: '/ws',
  login: 'user',
  passcode: 'pass',
  heartbeatIncoming: 10000,
  heartbeatOutgoing: 10000
};

Внутри адаптера выполняется трансформация:

function buildUrl(config) {
  switch (config.broker) {
    case 'rabbitmq':
      return `ws://${config.host}:${config.port}${config.endpoint}`;
    case 'activemq':
      return `ws://${config.host}:${config.port}/stomp`;
    case 'spring':
      return `ws://${config.host}${config.endpoint}`;
    default:
      throw new Error('Unsupported broker');
  }
}

Инкапсуляция STOMP.js клиента

Адаптер обычно содержит экземпляр STOMP.js клиента и управляет его жизненным циклом.

import { Client } from '@stomp/stompjs';

class StompBrokerAdapter {
  constructor() {
    this.client = null;
    this.subscriptions = new Map();
  }

  connect(config) {
    const url = buildUrl(config);

    this.client = new Client({
      brokerURL: url,
      connectHeaders: {
        login: config.login,
        passcode: config.passcode
      },
      heartbeatIncoming: config.heartbeatIncoming,
      heartbeatOutgoing: config.heartbeatOutgoing,
      reconnectDelay: 5000
    });

    this.client.activate();
  }
}

Унификация destination-логики

Разные брокеры используют разные соглашения:

  • /queue/...
  • /topic/...
  • /exchange/... (RabbitMQ специфично)
  • /app/... (Spring)

Адаптер вводит логический уровень:

const DestinationType = {
  QUEUE: 'queue',
  TOPIC: 'topic',
  EVENT: 'event'
};

И преобразует:

function normalizeDestination(type, name, broker) {
  if (broker === 'rabbitmq') {
    if (type === 'QUEUE') return `/queue/${name}`;
    if (type === 'TOPIC') return `/exchange/amq.topic/${name}`;
  }

  if (broker === 'activemq') {
    if (type === 'QUEUE') return `/queue/${name}`;
    if (type === 'TOPIC') return `/topic/${name}`;
  }

  if (broker === 'spring') {
    if (type === 'QUEUE') return `/user/queue/${name}`;
    if (type === 'TOPIC') return `/topic/${name}`;
  }

  return `/${type.toLowerCase()}/${name}`;
}

Абстракция подписок

STOMP.js возвращает объект subscription, который зависит от реализации. Адаптер скрывает это:

subscribe(destination, handler, options = {}) {
  const subId = crypto.randomUUID();

  const stompSub = this.client.subscribe(destination, (message) => {
    handler({
      body: message.body,
      headers: message.headers,
      ack: () => message.ack(),
      nack: () => message.nack()
    });
  }, options);

  this.subscriptions.set(subId, stompSub);

  return subId;
}

Теперь клиентский код работает только с subId, не зная о STOMP деталях.


Унификация отправки сообщений

Разные брокеры по-разному обрабатывают headers и content-type.

Адаптер приводит всё к единому виду:

publish(destination, body, headers = {}) {
  const normalizedHeaders = {
    'content-type': 'application/json',
    ...headers
  };

  const payload =
    typeof body === 'string'
      ? body
      : JSON.stringify(body);

  this.client.publish({
    destination,
    body: payload,
    headers: normalizedHeaders
  });
}

Управление переподключением

Разные брокеры имеют разную стабильность соединения. Адаптер централизует стратегию reconnect:

this.client = new Client({
  brokerURL: url,
  reconnectDelay: 3000,
  onDisconnect: () => {
    console.log('Disconnected, reconnecting...');
  },
  onStompError: (frame) => {
    console.error('Broker error:', frame.headers['message']);
  }
});

Для более сложных сценариев добавляется внешняя логика:

  • экспоненциальная задержка
  • ограничение числа попыток
  • переключение fallback брокера

Heartbeat унификация

Heartbeat работает по-разному в брокерах: где-то обязателен, где-то игнорируется.

Адаптер стандартизирует:

const heartbeat = {
  incoming: config.heartbeatIncoming ?? 10000,
  outgoing: config.heartbeatOutgoing ?? 10000
};

И применяет политику:

  • если брокер не поддерживает heartbeat → отключается автоматически
  • если поддерживает частично → выравнивается до минимального интервала

Поддержка нескольких брокеров (failover)

Адаптер может поддерживать список брокеров:

const brokers = [
  'ws://broker1:61614/stomp',
  'ws://broker2:61614/stomp'
];

Логика подключения:

async connectWithFailover() {
  for (const url of brokers) {
    try {
      await this.tryConnect(url);
      return;
    } catch (e) {
      continue;
    }
  }

  throw new Error('All brokers are unavailable');
}

Изоляция бизнес-логики от инфраструктуры

Ключевая архитектурная цель адаптера — отделить доменную часть приложения от транспортного уровня.

Без адаптера:

  • компоненты знают о /topic/
  • компоненты знают о RabbitMQ или ActiveMQ
  • изменение брокера ломает клиентский код

С адаптером:

  • компоненты работают с publish("chat.message")
  • подписки выражаются логически: subscribe(EVENT_CHAT)
  • транспорт скрыт полностью

Расширяемость адаптера

Адаптер проектируется как расширяемая система:

class BaseBrokerAdapter {
  transformDestination() {}
  transformMessage() {}
  transformError() {}
}

Пример расширения под RabbitMQ:

class RabbitMQAdapter extends BaseBrokerAdapter {
  transformDestination(type, name) {
    return `/exchange/amq.topic/${name}`;
  }
}

Централизованная обработка ошибок

STOMP ошибки приходят в разных форматах:

  • frame.error
  • broker disconnect
  • subscription failure

Адаптер приводит их к единому виду:

function normalizeError(frame) {
  return {
    message: frame.headers?.message || 'Unknown error',
    code: frame.headers?.code || 'STOMP_ERROR',
    raw: frame
  };
}

Логирование и трассировка

Адаптер часто включает слой observability:

  • логирование connect/disconnect
  • логирование отправки сообщений
  • трассировка подписок
  • измерение latency сообщений
publish(destination, body) {
  const start = performance.now();

  this.client.publish({ destination, body });

  const end = performance.now();
  this.logger.log(`Publish latency: ${end - start}ms`);
}

Итоговая структура адаптера

Типичная архитектура выглядит следующим образом:

  • STOMP.js Client (низкий уровень)
  • Broker Adapter (унификация)
  • Messaging Service (бизнес-логика)
  • UI / Application Layer

Такой подход делает транспортный слой заменяемым компонентом без влияния на остальную систему.