Предыстория

Все началось с рабочей задачи по реализации весьма специфичной web-панели мониторинга и управления различным оборудованием, в которой нужно было получать данные и отправлять команды по разнообразным сценариям (периодический опрос, подписка на websocket-события, https-запросы и т.д.), а также выводить состояние и графики в реальном времени, динамически комбинируя различные источники данных в разных виджетах.

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

В первой итерации с использованием RxJS получилось много неструктурированного и сложно поддерживаемого кода, так как архитектура на RxJS вынуждала императивно описывать перестроение топологии при изменении правил маршрутизации «на лету». В нашем специфичном кейсе с динамическими виджетами это приводило к сайд-эффектам и сложностям в отладке. Кроме того, местами размывалась строгая типизация и мы лишались compile-time гарантий. Так появилась идея разработать собственное решение на основе альтернативной концепции — модели графа потоков данных.

В результате получившаяся модель продемонстрировала предсказуемое поведение и низкую связность компонентов, а кодовая база сократилась на ~30% и приобрела более декларативный вид. Убедившись в эффективности решения, мы решили оформить его в виде отдельной библиотеки с открытым исходным кодом — Transferum.

И что, получилась просто еще одна реактивная библиотека?

Не совсем. Классические FRP-библиотеки (RxJS, Bacon, Most) построены вокруг единственного примитива (Observable). Transferum же основан на композиции различных типов узлов с явно определенным поведением.

Каждый узел в графе потоков распространения данных явно декларирует свои способности: может ли он принимать данные через push, отдавать через pull, распространять полученный сигнал подписчикам, опрашивать источник, фильтровать, блокировать поток и т.д.

Объявленные узлом возможности являются одновременно флагами для использования в runtime и compile-time гарантиями наличия соответствующих методов, определяющих его поведение.

Ключевая идея: поведение системы описывается как композиция независимых возможностей, которые одновременно определяют тип, реализацию и правила взаимодействия.

Transferum предоставляет четыре слоя абстракции:

  1. Трансферы — узлы графа (каналы, поллеры, мапперы, буферы, разветвители и концентраторы, реализации debounce, throttle, switchMap и т.д.).

  2. Мосты — ребра графа — управляемые вентили между узлами с динамической маршрутизацией и гейтингом.

  3. Операторы — stateless-трансформаторы и фильтры данных (используются трансферами, отвечающими за конвертацию данных).

  4. Билдеры — fluent-конструкторы композитных трансферов из цепочечных трансферов-примитивов.

Концептуальная и архитектурная основа — capability flags system. Каждый трансфер реализует CommunicationContractInterface — набор булевых флагов, определяющих его возможности.

Флаги isPushable, isPullable, isSubscribable, isGate и другие — это не просто свойства объекта. Это метаданные, которые:

  • Определяют TypeScript-интерфейс трансфера на этапе компиляции.

  • Управляют стратегией связывания с другими трансферами в рантайме (с помощью функции linkTransfers()).

  • Обеспечивают совместимость в билдерах без приведений типов.

Один набор флагов — три потребителя. Это единый источник истины для всей системы.

Когда флаг равен true, соответствующий метод входит в TypeScript-интерфейс трансфера. Это позволяет предоставлять трансфер пользователю вот так:

Transfer<T, [Pushable, Pullable, Subscribable, Triggerable]>
// гомогенный трансфер: T -> [ Transfer ] -> T

Или вот так:

Transfer<TInput, TOutput, [Pushable, Pullable, Subscribable, Triggerable]>
// гетерогенный трансфер: TInput -> [ Transfer ] -> TOutput

Именно в таком формате типов фабрики в библиотеке возвращают трансферы. Например:

const pushChannel:  Transfer<number, [Pushable, Subscribable]>;
const converter:    Transfer<number, string, [Pushable, Subscribable]>;
// содержат методы push(), subscribe()

const manualBuffer: Transfer<number, [Pushable, Pullable, Triggerable]>;
// содержит методы push(), pull(), trigger()
// trigger() в данном случае нужен, чтобы сделать запушенное значение 
// доступным для чтения через pull()

