Skip to main content

Референс MorphCluster NATS

Данный модуль расширяет MorphCluster Core поддержкой распределённого взаимодействия через NATS. Он предоставляет сервисы для публикации схем, маршрутизации запросов, событий, каналов данных, мониторинга здоровья и таймеров между узлами кластера.


Обзор

Модуль реализует транспортный слой на базе NATS, позволяя локальным сервисам прозрачно вызывать удалённые сервисы и подписываться на события других узлов. Основные компоненты:

  • NatsConnection – управление подключением, отправка запросов, подписки, сериализация.
  • Сервисы публикации – экспортируют локальные схемы, запросы, события, таймеры, данные каналов в NATS.
  • Сервисы подписки – получают объявления от других узлов и регистрируют их в локальном GlobalServices, делая удалённые сервисы доступными для локального кода.
  • Вспомогательные классы – ServiceNats (устаревший), NatsEvent (устаревший), NatsChannels и т.д.

Все сервисы строятся на базе ServiceRequire и ожидают наличия в хосте NatsConnection и GlobalServices.


NatsConnection

Базовый сервис, устанавливающий соединение с NATS-сервером и предоставляющий методы для обмена сообщениями. Наследуется от ServiceRequire, поэтому сам является сервисом.

import { NatsConnection } from '@morphcluster/nats';
// или
const { NatsConnection } = require('@morphcluster/nats');

Конструктор

new NatsConnection(host, natsConfig)
  • host – экземпляр ServiceHost.
  • natsConfig – объект конфигурации:
    • server – адрес NATS-сервера (например, "localhost:4222").
    • timeout – таймаут запросов в мс (по умолчанию из конфигурации).
    • prefix – префикс для всех subject'ов (по умолчанию "").
    • requestMode – режим запросов: "auto" (автовыбор короткого/длинного запроса), "long" (принудительно длинные запросы), "simple" (короткие, устаревший).
    • format – формат сериализации: "json" или "msgpack".
    • queue – имя группы подписчиков для балансировки (используется в NatsSubscriber).

Свойства

  • connection – объект соединения NATS (Nats.NatsConnection) или null до подключения.
  • serializer – экземпляр сериализатора (SerializerJson или SerializerMsgpack).
  • requestMode – режим запросов.
  • prefix – префикс subject'ов.

Методы

connect()

Устанавливает соединение с сервером. Вызывается автоматически при start().

disconnect()

Корректно закрывает соединение (drain).

request(subject, params)

Отправляет запрос и ожидает ответ. В зависимости от requestMode использует короткий или длинный протокол.

  • subject – строка (автоматически добавляется префикс).
  • params – объект данных.
  • Возвращает Promise с ответом. Если ответ содержит поле error, выбрасывается ComplexError.

subscribe(subject, callback)

Подписывается на сообщения. Возвращает объект подписки (Nats.Subscription), который можно использовать для отмены.

  • callback(msg, subject) – вызывается при получении сообщения. Данные уже десериализованы.

unsubscribe(subscription)

Отменяет подписку.

subscribeEvent(subject, callback)

Подписка на события (см. NatsEventer). Отличие от subscribe в том, что не ожидается ответов и используется отдельная коллекция подписок.

unsubscribeEvent(subject, callback)

Отписывается от события.

publish(subject, params)

Публикует сообщение без ожидания ответа.

Внутренние классы

SerializerJson / SerializerMsgpack

Реализуют кодирование/декодирование данных с помощью JSON или msgpackr. Используются внутри NatsConnection.

NatsRequester

Отвечает за логику запросов. Методы:

  • simpleRequest(subject, rawRequest) – короткий запрос-ответ.
  • longRequest(subject, params) – отправляет запрос, при необходимости создаёт канал для больших данных, получает ответ.
  • autoRequest(subject, params) – автоматически выбирает между коротким и длинным протоколом в зависимости от размера запроса.
  • fetchLongResponse(channel) – собирает ответ по частям через созданный канал.

NatsSubscriber

Управляет подписками на входящие запросы. Использует очередь, если задан queue. Для каждого сообщения с полем reply автоматически формирует ответ (в том числе через канал, если ответ большой).

NatsEventer

Упрощённая подписка на события (без ответов). Агрегирует несколько колбэков на один subject.

RequestChannel

Временный канал для передачи больших запросов/ответов. Получает данные по частям, собирает их, затем передаёт ответ также по частям.

  • Свойства: subject, status (этапы: request, process, responseHeader, response, idle).
  • Методы: subscribe(), waitRequest(), sendResponse(data), unsubscribe().

ServiceNats (УСТАРЕВШИЙ)

import { ServiceNats } from '@morphcluster/nats';

