8000
Skip to content

Repository files navigation

Event Sourcing Playground

Практическое пошаговое изучение паттерна Event Sourcing на примере предметной области морских грузоперевозок. Каждый пример усложняет предыдущий, постепенно вводя новые концепции: от базового apply/replay до ретроактивных событий с откатом побочных эффектов во внешних системах.

Содержание


Стек технологий

Компонент Версия Назначение
TypeScript ^5.8 Язык
tsx ^4.19 Runtime (запуск .ts напрямую, без предкомпиляции)
Node.js Платформа
MongoDB* ^6.16 Персистенция (event-based-CRUD)
Express* ^5.1 HTTP-сервер (event-based-CRUD)
fp-ts* ^2.16 Функциональные утилиты (event-based-CRUD)
pg / knex* ^8.16 / ^3.1 PostgreSQL-клиент (не используется в текущих примерах)
uuid* ^11.1 Генерация ID (не используется в текущих примерах)

* MongoDB, Express, fp-ts, pg, knex и uuid установлены как зависимости, но используются только в бонусном примере event-based-CRUD (или зарезервированы для будущих примеров). Пять основных примеров работают полностью in-memory без внешних сервисов.

Примечание о модульной системе: package.json объявляет "type": "module" (ESM), а tsconfig.json указывает "module": "CommonJS". Это не является конфликтом — tsx прозрачно транспилирует TypeScript-модули и обрабатывает разрешение модулей самостоятельно, минуя стандартный tsc + node pipeline.

Установка и запуск

# Клонирование репозитория
git clone <url>
cd event-sourcing-playground

# Установка зависимостей
pnpm install

# Запуск конкретного примера
pnpm run 1-tracking-ships
pnpm run 2-reverse-event
pnpm run 3-external-system
pnpm run 4-retroactive-cargoes-with-rewind
pnpm run 5-retroactive-external-systems

У бонусного примера event-based-CRUD нет скрипта запуска в package.json — это незавершённый набросок без точки входа index.ts.

Отладка

В .vscode/launch.json настроена конфигурация для отладки через tsx:

  • Открыть index.ts нужного примера
  • Нажать F5 — запустится текущий файл через tsx с подключённым отладчиком
  • node_internals и node_modules исключены из стека вызовов

Архитектура примеров

Каждый нумерованный каталог следует единой структуре (с вариациями):

N-example-name/
├── domain/
│   ├── events.ts              # Доменные события (ArrivalEvent, LoadEvent, ...)
│   ├── ship.ts                # Агрегат Ship
│   ├── value-objects.ts       # Объекты-значения (Port, Cargo)
│   ├── types.ts               # Базовый интерфейс/абстрактный класс DomainEvent
│   ├── cargo.ts               # Агрегат Cargo (примеры 5)
│   ├── cargo-events.ts        # События с грузом (примеры 4, 5)
│   └── replacement-event.ts   # ReplacementEvent (примеры 4, 5)
├── infra/
│   ├── event-processor.ts     # Процессор событий (apply, log, replay, rewind)
│   ├── registry.ts            # Service Locator (примеры 3, 5)
│   └── ...                    # Шлюзы, буферы, кеши
└── index.ts                   # Исполняемый сценарий

В примере 3 агрегаты вынесены в подкаталог domain/aggregates/ (ship.ts, cargo.ts).

Прогрессия сложности

[1] Apply + Replay (интерфейс DomainEvent)
       ↓
[2] + Reverse/undo (интерфейс DomainEvent + reverse())
       ↓
[3] + External system calls + QueryLog (кеш для replay)
       ↓
[4] + Retroactive events (абстрактный класс DomainEvent, rewind/replay, replacement, rejection)
       ↓
[5] + External side-effects buffering (cancel/re-notify, ReplayBuffer)

Важное архитектурное изменение: в примерах 1–3 DomainEvent — это интерфейс с асинхронными методами (Promise<void>). Начиная с примера 4, DomainEvent становится абстрактным классом с синхронными методами (void), несущим общее состояние (processingError, rejected) и поведение (shouldIgnoreOnReplay, after(), isConsequenceOf()).


Пример 1: Tracking Ships — базовый Event Sourcing

Каталог: 1-tracking-ships/

Концепция

Самый простой пример: доменные события применяются к агрегату Ship, сохраняются в лог, а затем могут быть воспроизведены с нуля для восстановления состояния.

