Раздел 25 · Effect-TS

CQRS-проекции: Daemon, SubscriptionRef, live read-model

senior~70 мин

открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти

CQRS-проекции: Daemon, SubscriptionRef, live read-model

Сцена · команды и запросы говорят на разных языках

У отеля две роли. Менеджер ресепшен меняет состояние: бронирует, отменяет, оформляет заезд. Менеджер ресторана читает состояние: “сколько гостей в отеле сейчас, кто из них на полу 1, кто на полу 2”. Их интересы не пересекаются: первому важна корректность транзакций, второму нужна свежая красивая таблица.

CRUD-подход пытается обслужить обе роли одной моделью. Получается компромисс: схема, которая нормализована для писателя, но плохо читается, или денормализована для читателя, но мучительна для писателя. CQRS говорит: разведём write-сторону и read-сторону по разным моделям. Писатель кладёт события в Event Store. Читатель смотрит на read-model, специально заточенную под его экран. Между ними стоит проекция, которая трансформирует поток событий в read-model.

В этом уроке мы строим живую проекцию “occupancy”, “кто сейчас в отеле”. Она читает поток событий из Event Store, обновляется при каждом GuestCheckedIn/GuestCheckedOut, и UI подписывается на её изменения через SubscriptionRef. Без polling, без своего websocket, без задержки.

Карта урока · что заберёшь домой

  1. CQRS как разделение языков. Команды vs запросы, write-модель vs read-модель.
  2. Проекция как чистая функция. (state, event) -> state. То же, что evolve, но другой state.
  3. SubscriptionRef. Ref с встроенной подпиской на изменения. Подписчики получают каждое новое значение.
  4. Daemon под Effect.forkScoped. Фоновый файбер, который держит read-model в актуальном состоянии.
  5. Идемпотентность. Почему проекция должна выдерживать повторное применение события.
  6. Last_seq tracking. Как пережить рестарт без потери позиции.
  7. Polling против LISTEN. Когда хватит первого, когда нужен второй.

Раздел 1 · CQRS, кратко

Сцена

CQRS (Command Query Responsibility Segregation) сформулирован Грегом Янгом в 2010 году. Идея простая: команды и запросы это разные операции, и им можно дать разные модели данных.

        Command path                              Query path
        ─────────────                             ──────────

  POST /reservations                       GET /occupancy
        │                                       │
        ▼                                       ▼
   CommandHandler                          ReadModel
        │                                       ▲
        ▼                                       │
   reservationDecider             ┌─── Projection (Daemon)
        │                         │
        ▼                         │
   EventStore (Postgres)  ────────┘  events stream

Левая ветка это уроки 14-16: команда -> decider -> event store. Правая ветка это сегодняшний урок: event store -> проекция -> read model -> UI.

Между ветками есть один общий ресурс: поток событий. Write-сторона его пишет, read-сторона его читает. Read-модели не обновляются напрямую командами; они материализуются исключительно из событий, и в этом залог их корректности.

Что взять с собой

  • CQRS разделяет write и read модели, не репликацию БД.
  • Источник истины это events; read-модели производны от него.
  • Read-модель можно пересоздать с нуля, прокрутив все события.

Раздел 2 · Проекция как чистая функция

Сцена

Проекция, по форме, это та же функция, что evolve в Decider-е. Берёт текущее состояние read-модели и событие, возвращает новое состояние. Но state другой: не агрегат, а денормализованная view.

Шаг 1 · Что в read-model

Хотим показать таблицу “кто сейчас в отеле”. Каждая строка это { reservationId, room, guest }. Гость попадает в таблицу при GuestCheckedIn, выходит при GuestCheckedOut или ReservationCancelled (после check-in отмены не бывает, но защищаемся на всякий).

// examples/ddd-hotel/effect/src/projections/occupancy.ts
export type OccupancyRow = {
  readonly reservationId: ReservationId;
  readonly room: RoomNumber;
  readonly guest: GuestEmail;
};

export type Occupancy = ReadonlyMap<ReservationId, OccupancyRow>;

Шаг 2 · Apply

import { Match } from 'effect';

type Pending = ReadonlyMap<ReservationId, OccupancyRow>;

