Skip to main content

MorphCluster Core

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

Основные концепции

Фреймворк строится вокруг следующих сущностей:

  • Service – минимальная единица бизнес-логики. Сервис описывает свою схему (ServiceSchema), может принимать запросы и генерировать события.
  • ServiceHost – контейнер, управляющий жизненным циклом группы сервисов (запуск, остановка, перезапуск).
  • GlobalService – прокси-объект, который объединяет локальные и удалённые реализации одного и того же сервиса, скрывая разницу между ними.
  • GlobalServices – реестр всех глобальных сервисов в системе, отвечает за их связывание.
  • ServiceClient – декларативное описание потребителя сервиса. Позволяет сервису подключаться к другому сервису по имени.
  • Logger – иерархическая система логирования с поддержкой разных бэкендов (консоль, воркер, файл).
  • Event / Trigger / Timer – механизмы для реактивного взаимодействия, планирования задач и обработки событий.
  • ChannelSender / ChannelAggregator – средства публикации и сбора данных между сервисами без прямых вызовов.

Ядро (Core)

Config

Config управляет конфигурацией приложения.

const config = new Config(defaultConfig);
await config.load();
  • defaultConfig – объект с параметрами по умолчанию.
  • load() загружает глобальный конфиг из файла (MCL_CONFIG или global-config.json), объединяет через utils.overlay и создаёт директорию данных.
  • Свойства:
    • values – итоговая конфигурация.
    • packageRoot – корень пакета.
    • dataDir – директория для данных.

Logger и бэкенды

Logger обеспечивает иерархическую запись логов с возможностью создания дочерних "писателей" для отслеживания цепочек вызовов.

const logger = new Logger({ host: 'main', service: 'MyService' });
logger.backends = [new LoggerBackendConsole()];

const logId = logger.write('Старт', { someData: 1 });
const subLogger = logger.createWriter('Вложенная операция', payload);
// ...
subLogger.close();

Методы:

  • write(message, payload, options) – записать простое сообщение.
  • writeRaw(message, payload, options) – записать "сырое" сообщение (не закрытое).
  • createWriter(message, payload, options) – создать дочерний Logger, связанный с текущей записью.
  • attachWriter(parentId, options) – прикрепиться к существующему родительскому логу.
  • writeException(error, options) – записать ошибку и закрыть логгер.
  • writeExceptionOnly(error, options) – записать ошибку без закрытия.
  • wrap(callback) – выполнить функцию, перехватить исключения и закрыть логгер.
  • createWrapped(message, callback, payload, options) – выполнить функцию внутри нового логгера.

Бэкенды логов:

  • LoggerBackend – абстрактный класс.
  • LoggerBackendConsole – вывод в консоль (info / error в зависимости от уровня).
  • LoggerBackendWorker – отправка логов в родительский поток через parentPort.

ComplexError

Расширенный класс ошибки с дополнительной информацией.

throw new ComplexError('Сообщение', 'КодОшибки', { detail: '...' }, { showUser: true, httpStatus: 422 });

Поля: name, message, payload, options (может содержать showUser, httpStatus, logId).

Event

Паттерн "наблюдатель" с поддержкой вложенных событий и подсчёта подписчиков.

const event = new Event();

event.on((workspace, ...args) => { ... });
event.off(callback);
event.emit('workspace', data);

// Вложенные события (subEvents)
const sub = new Event();
event.addSubEvent(sub);   // подписчики sub учитываются в event.isSubscribed()
event.removeSubEvent(sub);

Свойства:

  • onSubscribe / onUnsubscribe – колбэки при появлении/исчезновении подписчиков.
  • isSubscribed() – есть ли активные слушатели (свои или subEvents).

Timer

Периодическое выполнение задач с обработкой ошибок и логированием.