Доменные события

Событие Поля Эффект
ArrivalEvent ship, port, type='Arrival' Устанавливает ship.currentPort = port
DepartureEvent ship, type='Departure' Сбрасывает ship.currentPort = null
LoadEvent ship, cargo, type='Load' Добавляет груз в ship.loadedCargo
UnloadEvent ship, cargo, type='Unload' Убирает груз из ship.loadedCargo

Каждое событие имеет readonly-поле type (строковый дискриминатор) и реализует интерфейс:

interface DomainEvent {
  occurred: Date
  apply(): Promise<void>
}

Агрегат Ship

class Ship {
  name: string              // имя корабля (через конструктор)
  currentPort: Port | null = null
  loadedCargo: Cargo[] = []

  arriveAt(port: Port): void
  depart(): void
  load(cargo: Cargo): void
  unload(cargo: Cargo): void
  reset(): void             // сброс состояния для replay
}

Value Objects (классы)

  • Portclass Port { name: string, country: string } — порт с указанием страны
  • Cargoclass Cargo { description: string } — груз с описанием

EventProcessor

Ядро системы — процессор хранит лог событий и управляет их жизненным циклом:

  • process(event) — вызывает event.apply(), добавляет в лог
  • replay() — сбрасывает все агрегаты из лога через reset(), заново применяет все события в хронологическом порядке

Сценарий

Корабль «King Roy» совершает рейс: прибывает в Ванкувер → загружает грузы (Candies, Tomatoes) → отплывает → прибывает в Сан-Франциско → разгружает Candies. После полного replay состояние совпадает с оригиналом.


Пример 2: Reverse Event — отмена событий

Каталог: 2-reverse-event/

Концепция

Расширяет модель событий методом reverse(), позволяя отменять последнее событие. Событие сохраняет снимок состояния перед применением (паттерн Memento), чтобы корректно его восстановить при откате.

Интерфейс DomainEvent

interface DomainEvent {
  readonly occurred: Date
  apply(): Promise<void>
  reverse(): Promise<void>  // откат эффекта
}

Агрегат Ship (упрощённая модель)

В отличие от примера 1, здесь Ship содержит только счётчик груза:

class Ship {
  cargoCount: number = 0
  reset(): void  // сброс cargoCount в 0
}

LoadEvent — пример реверсивного события

class LoadEvent implements DomainEvent {
  private previousCount: number = 0 // снимок до apply (Memento)

  async apply() {
    this.previousCount = this.ship.cargoCount  // сохраняем
    this.ship.cargoCount += this.count          // применяем
  }

  async reverse() {
    this.ship.cargoCount = this.previousCount   // восстанавливаем
  }
}

EventProcessor

  • process(event) — apply + запись в лог (лог public)
  • undoLast() — извлекает последнее событие из лога, вызывает reverse()
  • replay() — полное воспроизведение (отменённые события отсутствуют в логе, т.к. undoLast удаляет их)

Сценарий

Корабль «Sea Queen»: загрузка 5 единиц → загрузка 3 единиц (итого 8) → undo (возврат к 5) → replay (подтверждает 5).


Пример 3: External System — интеграция с внешними системами

Каталог: 3-external-system/

Концепция

При обработке события может потребоваться вызов внешнего сервиса (расчёт стоимости груза через PricingGateway). При replay нельзя повторно вызывать внешний API — результаты кешируются в QueryLog.

Структура domain/aggregates/

В этом примере агрегаты вынесены в отдельный подкаталог domain/aggregates/:

3-external-system/domain/
├── aggregates/
│   ├── ship.ts    # Ship: loadedCargo, totalPrice, loadCargo(), setSumCargoPrice()
│   └── cargo.ts   # Cargo: id, weight, cargoPrice, handleArrival()
├── events.ts      # ArrivalEvent
└── types.ts       # DomainEvent (interface), Money (class)

ArrivalEvent (специфика примера 3)

В отличие от примера 1, ArrivalEvent здесь принимает cargo (а не port) и отвечает за загрузку груза и расчёт стоимости:

class ArrivalEvent implements DomainEvent {
  constructor(occurred: Date, ship: Ship, cargo: Cargo) {}