const apply = (
  state: Occupancy,
  pending: Pending,
  item: { readonly streamId: ReservationId; readonly event: ReservationEvent },
): { readonly state: Occupancy; readonly pending: Pending } =>
  Match.value(item.event).pipe(
    Match.tag('ReservationPlaced', (e) => {
      const nextPending = new Map(pending);
      nextPending.set(item.streamId, {
        reservationId: item.streamId, room: e.room, guest: e.guest,
      });
      return { state, pending: nextPending };
    }),
    Match.tag('ReservationCancelled', () => ({
      state: removeKey(state, item.streamId),
      pending: removeKey(pending, item.streamId),
    })),
    Match.tag('GuestCheckedIn', () => {
      const row = pending.get(item.streamId);
      if (row === undefined) return { state, pending };
      const next = new Map(state);
      next.set(item.streamId, row);
      return { state: next, pending: removeKey(pending, item.streamId) };
    }),
    Match.tag('GuestCheckedOut', () => ({
      state: removeKey(state, item.streamId),
      pending,
    })),
    Match.exhaustive,
  );

Что тут происходит:

  • ReservationPlaced, кладём ещё-не-заехавшего гостя в pending. В основной state он попадёт после GuestCheckedIn.
  • ReservationCancelled, удаляем из обоих контейнеров (вдруг успели зачеркнуть).
  • GuestCheckedIn, перекладываем из pending в state. Гость теперь “в отеле”.
  • GuestCheckedOut, удаляем из state. Гость выехал.

pending это вспомогательная структура самой проекции, не часть публичного state. Это нормально: read-model для UI знает только Occupancy, внутренняя книга projection-а это её приватный мир.

Шаг 3 · Идемпотентность

Apply должна выдерживать повторное применение одного и того же события. Это критично, потому что:

  1. После рестарта daemon может начать читать с уже обработанного места.
  2. Сетевые retries могут продублировать NOTIFY.
  3. Repair-режим (пересоздать read-model с нуля) гоняет все события заново.

Проверь: повторное ReservationPlaced перезаписывает row в pending (один и тот же row). Повторное GuestCheckedIn ищет row в pending, не находит (уже перемещён в state), возвращает state без изменений. Все ветки идемпотентны.

Что взять с собой

  • Проекция это (state, event) -> state, то же по форме что evolve, другой state.
  • Внутренние структуры (pending) не часть публичного API read-model.
  • Каждая ветка обязана быть идемпотентной: повторное применение не ломает state.

Покрути проекцию руками

Прежде чем читать про daemon, прогони apply сам. Жми события у броней и смотри, как строка едет pending -> в отеле -> прочь. Повтори Placed или CheckIn дважды подряд, увидишь, как загорается флаг идемпотентности: state не изменился.

Раздел 3 · SubscriptionRef

Сцена

Ref это атомарная переменная: get, set, update. Этого хватает для shared state, но не хватает для live UI. Когда state меняется, как сказать всем подписчикам “обновитесь”? Через timer и сравнение версий? Через ручной EventEmitter? Effect даёт прямой ответ: SubscriptionRef.

Идея словами

SubscriptionRef<A> это Ref<A> плюс встроенный Stream изменений. Каждый SubscriptionRef.set(ref, value) атомарно обновляет state и публикует новое значение в SubscriptionRef.changes(ref). Подписчик читает SubscriptionRef.changes(ref) как обычный Stream и получает все обновления.

Шаг 1 · Создание

import { SubscriptionRef } from 'effect';

const ref = yield* SubscriptionRef.make<Occupancy>(new Map());

SubscriptionRef.make(initial) это эффект, который создаёт ref. Работает внутри Effect.gen.

Шаг 2 · Чтение и запись

// чтение текущего значения
const current = yield* SubscriptionRef.get(ref);

// атомарная замена
yield* SubscriptionRef.set(ref, newOccupancy);

// функциональное обновление
yield* SubscriptionRef.update(ref, (current) => apply(current, ...));

API идентично Ref, плюс дополнительно SubscriptionRef.changes(ref) это Stream.

Шаг 3 · Подписка

import { Stream } from 'effect';

const consumer = Stream.runForEach(SubscriptionRef.changes(ref), (occupancy) =>
  Effect.sync(() => console.log('новая read-model:', occupancy)),
);