const timer = new Timer('MyTimer', 5000, async (log) => {
  // логика
});
timer.logger = serviceLogger;
timer.serviceName = 'MyService';
timer.start();
timer.stop();
  • При возникновении ошибок вызывается onTickError.
  • ignoreErrorCount – число допустимых ошибок подряд; при превышении таймер останавливается.
  • maxRunningTime – максимальное время выполнения одного тика (для детекта зависаний).

Trigger

Связывает события с выполнением callback-функции.

const trigger = new Trigger('Обработка события', async (log, ...args) => { ... });
trigger.logger = serviceLogger;
trigger.serviceName = 'MyService';
trigger.connect(someEvent);   // подписаться на событие
trigger.disconnect();         // отписаться от всех

При срабатывании события автоматически создаётся лог-писатель.


Сервисная архитектура

ServiceSchema

Описывает контракт сервиса: запросы и события.

const schema = new ServiceSchema();
schema.name = 'MyService';
schema.addRequest({ name: 'doSomething', ... });
schema.addEvent({ name: 'onChanged', ... });

// Сериализация
const json = schema.getJson();
schema.setJson(json);

RequestSchema – описание запроса, включает name, request (JSON-схема параметров), response, флаги anonymous, needAdmin, noLogs и привязку http. EventSchema – описание события, содержит name, description, structure.

ServiceClient

Клиент для подключения к глобальному сервису. Используется сервисами для декларативного указания зависимостей.

class MyService extends Service {
  constructor() {
    super();
    this.clients.push(new ServiceClient({ name: 'OtherService', requests: [...], events: [...] }));
  }
}
  • connected – флаг подключения.
  • waitConnect() – ожидание подключения к сервису.
  • requestHandlers и eventHandlers – заполняются при подключении к реальному сервису.

Service

Базовый класс для любого сервиса.

class MyService extends Service {
  constructor() {
    super();
    this.name = 'MyService';
    this.addRequest({ name: 'ping', ... });
  }

  async ping(params, workspace, log) {
    return { pong: true };
  }

  async start(log) {
    await super.start(log);
    // инициализация, запуск таймеров
  }
  async stop() { ... }
}

Основные свойства и методы:

  • schema (ServiceSchema)
  • name, local (приватный), remote (виртуальный)
  • clients, timers, triggers, senders, aggregators
  • start(log), stop()
  • onStarted, onFailed, onRestored – события жизненного цикла.

ServiceHost

Управляет набором сервисов.

const host = new ServiceHost('main');
host.Config = config;
host.Logger = logger;

host.addService(new MyService());
await host.start(log);
  • addService(service, name?) – добавляет сервис, создаёт для него Logger и Config.
  • startService(service), stopService(service), restartService(service) – управление.
  • requireService(name, workspace?) – получить запущенный сервис по имени.
  • waitService(name) – ожидать появления сервиса.
  • services – массив всех сервисов.
  • События: onServiceAdded, onServiceStarted.

ServiceRequire

Сервис, который ожидает доступности других сервисов перед собственным запуском.

class MyConsumer extends ServiceRequire {
  constructor(host) {
    super(host);
    this.requirements = ['DataService'];      // обязательные зависимости
    this.optionalServices = ['LoggerService']; // необязательные
  }

  async start(log) {
    await super.start(log);
    // this.DataService уже доступен
    const result = await this.DataService.query(...);
  }
}

После выполнения super.start() в свойствах сервиса появятся ссылки на требуемые сервисы (имена из requirements). Необязательные сервисы подтягиваются асинхронно.


Глобальные сервисы и взаимодействие

GlobalServiceInterface

Абстрактный интерфейс для удалённого взаимодействия с сервисом.

class RemoteInterface extends GlobalServiceInterface {
  async request(name, params, parentLog) { ... }
  async subscribeEvent(name, callback) { ... }
  async unsubscribeEvent(name, callback) { ... }
}

GlobalService

Прокси-объект, скрывающий разницу между локальной и удалённой реализацией сервиса.