Deprecated. Вместо него рекомендуется использовать глобальные сервисы и CoreClient. ServiceNats реализует прокси для удалённого сервиса через NATS, динамически создавая методы запросов и событий.

Конструктор

new ServiceNats(host, svcSchema, dontWait = false)
  • host – ServiceHost.
  • svcSchema – объект схемы (ServiceSchema).
  • dontWait – если true, не ждать появления схемы в NATS при старте.

Особенности

  • При start() создаёт динамические методы на основе схемы: для каждого запроса reqS.name становится методом, для каждого события – свойством типа NatsEvent.
  • Отправляет запросы через NatsConnection.request.
  • Может ожидать появления схемы от других узлов через waitSchema.

NatsEvent (УСТАРЕВШИЙ)

import { NatsEvent } from '@morphcluster/nats';

Deprecated. Используйте стандартный Event из Core в связке с глобальными сервисами. NatsEvent расширяет Event, автоматически подписываясь на соответствующий subject в NATS при появлении слушателей.

Методы

  • on(callback, workspace) – добавляет слушателя и при необходимости инициирует NATS-подписку.
  • off(callback) – удаляет слушателя и, если больше нет подписчиков, отписывается от NATS.
  • setNats(natsConnection) – задаёт соединение и активирует подписку, если уже есть слушатели.

Сервисы NATS-интеграции

Все перечисленные ниже сервисы должны быть добавлены в ServiceHost и правильно сконфигурированы. Они используют NatsConnection и GlobalServices как зависимости.

NatsSchemaPublisher

Публикует схемы локальных сервисов в NATS, чтобы другие узлы могли их обнаружить.

import { NatsSchemaPublisher } from '@morphcluster/nats';

Зависимости

  • GlobalServices
  • NatsConnection
  • NatsRequestListener

Поведение

  • При старте подписывается на события Services.broadcast.refresh и Services.broadcast.refreshOne – по ним отправляет все схемы или конкретную.
  • Слушает onServicePublished от NatsRequestListener и автоматически публикует новую схему.
  • Если задан config.interval, периодически рассылает все схемы.
  • Публикует в subject Services.broadcast.schema.

Конфигурация

  • interval – интервал в мс для повторной публикации (необязательно).

NatsSchemaListener

Слушает объявления схем от других узлов и регистрирует их в GlobalServices как удалённые сервисы.

import { NatsSchemaListener } from '@morphcluster/nats';

Зависимости

  • GlobalServices
  • NatsConnection

Поведение

  • Подписывается на Services.broadcast.schema.
  • При получении схемы создаёт NatsServiceInterface и вызывает GlobalServices.addRemoteService.
  • При старте отправляет запрос Services.broadcast.refresh, чтобы получить схемы уже работающих узлов.

NatsServiceInterface (внутренний)

Реализует RemoteServiceInterface для NATS.

  • request(requestName, params, log) – отправляет запрос через NatsConnection.request.
  • subscribeEvent(eventName, callback) – подписывается на событие.
  • unsubscribeEvent(eventName, callback) – отписывается.

NatsRequestListener

Принимает входящие запросы из NATS и передаёт их локальным сервисам.

import { NatsRequestListener } from '@morphcluster/nats';

Зависимости

  • GlobalServices
  • NatsConnection

Поведение

  • Отслеживает локальные сервисы (не удалённые и не виртуальные), для каждого запроса в схеме создаёт NATS-подписку вида <serviceName>.<requestName>.
  • При получении запроса создаёт дочерний логгер (если разрешено), вызывает метод сервиса и возвращает ответ.
  • Поддерживает специальный subject для каналов (большие запросы).
  • Публикует событие onServicePublished при добавлении схемы сервиса.

Важно

  • Использует host.onServiceStarted, чтобы подписывать сервисы, запущенные позже.

NatsEventPublisher

Публикует события локальных сервисов в NATS.

import { NatsEventPublisher } from '@morphcluster/nats';

Зависимости

  • NatsConnection

Поведение

  • Для каждого локального сервиса находит события (из schema.events), подписывается на них и при срабатывании публикует данные в NATS с subject <serviceName>.<eventName>.
  • Использует host.onServiceStarted для автоматической обработки новых сервисов.

NatsServicePublishers

Сервис для публикации данных от Publisher (вероятно, имеются в виду ChannelSender) в NATS. Аналогичен NatsChannels, но использует другой формат subject'ов: Channels.<channelName>.

Примечание: В предоставленном коде используется свойство svc.publishers и метод pub.onDataChanged, но в Core таких сущностей нет (есть ChannelSender). Возможно, это устаревший компонент. Рекомендуется использовать NatsChannels.