SubscriptionRef.changes(ref) начинает с текущего значения, потом отдаёт каждое последующее set. Если поставить два set-а быстро подряд, подписчик гарантированно увидит последнее, но не обязательно промежуточные (стратегия “skip if behind”). Это правильное поведение для UI: важно показать актуальное, а не воспроизвести анимацию.

Что взять с собой

  • SubscriptionRef это Ref плюс Stream изменений.
  • Подписчик читает SubscriptionRef.changes(ref) без polling.
  • Если подписчик не успевает, он получает последнее значение, не все промежуточные.

Раздел 4 · Daemon под Effect.forkScoped

Сцена

Кто-то должен крутить цикл “прочитал событие -> применил к state -> обновил SubscriptionRef”. Это фоновый процесс, который должен жить, пока живёт владеющий Scope (например, всё время работы программы или время жизни одного агрегата).

Шаг 1 · forkScoped, не fork

В уроке 06 мы видели Effect.forkChild. Он привязывает файбер к родительскому файберу, и имя это прямо говорит: умирает родитель, умирает и ребёнок. Это не то, что нам нужно: родитель это конструктор Layer, он умирает сразу после return {...}. Daemon должен жить дольше.

Effect.forkScoped привязывает файбер к Scope из Context. Конструктор Layer.effect даёт нам такой Scope: файбер живёт, пока живёт Scope, и корректно прерывается при его закрытии. Отдельного Layer.scoped не существует, Layer.effect сам принимает эффект, которому нужен Scope, и убирает эту потребность из требований слоя.

Шаг 2 · Полная реализация

import { Effect, Ref, Stream, SubscriptionRef } from 'effect';

export const buildOccupancyProjection = Effect.gen(function* () {
  const store = yield* InMemoryEventStore;       // или PgEventStore
  const ref = yield* SubscriptionRef.make<Occupancy>(new Map());
  const pendingRef = yield* Ref.make<Pending>(new Map());

  const stream = yield* store.subscribe;  // Stream<AllEventsItem>

  yield* Effect.forkScoped(
    Stream.runForEach(stream, (item) =>
      Effect.gen(function* () {
        const pending = yield* Ref.get(pendingRef);
        const state = yield* SubscriptionRef.get(ref);
        const next = apply(state, pending, item);
        yield* SubscriptionRef.set(ref, next.state);
        yield* Ref.set(pendingRef, next.pending);
      }),
    ),
  );

  return ref;
});

Шаги:

  1. Получаем store и создаём refs.
  2. store.subscribe это Effect, который возвращает Stream<AllEventsItem>. Важно: подписка к PubSub устанавливается синхронно в момент yield, поэтому все события после этого момента гарантированно дойдут. Внутри store это две строки: PubSub.subscribe(pubsub) отдаёт Subscription, а Stream.fromSubscription(subscription) превращает подписку в стрим. Подписка это отдельный тип, а не очередь, поэтому и конструктор стрима для неё свой.
  3. Effect.forkScoped запускает фоновый файбер, который крутит Stream.runForEach. Для каждого события мы читаем state, применяем apply, пишем обратно.
  4. Возвращаем SubscriptionRef. Caller подпишется через SubscriptionRef.changes(ref).

Шаг 3 · Использование

const program = Effect.gen(function* () {
  const occupancy = yield* buildOccupancyProjection;

  // подписка
  yield* Effect.forkScoped(
    Stream.runForEach(SubscriptionRef.changes(occupancy), (current) =>
      Effect.sync(() => console.log('гостей сейчас:', current.size)),
    ),
  );

  // ... тем временем кто-то пишет события в store
}).pipe(Effect.scoped, Effect.provide(InMemoryEventStoreLive));

Effect.scoped создаёт корневой Scope. Когда программа завершается, Scope закрывается, и оба наших forkScoped-файбера корректно прерываются.

Что взять с собой

  • Effect.forkScoped привязывает файбер к Scope, не к родительскому файберу.
  • В Layer.effect или в Effect.gen плюс Effect.scoped это правильный выбор для daemon-а.
  • Daemon крутит Stream.runForEach, на каждом событии обновляет state.

Раздел 5 · Тест проекции

Сцена

