Angular и сервисы

Angular-сервисы представляют собой основной механизм инкапсуляции бизнес-логики и взаимодействия с внешними источниками данных, включая WebSocket-соединения. При работе с STOMP.js в Angular сервис становится центральным слоем, управляющим жизненным циклом соединения, подписками на каналы и распределением сообщений между компонентами через реактивные потоки RxJS.

STOMP-протокол поверх WebSocket используется для обмена сообщениями в режиме publish/subscribe. STOMP.js реализует клиентскую часть протокола, обеспечивая подключение к брокерам сообщений (RabbitMQ, ActiveMQ, Spring WebSocket и другие). В Angular это обычно оформляется в виде singleton-сервиса, предоставляемого на уровне root, что гарантирует единственное соединение на всё приложение.


Сервис выступает посредником между WebSocket-слоем и UI-компонентами. Его задачи включают:

  • создание и управление STOMP-клиентом
  • установление соединения через WebSocket или SockJS
  • подписка на топики (destinations)
  • обработка входящих сообщений
  • трансляция данных в RxJS-потоки
  • восстановление соединения при сбоях
  • централизованная обработка ошибок

Такой подход исключает прямую работу компонентов с низкоуровневым API STOMP.js, снижая связность и упрощая тестирование.


Базовая структура сервиса

Инкапсуляция STOMP-клиента обычно строится вокруг класса сервиса, использующего BehaviorSubject для хранения состояния соединения и потоков сообщений.

import { Injectable } from '@angular/core';
import { BehaviorSubject, Observable } from 'rxjs';
import { Client, IMessage, StompSubscription } from '@stomp/stompjs';

@Injectable({ providedIn: 'root' })
export class StompService {

  private client: Client | null = null;

  private connectionState$ = new BehaviorSubject<boolean>(false);

  get connectionStatus(): Observable<boolean> {
    return this.connectionState$.asObservable();
  }

}

Состояние соединения отделяется от логики доставки сообщений, что позволяет UI реагировать на подключение независимо от бизнес-данных.


Инициализация STOMP-клиента

STOMP.js использует объект Client, который настраивается перед активацией соединения.

import { Client } from '@stomp/stompjs';
import * as SockJS from 'sockjs-client';

private createClient(): Client {
  return new Client({
    webSocketFactory: () => new SockJS('https://example.com/ws'),
    reconnectDelay: 5000,
    heartbeatIncoming: 4000,
    heartbeatOutgoing: 4000,
    debug: () => {}
  });
}

Ключевые параметры:

  • webSocketFactory — обеспечивает совместимость с SockJS
  • reconnectDelay — автоматическое восстановление соединения
  • heartbeatIncoming / heartbeatOutgoing — контроль “живости” канала
  • debug — отладочный вывод

Управление жизненным циклом соединения

Сервис управляет активацией и деактивацией клиента, синхронизируя состояние с RxJS.

connect(): void {
  if (this.client?.active) {
    return;
  }

  this.client = this.createClient();

  this.client.onConn ect = () => {
    this.connectionState$.next(true);
  };

  this.client.onDisconn ect = () => {
    this.connectionState$.next(false);
  };

  this.client.onStompEr ror = () => {
    this.connectionState$.next(false);
  };

  this.client.activate();
}

Данная схема обеспечивает централизованную реакцию на события транспортного уровня.


Подписка на каналы и управление подписками

STOMP использует модель destinations (например, /topic/chat, /queue/events). Каждая подписка возвращает объект StompSubscription, который необходимо хранить для последующей отписки.

private subscriptions: Map<string, StompSubscription> = new Map();
subscribe<T>(destination: string): Observable<T> {
  return new Observable<T>(observer => {

    if (!this.client?.connected) {
      observer.error('STOMP not connected');
      return;
    }

    const sub = this.client.subscribe(destination, (message: IMessage) => {
      const body = JSON.parse(message.body);
      observer.next(body as T);
    });

    this.subscriptions.set(destination, sub);

    return () => {
      sub.unsubscribe();
      this.subscriptions.delete(destination);
    };
  });
}

Такой подход связывает STOMP-подписку с жизненным циклом Observable, автоматически освобождая ресурсы при отписке компонентов.


Отправка сообщений

Отправка данных в брокер осуществляется через метод publish или send (в зависимости от версии STOMP.js).

send(destination: string, payload: unknown): void {
  if (!this.client?.connected) {
    return;
  }

  this.client.publish({
    destination,
    body: JSON.stringify(payload)
  });
}

При необходимости можно добавлять заголовки:

this.client.publish({
  destination,
  body: JSON.stringify(payload),
  headers: {
    priority: 'high',
    authorization: 'Bearer token'
  }
});

