Архитектура бэкенда на Effect: слайсы, очередь на Postgres, outbox, воркер
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Архитектура бэкенда на Effect: слайсы, очередь на Postgres, outbox, воркер
Сцена · тебе заказали платёжный шлюз
Двадцать пять уроков ты собирал Pulse: сервисы, слои, стримы, телеметрия, HTTP. Каждый примитив ты видел на маленьком, честном примере. Сегодня другой жанр: разбор целого боевого сервиса, где эти примитивы встречаются друг с другом и начинают спорить.
Кейс настоящий. billing-service, открытый платёжный шлюз для подписок: чекаут, вебхуки платёжного провайдера, собственный цикл списаний с ретраями и события наружу, в CRM. Требования к нему формулируются жёстче, чем к мониторингу адресов:
- Ни один платёж не теряется. Упал процесс, оборвалась сеть, база перезапустилась: после рестарта состояние полное.
- Дубликат вебхука не списывает деньги дважды. Провайдер шлёт колбэк повторно, это норма, а не авария.
- Событие об оплате доходит до CRM за минуту. Иначе автоматика на той стороне успеет отработать по устаревшему состоянию и отключит оплатившему доступ.
- Новый провайдер или новый получатель событий не трогает ядро. Подключение криптоплатежей не должно превращаться в правки по всем модулям.
Знакомый рефлекс, роутер плюс ORM плюс пара setInterval, здесь не выдерживает ни одного пункта. Не потому что технологии плохие, а потому что каждый пункт это архитектурное решение, которое надо принять явно. Урок про эти решения: как нарезать код на модули, где провести границы, как подружить Effect с чужим фреймворком, как построить очередь и outbox на голом Postgres и почему воркер живёт отдельным процессом.
Оригинальный сервис написан на Effect 3 и держит HTTP на Fastify. Мы читаем его архитектуру, а код в листингах переводим на идиомы v4, которыми ты пользуешься весь раздел. Разница косметическая: Layer, ManagedRuntime, tagged-ошибки и SQL-шаблоны работают одинаково, а архитектурные решения от версии не зависят вовсе.
Карта урока · что заберёшь домой
Восемь разделов:
- Раздел 1, вертикальные слайсы внутри горизонтальных границ. Как нарезать сервис, чтобы каждая итерация давала ценность.
- Раздел 2, типы держат границы: домен принимает команды, ошибки несут свой HTTP-маппинг, роут объявляется в одном месте, а SQL типизирован kysely-билдером с типами из живой базы.
- Раздел 3, сервис или функция. Когда нужен
Layer, а когда достаточно передать функцию в composition root. - Раздел 4, один рантайм на процесс:
ManagedRuntimeкак мост в не-Effect-фреймворк, тонкие bin-точки входа, воркер отдельным процессом и REPL к живому графу. - Раздел 5, надёжная очередь на голом Postgres: три слоя идемпотентности,
FOR UPDATE SKIP LOCKED, reaper. - Раздел 6, transactional outbox: событие и доставки одной транзакцией, at-least-once до каждого получателя.
- Раздел 7, собственный биллинг-цикл: шедулер, лестница ретраев, якорь даты против дрейфа.
- Раздел 8, границы с внешним миром: модуль провайдера, терпимый парсинг, карантин вместо потери данных.
К концу урока у тебя будет рабочий чертёж бэкенда на Effect, который ты сможешь приложить к любому своему сервису.
Раздел 1 · Вертикальные слайсы внутри горизонтальных границ
Две оси, которые все путают
Про слои ты слышал: api, application, domain, infrastructure, зависимости направлены внутрь, домен не знает про Postgres. Это ось зависимостей. Про слайсы тоже слышал: фича режется от кнопки до базы, каждая итерация даёт видимый результат. Это ось поставки.
Классическая ошибка, выбрать одну ось и забыть вторую. Команда, живущая только слоями, месяц пишет репозитории, потом месяц юзкейсы, и первый работающий сценарий появляется к концу квартала, когда половина требований уже устарела. Команда, живущая только слайсами, быстро получает демо, а через полгода обнаруживает, что бизнес-правила размазаны по роутам и SQL-запросам, и изменить политику ретраев значит перечитать весь проект.
billing-service явно фиксирует ответ: слайсы как модель поставки, слои как модель зависимостей. Каждый значимый сценарий (создание чекаута, колбэк провайдера, продление подписки) реализуется целиком, от HTTP до базы и событий. Но внутри каждого слайса действуют одни и те же горизонтальные границы, и направление зависимостей всегда внутрь: роут зависит от юзкейса, юзкейс от доменных правил, и никогда наоборот.
Вертикальный слайс отвечает на вопрос, что делать следующим. Горизонтальная граница отвечает на вопрос, кому можно импортировать кого.
Модуль как единица кода
На диске это выглядит буднично. Каждый слайс живёт в своём модуле с одинаковой анатомией:
src/modules/checkout/
routes.ts транспорт: HTTP внутрь и наружу
domain.ts правила: чистые функции и юзкейсы
data-access.ts хранение: SQL за интерфейсом
contracts.ts формы, которые нельзя выносить наружу
test/ юниты рядом с кодом
src/infra/
db/ общий доступ к Postgres
queue/ очередь (Раздел 5)
http/ роут-DSL и маппинг ошибок (Раздел 2)
Правила простые и скучные, в этом их сила:
routes.tsзнает про HTTP и не знает про SQL.domain.tsне импортирует ни HTTP, ни SQL: только типы команд, доменные ошибки и интерфейсы хранилищ.data-access.tsреализует интерфейсы домена поверхSqlClientи не содержит правил: если в SQL-функции появилсяifпро бизнес, он переехал не туда.- Общая инфраструктура (маппинг ошибок Postgres, роут-DSL, очередь) живёт в
infra/один раз, а не копипастой по модулям.
Сравни с Pulse: там масштаб позволял держать всё в паре файлов. Здесь модулей девять, и одинаковая анатомия означает, что открыв любой из них, ты заранее знаешь, где искать правило, где запрос, а где транспорт.
Правило, которое дороже всех: база не владеет правилами
Самое коварное следствие подхода роутер плюс ORM: схема базы незаметно становится доменной моделью, а правила расползаются по запросам. Политика я вижу дубликат, значит не списываю живёт в одном UPDATE, политика ретраев в другом, и связать их некому.
В billing-service это запрещено явно: изменение платёжного состояния всегда проходит через доменные функции, а провайдерские колбэки входят через юзкейс, а не через прямой UPDATE таблицы. База хранит факты и ограничения целостности (уникальные индексы нам ещё пригодятся), но решения принимает домен.
Раздел 2 · Типы держат границы
Границы из Раздела 1 легко провозгласить и легко потерять: через полгода кто-нибудь протащит объект запроса в домен, и слои снова срастутся. Надёжнее, когда границу охраняет компилятор.
Домен принимает команды, а не транспорт
Первое правило: доменная функция никогда не видит HTTP-запрос, заголовки или тело в сыром виде. Роут декодирует вход схемой в команду, плоский доменный тип, и передаёт её:
import { Schema } from 'effect';
export const CreateCheckout = Schema.Struct({
externalUserId: Schema.String,
amount: Schema.Int,
currency: Schema.Literals(['UAH', 'USD']),
period: Schema.Literals(['month', 'year']),
});
export type CreateCheckout = typeof CreateCheckout.Type;
export const createCheckout = (repo: CheckoutRepo) => (command: CreateCheckout) =>
Effect.gen(function* () {
// только правила: никакого request, headers, reply
});
Выгода не только в чистоте. Команду можно построить откуда угодно: из HTTP, из CLI, из обработчика очереди, из теста. Домен, принимающий команды, тестируется без единого мока транспорта, ты видел этот приём в 13 · Testing.
Сюда же относится авторизация. Проверка роли это доменное правило, а не транспортное: транспорт достаёт учётные данные и превращает их в значение Actor, а requireRole(actor, 'operator') живёт в домене и падает типизированной ошибкой Forbidden. Перенеси проверку в роут, и обработчик очереди, зовущий тот же юзкейс, пройдёт мимо неё.
Ошибки несут свой HTTP-маппинг
Ты уже строишь семейства ошибок через Data.TaggedError с 03 · Tagged-ошибки. Вопрос масштаба: когда роутов тридцать, где живёт таблица ошибка, статус? Вариант со switch в каждом роуте умирает первым же рефакторингом.
Решение billing-service: ошибка сама знает свой статус.
import { Data } from 'effect';
export interface HttpReply {
readonly status: number;
readonly body: unknown;
}
export class NotFound extends Data.TaggedError('NotFound')<{ entity: string }> {
toHttp(): HttpReply {
return { status: 404, body: { error: `${this.entity} not found` } };
}
}
export class Forbidden extends Data.TaggedError('Forbidden')<{}> {
toHttp(): HttpReply {
return { status: 403, body: { error: 'forbidden' } };
}
}
const hasToHttp = (e: unknown): e is { toHttp: () => HttpReply } =>
typeof e === 'object' && e !== null && 'toHttp' in e;
export const toHttp = (error: unknown): HttpReply =>
hasToHttp(error) ? error.toHttp() : { status: 500, body: { error: 'internal error' } };
Одна функция toHttp в транспортном слое, и любая непредусмотренная ошибка (SqlError, дефект, забытый кейс) честно схлопывается в 500, а не протекает наружу стектрейсом. Добавил новую доменную ошибку, дописал метод, и все роуты автоматически отвечают правильно.
Роут объявляется в одном месте
Остался бойлерплейт роутов: декодировать вход, вызвать юзкейс, закодировать выход, замапить ошибки. Тридцать роутов, тридцать одинаковых обвязок. billing-service сворачивает их в маленький DSL:
export interface RouteDef<In, Out, R> {
readonly method: 'GET' | 'POST' | 'PUT' | 'DELETE';
readonly path: string;
readonly input: Schema.Codec<In>;
readonly output: Schema.Codec<Out>;
readonly handler: (input: In) => Effect.Effect<Out, unknown, R>;
}
Функция route(def) одна на проект: декодирует тело (при несовпадении схемы сразу 400), запускает обработчик на рантайме процесса, кодирует результат, прогоняет любой провал через toHttp. Модуль объявляет роуты декларативно: метод, путь, две схемы, обработчик. Всё остальное, включая единообразные ошибки, достаётся бесплатно.
Узнаёшь идею? Это ручная, минимальная версия HttpApi из 12 · CLI и терминал и 19 · HttpClient и обвязка API. Когда сервер целиком на Effect, бери HttpApi: он даёт то же самое плюс OpenAPI и производный клиент. DSL руками нужен, когда транспорт чужой, и об этом следующий раздел.
Формы с дискриминантом это объединение
Последний типовой приём. Доменное событие бывает шести видов, и у каждого свой набор полей. Соблазн завести широкий тип с опциональными полями и полем payload: Record<string, unknown> заканчивается тем, что никто не помнит, какие поля у какого вида заполнены. Правильный ответ ты знаешь по 14 · DDD-типы: Schema.Union из структур с общим полем name, и компилятор сам заставит обработать каждый вариант.
Единственное законное исключение: намеренно непрозрачный груз. Сырой колбэк провайдера хранится как Record<string, unknown> осознанно, это боковой карман для аудита, а не доменная форма. Различай эти два случая: форма, по которой ветвится логика, всегда объединение; форма, которую мы только храним и пересылаем, может остаться мешком.
Типы доезжают до SQL: kysely как билдер поверх SqlClient
Осталась последняя граница, где типы обычно умирают: сама база. Сырые sql-теги, которыми написана очередь в Разделе 5, честны и прозрачны, но строково-типизированы: опечатка в имени колонки обнаружится в рантайме, а форму результата ты объявляешь руками и ничем не проверяешь. На трёх таблицах это терпимо. На тридцати хочется, чтобы компилятор знал схему.
Первый рефлекс, взять полноценный ORM, тянет за собой чужой пул соединений, чужие транзакции и чужой рантайм, и всё это начнёт воевать с Layer-графом. Есть ход точнее, и ты его уже видел в основном репозитории BatSchool с drizzle: билдер компилирует, SqlClient исполняет. Здесь покажу его на kysely, типобезопасном SQL-билдере без ORM-амбиций.
Kysely умеет работать в холодном режиме: инстанс без драйвера, который не может исполнить ни одного запроса, зато компилирует их в пару из строки и параметров:
import {
DummyDriver,
Kysely,
PostgresAdapter,
PostgresIntrospector,
PostgresQueryCompiler,
} from 'kysely';
import type { DB } from './schema.ts'; // сгенерирован kysely-codegen, см. ниже
// холодный инстанс: исполнять не умеет, умеет только компилировать
export const qb = new Kysely<DB>({
dialect: {
createAdapter: () => new PostgresAdapter(),
createDriver: () => new DummyDriver(),
createIntrospector: (db) => new PostgresIntrospector(db),
createQueryCompiler: () => new PostgresQueryCompiler(),
},
});
Исполнение остаётся за Effect, через один хелпер на проект:
import type { CompiledQuery } from 'kysely';
import { Effect } from 'effect';
import { SqlClient } from 'effect/unstable/sql';
export const runQuery = <T>(query: CompiledQuery<T>) =>
Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
const rows = yield* sql.unsafe(query.sql, query.parameters as never[]);
return rows as ReadonlyArray<T>;
});
И data-access начинает читаться как обещали рекламные проспекты, только без обмана:
export const findDue = (now: Date, limit: number) =>
runQuery(
qb
.selectFrom('payments')
.selectAll()
.where('status', 'in', ['active', 'past_due'])
.where('nextPaymentDate', '<=', now)
.orderBy('nextPaymentDate')
.limit(limit)
.compile(),
);
Имена таблиц и колонок проверяет компилятор, тип результата выводится из запроса, а не написан рядом с надеждой, что не разъедется. Обрати внимание на два свойства связки. Оба каста, as never[] и as ReadonlyArray<T>, живут в одном месте, в runQuery: это честная граница между мирами типов kysely и Effect SQL, и она не размазана по коду. И второе, важнее: запрос исполняется на SqlClient текущего файбера, поэтому он видит sql.withTransaction и спокойно участвует в транзакциях очереди и outbox. Если бы билдер исполнял запросы сам, через собственный пул, транзакционные гарантии Разделов 5 и 6 молча рассыпались бы.
Откуда берётся тип DB? Из живой базы. kysely-codegen подключается к Postgres, читает каталог и генерирует файл с типами всех таблиц. Цикл закольцовывается в скриптах:
{
"db:codegen": "kysely-codegen --camel-case --out-file src/db/schema.ts",
"migrate": "kysely migrate:latest && pnpm db:codegen"
}
Прогнал миграцию, типы перегенерировались, и компилятор мгновенно подсвечивает каждое место, которое переименованная колонка сломала. Заметь, как это соотносится с правилом производные типы, а не дублированные из начала раздела. Направление вывода другое: в billing-service источник истины это Schema, и из её полей выводится список SQL-колонок; здесь источник истины это миграции и сама база, и из неё выводятся типы. Оба варианта соблюдают главный принцип, источник один, остальное выводится. Выбирай направление по тому, кто владеет схемой: если базой правят только твои миграции, кодоген из базы дешевле.
Для колонок-перечислений, которые в базе хранятся как text, у кодогена есть overrides: колонке payments.status назначается литеральное объединение вместо string, и дискриминированные объединения из прошлого подраздела доезжают до самого SQL.
Два примечания на полях. У kysely есть CamelCasePlugin, который транслирует camelCase кода в snake_case базы: это альтернатива решению billing-service с camelCase-колонками в кавычках, выбери одно и не смешивай. И для Effect 3 существует готовый пакет @effect/sql-kysely, где тот же приём упакован в слой; на v4 связка из тридцати строк выше делает то же самое и не прячет механику.
Раздел 3 · Сервис или функция
В 04 · Services и Layer ты научился заводить сервисы, и есть искушение применять этот молоток ко всему: каждый шаг логики оборачивать в Context.Service и собирать пирамиду слоёв. billing-service даёт трезвое правило, которое стоит унести с собой:
Layer заслуживает настоящий ресурс или настоящая граница подмены. Шаг домена это функция.
Настоящий ресурс: соединение с Postgres, очередь, outbox, HTTP-клиент провайдера с лимитером запросов. У него есть жизненный цикл, его надо открыть и закрыть, он один на процесс. Настоящая граница подмены: место, где прод и тест обязаны отличаться реализацией.
А вот шаги конвейера платежей, матчер (какой сущности принадлежит входящий платёж) и эпплаер (что сделать с совпавшим), это чистые функции. Их не заводят в контекст, их передают параметрами и собирают в одном месте, в composition root:
// composition root воркера: единственное место, где всё встречается
const pipeline = makeChargePipeline({
matcher: composeMatchers([
makeCheckoutMatcher(checkoutRepo),
makeRecurringMatcher(paymentRepo),
]),
applier: makeCheckoutApplier(paymentRepo, checkoutRepo),
});
Посмотри, что здесь произошло с расширяемостью. Требование новый провайдер не трогает ядро выполнено без единого сервиса: подключить нового провайдера значит написать ещё один матчер-функцию и добавить её в массив. Пирамида из Context.Service дала бы то же самое, но ценой графа слоёв, который надо собирать в каждом тесте.
Проверка, которой я пользуюсь: если у кандидата в сервисы нет ни жизненного цикла, ни второй реализации, это функция. Появится вторая реализация, рефакторинг в сервис занимает полчаса. Обратный рефакторинг, разбор ненужной пирамиды слоёв, занимает неделю.
Сам граф слоёв воркера при этом остаётся обычным, послойным, как ты привык:
const base = Layer.mergeAll(SqlLive, TaskRegistryLive, wayForPayLayer);
const withQueue = Layer.provideMerge(QueueLive, base); // очереди нужен Sql
const withOutbox = Layer.provideMerge(OutboxLive, withQueue); // outbox нужны Sql + очередь
export const WorkerLive = Layer.provideMerge(pipelineLayer, withOutbox);
Layer.provideMerge вместо Layer.provide в каждой ступени не случайность: SqlClient должен остаться видимым и очереди, и outbox, и конвейеру, поэтому каждый этаж отдаёт наружу и себя, и всё, что под ним.
Раздел 4 · Один рантайм на процесс
Мост в чужой фреймворк
Pulse целиком на Effect, поэтому вопроса не возникало. В реальности чаще наоборот: HTTP-сервер уже есть (Fastify, Express, Nest), команда его знает, экосистема плагинов нужна, и переписывать вход на HttpApi никто не даст. Как врастить Effect в такой проект?
Ответ billing-service: один ManagedRuntime на процесс, поднятый при старте и проброшенный в обработчики. Ты собираешь Layer-граф как обычно, а фреймворк видит лишь объект с методом runPromise:
import { Layer, ManagedRuntime } from 'effect';
const AppLive = Layer.mergeAll(
HasherLive,
AuthConfigLive,
AuthRepoLive.pipe(Layer.provideMerge(SqlLive)),
);
export const runtime = ManagedRuntime.make(AppLive);
// плагин Fastify: рантайм один, живёт вместе с процессом
fastify.decorate('runtime', runtime);
fastify.addHook('onClose', () => runtime.dispose());
// любой роут: одна строка моста
fastify.post('/api/checkout-sessions', async (request, reply) => {
const result = await fastify.runtime.runPromise(
handle(request.body).pipe(Effect.catchAll((e) => Effect.succeed(toHttp(e)))),
);
reply.status(result.status).send(result.body);
});
Три свойства этой схемы, ради которых всё затевалось.
Один экземпляр сервисов на процесс. Пул соединений с Postgres, лимитер запросов к провайдеру, конфиг: всё собирается однажды. Ошибка, которую ловят на код-ревью в каждом втором проекте с Effect, это Effect.runPromise(program.pipe(Effect.provide(AppLive))) прямо в обработчике: каждый запрос собирает граф заново, открывает свой пул и через час сервис упирается в лимит соединений. Ты знаешь причину из 04 · Services и Layer: у каждого provide своя сборка, мемоизация не сшивает независимые запуски.
Ленивость как бесплатный подарок. ManagedRuntime.make не строит Layer немедленно: первое обращение к базе случится при первом runPromise. Значит, сборка приложения в тестах не требует поднятого Postgres, а модульные тесты роутов остаются герметичными.
Явный dispose. Финализаторы из 05 · Resources срабатывают на закрытии рантайма: пул закрыт, фоновые файберы прерваны, procesо завершается чисто.
Ленивость означает, что порт может слушать раньше, чем база доступна. Поэтому /health здесь не заглушка с ok, а readiness-проба: она запускает SELECT 1 на том же рантайме с таймаутом и отвечает 503, пока база недоступна. Балансировщик не пошлёт трафик на процесс, который не сможет его обслужить.
Воркер это отдельный процесс
Второе решение: фоновая работа (очередь, шедулер, доставка событий) живёт в отдельном процессе со своим рантаймом и своим Layer-графом. Не в setInterval рядом с HTTP-сервером.
Причины прозаичные. Деплой и масштабирование разные: HTTP-процессов может быть пять за балансировщиком, воркер один или два. Падение изолировано: OOM в обработчике очереди не убивает приём вебхуков. Граф зависимостей разный: воркеру нужны очередь, outbox и клиент провайдера, HTTP-процессу нет.
При этом кода два процесса почти не добавляют: оба entrypoint-а импортируют одни и те же модули и отличаются только тем, какой Layer-граф собирают. Это ровно та же идея две поверхности, один граф, что была у Pulse с CLI и HTTP в 12 · CLI и терминал, только теперь поверхности разнесены по процессам.
Каталог bin: тонкие точки входа
Раз точек входа стало несколько, им нужно место и дисциплина. Конвенция, подсмотренная у зрелых сервисов: каталог src/bin/, по одному файлу на процесс, и жёсткое правило, bin не содержит логики:
src/bin/
server.ts HTTP-процесс
worker.ts фоновый процесс
repl.ts операторская консоль (следующий подраздел)
Каждый файл это несколько строк: собрать конфиг, собрать приложение, запустить, отчитаться об ошибке:
// src/bin/server.ts, целиком
import { loadConfig } from '../config.ts';
import { buildApp } from '../app.ts';
const config = loadConfig();
const app = await buildApp(config);
await app.listen({ port: config.port });
// src/bin/worker.ts, целиком
import { loadConfig } from '../config.ts';
import { makeWorkerRuntime } from '../runtime.ts';
import { bootWorker } from '../worker-boot.ts';
const config = loadConfig();
const runtime = makeWorkerRuntime(config);
await runtime.runPromise(bootWorker(config));
Проверка на правильность нарезки: всё, что умеет bin, должен уметь и тест. buildApp и bootWorker импортируемы, значит интеграционный тест поднимает то же самое приложение без спауна процесса, а bin остаётся настолько тонким, что в нём нечему ломаться. В package.json каждая точка входа получает свой скрипт (dev, dev:worker, repl), и с Node 24 файлы запускаются напрямую, node --experimental-strip-types src/bin/server.ts, либо через tsx на старших версиях.
Третий bin: REPL к живому графу
У rails-разработчиков есть привычка, которой энтерпрайзный TypeScript незаслуженно лишён: консоль, где доступно всё приложение. Не psql с голыми таблицами, а именно приложение: домен, репозитории, правила. С ManagedRuntime она строится в двадцать строк:
// src/bin/repl.ts
import repl from 'node:repl';
import { Effect } from 'effect';
import { runtime, type AppServices } from '../runtime.ts';
import * as payments from '../modules/payment/data-access.ts';
import * as outbox from '../modules/outbox/domain.ts';
const r = repl.start({ prompt: 'billing> ', useColors: true });
r.setupHistory('.repl_history', () => {});
r.context.Effect = Effect;
r.context.payments = payments;
r.context.outbox = outbox;
r.context.run = <A, E>(program: Effect.Effect<A, E, AppServices>) =>
runtime.runPromise(program);
r.on('exit', () => {
void runtime.dispose().then(() => process.exit(0));
});
Сессия выглядит так:
$ pnpm repl
billing> const sub = await run(payments.findByExternalUserId('u_42'))
billing> sub.status
'past_due'
billing> await run(outbox.publish(paymentCancelled(sub, new Date())))
Разбери, почему каждая строка файла на своём месте. runtime тот же, что у сервера: ленивость ManagedRuntime означает, что консоль открывается мгновенно, а соединение с базой случится при первом run. Хелпер run прячет весь бойлерплейт запуска: дальше оператор работает с обычными доменными функциями. setupHistory сохраняет историю между сессиями, вчерашняя команда достаётся стрелкой вверх. А dispose на выходе закрывает пул через финализаторы, консоль не оставляет висящих соединений.
Главная ценность не в удобстве, а в безопасности разовых операций. Инцидент: платёж завис, оператору надо перевыпустить событие. Вариант psql, руками написать UPDATE и INSERT, обойдёт все правила Раздела 2 и все журналы Раздела 5: ни идемпотентности, ни аудита, ни outbox. Вариант REPL зовёт ту же доменную функцию, что и продакшен-код, со всеми проверками, транзакциями и журналами. Консоль это операторский инструмент с полным доступом, относись к её запуску на проде так же серьёзно, как к psql, но в отличие от psql она проводит оператора через домен, а не мимо него.
Три bin-файла, один граф. Сервер, воркер и консоль собирают одни и те же слои и зовут одни и те же функции: добавив модуль, ты автоматически получаешь его во всех трёх мирах.
Раздел 5 · Надёжная очередь на голом Postgres
Зачем очередь и почему не брокер
Всё асинхронное в сервисе (доставка событий, обработка входящих платежей, повторные попытки) должно переживать рестарт. В памяти хранить нельзя, значит нужна надёжная очередь. Первый рефлекс, поставить Kafka или RabbitMQ, здесь осознанно отвергнут: до ста тысяч клиентов и одного списания на клиента в месяц нагрузка смешная, а каждый брокер это ещё одна система, которую надо деплоить, мониторить и уметь чинить в три ночи. Postgres уже есть, транзакции уже есть.
Очередь на Postgres это четыре таблицы с ясным разделением ролей:
| Таблица | Роль | Характер |
|---|---|---|
messages | текущее состояние задачи: status, попытки, блокировка | изменяемая, одна строка на задачу |
raw_events | всё, что когда-либо приходило, включая дубликаты | только append |
attempts | каждая попытка обработки с результатом или ошибкой | только append |
message_status_events | журнал переходов статуса | только append |
Правило чтения: текущее состояние берём из кэша messages, историю из журналов. Никогда не восстанавливаем текущий статус сканом журнала на горячем пути и никогда не правим журнал задним числом. Ты уже видел эту пару в 16 · Event Store: append-only лог как истина, проекция как кэш. Здесь та же схема в миниатюре, и каждый переход статуса пишется в журнал тем же SQL-стейтментом, что меняет кэш, чтобы они не могли разъехаться.
Три слоя идемпотентности
Сердце очереди: идемпотентность на трёх этажах. Запомни эту тройку, она переносится на любую систему доставки сообщений.
Этаж 1, дедупликация на входе. У каждой задачи есть детерминированный idemKey, и пара (messageType, idemKey) уникальна в базе:
export const enqueue = (sql: SqlClient) => (input: EnqueueInput) =>
sql.withTransaction(
Effect.gen(function* () {
const inserted = yield* sql<{ id: string }>`
INSERT INTO messages ("messageType", "idemKey", payload)
VALUES (${input.messageType}, ${input.idemKey}, ${JSON.stringify(input.payload)}::jsonb)
ON CONFLICT ("messageType", "idemKey") DO NOTHING
RETURNING id
`;
const enqueued = inserted.length > 0;
yield* sql`
INSERT INTO raw_events ("messageType", "idemKey", payload, "wasDuplicate")
VALUES (${input.messageType}, ${input.idemKey},
${JSON.stringify(input.payload)}::jsonb, ${!enqueued})
`;
return { enqueued, messageId: inserted[0]?.id ?? null };
}),
);
Дубликат вебхука упирается в ON CONFLICT DO NOTHING и не рождает вторую задачу. Заметь: сырое событие в raw_events пишется всегда, дубликат оно или нет, просто с флагом. Аудит полный, обработка одна.
Этаж 2, эксклюзивность взятия. Воркеров может быть несколько, и одну задачу не должны взять двое. Это решает Postgres одной конструкцией:
export const claimBatch = (sql: SqlClient) => (params: ClaimParams) =>
sql<ClaimedMessage>`
WITH claimed AS (
SELECT id FROM messages
WHERE status IN ('pending', 'retry')
AND ("retryAt" IS NULL OR "retryAt" <= now())
AND "messageType" IN ${sql.in(params.types)}
ORDER BY id
LIMIT ${params.limit}
FOR UPDATE SKIP LOCKED
),
upd AS (
UPDATE messages m
SET status = 'in_progress', "lockedBy" = ${params.workerId},
"lockedAt" = now(), "attemptCount" = m."attemptCount" + 1
FROM claimed WHERE m.id = claimed.id
RETURNING m.id, m."messageType", m."idemKey", m.payload, m."attemptCount"
),
logged AS (
INSERT INTO message_status_events ("messageId", status)
SELECT id, 'in_progress' FROM upd
)
SELECT id, "messageType", "idemKey", payload, "attemptCount" FROM upd
`;
FOR UPDATE SKIP LOCKED делает взятие кластер-безопасным без единого мьютекса в коде: второй воркер просто не увидит строк, которые держит первый. Обрати внимание на два штриха. Отбор фильтруется по типам, для которых у воркера есть обработчики: задача незнакомого типа не будет взята и потеряна. И перевод в in_progress, инкремент попытки и журнальная запись происходят в одном стейтменте с выборкой: между увидел и забрал нет щели, куда мог бы влезть конкурент.
Колонки в кавычках и в camelCase это не украшение: без кавычек Postgres приведёт имя к нижнему регистру, и messageType молча превратится в messagetype. Проект выбирает кавычки, чтобы имена колонок буква в букву совпадали со схемами и типами.
Этаж 3, идемпотентность обработчика. Первые два этажа дают доставку at-least-once: задача не потеряется, но может быть обработана повторно (сейчас увидим почему). Значит, сам обработчик обязан быть повторяемым: создавать сущности с детерминированными идентификаторами, писать через upsert, проверять а не сделано ли уже. Это не удаётся спрятать в инфраструктуру, это контракт для каждого, кто пишет обработчик.
Reaper и visibility timeout
Откуда берётся повторная обработка? Воркер взял задачу и умер: OOM, деплой, выдернутый шнур. Строка застряла в in_progress навсегда. Лечит это reaper, шаг в начале каждого тика: строки, чей lockedAt старше visibility timeout, возвращаются в retry и становятся доступны снова. Отсюда и возможность повтора: воркер мог успеть сделать работу и умереть за миллисекунду до отметки о завершении. Reaper вернёт задачу, второй воркер выполнит её ещё раз, и спасёт всех только идемпотентный обработчик.
Диспетчер: цикл, который не умирает
Сам цикл воркера собирается из знакомых тебе кубиков 11 · Runtime и Schedule:
const processOne = (message: ClaimedMessage, handler: TaskHandler) =>
Effect.gen(function* () {
const startedAt = new Date(yield* Clock.currentTimeMillis);
// exit, а не either: ловим и типизированные провалы, И дефекты
const exit = yield* Effect.exit(handler(message.payload));
const failed = Exit.isFailure(exit);
const outcome = decideOutcome({ failed, attemptCount: message.attemptCount });
yield* complete({
messageId: message.id,
startedAt,
error: failed ? describeError(Cause.squash(exit.cause)) : null,
outcome,
});
});
export const runDispatcher = (options: DispatcherOptions) =>
tick(options).pipe(
Effect.catchAllCause((cause) =>
Effect.logError('queue dispatcher tick failed').pipe(
Effect.annotateLogs('cause', Cause.pretty(cause)),
),
),
Effect.andThen(Effect.sleep(Duration.millis(options.pollIntervalMillis))),
Effect.forever,
);
Две строки здесь несут всю живучесть. Effect.exit вокруг обработчика: обработчик, который бросил исключение (дефект, не типизированный провал), записывается как неудачная попытка со всеми подробностями из Cause, а не роняет процесс воркера. И Effect.catchAllCause вокруг тика: недоступная база на одном тике это строка в логе, а не смерть цикла, Effect.forever продолжит со следующего тика.
Решение о судьбе неудачной попытки вынесено в чистую функцию decideOutcome: успех терминален, провал уходит в retry с экспоненциальным бэкоффом, после исчерпания попыток терминальный fail. Чистую политику можно исчерпывающе покрыть юнит-тестами без базы и без часов, приём из 13 · Testing.
Раздел 6 · Transactional outbox
Проблема двух миров
Сценарий: платёж прошёл, надо обновить подписку в базе и сообщить об этом в CRM. Два хранилища, одной транзакции на двоих не существует. Отправишь событие после коммита, процесс может умереть между, и CRM никогда не узнает об оплате. Отправишь до коммита, транзакция может откатиться, и CRM узнает о платеже, которого не было. Это не экзотика, это главный источник рассинхрона в интеграциях.
Классическое решение называется transactional outbox: событие не отправляется, а записывается в ту же базу той же транзакцией, что и изменение состояния. Раз база одна, атомарность даёт Postgres. Доставкой занимается отдельный процесс, у нас он уже есть: воркер с очередью из Раздела 5.
Публикация: событие плюс веер доставок
Получателей у события несколько (CRM, аналитика, вебхук клиента), и у каждого своя судьба доставки: CRM может лежать, пока аналитика отвечает. Поэтому outbox хранит не только событие, но и по одной записи доставки на каждого получателя:
export const publish = (deps: PublishDeps) => (event: DomainEvent) =>
deps.repo.transaction(
Effect.gen(function* () {
const isNew = yield* deps.repo.insertEvent(event);
if (!isNew) {
return; // уже публиковали: идемпотентный no-op
}
const sinks = yield* deps.sinks.all();
yield* Effect.forEach(
sinks,
(sink) =>
deps.repo.insertDelivery(event.id, sink.name).pipe(
Effect.flatMap((deliveryId) =>
deps.enqueue({
messageType: 'deliver_event',
idemKey: deliveryId,
payload: { deliveryId },
}),
),
),
{ discard: true },
);
}),
);
Разбери, как здесь сцеплены гарантии. Идентификатор события детерминированный и приходит от вызывающего, поэтому повторный вызов publish (например, из переигранного обработчика очереди, Раздел 5 обещал такие повторы) увидит isNew === false и молча выйдет: двойной публикации не будет. Каждая доставка становится задачей в той же очереди, а значит наследует всю её механику: ретраи с бэкоффом, журнал попыток, reaper. Мы не строим вторую систему надёжности, мы переиспользуем первую.
Сама доставка симметрично проста: прочитать доставку с событием, если уже delivered, выйти (идемпотентность), иначе позвать получателя. Успех помечает доставку. Провал записывает ошибку для оператора и перебрасывает её дальше, чтобы очередь запланировала ретрай:
yield* sink.deliver(event).pipe(
Effect.matchCauseEffect({
onSuccess: () => deps.repo.markDelivered(deliveryId),
onFailure: (cause) =>
deps.repo
.markFailed(deliveryId, describeError(Cause.squash(cause)))
.pipe(Effect.andThen(Effect.failCause(cause))),
}),
);
Строка с Effect.failCause(cause) легко теряется при чтении, а она здесь главная: записать ошибку и проглотить её значило бы отчитаться очереди об успехе, и ретраев бы не было. Записать и переупасть, вот полный жест.
Итоговая гарантия: каждое событие дойдёт до каждого получателя как минимум один раз, обычно за секунды, а не дойдя, останется видимым оператору со всей историей попыток. Требование про минуту до CRM закрывается архитектурой, а не надеждой.
Раздел 7 · Собственный биллинг-цикл
Шедулер как третья поверхность
Сервис не только реагирует на входящие, он сам инициирует списания: настал день оплаты, надо списать с сохранённого токена. Это третья поверхность процесса-воркера, рядом с диспетчером очереди: цикл на Effect.forever, который каждый тик выбирает подписки с наступившей датой и пытается списать.
Внутри тика два приёма, оба тебе знакомы, но здесь они несут деньги.
Изоляция элементов. Одна неудачная подписка не должна прервать пачку: каждое списание завёрнуто в собственный catchAllCause с логом, и цикл идёт дальше. Тот же жест, что в диспетчере, уровнем ниже.
Детерминированный идентификатор заказа. Каждой попытке списания нужен orderReference для провайдера. Он строится из идентификатора подписки и запланированной даты списания, не из текущего времени. Если процесс упал после запроса к провайдеру, но до записи результата, повторная попытка построит тот же orderReference, и провайдер отбросит дубль на своей стороне. Идемпотентность, этаж четвёртый: на границе с чужой системой ключ должен быть воспроизводимым.
Лестница ретраев это домен, а не очередь
Списание не прошло (карта просрочена, нет денег). Что дальше? У очереди из Раздела 5 есть свой ретрай с экспоненциальным бэкоффом, но он здесь не подходит, и важно понять почему. Бэкофф очереди это техническая политика: сеть мигнула, попробуй через секунду, потом через две. А повтор списания это бизнес-политика: продукт решил, что пробуем в день неудачи, потом через день, на третий, пятый и седьмой день, после чего сдаёмся и шлём терминальное событие о неудачном продлении. Эти цифры видны пользователю в письмах, зафиксированы в требованиях и не имеют отношения к сетевым сбоям.
Поэтому лестница живёт в домене, в чистой функции:
export const RETRY_SCHEDULE_DAYS = [0, 1, 3, 5, 7] as const;
export const planRetry = (attempt: number, firstFailureAt: Date): RetryPlan => {
const nextAttempt = attempt + 1;
const day = RETRY_SCHEDULE_DAYS[nextAttempt];
if (day === undefined) {
return { final: true, attempt: nextAttempt, nextPaymentDate: null };
}
return {
final: false,
attempt: nextAttempt,
nextPaymentDate: addDays(firstFailureAt, day),
};
};
Дата каждой следующей попытки считается от firstFailureAt, зафиксированного момента первой неудачи, а не от предыдущей попытки. Если воркер лежал сутки, лестница не растянется: пропущенные ступени просто наступят разом.
Заметь разделение труда: домен решает, когда и что (следующая дата, терминальность, какое событие публиковать), очередь решает, как надёжно исполнить решение. Два вида ретраев сосуществуют, не подменяя друг друга.
Якорь даты против дрейфа
Самая дорогая строка урока. Подписка на месяц, оплата второго числа. Списание не прошло, ретраи шли пять дней, седьмого числа деньги списались. Когда следующее списание?
Если следующая дата считается как успешная попытка плюс месяц, то седьмого. Ещё одна неудача, и уже двенадцатого. Подписка дрейфует: каждый инцидент навсегда сдвигает день оплаты, через год биллинг размазан по всему месяцу, сверка с бухгалтерией превращается в ад. Эта ошибка кочует из проекта в проект, потому что локально каждая строка выглядит разумно.
Лекарство: даты периода образуют якорь. У подписки хранятся currentPeriodStart и currentPeriodEnd, и следующая дата списания всегда выводится из конца периода, никогда из даты фактической попытки:
if (response.transactionStatus === 'Approved') {
const newEnd = addPeriod(sub.currentPeriodEnd, sub.period);
yield* deps.subs.advanceAfterSuccess(sub.id, {
currentPeriodStart: sub.currentPeriodEnd,
currentPeriodEnd: newEnd,
nextPaymentDate: newEnd,
});
}
Оплатили седьмого за период со второго, следующий период всё равно начинается второго. Дрейф не подавлен дисциплиной, он структурно невозможен: в формуле следующей даты просто нет фактической даты платежа. Осталась одна тонкость, календарь: месяц от 31 января это 28 февраля, поэтому addPeriod прижимает день к последнему дню целевого месяца и живёт в одном модуле с тестами на високосные годы.
Раздел 8 · Границы с внешним миром
Провайдер это модуль, а не россыпь утилит
Всё специфичное для платёжного провайдера (подпись запросов, разбор колбэков, форма оплаты, HTTP-клиент, его лимитер) живёт в одном модуле wayforpay/. Чекаут и конвейер платежей провайдера не знают: они оперируют нейтральными формами и получают функции-шаги через composition root из Раздела 3. Подключить второго провайдера значит написать второй модуль с той же поверхностью, не тронув ядро. Это ровно anti-corruption layer из 09-backend · DDD-обзор, доведённый до механической проверяемости: грепни импорты wayforpay вне модуля, их не должно быть нигде, кроме composition root.
Внутри модуля есть чему поучиться и без архитектуры. Подпись запросов по таблице полей в зафиксированном порядке, где порядок это протокол, а не стиль. Клиент, завёрнутый в лимитер запросов из 19 · HttpClient, потому что у провайдера свои квоты. Ответ отклонено, которое не ошибка Effect, а законная ветка результата: шедулер по ней идёт в лестницу ретраев, и путать её с сетевым сбоем нельзя.
Терпимый парсинг: не потерять ни одного события
Провайдерские колбэки это дикая природа: суммы приходят то числом, то строкой, поля появляются без предупреждения. Строгая схема, которая отбрасывает непонятное, здесь означала бы потерю платежа. Правило сервиса: сырое сохраняется всегда, непонятное никогда не выбрасывается.
Каждое входящее событие сначала целиком пишется в raw_events (Раздел 5), и только потом разбирается. Схемы разбора терпимые: числовые поля принимают строку и число (Schema.decodeTo с трансформацией, ты строил такие в 02 · Schema и 21 · Schema вглубь). А если событие разобралось, но не совпало ни с одной сущностью, оно уходит не в /dev/null, а в карантин: отдельную таблицу с сырым телом, алертом оператору и ручкой привязки. Оператор привязывает платёж к подписке, и событие переигрывается по обычному конвейеру, как будто совпало сразу. Метрика вечно неопознанных платежей должна быть нулём, и она видима.
Сравни с DLQ из 23 · Sink и конвейеры: та же философия (не терять, а откладывать и показывать), но карантин это не свалка ошибок, а рабочая очередь человека с кнопкой переиграть.
Опаковый идентификатор пользователя
Последний штрих границы. Сервис не хранит пользователей: снаружи приходит externalUserId, и он опаковый: возвращается в каждом событии буква в букву, без попыток сматчить по почте или телефону. Матчинг по нечётким признакам это чужая ответственность и бесконечный источник инцидентов. Урок шире платежей: если твоему сервису не нужна собственная идентичность пользователя, не заводи её, таскай чужой идентификатор как непрозрачный жетон.
Чеклист · собери свой бэкенд
ДЗ
Все задания делаются в pulse-<nick>. После каждого открой PR с тегом lesson-26. Задания складываются в маленький бэкенд с надёжной доставкой алертов: делай их по порядку.
Дальше
- Открой сам billing-service и прочитай его как код коллеги: начни с
docs/(там решения записаны раньше кода), потомruntime.ts, потом любой модуль. Теперь ты знаешь там каждую конструкцию. - 16 · Event Store и 17 · CQRS-проекции, если хочешь увидеть, как append-only журнал становится главной моделью, а не боковым карманом.
- 09-backend · Надёжная доставка, тот же outbox и идемпотентность с высоты птичьего полёта, без привязки к Effect. А 24-microservices, когда один процесс перестанет вмещать команду: там эти решения становятся обязательной программой.
- 29-testing, раздел про тестирование: чистые политики вроде decideOutcome и planRetry это идеальные кандидаты для property-тестов.