Проекцию тестируем end-to-end: подкладываем InMemoryEventStore, append-аем события, читаем SubscriptionRef, проверяем содержимое. Никаких моков, реальная mini-инфраструктура.

Шаг 1 · Тестовый сценарий

// examples/ddd-hotel/effect/test/projections.test.ts
import { Effect, SubscriptionRef } from 'effect';
import { describe, expect, it } from 'vitest';

const placed: ReservationEvent = {
  _tag: 'ReservationPlaced', reservationId: id, guest, room, range, occurredAt: at,
};
const checkedIn: ReservationEvent = { _tag: 'GuestCheckedIn', reservationId: id, occurredAt: at };
const checkedOut: ReservationEvent = { _tag: 'GuestCheckedOut', reservationId: id, occurredAt: at };

const waitTick = Effect.sleep('20 millis');

describe('occupancy projection', () => {
  it('после CheckIn read-model содержит гостя, после CheckOut пустой', async () => {
    const program = Effect.gen(function* () {
      const ref = yield* buildOccupancyProjection;
      const store = yield* InMemoryEventStore;

      yield* store.append(id, 0, [placed, checkedIn]);
      yield* waitTick;
      const afterCheckIn = yield* SubscriptionRef.get(ref);
      expect(afterCheckIn.size).toBe(1);
      expect(afterCheckIn.get(id)?.room).toBe(room);

      yield* store.append(id, 2, [checkedOut]);
      yield* waitTick;
      const afterCheckOut = yield* SubscriptionRef.get(ref);
      expect(afterCheckOut.size).toBe(0);
    }).pipe(Effect.scoped, Effect.provide(InMemoryEventStoreLive));

    await Effect.runPromise(program);
  });
});

Заметь waitTick = Effect.sleep('20 millis'). Это не “плохой sleep в тесте”, а необходимая пауза на eventually-consistent системе. Append-аем мы из одного файбера, проекция читает из другого. Чтобы daemon успел обработать событие, нам нужно отдать ему квант времени. В реальной системе UI и так увидит обновление с задержкой.

Альтернатива sleep: подписаться на SubscriptionRef.changes(ref) и ждать первого значения с нужным размером. Это надёжнее, но многословнее. Для иллюстрации sleep понятнее.

Шаг 2 · Запуск

cd examples/ddd-hotel/effect
pnpm test

Если всё собрано правильно (decider, in-memory store, проекция), тест проходит за пятьдесят миллисекунд. Можешь повторить append с теми же событиями, проверить идемпотентность: после повторного ReservationPlaced row в pending перезатёрся, после повторного GuestCheckedIn state не изменился (потому что pending уже пуст).

Что взять с собой

  • Проекцию тестируем end-to-end с in-memory store.
  • Effect.sleep в тесте это пауза для daemon-а, не плохой паттерн.
  • Альтернатива: подписаться на SubscriptionRef.changes(ref) и ждать заданного значения.

Раздел 6 · Идемпотентность и last_seq

Сцена

В реальной системе daemon может рестартнуться (CI deploy, OOM, плановая перезагрузка). Что произойдёт с проекцией?

Два сценария:

  1. Naive: проекция стартует с пустой read-model, и каждое чтение начинается с id = 0. Это работает для маленьких систем, но на миллионе событий запуск превращается в часы.
  2. Snapshot: проекция периодически сохраняет state в Postgres вместе с last_seq. После рестарта читает snapshot, потом дочитывает события с id > last_seq.

Шаг 1 · Таблица projection_state

create table if not exists projection_state (
  name      text primary key,
  last_seq  bigint not null,
  state     jsonb not null,
  updated_at timestamptz not null default now()
);

name это имя проекции (occupancy, daily_occupancy, и так далее). last_seq это id последнего обработанного события из events. state это сама материализованная read-model.

Шаг 2 · Цикл с восстановлением

Псевдокод (полную реализацию делаешь в ДЗ):