  async apply() {
    this.ship.loadCargo(this.cargo)                        // 1. Загрузка
    const prices = await Promise.all(                      // 2. Запрос цен
      this.ship.loadedCargo.map(c => c.handleArrival())    //    для каждого груза
    )
    const total = prices.reduce((acc, p) => acc + p.amount, 0)
    this.ship.setSumCargoPrice(new Money(total))           // 3. Итоговая цена
  }
}

Cargo.handleArrival() обращается к Registry.pricingGateway.getPrice(this).

Архитектура (паттерн Decorator)

ArrivalEvent.apply()
    → ship.loadCargo(cargo)
    → cargo.handleArrival()
        → Registry.pricingGateway.getPrice(cargo)
            → LoggedPricingGateway (Decorator)
                ├── [cache miss] → RealPricingGateway → сохранить в QueryLog
                └── [cache hit]  → вернуть из QueryLog

PricingGateway

interface PricingGateway {
  getPrice(cargo: Cargo): Promise<Money>
}

// Имитация внешнего API (логирует вызов в консоль)
class RealPricingGateway implements PricingGateway {
  getPrice(cargo: Cargo): Promise<Money> {
    console.log('[External pricing service]: Do a real job...')
    return Promise.resolve(new Money(cargo.weight * 100 + 500))
  }
}

LoggedPricingGateway (паттерн Decorator)

Оборачивает реальный шлюз, реализуя тот же интерфейс PricingGateway. При вызове getPrice():

  1. Проверяет QueryLog по составному ключу (тип запроса, класс доменного события, ID груза)
  2. Cache miss — вызывает настоящий шлюз, сохраняет результат как QueryEvent в QueryLog
  3. Cache hit — возвращает кешированный результат (при replay внешний сервис не вызывается)

Для корреляции кеш-записей с текущим событием используется указатель EventProcessor.currentEvent, который процессор устанавливает перед apply() и сбрасывает после.

EventProcessor (специфика примера 3)

class EventProcessor {
  currentEvent: DomainEvent | null = null  // указатель на обрабатываемое событие
  log: DomainEvent[] = []

  async process(event) {
    this.currentEvent = event    // установить указатель
    await event.apply()
    this.log.push(event)
    this.currentEvent = null     // сбросить
  }

  async replay() {
    // reset всех Ship из лога, затем повторный apply каждого события
    // currentEvent устанавливается и при replay — для корректной работы кеша
  }
}

QueryLog и QueryEvent

In-memory хранилище результатов внешних вызовов:

// query-log.ts
class QueryEvent {
  type: string           // тип запроса ('GetPriceReqForCargo')
  domainEvent: DomainEvent  // ссылка на обрабатываемое событие
  payload: string        // ID груза
  result: Money          // кешированный результат
}

class QueryLog {
  cacheRequest(req: QueryEvent): void
  findCachedVal(options: { key, event, id }): Money | null
  // Поиск использует instanceof по constructor класса события,
  // т.е. кеш-попадание определяется классом события, а не идентичностью
}

В каталоге также есть query-event.ts — альтернативное определение QueryEvent с полем cargoId вместо payload и без result в конструкторе. Используется в LoggedPricingGateway.

Registry (паттерн Service Locator)

Связывает компоненты при инициализации:

Registry.init()
  → EventProcessor
  → QueryLog
  → LoggedPricingGateway(RealPricingGateway, QueryLog, EventProcessor)

Сценарий

Корабль «King Roy», два груза (REF123 вес=20, Uni123 вес=42). Обрабатываются два ArrivalEvent. При replay цены берутся из кеша — в консоли отсутствуют повторные сообщения [External pricing service]: Do a real job....


Пример 4: Retroactive Events — ретроактивные события с перемоткой

Каталог: 4-retroactive-cargoes-with-rewind /

Концепция

Ретроактивные события — ключевой паттерн данного примера. События могут поступать не в хронологическом порядке, могут быть заменены исправленной версией или полностью отклонены. Процессор обрабатывает это через механизм rewind (обратная перемотка) и replay (повторное воспроизведение).

DomainEvent — переход от интерфейса к абстрактному классу

Начиная с этого примера DomainEventабстрактный класс (а не интерфейс, как в примерах 1–3). Это позволяет нести общее состояние и поведение для управления отклонением и ошибками:

abstract class DomainEvent {
  processingError: Error | null = null  // ошибка при apply()
  private rejected: boolean = false     // отклонено