const pollingProxy: Transfer<T, [PollingProxy, Pullable, Subscribable, Triggerable, Gate]>
// содержит методы pull(), subscribe(), trigger(), activate(), deactivate(), toggle()
// может быть связан с pullable-источником благодаря поддержке поллинга
Эта «магия» работает в compile-time благодаря несколько замысловатой системе вычислимых типов
type FeatureMap<TIn, TOut, F> =
  F extends { readonly isPushable: true } ? Pushable<TIn> & { readonly isInput: true } :
  F extends { readonly isPollingProxy: true } ? PollingProxy<TIn> & { readonly isInput: true, readonly isOutput: true } :
  F extends { readonly isPullable: true } ? Pullable<TOut> & { readonly isOutput: true } :
  F extends { readonly isSubscribable: true } ? Subscribable<TOut> & { readonly isOutput: true } :
  F extends { readonly isTriggerable: true } ? Triggerable :
  F extends { readonly isGate: true } ? Gate :
  F extends { readonly isAsyncPushable: true } ? AsyncPushable<TIn> & { readonly isInput: true } :
  F extends { readonly isAsyncPollingProxy: true } ? AsyncPollingProxy<TIn> & { readonly isInput: true, readonly isOutput: true } :
  F extends { readonly isAsyncPullable: true } ? AsyncPullable<TOut> & { readonly isOutput: true } :
  F extends { readonly isAsyncTriggerable: true } ? AsyncTriggerable :
  unknown;

type UnionToIntersection<U> = (U extends any ? (k: U) => void : never) extends ((k: infer I) => void) ? I : never;

type ResolveFeatures<TIn, TOut, Features extends any[]> = UnionToIntersection<{
  [K in keyof Features]: FeatureMap<TIn, TOut, Features[K]>
}[number]>;

type IsDuplexFeatures<Features extends any[]> =
  ResolveFeatures<any, any, Features> extends { readonly isInput: true, readonly isOutput: true }
    ? { readonly isDuplex: true }
    : unknown;

type ExplicitTransfer<TIn, TOut, Features extends any[]> = BaseTransferInterface
  & ResolveFeatures<TIn, TOut, Features>
  & IsDuplexFeatures<Features>;

export type Transfer<
  TInOrAll,
  TOutOrFeatures extends any[] | any,
  TFeaturesOrUndefined extends any[] | undefined = undefined,
> = TFeaturesOrUndefined extends any[]
  ? ExplicitTransfer<TInOrAll, TOutOrFeatures, TFeaturesOrUndefined>
  : TOutOrFeatures extends any[]
    ? ExplicitTransfer<TInOrAll, TInOrAll, TOutOrFeatures>
    : never;

Архитектурные инварианты

1. Трансферы не знают своих соседей

Трансфер определяет своё поведение (push, pull, subscribe и др.), но никогда не ссылается и не проверяет класс другого трансфера. Он не знает, что является upstream или downstream — лишь выполняет свой контракт. Пользователь может создать свой трансфер, объявить и реализовать его возможности — и он органично и бесшовно впишется в экосистему.

2. Мосты не знают конкретных реализаций

Мост инспектирует capability flags, а не имена классов. Нет цепочки instanceof, нет переключения по имени класса. Любой output-трансфер может быть соединен с любым input-трансфером — при условии совместимости их флагов, о чем мы поговорим чуть ниже. Это применимо и к тем узлам, которые еще не существуют и будут созданы пользователем.

3. Значение undefined никогда не распространяется

В Transferum undefined означает «нет данных», а не «пустое значение». Оно подавляется на уровне внутренней реализации менеджера подписок — подписчики никогда не уведомляются с undefined. При этом для явных маркеров пустых значений можно использовать null. Мы сознательно пошли на этот компромисс, чтобы избежать runtime-оверхеда и сохранить нативную скорость работы на плотных потоках данных.

Связывание трансферов

Функция linkTransfers(lhs, rhs) соединяет output-трансфер (lhs) с input-трансфером (rhs) с автоматическим выбором стратегии связывания на основе возможностей этих трансферов:

  • isSubscribableisPushable (реактивная подписка);

  • isPullableisPollingProxy (активный опрос);

  • isSubscribableisAsyncPushable (реактивная подписка + асинхронный push);

  • isAsyncPullableisAsyncPollingProxy (активный асинхронный опрос асинхронного pull-источника);

  • isPullableisAsyncPollingProxy (активный асинхронный опрос синхронного pull-источника).