const buildOccupancyProjection = Effect.gen(function* () {
  const sql = yield* PgClient.PgClient;
  const ref = yield* SubscriptionRef.make<Occupancy>(/* загрузить из projection_state.state */);

  let lastSeq = /* загрузить из projection_state.last_seq */;

  // catchup: прочитать всё, что прибавилось с момента сохранения
  const missed = yield* sql`select id, payload from events where id > ${lastSeq} order by id`;
  for (const row of missed) {
    /* apply, обновить state */
    lastSeq = row.id;
  }
  yield* SubscriptionRef.set(ref, newState);
  yield* saveState(name, lastSeq, newState);

  // live: подписаться на свежее
  yield* Effect.forkScoped(
    /* polling каждые 100ms или LISTEN/NOTIFY */
    Effect.gen(function* () {
      const newEvents = yield* sql`select id, payload from events where id > ${lastSeq} order by id`;
      for (const row of newEvents) {
        /* apply, обновить ref, persist */
      }
    }).pipe(Effect.repeat({ schedule: Schedule.spaced('100 millis') })),
  );

  return ref;
});

Идемпотентность тут критична: при рестарте мы можем дважды применить событие на грани last_seq, и проекция не должна сломаться.

Что взять с собой

  • Без persistence projection-state проекция начинает с нуля при каждом рестарте.
  • projection_state(name, last_seq, state) сохраняет место и снимок.
  • Идемпотентность apply страхует от дублирующего применения на границе catchup.

Раздел 7 · Polling vs LISTEN/NOTIFY

Сцена

Чтобы daemon знал, что в events появилась новая строка, есть два пути.

Polling

SELECT ... WHERE id > $last_seq каждые 100ms. Просто, переносимо, работает на любой БД. Минус: лишние запросы при отсутствии нагрузки, и задержка до 100ms на каждое событие.

LISTEN/NOTIFY

Postgres-специфичный механизм. Триггер на INSERT INTO events шлёт NOTIFY events_channel. Daemon делает LISTEN events_channel и получает push-уведомление мгновенно.

create or replace function notify_events() returns trigger as $$
begin
  perform pg_notify('events_channel', new.stream_id || ':' || new.id);
  return new;
end;
$$ language plpgsql;

create trigger events_notify after insert on events
for each row execute function notify_events();

Плюс: нулевая задержка. Минус: NOTIFY ограничен 8000 байт payload, требует постоянного соединения, не работает за PgBouncer в transaction-mode.

Гибрид

Хороший production-паттерн: LISTEN для уведомления “появилось новое”, catchup через SELECT. NOTIFY это сигнал (“проверь events”), не сами данные. Это снимает ограничения payload и работает корректно при пропусках уведомлений.

В нашем уроке мы оставляем in-memory PubSub для in-memory store и polling для Postgres. LISTEN/NOTIFY это ДЗ для тех, кому хочется глубже.

Что взять с собой

  • Polling прост, переносим, имеет задержку.
  • LISTEN/NOTIFY мгновенный, но Postgres-специфичный.
  • Гибрид: NOTIFY это сигнал, SELECT это данные. Дополнительно страхует от пропусков.

Чек-лист

Чек-листготово

ДЗ

Финал блока

Этот урок закрывает четвёрку Effect.ts-трека по functional DDD. По карте:

У тебя в examples/ddd-hotel/effect/ сейчас лежит полный мини-домен Hotel Booking: пять состояний, четыре команды, четыре события, Event Store, проекция. Около 350 строк прикладного кода, и тесты, которые гоняются за полсекунды.

Что дальше: тот же домен в Gleam-треке, потом в Haskell и Rust. После четырёх треков в 16-arch появится один сводный урок, где мы сравним четыре языка на одном и том же Decider-е и поймём, где какой язык удобнее, где проигрывает, и что общего во всех вариантах ФП-DDD.

Дальше

Раздел продолжается блоком прод-обвязки, и он начинается с 18 · Observability: логи, спаны, метрики и экспорт наружу. Дальше HttpClient и API, конфигурация и секреты, Schema вглубь, приёмники потоков, стандартная библиотека и инструменты.

Полезно перечитать:

  • 09-stream, про Stream и Sink. Проекции это classic stream-processing, а приёмники разбираются в 23 · Sink и продвинутый Stream.
  • 11-runtime-and-schedule, про Schedule.spaced и Schedule.exponential. Polling-цикл это Effect.repeat(spaced(100 millis)).
  • Greg Young, “8 Lines of Code”. Видео, в котором за десять минут показывается, как CQRS убирает бойлерплейт.