  constructor(readonly occurred: Date) {}

  get shouldIgnoreOnReplay(): boolean   // true если rejected или есть ошибка
  after(other: DomainEvent): boolean    // хронологическое сравнение по occurred
  isConsequenceOf(base: DomainEvent): boolean  // !shouldIgnoreOnReplay && after(base)
  reject(): void                        // пометить как отклонённое

  abstract apply(): void    // синхронный (не Promise!)
  abstract reverse(): void  // синхронный (не Promise!)
}

apply() и reverse() в примерах 4 и 5 синхронные (void), в отличие от асинхронных (Promise<void>) в примерах 1–3.

Доменные события

Событие apply() reverse()
ArrivalEvent ship.arriveAt(port) ship.depart()
LoadEvent Сохраняет priorPort (Memento), вызывает ship.load(cargo) ship.unload(cargo), восстанавливает cargo.currentPort из priorPort

ReplacementEvent

Обёртка над парой (оригинал, замена):

class ReplacementEvent extends DomainEvent {
  original: DomainEvent         // исходное событие
  replacement: DomainEvent | null  // замена (null = отклонение)

  get hasPriorReplacement(): boolean  // replacement.occurred < original.occurred
}

apply() и reverse() выбрасывают ошибку — ReplacementEvent обрабатывается процессором напрямую через processReplacement().

Агрегат Ship

class Ship {
  name: string
  cargoes: Cargo[] = []
  currentPort: Port | null = null

  arriveAt(port: Port): void
  depart(): void
  load(cargo: Cargo): void    // валидация: корабль пришвартован И груз в том же порту
  unload(cargo: Cargo): void  // валидация: корабль пришвартован И груз на борту
  unloadAll(): void           // разгрузка всех грузов
  getCargoes(): Cargo[]
}

Валидация в load():

  • Если корабль в море (currentPort === null) — бросает ошибку
  • Если порт груза не совпадает с портом корабля (cargo.currentPort.name !== ship.currentPort.name) — бросает ошибку

Value Objects

  • Portclass Port { name: string } (без country, в отличие от примера 1)
  • Cargoclass Cargo { id: string, weight: number, currentPort: Port | null } — груз отслеживает свой текущий порт; метод arriveAt(port) устанавливает currentPort

EventProcessor — ядро ретроактивной логики

Ключевые изменения относительно примеров 1–2:

  • Лог всегда отсортированinsertToLog() после push сортирует по occurred.getTime()
  • Обработка ошибокprocess() оборачивает apply в try/catch: ошибка сохраняется в event.processingError, событие с ошибкой игнорируется при replay (shouldIgnoreOnReplay). Флаг shouldRethrow управляет всплытием ошибок

Маршрутизация в process():

process(event)
├── instanceof ReplacementEvent?  → processReplacement()
├── outOfOrder(event)?            → processOutOfOrder()
└── default                       → basicProcess() (простой apply)

Обработка замены/отклонения (processReplacement):

  1. Отклонение (replacement = null):

    • Rewind всех последствий оригинала (reverse в обратном порядке)
    • Reverse оригинала, пометить как rejected
    • Replay последствий
  2. Замена с более ранней датой (replacement.occurred < original.occurred):

    • Rewind до даты замены
    • Отклонить оригинал
    • Apply замену
    • Replay всех последующих событий
  3. Замена с более поздней датой (replacement.occurred > original.occurred):

    • Rewind последствий оригинала
    • Reverse оригинала, отклонить
    • Replay событий между оригиналом и заменой (replayBetween)
    • Apply замену
    • Replay событий после замены

Обработка события вне хронологии (processOutOfOrder):

  • Определяется через outOfOrder(): last.after(newEvent) — последнее событие в логе позже нового
  • Rewind всех последствий → Apply опоздавшее → Replay последствий

В примере 4 нет lifecycle-хуков (onRewindStarted/Finished/onReplayFinished) — они появляются только в примере 5 для интеграции с ReplayBuffer.

Сценарий (index.ts)

Файл содержит 4 примера (3 закомментированных + 1 активный):