Реактивная модель данных через RxJS

Для распространения сообщений между компонентами используется Subject или BehaviorSubject.

private messages$ = new BehaviorSubject<any[]>([]);
listenToTopic(): void {
  this.client?.subscribe('/topic/updates', (msg) => {
    const data = JSON.parse(msg.body);

    const current = this.messages$.value;
    this.messages$.next([...current, data]);
  });
}

Такой подход позволяет реализовать глобальное состояние без дополнительного state-management слоя.


Разделение потоков по типам данных

В сложных приложениях используется разделение потоков:

  • системные события
  • пользовательские сообщения
  • уведомления
  • ошибки
private userMessages$ = new Subject<any>();
private systemEvents$ = new Subject<any>();
private errorEvents$ = new Subject<any>();

Роутинг сообщений выполняется на уровне сервиса:

private routeMessage(destination: string, payload: any): void {
  if (destination.includes('/topic/system')) {
    this.systemEvents$.next(payload);
  } else if (destination.includes('/topic/user')) {
    this.userMessages$.next(payload);
  }
}

Обработка переподключений

STOMP.js поддерживает автоматическое переподключение, однако состояние подписок требует восстановления вручную.

private restoreSubscriptions(): void {
  this.subscriptions.forEach((_, destination) => {
    this.subscribe(destination).subscribe();
  });
}

Вызывается после успешного reconnection:

this.client.onConn ect = () => {
  this.connectionState$.next(true);
  this.restoreSubscriptions();
};

Интеграция с Angular Zone

WebSocket-события происходят вне Angular Zone, что может приводить к отсутствию обновления UI.

constructor(private ngZone: NgZone) {}
this.client.subscribe('/topic/data', (msg) => {
  this.ngZone.run(() => {
    const data = JSON.parse(msg.body);
    this.messages$.next(data);
  });
});

Это гарантирует корректный запуск change detection.


Авторизация и передача токенов

При использовании защищённых брокеров сообщений токен передаётся через headers CONNECT:

this.client = new Client({
  webSocketFactory: () => new SockJS('/ws'),
  connectHeaders: {
    Authorization: `Bearer ${token}`
  }
});

При обновлении токена требуется пересоздание соединения, так как STOMP не всегда поддерживает динамическое обновление headers.


Типизация сообщений

Для повышения надёжности используется строгая типизация payload:

interface ChatMessage {
  id: string;
  text: string;
  author: string;
  timestamp: number;
}
subscribe<ChatMessage>('/topic/chat');

Это снижает количество runtime-ошибок при обработке JSON.


Тестирование сервиса

При тестировании STOMP-сервиса используется мокирование Client:

const mockClient = {
  activate: jasmine.createSpy(),
  deactivate: jasmine.createSpy(),
  subscribe: jasmine.createSpy().and.returnValue({
    unsubscribe: jasmine.createSpy()
  })
};

Сервис проверяется как изолированная единица без реального WebSocket.


Управление памятью и утечки

Критическим аспектом является корректное освобождение ресурсов:

  • отписка от STOMPSubscription
  • завершение Observable
  • деактивация клиента при destroy приложения
disconnect(): void {
  this.subscriptions.forEach(sub => sub.unsubscribe());
  this.subscriptions.clear();

  this.client?.deactivate();
  this.connectionState$.next(false);
}

Масштабирование архитектуры

При росте приложения STOMP-сервис часто расширяется до нескольких уровней:

  • базовый WebSocketService
  • STOMPService (протокольный слой)
  • DomainMessagingService (бизнес-логика)

Такое разделение позволяет:

  • изолировать транспортный слой
  • переиспользовать сервисы
  • тестировать бизнес-логику без сети

Работа с несколькими соединениями

В некоторых системах требуется несколько брокеров или пространств:

private clients: Map<string, Client> = new Map();

Каждый клиент управляется отдельно, а маршрутизация сообщений выполняется по ключу соединения.


Потоковая обработка и backpressure

При высокой частоте сообщений используется буферизация:

private buffered$ = new Subject<any>();

this.client.subscribe('/topic/high-load', msg => {
  this.buffered$.next(JSON.parse(msg.body));
});

Далее применяются RxJS операторы bufferTime, throttleTime, auditTime для контроля нагрузки UI.


Итоговая модель поведения сервиса

STOMP-сервис в Angular формирует слой:

  • между WebSocket-транспортом и UI
  • между брокером сообщений и RxJS-архитектурой
  • между событиями реального времени и состоянием приложения

Он объединяет асинхронный поток сообщений с реактивной моделью Angular, обеспечивая предсказуемую обработку событий, управляемые подписки и централизованный контроль соединения.