const gsvc = new GlobalService(schema);
gsvc.localService = localInstance;   // связать с локальным сервисом
// или
gsvc.interface = remoteInterface;    // связать с удалённым интерфейсом

// Отправка запроса
const result = await gsvc.sendRequest('method', params, log);

// Подписка на событие
gsvc.getEvent('onChanged').on(callback);

// Подключение клиента
gsvc.connectClient(client);

GlobalService управляет маршрутизацией запросов и событий через requestHandlers и eventHandlers. При изменении источника (локальный/удалённый) он корректно переключает связи.

GlobalServices

Центральный реестр, объединяющий все глобальные сервисы в системе.

class MyGlobalServices extends GlobalServices {
  async start(log) {
    await super.start(log);
    // Автоматически подписывается на host.onServiceStarted для добавления локальных сервисов
  }
}

Методы:

  • addLocalService(service) – зарегистрировать локальный сервис, создать для него GlobalService.
  • addRemoteService(schema, iface, hostName) – добавить удалённый сервис.
  • getServiceBySchema(schema) – найти или создать GlobalService по схеме.
  • Событие onServiceChanged – оповещает об изменениях.

ServiceFactory

Упрощает создание сервисов из конфигурации.

const factory = new ServiceFactory();
factory.addServiceClass(MyService);
factory.addMicroservice(pack); // объект со свойством services
factory.run(host, config);

config.services – массив объектов вида { name: '...', className: '...', config: {...} }.

ServiceBridge

Сервис, предоставляющий HTTP-доступ к любому зарегистрированному сервису.

const bridge = new ServiceBridge(host, config);
bridge.GlobalServices = globalServices;

Запросы:

  • list – список всех сервисов с их схемами.
  • getSchema – схема конкретного сервиса.
  • request – выполнить произвольный запрос: { service, request, payload }. Автоматически логируется, ошибки оборачиваются в ComplexError.

Встроенные сервисы

LocalHealth

Предоставляет информацию о здоровье хоста.

const health = new LocalHealth(host, config);

Запрос list возвращает объект:

{
  "result": [{
    "hostName": "main",
    "allStarted": true,
    "noRequirements": [...],
    "fullServices": [...]
  }]
}

Sessions

Хранилище сессий в памяти (неперсистентное).

const sessions = new MemorySessions(host, config);

Основные запросы:

  • create({ userId, login, isAdmin }){ token }
  • delete({ token })
  • get({ token }) → объект сессии (или ошибка InvalidToken)
  • list, deleteMany – работа с фильтрами.
  • validateHttp(request, reply, log, needAdmin) – извлечение сессии из заголовка Authorization.

Auth

Базовая аутентификация администратора.

const auth = new Auth(host, { admin: { login: 'admin', password: '...' } });

Запросы:

  • admin({ login, password }) – возвращает токен сессии.
  • logout – удаляет текущую сессию.

Требует Sessions в зависимостях.


Каналы данных

ChannelSender

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

class Producer extends Service {
  constructor() {
    super();
    this.senders.push(new ChannelSender('updates'));
  }

  async onNewData(data) {
    this.senders[0].data = data;  // автоматически вызовет onDataChanged
  }
}

ChannelAggregator

Собирает данные от нескольких сервисов в один канал.

class Consumer extends Service {
  constructor() {
    super();
    const agg = new ChannelAggregator('updates');
    this.aggregators.push(agg);
    agg.onDataChanged.on((ws, allData) => {
      // allData: [{ serviceName, data }, ...]
    });
  }
}

Для связи Sender и Aggregator используется GlobalService, который при обнаружении ChannelSender начинает передавать данные в соответствующий ChannelAggregator (детали реализации не входят в данную документацию, но являются частью внутренней логики).


Утилиты

Кэширование

MemoryCache

In-memory кэш с TTL, LRU-подобным вытеснением и событиями.

const cache = new MemoryCache({ stdTTL: 60, checkperiod: 30, maxKeys: 1000 });
cache.set('key', value, ttl);
const val = cache.get('key');
cache.del('key');
cache.flushAll();