Protocol-oriented design: механизм не спрашивает «какой это класс?» — он выясняет, какие у него есть возможности. Любая пара трансферов с совместимыми возможностями является linkable. Добавление нового класса трансфера требует только объявления его флагов и реализации соответствующих методов — как связать его с другим трансфером, связующий алгоритм разберется сам.

import { createPushChannelTransfer, createSinkTransfer, linkTransfers } from 'transferum';
 
const source = createPushChannelTransfer<number>();
const target = createSinkTransfer<number>({ callback: (x) => console.log(x) });
 
const link = linkTransfers(source, target);
source.push(22);
 
link.unsubscribe(); // разорвать связь

Sync и async в одной экосистеме

Синхронные и асинхронные трансферы сосуществуют и могут быть связаны между собой. linkTransfers() предпочитает sync-связывание, когда это возможно, а async-стратегии применяет только когда sync неприменим. Нет отдельного «асинхронного мира».

Поддержка backpressure

Ряд асинхронных трансферов (AsyncSinkTransfer, AsyncWriteTransfer, AsyncConvertTransfer, AsyncConditionTransfer) поддерживают необязательные поля в конфигурации: maxConcurrency, bufferSize и onBufferOverflow — для ограничения параллельных async-операций, очереди избыточных данных и graceful-обработки переполнения. По умолчанию — неограниченная обработка, без буферизации.

Локальная обработка ошибок

Transferum использует единую модель обработки ошибок для всех трансферов. Каждый трансфер, который может столкнуться с runtime-ошибкой, принимает опциональный onError-хэндлер в своей конфигурации.

А теперь — к примерам использования

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

import { 
  createMergeTransfer, 
  createConditionTransfer, 
  createAsyncWriteTransfer,
  createAsyncSinkTransfer,
  createAsyncPollingSourceTransfer, 
  OutputPipelineBuilder, 
} from 'transferum';

const sensor1 = createAsyncPollingSourceTransfer<SensorData>({
  fetcher: () => Promise.resolve({ sensorId: 1, temperature: 25, humidity: 50 }),
  interval: 50,
  activated: true,
});

const sensor2 = createAsyncPollingSourceTransfer<SensorData>({
  fetcher: () => Promise.resolve({ sensorId: 2, temperature: 26, humidity: 55 }),
  interval: 50,
  activated: true,
});

const sensorsAggregator = createMergeTransfer<SensorData>({
  sources: [sensor1, sensor2],
});

const alertPipeline = OutputPipelineBuilder
  .start(sensorsAggregator)
  .to(createConditionTransfer<SensorData>({
    shouldAccept: (d) => d.temperature > 95,
  }))
  .finish(createAsyncSinkTransfer<SensorData>((sensorData) => console.warn(`Alert! Sensor ${sensorData.sensorId}: temperature is ${sensorData.temperature}`)));

sensorsAggregator.subscribe((sensorData) => console.log('Debugging', sensorData));

А вот как можно организовать динамический роутинг:

import { 
  createPushStoredChannelTransfer,
  createBridgeMultiSelector, 
  createBridgeSelector, 
  createPassBridge, 
} from 'transferum';

// Вариант 1. Выбор одного реципиента
const commandRouter = createBridgeSelector({
  bridges: {
    light: createPassBridge({ source: commandChannel, target: lightController, activated: false }),
    thermostat: createPassBridge({ source: commandChannel, target: thermostatController, activated: false }),
    lock: createPassBridge({ source: commandChannel, target: lockController, activated: false }),
  },
  initialKey: 'thermostat', // выбор по умолчанию
  activated: true,
  owned: true,
});

commandRouter.select('lock'); // выберем другой вариант

// Вариант 2. Выбор нескольких реципиентов
const loggersRouter = createBridgeMultiSelector({
  bridges: {
    prometheus: createPassBridge({ source: metricsChannel, target: prometheusWriter, activated: false }),
    elk: createPassBridge({ source: metricsChannel, target: elkWriter, activated: false }),
    sentry: createPassBridge({ source: metricsChannel, target: sentryWriter, activated: false }),
  },
  initialKeys: ['prometheus', 'elk'], // будут активированы по умолчанию
  activated: true,
  owned: true,
});