Пример Описание Результат
Example 1 (закомментирован) Замена с поздней датой, но нарушает последовательность load→arrival Не работает (ошибка)
Example 2 (закомментирован) Замена с более поздней датой (corrected > original) Работает (processReplacementAfterOriginal)
Example 3 (закомментирован) Замена с более ранней датой (corrected < original) Работает (processReplacementBeforeOriginal)
Example 4 (активный) Out-of-order: LoadEvent (2 янв) → ArrivalEvent (1 янв) Работает: rewind Load → apply Arrival → replay Load

Пример 5: Retroactive External Systems — ретроактивность + побочные эффекты

Каталог: 5-retroactive-external-systems/

Концепция

Самый сложный пример. Объединяет ретроактивные события (rewind/replay/replacement) с побочными эффектами во внешних системах (таможенные уведомления). При ретроактивном изменении событий внешние уведомления должны быть отменены и переотправлены корректно.

Бизнес-правила (инварианты)

  • Если корабль прибывает в US и груз ранее был в Canada — уведомить таможню со специальным сообщением «Hello from Canada to US!»
  • Если корабль прибывает в China и груз ранее был в Canada — уведомить китайскую таможню
  • При прибытии в любой порт без выполнения условий выше — уведомить со стандартным сообщением

Агрегат Cargo — отслеживание флагов стран

class Cargo {
  readonly id: string
  readonly weight: number
  currentPort: Port | null

  // Приватные флаги-маркеры посещённых стран
  private inCanada: boolean = false
  private inChina:  boolean = false
  private inUS:     boolean = false

  arriveAt(port: Port): void
  // Устанавливает currentPort и переключает соответствующий флаг в true.
  // Флаги "липкие" — устанавливаются в true и не сбрасываются естественным путём.
  // Сброс возможен только через restoreFlags().

  needsUSCustomsNotification(): boolean     // inCanada && currentPort?.name === "US"
  needsChinaCustomsNotification(): boolean  // inCanada && currentPort?.name === "China"

  snapshotFlags(): { inCanada, inChina, inUS }  // Memento: сохранение
  restoreFlags(flags): void                      // Memento: восстановление
}

ArrivalEvent — apply/reverse с snapshot

Уведомления таможни не вызываются внутри ArrivalEvent.apply(). Событие отвечает только за мутацию состояния агрегатов и сохранение снимков. Уведомления инициируются внешним кодом (из index.ts) через CustomsGatewayBuffer:

class ArrivalEvent extends DomainEvent {
  priorFlags = new Map<Cargo, FlagSnapshot>()  // снимки до apply

  apply(): void {
    this.ship.arriveAt(this.port)
    for (const cargo of this.ship.getCargoes()) {
      this.priorFlags.set(cargo, cargo.snapshotFlags())  // Memento
      cargo.arriveAt(this.port)
    }
  }

  reverse(): void {
    this.ship.depart()
    for (const [cargo, flags] of this.priorFlags) {
      cargo.restoreFlags(flags)  // восстановление из Memento
    }
  }
}

EventProcessor — с lifecycle-хуками

В отличие от примера 4, процессор реализует интерфейс IRewindable с тремя хуками:

class EventProcessor implements IRewindable {
  onRewindStarted:  () => void   // вызывается в начале перемотки
  onRewindFinished: () => void   // вызывается после завершения реверсов
  onReplayFinished: () => void   // вызывается после завершения повторного воспроизведения
}

Хуки вызываются вручную внутри processRejection, processReplacementBeforeOriginal, processReplacementAfterOriginal и processOutOfOrder в определённые моменты. Подписка на хуки — прямое присвоение функций (паттерн Observer в минималистичной форме, один подписчик).

Инфраструктура побочных эффектов

Архитектура состоит из четырёх слоёв:

EventProcessor (IRewindable)
    ↓ lifecycle hooks (onRewindStarted / Finished / onReplayFinished)
ReplayBuffer
    ↓ send() / determineChange()
CustomsGatewayBuffer (IAdjustable)
    ↓ adjust() → apply()
CustomsGatewayFront (ICustomsGatewayFront — реальный шлюз)

Интерфейсы (replay-buffer.ts)

interface IGatewayEvent {
  readonly key: string
  readonly gateway: ICustomsGatewayFront
  readonly arrivalDate: Date
  readonly ship: Ship
  readonly port: Port
  readonly customMessage: string | null
  apply(): void
}

interface IRewindable {
  onRewindStarted: Function
  onRewindFinished: Function
  onReplayFinished: Function
}

