Event Store на @effect/sql-pg: append-only, optimistic concurrency
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Event Store на @effect/sql-pg
Сцена · события первого класса
В предыдущем уроке Decider жил в массиве. Его легко тестировать, нечего сохранять. В реальной системе события нужно куда-то класть, и место это особенное: append-only, никаких update, никаких delete. История фактов не переписывается; если факт оказался “ошибочным”, его компенсирует новое событие, а не “отмена старого”.
Это и есть Event Sourcing в чистом виде: вместо снимка состояния храним полный лог фактов. State это производное от лога; его всегда можно пересчитать. Эта инверсия меняет три вещи:
- Любое прошлое восстанавливается. Spotify хочет увидеть, что было в корзине пользователя на момент
2026-01-15T12:00. Replay до этой временной точки, и состояние ровно то, что было. - Аудит бесплатен. Каждая мутация это событие, у каждого события есть
occurredAt. Регулятор спрашивает “почему этот гость попал в check-in”, читаем стрим. - Новые read-models дёшевы. Появилась задача “статистика загрузки по этажам”, не надо мигрировать БД, прокручиваем стрим существующих событий и материализуем новый view.
В этом уроке мы складываем ReservationEvent в Postgres через @effect/sql-pg. На выходе: типизированный EventStore-сервис, оптимистичная конкуренция, и handler, который связывает Event Store с decider-ом из урока 15.
Карта урока · что заберёшь домой
- Схема events. Какие колонки нужны и почему. Уникальный ключ для оптимистики.
- @effect/sql-pg, базовый Layer. Чем
PgClient.layerотличается отPgClient.layerConfig. - EventStore сервис. Контракт: load, append, ConcurrencyConflict.
- load: SELECT по stream_id, decode payload. Schema на чтении.
- append: INSERT с уникальным версионным ключом. Нарушение уникальности как сигнал, и как его прочитать типизированно.
- commandHandler: load -> replay -> decide -> append. Полный цикл для одной команды.
- Что не покрыто: snapshot, idempotency, projection (урок 17).
Раздел 1 · Схема events
Сцена
Под Event Store нужна одна таблица. Не двадцать таблиц на каждый тип события, не нормализованная схема “reservations + check_ins + cancellations”. Одна. Это контринтуитивно, если ты приходишь из CRUD-мира, но это и есть суть подхода: события самодостаточны, их форма выводится из type плюс payload.
Шаг 1 · DDL
create table if not exists events (
id bigserial primary key,
stream_id text not null,
version integer not null,
type text not null,
payload jsonb not null,
created_at timestamptz not null default now(),
unique (stream_id, version)
);
create index if not exists events_stream_id_idx on events (stream_id, version);
create index if not exists events_id_idx on events (id);
Разбор колонок:
idэто глобальный монотонный счётчик. По нему проекции (урок 17) читают поток: “дай мне все события сid > last_seen”.stream_idэто id агрегата. Для нашего домена этоReservationId.versionэто локальный счётчик внутри одногоstream_id. Первая запись имеетversion = 1, втораяversion = 2, и так далее.typeэто_tagсобытия (ReservationPlaced,GuestCheckedIn, …). Полезно для индексирования и фильтрации в SQL без декода payload.payloadэто сам факт вjsonb. Что внутри, целиком определяется типом и формойReservationEvent.created_atэто серверное время вставки. Не то же, чтоoccurredAtв payload (которое могло быть установлено decider-ом в момент решения).
Шаг 2 · Почему unique (stream_id, version)
Это главный механизм optimistic concurrency. Когда два процесса одновременно пытаются записать второе событие в один и тот же стрим, оба хотят version = 2. Один успевает первым, второй получает unique violation. БД сама гарантирует, что одна и та же версия не запишется дважды. Никакой ручной блокировки, никакого SELECT ... FOR UPDATE.
В классической CRUD-таблице такой инвариант пришлось бы держать на уровне приложения через явный lock. В event store он живёт в схеме и стоит копеек по производительности.
Шаг 3 · Почему bigserial id
Глобальный id нужен проекциям. Когда мы строим read-model, мы хотим читать события в порядке вставки независимо от stream_id. SELECT * FROM events WHERE id > $1 ORDER BY id это однородный indexed scan, который работает для любой read-model. Без id пришлось бы выдумывать комбинированный курсор (created_at, stream_id, version), и каждый пропуск был бы дорогой операцией.
Что взять с собой
- Event Store это одна таблица, не семейство.
unique (stream_id, version)даёт optimistic concurrency бесплатно.id bigserialнужен проекциям для глобального порядка чтения.payloadэтоjsonb; схема живёт в коде, а не в БД.
Раздел 2 · @effect/sql-pg, базовый Layer
Сцена
В основном проекте BatSchool мы используем @effect/sql-pg, а drizzle держим рядом как чистый сборщик запросов. В examples-проекте мы возьмём только @effect/sql-pg, чтобы показать минимальный путь без ORM. Это та же машинерия, что используется в our-db/effect.
Ещё одна деталь про экосистему. Ядро Effect SQL (SqlClient, SqlError, шаблоны запросов) живёт прямо внутри пакета effect, в точке входа effect/unstable/sql. Отдельными пакетами остались только драйверы под конкретную СУБД, и @effect/sql-pg как раз такой драйвер.
Шаг 1 · Layer для PgClient
// examples/ddd-hotel/effect/src/event-store/layer.ts
import { PgClient } from '@effect/sql-pg';
import { Config } from 'effect';
export const PgLive = PgClient.layerConfig({
url: Config.redacted('DATABASE_URL'),
});
Здесь два родственных конструктора, и выбирать между ними надо осознанно:
PgClient.layer(config)берёт готовые значения (url: Redacted.make(process.env.DATABASE_URL)) и отдаётLayer<PgClient | SqlClient, SqlError>.PgClient.layerConfig(config)берёт описания значений (Config) и сам их читает на старте, добавляяConfigErrorв канал ошибок слоя.
Второй вариант удобнее, когда строка подключения приходит из окружения. Config.redacted('DATABASE_URL') это механизм Config из Effect: читает переменную окружения, валидирует на старте и заворачивает значение в Redacted, чтобы пароль не утёк в лог при печати объекта. Если переменной нет, программа падает на инициализации с понятной ошибкой.
Заметь ещё, что слой отдаёт сразу два сервиса: PgClient (специфика Postgres) и SqlClient (общий интерфейс). Код, которому хватает обычных запросов, зависит от SqlClient и остаётся переносимым на другую СУБД.
Шаг 2 · Скрипт инициализации схемы
В тестах и dev-сборке удобно создать таблицу один раз при старте программы:
// examples/ddd-hotel/effect/src/event-store/pg.ts
export const SCHEMA_DDL = `
create table if not exists events (
id bigserial primary key,
stream_id text not null,
version integer not null,
type text not null,
payload jsonb not null,
created_at timestamptz not null default now(),
unique (stream_id, version)
);
create index if not exists events_stream_id_idx on events (stream_id, version);
create index if not exists events_id_idx on events (id);
`;
Запуск:
import { Effect } from 'effect';
import { PgClient } from '@effect/sql-pg';
const setup = Effect.gen(function* () {
const sql = yield* PgClient.PgClient;
yield* sql.unsafe(SCHEMA_DDL);
});
Effect.runPromise(setup.pipe(Effect.provide(PgLive)));
sql.unsafe(text) это escape hatch: выполнить произвольный SQL без интерполяции. В production миграции мы кладём в отдельный механизм (drizzle-kit, raw-migrate), но для examples-проекта одна точка сборки удобнее.
Что взять с собой
PgClient.layerConfig({ url: Config.redacted(...) })даёт типобезопасный Layer для Postgres, аPgClient.layer(...)тот же слой из уже готовых значений.Configвалидирует env на старте, не в runtime.- Ядро Effect SQL лежит в
effect/unstable/sql, отдельный пакет нужен только под драйвер. sql.unsafe(text)для DDL и других “немодельных” запросов.
Раздел 3 · Контракт EventStore
Сцена
Перед реализацией договоримся об интерфейсе. EventStore это два метода: load (получить полный стрим событий) и append (записать новые события с проверкой версии). Опциональный третий, subscribe, будет в уроке 17, но он не часть основного контракта Event Store.
Шаг 1 · Тип
// examples/ddd-hotel/effect/src/event-store/types.ts
import type { Effect } from 'effect';
import { ConcurrencyConflict } from '../domain/errors.ts';
import type { ReservationEvent } from '../domain/events.ts';
import type { ReservationId } from '../domain/types.ts';
export type StreamSlice = {
readonly version: number;
readonly events: ReadonlyArray<ReservationEvent>;
};
export type EventStore = {
readonly load: (streamId: ReservationId) => Effect.Effect<StreamSlice>;
readonly append: (
streamId: ReservationId,
expectedVersion: number,
events: ReadonlyArray<ReservationEvent>,
) => Effect.Effect<void, ConcurrencyConflict>;
};
Разберём по полям:
load(streamId)возвращаетStreamSlice, пару из текущей версии и массива событий. Пустой стрим даёт{ version: 0, events: [] }.append(streamId, expectedVersion, events).expectedVersionэто та версия, которую видел писатель приload. Если в БД уже больше (кто-то записал в фоне), append падает сConcurrencyConflict. Если меньше или равно, события записываются, и БД сама присвоит имexpectedVersion + 1,expectedVersion + 2, …
ConcurrencyConflict это уже знакомая нам tagged-ошибка:
export class ConcurrencyConflict extends Data.TaggedError('ConcurrencyConflict')<{
readonly reservationId: ReservationId;
readonly expected: number;
readonly actual: number;
}> {}
Что взять с собой
- EventStore это два метода:
loadиappend. expectedVersionэто контракт optimistic concurrency.ConcurrencyConflictэто специфичная tagged-ошибка, видна в сигнатуре.
Раздел 4 · load · SELECT и decode
Сцена
load читает все события стрима и декодит payload через Schema. На большом стриме это становится медленно (тысячи декодов), и тогда подключается snapshot. Сегодня обходимся без snapshot.
Шаг 1 · Реализация
// examples/ddd-hotel/effect/src/event-store/pg.ts
import { PgClient } from '@effect/sql-pg';
import { Context, Effect, Layer, Schema } from 'effect';
import { ConcurrencyConflict } from '../domain/errors.ts';
import { ReservationEvent } from '../domain/events.ts';
import type { ReservationId } from '../domain/types.ts';
import type { EventStore, StreamSlice } from './types.ts';
export class PgEventStore extends Context.Service<PgEventStore, EventStore>()(
'ddd-hotel/PgEventStore',
) {}
const decodeEvent = Schema.decodeUnknownEffect(ReservationEvent);
export const PgEventStoreLive = Layer.effect(
PgEventStore,
Effect.gen(function* () {
const sql = yield* PgClient.PgClient;
const load = (streamId: ReservationId): Effect.Effect<StreamSlice> =>
Effect.gen(function* () {
const rows = yield* sql<{
readonly version: number;
readonly payload: unknown;
}>`select version, payload from events
where stream_id = ${streamId}
order by version asc`.pipe(Effect.orDie);
const events: ReservationEvent[] = [];
let version = 0;
for (const row of rows) {
const parsed = yield* decodeEvent(row.payload).pipe(Effect.orDie);
events.push(parsed);
version = row.version;
}
return { version, events } satisfies StreamSlice;
});
// ... append дальше
return { load, /* append, */ };
}),
);
Заметные моменты:
Context.Service<Self, Shape>()('id')объявляет тег сервиса: сначала типы, потом строковый идентификатор. Класс-наследник и есть тег, его же мы дальшеyield*-им, чтобы достать реализацию.sql<RowType>`...`это template literal с типизацией строки результата. Параметры (${streamId}) автоматически параметризуются, никакой инъекции.Effect.orDieпревращает ошибки БД в defect. Это сознательно: ошибка БД на read это инфраструктурный сбой, не доменная ошибка. Если БД упала, разработчик должен увидеть defect и поднять алерт, а не переводить это в “что-то невалидное”.decodeEvent(row.payload)валидирует payload по Schema-у. Если в БД случайно лежит payload неизвестного_tag(например, после неудачной миграции), мы поймём это сразу при чтении.
Шаг 2 · Что если payload устарел
Допустим, ты выпустил версию 1.0 с событием GuestCheckedIn { reservationId }, потом версия 1.1 добавила roomNumber в payload. Старые события в БД не имеют roomNumber. Decode по новой Schema упадёт.
Решений два:
- Backwards-compatible изменения: добавляй только optional поля в Schema (
Schema.optional, а если нужно подставить значение при декоде,Schema.withDecodingDefaultType). Тогда старые payload-ы декодятся, новые тоже. Это дефолтная стратегия. - Версия события: добавь поле
versionв payload и сделай Union по версии. Это для редких ломающих миграций.
Этой темы (event schema evolution) хватает на целый отдельный курс, мы её не разбираем глубоко. Главное правило: поле добавляй optional, никогда не меняй смысл существующего поля, никогда не удаляй.
Что взять с собой
sql<RowType>это типизированный template, параметризованный.Effect.orDieпревращает БД-ошибки в defect (инфраструктурный сбой, не expected).- Schema decode на чтении ловит устаревшие payload-ы. Эволюцию схемы делаем через optional поля.
Раздел 5 · append · INSERT с unique violation
Сцена
Append это главная ответственность Event Store. Здесь живёт optimistic concurrency, потому что одновременных писателей в один стрим может быть несколько.
Шаг 1 · Реализация
import type { SqlError } from 'effect/unstable/sql';
const isUniqueViolation = (error: SqlError.SqlError): boolean =>
error.reason._tag === 'UniqueViolation';
const append = (
streamId: ReservationId,
expectedVersion: number,
events: ReadonlyArray<ReservationEvent>,
): Effect.Effect<void, ConcurrencyConflict> =>
Effect.gen(function* () {
if (events.length === 0) return;
const rows = events.map((event, index) => ({
stream_id: streamId,
version: expectedVersion + index + 1,
type: event._tag,
payload: event as unknown,
}));
yield* sql`insert into events ${sql.insert(rows)}`.pipe(
Effect.catch((error) =>
isUniqueViolation(error)
? Effect.fail(
new ConcurrencyConflict({
reservationId: streamId,
expected: expectedVersion,
actual: expectedVersion + events.length,
}),
)
: Effect.die(error),
),
);
});
Шаги:
- Если событий нет, ничего не делаем (no-op append законен).
- Готовим rows: каждое событие получает version =
expectedVersion + index + 1. Если в стриме уже 5 событий (expectedVersion=5), новые получат 6, 7, 8. sql.insert(rows)это helper, генерирующийINSERT INTO events (...) VALUES (...), (...), (...)с правильным escape.- Главный трюк:
Effect.catchловит ошибку из каналаE, разбирает её причину и превращает нарушение уникальности вConcurrencyConflict. Все остальные ошибки (соединение, синтаксис, диск) идут вEffect.dieкак defect.
Шаг 2 · Почему причина ошибки, а не код драйвера
Присмотрись к isUniqueViolation. Мы не заглядываем в сырую ошибку драйвера и не сверяем строковый код. В канал ошибки прилетает SqlError, и у него есть поле reason, это размеченное объединение причин: UniqueViolation, ConstraintError, DeadlockError, SerializationError, SqlSyntaxError, ConnectionError и ещё несколько. Драйвер уже разобрал свой код возврата и назвал причину словом.
Это тот случай, когда стоит остановиться и посмотреть на приём, а не только на строчку кода. Сравни две проверки:
// плохо: гадаем по коду драйвера, ошибка типизирована как unknown
const byCode = (cause: unknown): boolean =>
(cause as { code?: unknown }).code === '23505';
// хорошо: спрашиваем у типа, компилятор знает все варианты
const byReason = (error: SqlError.SqlError): boolean =>
error.reason._tag === 'UniqueViolation';
Первая версия опирается на знание, которого нет в типах: 23505 это код unique violation в Postgres, и надо помнить его наизусть. Опечатался в цифре, поменял СУБД, обернул ошибку по дороге, и проверка молча начинает возвращать false. Ничего не сломается громко: optimistic concurrency просто перестанет работать, конфликт вместо ConcurrencyConflict улетит в defect, а Effect.retry его не поймает, потому что defect не ошибка канала E.
Вторая версия спрашивает у типа. Опечатку в 'UniqueViolation' компилятор подсветит сразу, потому что список тегов конечен и известен. Плюс у самой причины есть полезные поля: у UniqueViolation лежит constraint с именем нарушенного ограничения, а SqlError.isRetryable сразу отвечает, имеет ли смысл повтор.
Мораль шире одного драйвера: если библиотека уже разметила свои причины отказа типами, не разбирай их руками по строкам. Строковая проверка компилируется всегда, но проверить её может только рантайм и только на живой базе.
Шаг 3 · Что значит expected vs actual
В ConcurrencyConflict мы пишем expected (что писатель видел) и actual (что в БД сейчас, точнее, что должно быть после нашего append). actual это не точно то, что в БД (там может быть и больше, если за время нашей операции кто-то ещё записал). Но это разумный сигнал для вызывающего: “ты пытался писать с версии X, а БД уже впереди”. На практике вызывающий сделает retry: load заново, replay, decide, append.
Переключи таб на “по очереди” и увидишь контраст: те же два писателя, но без перекрытия во времени конфликта нет, версии растут 1, 2, 3. Конфликт рождается ровно тогда, когда оба прочитали одну версию и целят одно число.
Шаг 4 · Транзакционность
Один INSERT это одна транзакция Postgres. Если у нас несколько событий в одной команде (PlaceReservation плюс RoomLocked), они вставляются атомарно: либо все, либо ни одного. Это и есть то, что нам нужно: одна команда даёт согласованный набор фактов, частичных записей не бывает.
Если нужна транзакция на несколько разных стримов (например, перенести бронь с одного агрегата на другой), мы оборачиваем оба append в sql.withTransaction. В сегодняшнем уроке такого нет.
Что взять с собой
appendверсионирует события сам, на основеexpectedVersion.- Нарушение уникальности превращаем в
ConcurrencyConflict, узнаём его поerror.reason._tag === 'UniqueViolation'. - Причины отказа SQL размечены типами, поэтому проверку по коду драйвера писать не нужно и вредно.
- Один INSERT, одна транзакция. Несколько событий на команду = атомарно.
Раздел 6 · commandHandler
Сцена
Decider и Event Store собираются в одну точку: command handler. Он принимает команду, грузит стрим, replay-ит state, дёргает decide, складывает события через append. Это полный цикл одной операции.
Шаг 1 · Реализация
// examples/ddd-hotel/effect/src/event-store/handler.ts
import { Effect } from 'effect';
import type { ReservationCommand } from '../domain/commands.ts';
import { reservationDecider } from '../domain/decider.ts';
import type { DomainError } from '../domain/errors.ts';
import type { ReservationId } from '../domain/types.ts';
import { PgEventStore } from './pg.ts';
export const handleCommand = (streamId: ReservationId, command: ReservationCommand) =>
Effect.gen(function* () {
const store = yield* PgEventStore;
const slice = yield* store.load(streamId);
const state = slice.events.reduce(reservationDecider.evolve, reservationDecider.initial);
const events = yield* reservationDecider.decide(state, command);
yield* store.append(streamId, slice.version, events);
return events;
});
Этот код читается как короткий рассказ:
- load, прочитали текущий стрим.
- replay, восстановили состояние.
- decide, спросили decider, что делать.
- append, записали новые события с правильной версией.
- Вернули список событий (полезно вызывающему: HTTP-обработчик отправит их в WebSocket или SSE).
Тип:
Effect<ReadonlyArray<ReservationEvent>, DomainError, PgEventStore>
В канале E стоит DomainError (union из InvalidStateTransition, ConcurrencyConflict и так далее). Это уровень, на котором HTTP-обработчик пишет ответ:
.catchTag('InvalidStateTransition', (e) => Effect.succeed(toHttp422(e)))
.catchTag('ConcurrencyConflict', (e) => Effect.succeed(toHttp409(e)))
Шаг 2 · Retry на conflict
В реальности при ConcurrencyConflict мы хотим перезагрузить стрим и попробовать заново. Эта политика выражается через Effect.retry:
import { Schedule } from 'effect';
const handleWithRetry = (streamId: ReservationId, command: ReservationCommand) =>
handleCommand(streamId, command).pipe(
Effect.retry({
schedule: Schedule.exponential('20 millis'),
times: 5,
while: (e) => e._tag === 'ConcurrencyConflict',
}),
);
Ограничение числа попыток это отдельное поле times, а не второе расписание, скрещенное с первым. Расписание отвечает на вопрос “как долго ждать между попытками”, times на вопрос “сколько раз пробовать”, и разводить их по разным полям честнее, чем комбинировать два расписания оператором пересечения.
Шесть попыток, экспоненциальный backoff. Если все шесть провалились, выдаём ConcurrencyConflict наверх. На практике этого хватает для всех реалистичных сценариев multi-writer.
Шаг 3 · Где живёт command handler
В нашем examples-проекте это Effect-функция. В реальном приложении она будет вызываться из HTTP-ручки:
// в каком-нибудь Astro API route или HttpApi эндпоинте
export const POST = (request) =>
Effect.gen(function* () {
const command = yield* Schema.decodeUnknownEffect(ReservationCommand)(yield* request.json());
const events = yield* handleCommand(command.reservationId, command);
return Response.json({ events }, { status: 201 });
});
Это и есть весь backend для одного агрегата. Никаких репозиториев, никакого ORM-маппинга, никакой логики в HTTP-слое. Слой ровно один: чистый decider плюс event store.
Что взять с собой
commandHandler = load -> replay -> decide -> append. Пять строк, цикл одной команды.- На
ConcurrencyConflictоборачиваем вEffect.retryс exponential backoff. - HTTP-слой это тонкий перевод между HTTP и Command, бизнеса в нём нет.
Раздел 7 · Что не покрыто
В этом уроке мы не делаем:
- Snapshot. Если у одного стрима тысячи событий, replay тормозит. Snapshot это периодически сохраняемый state, чтобы load начинал не с initial, а с последнего snapshot. Опциональное ДЗ в этом уроке.
- Idempotency. При retry на сетевой ошибке клиент может прислать одну и ту же команду дважды. Без защиты получим два события
ReservationPlaced. Решается полемidempotency_keyв командах и таблицей processed_commands. Тоже ДЗ. - Cross-aggregate transactions. Перенос состояния между двумя стримами одной транзакцией. Решается через
sql.withTransaction, но это редкая ситуация: обычно cross-aggregate операции делаются через saga, а не одной транзакцией. - Подписка на стрим. Read-models нужно строить из тех же событий, и эта тема целиком вынесена в урок 17.
Чек-лист
ДЗ
Дальше
Следующий урок · 17. CQRS-проекции: Daemon и SubscriptionRef. Берём поток событий из Event Store, материализуем read-model для UI. Daemon-файбер живёт в фоне, SubscriptionRef даёт live-подписку на изменения. К концу урока 17 у тебя в Hotel Booking появится “список текущих гостей в отеле”, обновляемый без polling-а.
Параллельно полезно перечитать:
- 05-resources, про Scope и acquireRelease. PgClient это classic scoped resource, открывается на старте, закрывается на SIGTERM.
- 11-runtime-and-schedule, про Schedule и retry.
Effect.retry({ schedule, while })это идиома, которую мы использовали для оптимистического conflict-handling.