loggersRouter.check('elk');                     // теперь выбраны все три
loggersRouter.uncheck('prometheus');            // выбраны elk, sentry
loggersRouter.select(['prometheus', 'sentry']); // выбраны prometheus, sentry
Еще один более длинный пример c организацией игровой механики
import {
  createDebounceTransfer, 
  createConvertTransfer, 
  createMapOperator,
  createBridgeSelector, 
  createPassBridge, 
  linkTransfers,
} from 'transferum';

// Источник: debounce нажатий клавиш (16ms ~ 1 кадр при 60 FPS)
const rawInputSource = createDebounceTransfer<KeyboardEvent>({ delay: 16 });

// Конвертер: KeyboardEvent -> код клавиши
const inputConverter = createConvertTransfer<KeyboardEvent, string>({
  operator: createMapOperator((event) => event.code),
});

// Целевые системы-потребители пользовательских событий
const carSystem = { push: (cmd: string) => console.log(`🚗 Car executing: ${cmd}`) };
const planeSystem = { push: (cmd: string) => console.log(`✈️ Plane executing: ${cmd}`) };
const menuSystem = { push: (cmd: string) => console.log(`📋 Menu processing: ${cmd}`) };

// Роутер: динамическое переключение между игровыми контекстами
const gameplayRouter = createBridgeSelector({
  bridges: {
    driving: createPassBridge({ source: inputConverter, target: carSystem }),
    flying: createPassBridge({ source: inputConverter, target: planeSystem }),
    ui: createPassBridge({ source: inputConverter, target: menuSystem }),
  },
  initialKey: 'driving', // игрок начинает в машине
  activated: true,
});

// Связывание источника с конвертером с автовыбором стратегии
linkTransfers(rawInputSource, inputConverter);

// Демонстрация:

// Игрок в машине нажимает "W"
rawInputSource.push(new KeyboardEvent('keydown', { code: 'KeyW' }));
// Через 16ms: 🚗 Car: KeyW

// ==================

// Переключаемся на самолёт
gameplayRouter.select('flying');

// Та же клавиша теперь управляет самолётом
rawInputSource.push(new KeyboardEvent('keydown', { code: 'KeyW' }));
// Через 16ms: ✈️ Plane: KeyW

Когда имеет смысл попробовать Transferum

Библиотека подойдет для:

  • TypeScript-first проектов — благодаря максимально строгой типизации и compile-time вычислению доступных методов любого трансфера на основе объявленных у него флагов возможностей.

  • Работы с pull-based источниками данных — polling API, датчиков, хранилищ с PollingProxy.

  • Смешанных sync/async пайплайнов — в единой модели без ручного преобразования.

  • Явного flow control — gates, bridges, selectors для runtime-маршрутизации.

  • Game development / IoT — frame-aligned tickers, idle polling, sensor aggregation.

  • Устойчивой обработки ошибок — локальная, non-fatal обработка: одна стадия не убивает пайплайн при условии переданного в конфиге обработчика ошибок, ничего не подавляется молча.

Результаты и планы

  • Библиотека уже используется в двух наших внутренних проектах и показывает свою эффективность. Код доступен на GitHub под лицензией MIT, библиотека не имеет внешних зависимостей и поставляется с подробной документацией (README + API Reference).

  • В одном из проектов граф состоит из ~80 узлов и стабильно обрабатывает несколько сотен событий в секунду без деградации. На основе этих данных в том числе рендерится 3D-сцена в Babylon.js со стабильным фреймрейтом ~60 FPS без микрофризов.

  • Тесты библиотеки покрывают не только отдельные трансферы, но и поведение системы в динамике: переподключение мостов, обработку ошибок в длинных асинхронных цепочках, а также разнообразные граничные случаи. Покрытие — 100%.

  • В дальнейшем планируем реализовать хуки и утилиты для более удобного и нативного использования Transferum с Vue и React. Если они окажутся в достаточной мере переиспользуемыми, оформим в отдельные пакеты-адаптеры.

Буду рад, если вы заглянете в репозиторий, попробуете библиотеку в деле и поделитесь замечаниями — обратная связь поможет сделать Transferum лучше.

P. S. Если вы сталкивались с похожими задачами и решили их как-то иначе — буду рад прочитать о вашем опыте в комментариях.