interface IAdjustable {
  adjust(toCancelEvents: IGatewayEvent[], toNotifyEvents: IGatewayEvent[]): void
}

ReplayBuffer — буфер отслеживания побочных эффектов

Отслеживает побочные эффекты на трёх стадиях:

Стадия Массив isActive Поведение
Нормальная работа sent[] true send()current.push(ev) + ev.apply()
Rewind rewound[] false send()current.push(ev) (без apply)
Replay replayed[] false send()current.push(ev) (без apply)

Жизненный цикл (подписка через subscribeToProcessorEvents()):

onRewindStarted()  → isActive = false, current = rewound[] (новый пустой)
onRewindFinished() → current = replayed[] (новый пустой)
onReplayFinished() → current = sent[], вызов determineChange(), isActive = true

Два режима расчёта изменений

Базовый режим (determineChange):

  • Передать ВСЕ rewound[] как cancellations и ВСЕ replayed[] как notifications в adjuster.adjust()

Точный режим (determineChangeAccurately):

  • Вычислить разницу по ключам событий через Set:
    • toCancel = rewound-события, чьих ключей нет в replayed
    • toNotify = replayed-события, чьих ключей нет в rewound
  • Передать только дифференциальный набор в adjuster.adjust()

Режим выбирается при инициализации: retroactiveEventsAccurateCalculationMode: true|false. По умолчанию в index.tsfalse.

CustomsGatewayBuffer

  • notify(arrivalDate, ship, port, customMessage) — создаёт CustomsNotificationEvent с вычисленным ключом, отправляет через buffer.send()
  • cancel(arrivalDate, ship, port, customMessage) — создаёт CustomsCancellationEvent, отправляет через buffer.send()
  • adjust(toCancelEvents, toNotifyEvents) — вызывается буфером; применяет .apply() на каждом событии
  • calculateKey(options) — генерирует строковый ключ в формате ${arrivalDate}-${shipName}-${portName}[-${customMessage}] (разделитель -, дата — Date.toString())

CustomsNotificationEvent и CustomsCancellationEvent — реализации IGatewayEvent, делегирующие apply() в gateway.notify() / gateway.cancel().

CustomsGatewayFront

Реальный шлюз (ICustomsGatewayFront). В данном случае — логирование в консоль:

[NOTIFY_EVENT]: Mercury at Canada on 2024-01-02T00:00:00.000Z
[CANCEL_EVENT]: Mercury at US on 2024-01-04T00:00:00.000Z. <<Hello from Canada to US!>>

Сценарий

  1. Конфигурация: retroactiveEventsAccurateCalculationMode: false (можно переключить на true)
  2. Создание: порты US, China, Canada; корабль «Mercury» (стартует в US); груз «GoldBars» (1000, в US)
  3. Обработка цепочки: E865 LoadEvent (1 янв) → ArrivalEvent Canada (2 янв) → ArrivalEvent China (3 янв) → ArrivalEvent US (4 янв)
  4. При прибытии в US груз удовлетворяет needsUSCustomsNotification() (был в Canada) → уведомление с «Hello from Canada to US!»
  5. Ретроактивное отклонение: ReplacementEvent(arrivalCa, null) — отклонение прибытия в Canada → rewind → cancel уведомления о US с сообщением → replay → re-notify о US без специального сообщения (т.к. inCanada восстановлен в false)
  6. Ретроактивная замена: ReplacementEvent(arrivalUs, fixedUS) — замена даты прибытия в US с 4 на 5 января → cancel старого → notify нового

Вывод: точный режим расчёта

[NOTIFY_EVENT]: Mercury at Canada on 2024-01-02T00:00:00.000
[NOTIFY_EVENT]: Mercury at China on 2024-01-03T00:00:00.000
[NOTIFY_EVENT]: Mercury at US on 2024-01-04T00:00:00.000Z. <<Hello from Canada to US!>>

Rewind process...
Reject ArrivalEvent event at 2024-01-02

[Retroactive events accurate calculation mode] is enabled...

[CANCEL_EVENT]: Mercury at US on 2024-01-04T00:00:00.000Z. <<Hello from Canada to US!>>
[CANCEL_EVENT]: Mercury at Canada on 2024-01-02T00:00:00.000Z
[NOTIFY_EVENT]: Mercury at US on 2024-01-04T00:00:00.000Z