import { NatsServicePublishers } from '@morphcluster/nats';

NatsChannels

Обеспечивает прозрачную синхронизацию ChannelSender и ChannelAggregator между узлами через NATS.

import { NatsChannels } from '@morphcluster/nats';

Зависимости

  • NatsConnection

Поведение

  • При старте подписывается на host.onServiceStarted и обрабатывает все существующие сервисы.
  • Для каждого ChannelSender подписывается на onDataChanged и публикует данные в Channels.<channelName>, добавляя hostName и serviceName.
  • Для каждого ChannelAggregator подписывается на Channels.<channelName> и при получении данных вызывает aggregator.setData(serviceName, data).
  • Также подписывается на Channels.Request.* для ответа на запросы актуальных данных от других узлов.
  • При запуске отправляет запросы данных для всех локальных агрегаторов.

NatsHealthPublisher

Периодически публикует состояние здоровья локального хоста (какие сервисы запущены, есть ли невыполненные зависимости) в NATS.

import { NatsHealthPublisher } from '@morphcluster/nats';

Зависимости

  • NatsConnection

Конфигурация

  • publishTime – интервал публикации в мс (по умолчанию 30000). Если 0, то автоматическая периодическая публикация отключается (только при старте и изменении состояния).

Публикуемые данные

{
  "hostName": "main",
  "allStarted": false,
  "noRequirements": ["SomeService"],
  "fullServices": [
    {
      "name": "MyService",
      "starting": false,
      "started": true,
      "requirements": [...]
    }
  ]
}

Отправляется в Health.broadcast.status.


NatsHealthSubscriber

Принимает данные о здоровье от других узлов и предоставляет их через запрос list.

import { NatsHealthSubscriber } from '@morphcluster/nats';

Зависимости

  • NatsConnection

Запросы

  • list() – возвращает массив объектов с информацией о хостах, включая поле updated (секунд с последнего обновления). Требует needAdmin: true.

Поведение

  • Подписывается на Health.broadcast.status, сохраняет данные в this.healthData.

NatsTimersPublisher

Периодически публикует состояние всех таймеров локального хоста (запущены, ошибки) в NATS.

import { NatsTimersPublisher } from '@morphcluster/nats';

Зависимости

  • NatsConnection

Конфигурация

  • publishInterval – интервал в мс (по умолчанию 5000).

Данные

Публикует в Timers.broadcast.status объект с полями hostName и timers[], где каждый таймер содержит name, serviceName, enabled, running, lastError, errorCounter.


NatsTimersAggregator

Собирает данные о таймерах со всех узлов и предоставляет их через запрос list.

import { NatsTimersAggregator } from '@morphcluster/nats';

Зависимости

  • NatsConnection

Запросы

  • list() – возвращает агрегированный список таймеров с отметкой времени updated.

Тестовый сервис NatsConnectionTest

import { NatsConnectionTest } from '@morphcluster/nats';

Предназначен для проверки корректности работы NatsConnection, включая передачу больших сообщений.

Запросы (все требуют needAdmin: true)

  • testShort – короткий запрос-ответ.
  • testLargeRequest – большой запрос, короткий ответ.
  • testLargeResponse – короткий запрос, большой ответ.
  • testLarge – оба больших.
  • reqForTest – внутренний, используется для эхо-тестов.

Конструктор

new NatsConnectionTest(host, config)
  • config.largeSize – размер больших данных (по умолчанию 10 МБ).

Интеграция в хост

Типичная конфигурация хоста с NATS:

import { ServiceHost, Config, Logger } from '@morphcluster/core';
import {
  NatsConnection,
  NatsSchemaPublisher,
  NatsSchemaListener,
  NatsRequestListener,
  NatsChannels,
  NatsHealthPublisher,
  NatsHealthSubscriber,
  // ... другие
} from '@morphcluster/nats';

const host = new ServiceHost('main');
const config = new Config({ /* ... */ });

// Добавляем базовые сервисы Core (GlobalServices, Logger и т.д.)
// ...

// NATS-соединение
host.addService(new NatsConnection(host, { server: 'localhost:4222', format: 'msgpack' }));

// Обнаружение сервисов
host.addService(new NatsSchemaPublisher(host));
host.addService(new NatsSchemaListener(host));
host.addService(new NatsRequestListener(host));

// Каналы данных
host.addService(new NatsChannels(host));

// Мониторинг
host.addService(new NatsHealthPublisher(host, { publishTime: 10000 }));
host.addService(new NatsHealthSubscriber(host));

await host.start(logger);

После запуска локальные сервисы становятся доступны удалённым узлам, а удалённые сервисы – локальным через GlobalServices.