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.
No Comments