Rewind process...
Replace ArrivalEvent event at 2024-01-04 → 2024-01-05

[CANCEL_EVENT]: Mercury at US on 2024-01-04T00:00:00.000Z
[NOTIFY_EVENT]: Mercury at US on 2024-01-05T00:00:00.000Z

В точном режиме уведомление о China не отменяется и не переотправляется — оно не затронуто изменениями.

Вывод: базовый режим расчёта

[NOTIFY_EVENT]: Mercury at Canada on 2024-01-02T00:00:00.000Z
[NOTIFY_EVENT]: Mercury at China on 2024-01-03T00:00:00.000Z
[NOTIFY_EVENT]: Mercury at US on 2024-01-04T00:00:00.000Z. <<Hello from Canada to US!>>

Rewind process...
Reject ArrivalEvent event at 2024-01-02

[Retroactive events accurate calculation mode] is disabled...

[CANCEL_EVENT]: Mercury at US on 2024-01-04T00:00:00.000Z. <<Hello from Canada to US!>>
[CANCEL_EVENT]: Mercury at China on 2024-01-03T00:00:00.000Z              ← лишняя операция
[CANCEL_EVENT]: Mercury at Canada on 2024-01-02T00:00:00.000Z
[NOTIFY_EVENT]: Mercury at China on 2024-01-03T00:00:00.000Z              ← лишняя операция
[NOTIFY_EVENT]: Mercury at US on 2024-01-04T00:00:00.000Z

Rewind process...
Replace ArrivalEvent event at 2024-01-04 → 2024-01-05

[CANCEL_EVENT]: Mercury at US on 2024-01-04T00:00:00.000Z
[NOTIFY_EVENT]: Mercury at US on 2024-01-05T00:00:00.000Z

Финальное состояние Ship

Ship {
  name: 'Mercury',
  currentPort: Port { name: 'US' },
  cargoes: [
    Cargo {
      id: 'GoldBars',
      weight: 1000,
      currentPort: null,
      inCanada: false,  // Canada отклонена → флаг восстановлен
      inChina: true,
      inUS: true
    }
  ]
}

Поля inCanada, inChina, inUS — приватные. В реальном console.log они могут не отображаться (зависит от runtime). Показаны здесь для наглядности.


Event-based CRUD — событийный CRUD (бонус)

Каталог: event-based-CRUD/

Отдельный, незавершённый пример, показывающий более традиционный подход к Event Sourcing в контексте CRUD-приложения: корзина покупок с MongoDB-персистенцией и Express HTTP-эндпоинтами.

Статус: набросок. Нет скрипта запуска в package.json, нет index.ts. Ряд типов и зависимостей (EventBus, команды OpenShoppingCart/AddProductItemToShoppingCart/RemoveProductItemFromShoppingCart/ConfirmShoppingCart, функции from(), getShoppingCart(), getCollection()) используются, но не определены в кодовой базе. Маппер ShoppingCartMapper вызывает конструкторы new ShoppingCart(...) и new ShoppingCartModel(...), хотя оба определены как type, а не class — код не скомпилируется без доработки.

Структура

event-based-CRUD/
├── types.ts           # ShoppingCart (type), ShoppingCartEvent (discriminated union),
│                      #   ShoppingCartStatus, ProductItem
├── model.ts           # ShoppingCartModel (type — MongoDB-документ)
├── store.ts           # Проекция событий → MongoDB-операции (upsert, $push, $inc)
├── handlers.ts        # 4 командных обработчика с валидацией бизнес-правил
├── add-product-item-to-shopping-cart.handler.ts
│                      # Класс-обработчик: load → map → mutate → persist → publish
├── router.ts          # Express DELETE-роут для удаления товара из корзины
├── shopping-cart.mapper.ts    # Маппер: ShoppingCart ↔ ShoppingCartModel
└── mongodb.repository.ts     # Обобщённый CRUD-репозиторий (add, update, upsert, find)

Доменные типы

type ShoppingCart = {
  id: string
  status: ShoppingCartStatus
  productItems: Map<string, number>  // productId → quantity
}

const ShoppingCartStatus = { Opened: 'Opened', Confirmed: 'Confirmed' }

type ProductItem = { productId: string, quantity: number }

Доменные события (discriminated union)