MemoryCacheAsync

Расширяет MemoryCache для работы с асинхронной загрузкой данных, автоматически объединяя одновременные запросы одного ключа.

const cache = new MemoryCacheAsync({
  stdTTL: 120,
  callback: async (id, params) => {
    return await fetchFromDb(id);
  }
});

const data = await cache.fetch('someId', optionalParams);

AsyncCache (устаревший)

Аналогичный механизм без наследования от MemoryCache. Рекомендуется использовать MemoryCacheAsync.

AsyncCacheList

Кэш для списка элементов с индивидуальной асинхронной загрузкой.

const listCache = new AsyncCacheList();
// необходимо переопределить load(id)
const item = await listCache.get('id');

PromiseSingleton

Гарантирует однократное создание экземпляра.

const singleton = new PromiseSingleton(async (params) => await createResource(params));
const res = await singleton.getInstance(params);
singleton.reset(); // сброс для повторного создания

Работа с данными

JsonSchema

Утилиты для работы с JSON-схемами:

  • overlay(target, overlay) – рекурсивное наложение схем.
  • generateJson(schema) – создаёт пустой объект по схеме.
  • sanitize(schema, data) – очищает данные согласно схеме, подставляя значения по умолчанию.

JsonDirList

Хранение списка объектов в виде JSON-файлов в директории.

const list = new JsonDirList('./data/myList');
await list.reload();
await list.set({ id: '1', ... });
await list.remove('1');
list.changedEvent = () => { /* обновить UI или кэш */ };

ExpiringList

Список с автоматическим удалением устаревших элементов по TTL.

const list = new ExpiringList(3600, 60); // ttl=1 час, очистка раз в минуту
list.push({ id: 'token1', ... });
const item = list.get('token1');
list.remove('token1');

Дата и время

Функции для работы с dayjs:

  • removeTimezone(dateWithTz) – приводит к МСК и возвращает строку YYYY-MM-DDTHH:mm:ss.
  • strToDayjs(str) – парсит дату с поддержкой нескольких форматов.
  • dayjsToStr(djs) – обратное преобразование.
  • msTime() – high-resolution время в миллисекундах (process.hrtime).

Прочее

array-utils

  • groupBy(array, key) – группировка массива объектов.
  • unique(array) – удаление дубликатов.
  • compare(arr1, arr2) – поэлементное сравнение.
  • complement(arr1, arr2) – разность множеств.

object-utils

  • clone(obj) – глубокое клонирование.
  • compareRecursive(a, b) – глубокое сравнение.
  • overlay(a, b) – рекурсивное слияние объектов (shallow для необъектов).
  • arrayLimit(item, limit) – обрезает массивы в структуре данных до указанной длины (полезно для логирования).

sleep

  • sleep(ms) – пауза.
  • sleepM(minutes, abortSignal) – пауза в минутах с возможностью отмены.

stream-utils

  • stream2buffer(stream) – чтение ReadableStream в Buffer.

Экспорт

Библиотека экспортирует все основные классы и утилиты через src/index.js. Импорт может быть как деструктурированный, так и через объект utils:

import {
  Config, Logger, LoggerBackendConsole,
  ComplexError, Event, Timer, Trigger,
  ChannelSender, ChannelAggregator,
  ServiceSchema, ServiceClient, Service, ServiceHost, ServiceRequire,
  GlobalService, GlobalServices, GlobalServiceInterface,
  ServiceFactory, ServiceBridge,
  LocalHealth, Sessions, Auth,
  JsonDirList, Profiler, ExpiringList,
  MemoryCache, MemoryCacheAsync, PromiseSingleton, JsonSchema,
  utils
} from '@morphcluster/core';

Вспомогательные функции доступны через utils:

utils.sleep(1000);
utils.clone(data);
utils.groupBy(arr, 'category');
// и т.д.