Событие Данные
ShoppingCartOpened shoppingCartId, customerId, openedAt
ProductItemAddedToShoppingCart shoppingCartId, productItem: ProductItem, addedAt
ProductItemRemovedFromShoppingCart shoppingCartId, productItem: ProductItem, removedAt
ShoppingCartConfirmed shoppingCartId, confirmedAt

Командные обработчики (handlers.ts)

Четыре функции, каждая из которых принимает команду (+ текущее состояние агрегата для мутирующих операций), валидирует бизнес-правила и возвращает событие:

Обработчик Валидация
openShoppingCart(cmd)
addProductItemToShoppingCart(cmd, cart) cart.status === Opened
removeProductItemFromShoppingCart(cmd, cart) cart.status === Opened, quantity > 0
confirmShoppingCart(cmd, cart) cart.status === Opened

Проекция событий в MongoDB (store.ts)

Функция store(carts: Collection, event) выполняет MongoDB-операции в зависимости от типа события:

  • shopping-cart-openedupsert нового документа
  • product-item-added$push в массив + $inc quantity через arrayFilter
  • product-item-removed$inc с отрицательным quantity
  • shopping-cart-confirmed → обновление status на Confirmed

Поток обработки (паттерн Command)

HTTP Request
  → Router (парсинг в Command)
    → Handler (валидация бизнес-правил → Event)
      → Store (проекция Event → MongoDB)
        → EventBus (публикация Event подписчикам)

Используемые паттерны проектирования

Паттерн Где используется Описание
Event Sourcing Все примеры Состояние восстанавливается из последовательности событий
Memento Примеры 2, 4, 5 Событие сохраняет snapshot состояния до apply() для восстановления при reverse(): previousCount (пример 2), priorPort (пример 4), priorFlags (пример 5)
Decorator Пример 3 LoggedPricingGateway оборачивает RealPricingGateway, добавляя кеширование через QueryLog без изменения интерфейса
Service Locator Примеры 3, 5 Registry хранит синглтоны инфраструктурных компонентов и связывает их при init()
Observer Пример 5 ReplayBuffer подписывается на lifecycle-хуки EventProcessor через прямое присвоение callback-функций (onRewindStarted, onRewindFinished, onReplayFinished). Ограничение: один подписчик
Command event-based-CRUD Обработчики принимают Command-объект, валидируют инварианты, генерируют Event

Ключевые концепции

Event Sourcing

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

Apply / Replay

  • Apply — применение события к агрегату (мутация состояния)
  • Replay — сброс состояния (reset()) и повторное применение всех событий из лога

Reverse (Undo)

Каждое событие хранит снимок состояния до применения (паттерн Memento) и может откатить свой эффект через reverse().

QueryLog (кеширование внешних вызовов)

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

Retroactive Events (ретроактивные события)

Механизм обработки корректировок задним числом:

  • Rejection — полное отклонение события (ReplacementEvent(original, null))
  • Replacement — замена события исправленной версией (ReplacementEvent(original, corrected))
  • Out-of-order — вставка опоздавшего события в правильное место хронологии

Все три случая обрабатываются через rewind (обратная перемотка последствий в обратном хронологическом порядке) + replay (повторное воспроизведение в прямом порядке).

Sorted Log

В примерах 4–5 лог событий поддерживается в хронологическом порядке: insertToLog() после push вызывает sort() по occurred.getTime(). Это гарантирует корректную работу consequences() и replayAfter() независимо от порядка поступления событий.

Error Handling (примеры 4–5)

Процессор оборачивает apply() в try/catch:

  • Ошибка сохраняется в event.processingError
  • Событие с ошибкой автоматически игнорируется при replay (shouldIgnoreOnReplay)
  • Флаг shouldRethrow определяет, всплывает ли ошибка наверх или поглощается

ReplayBuffer (буферизация побочных эффектов)

При ретроактивных изменениях побочные эффекты (уведомления, вызовы API) не должны дублироваться. Буфер отслеживает эффекты на стадиях rewind и replay, а затем — в зависимости от режима — либо полностью отменяет/переотправляет все эффекты, либо вычисляет минимальный дифференциальный набор изменений (cancel/notify) по ключам событий.

About

Play around with event sourcing using Node.js

Topics

Resources

Stars

1 star

Watchers

1 watching

Forks

Releases

Packages

Used by

Contributors

Languages

0