Симптом механизма выглядит как «случайный» порядок: consumer сначала получает OrderCancelled v3, потом OrderPaid v2, а затем ещё один из этих event. Цена прямого применения первого входа — проекция перескакивает через факт, который объясняет отмену. Цена игнорирования версии — позднее v2 может вернуть состояние назад. Цена повторного вызова компенсации — два независимых решения по одному отказу. Слова duplicate, out-of-order и failed action описывают разные границы; одна проверка их не закрывает.
Ниже я держу модель предельно маленькой. Owner уже имеет paid v2. Контролируемый учебный отказ резерва создаёт у него новое решение cancelled v3 и один compensationKey. Затем отдельный consumer получает v3 раньше v2, откладывает её, применяет v2 и replay-ит отложенную v3. Повтор v2 или v3 не создаёт нового состояния. Это детерминированный Node-модуль с Map и Set. Он не реализует broker, database, HTTP, outbox, distributed transaction, реальный reserve API или exactly-once delivery.
Три ключа отвечают на три разных вопроса
Первый ключ — eventKey = source + id. Он нужен consumer-у, чтобы отличить уже применённую доставку от повтора того же конверта. Второй ключ — orderVersion. Он принадлежит owner-у одного заказа и говорит, какой переход должен быть следующим. Третий ключ — compensationKey. Он принадлежит одному компенсирующему решению, например «отменить paid v2 после подтверждённого отказа резерва». Если использовать один id для всех трёх задач, код станет короче только до первого повторного сообщения или повторной команды.
| Ключ | Что защищает | Где хранится в модели | Чего не гарантирует |
|---|---|---|---|
source + id | повтор одного event | set appliedEventIds consumer-а | порядок разных событий и бизнес-эффект |
orderVersion | переходы одного owner state | owner и projection | глобальный порядок всех заказов или всех topic |
compensationKey | повтор одного решения owner-а | ledger owner-а | автоматическое удаление внешнего действия |
reservationEvidence | разрешение локального шага | projection исполнения | что owner всё ещё в том же state без сверки версии |
reason | основание компенсации | вход решения owner-а | что любой unknown failure можно безопасно компенсировать |
Стабильный id полезен, но сам по себе не делает обработчик идемпотентным в доменном смысле. Kafka 2.7 описывает idempotent producer отдельно и подчёркивает, что прикладные повторные отправки остаются отдельной проблемой. Поэтому я не называю set event id решением «ровно один раз». В fixture он только запрещает второй переход проекции по тому же event. Внешний платёж, письмо или резерв требуют своего effect key, владельца и доказательства результата. Это продолжение февральского и мартовского контрактов: key полезен только в той границе, где его реально проверяют.
Gap не нужно превращать в отмену
Когда consumer с текущей версией 1 видит cancelled v3, у него нет права сразу сделать вывод, что v2 неважна. Возможно, v2 содержит обязательное решение, schema change или причину последующей компенсации. В учебной модели v3 попадает в deferredByVersion вместе с записью evidence: ожидалась v2, пришла v3. Проекция остаётся на v1. Это важный результат: не было «частично применённой отмены», которую потом придётся угадывать при разборе.
function applyEvent(projection, event) {
if (event.orderVersion > projection.version + 1) {
return { state: 'deferred-version-gap', expected: projection.version + 1 };
}
if (event.orderVersion === projection.version) {
return { state: 'same-version-event-rejected' };
}
if (event.orderVersion < projection.version) {
return { state: 'stale-event-ignored' };
}
projection.version = event.orderVersion;
projection.state = event.type === 'training.order.cancelled'
? 'cancelled'
: 'awaiting-reservation';
return { state: 'event-applied' };
}
После доставки paid v2 consumer применяет ровно следующий переход: accepted v1 → awaiting-reservation v2. Затем функция replay берёт v3 из отложенной коллекции и применяет её: awaiting-reservation v2 → cancelled v3. В этом месте порядок не «исправлен брокером». Его проверяет контракт projection. Если v2 так и не найдена, consumer не имеет автоматического перехода и должен удержать evidence либо направить запись в ручный маршрут. TTL, дополнительный retry или случайная сортировка payload не дают доказательства, что пропуск безопасен.
Компенсация — новое доменное решение
Слово «компенсация» легко создаёт ложное впечатление, что можно отменить всю распределённую операцию одним rollback. Здесь оно означает другое: owner получает проверяемый reason, проверяет исходную версию и создаёт следующий state. paid v2 не исчезает из истории; после него появляется cancelled v3 с причиной и ключом решения. Consumer не должен самостоятельно получить техническое исключение и поставить owner-у отмену. Иначе неизвестно, кто подтвердил причину, как выбран повтор и почему два consumer-а не создали две отмены.
function compensatePaidOrder(order, rejection, ledger) {
const key = 'compensation:' + order.id + ':reservation:v' + order.version;
if (ledger.has(key)) return { state: 'already-recorded', key };
if (order.status !== 'paid' || rejection.reason !== 'training-reservation-rejected') {
return { state: 'manual-review' };
}
const next = { ...order, status: 'cancelled', version: order.version + 1 };
ledger.set(key, next);
return { state: 'compensation-recorded', order: next, key };
}
В функции есть узкое условие: известен именно training-reservation-rejected и version отказа совпадает с owner version. Для unknown timeout это плохой автоматический путь. Он может означать, что внешний резерв уже создан, но ответ потерян. В таком случае нужны отдельные evidence и ручное или явно описанное сверочное действие. Уменьшить правила до «любая ошибка отменяет заказ» можно только ценой новых ложных отмен. Поэтому controlled compensation сначала ограничивает вход, а уже потом меняет state.
Локальная запись может быть повторяемой, но не распределённой
В одном PostgreSQL-хранилище можно сочетать update owner-а с записью compensation ledger и уникальным compensation_key. INSERT ... ON CONFLICT задаёт поведение при конфликте уникального ограничения, а transaction isolation — границу конкурентной работы этой базы. Это полезный строительный блок. Он не даёт одной командой атомарно изменить другую базу, отправить message и отменить внешний резерв. Такой вывод важнее красивой диаграммы: каждая надежда должна быть привязана к месту, где она действительно проверяется.
-- Локальная граница одной базы. Не делает две базы атомарными.
BEGIN;
UPDATE orders SET status = 'cancelled', version = version + 1
WHERE id = 'order-417' AND status = 'paid' AND version = 2;
INSERT INTO compensation_ledger (compensation_key, order_id, order_version)
VALUES ('compensation:order-417:reservation:v2', 'order-417', 2)
ON CONFLICT (compensation_key) DO NOTHING;
COMMIT;
Если запись owner-а повторилась, ledger не позволяет создать вторую такую компенсацию. Если доставка cancelled v3 повторилась, consumer set не применяет второй раз тот же event. Если пришла другая v3 с тем же orderVersion, это уже не duplicate того же конверта, а конфликт контракта: fixture возвращает same-version-event-rejected, оставляет projection на cancelled v3 и записывает evidence. Реальная система всё равно должна иметь owner, правило разбора и отдельный test для такого конфликта. Важно не спрятать нерешённый случай за названием ON CONFLICT.
Как fixture проверяет инвариант без инфраструктуры
Fixture сначала создаёт cancelled v3 через createControlledCompensation(); второй вызов с тем же входом получает compensation-already-recorded. Затем projection получает v3 первой, фиксирует gap и не меняет свой version. После v2 она применяет отложенную v3. В конце повторные v3 и v2 подавляются, а другая v3 отклоняется как same-version conflict; shippingIsAllowed возвращает false, потому что owner и projection находятся в cancelled state. Это не integration test, зато набор проверок делает условие изменения видимым.
const fixture = runDataConsistencyFixture();
if (!Object.values(fixture.assertions).every(Boolean)) {
throw new Error('training consistency contract failed');
}
console.log(fixture.deliveries);
// v3 сначала defer, v2 применяется, затем v3 replay; duplicate не меняет projection.
Здесь полезно заметить предел самого инварианта. Он говорит, что данная учебная projection не разрешает отгрузку после отмены и не двигается назад по version. Он не говорит, что любой настоящий сервис увидит события в этом порядке, что message никогда не пропадёт или что внешний резерв уже освобождён. После изменения fixture нужно тестировать реальные точки: запись owner-а, publish path, consumer storage, retry policy и внешний effect. Модель не заменяет эти проверки, но помогает сформулировать, чего именно от них ждать.
Маршрут механизма без лишних обещаний
- Выписать state, которым владеет один сервис, и последовательную версию этого state. Не строить version из времени получения сообщения.
- Отделить event id для duplicate от version объекта. Для обоих назвать место хранения и период жизни.
- При version gap сохранять evidence и не менять projection до следующей допустимой версии. Явно решить, сколько и где она ждёт.
- Определить узкий набор причин, для которых owner создаёт compensation. Unknown failure не переводить автоматически в отмену.
- Защитить одно решение compensation key в локальной границе owner-а; отдельно защитить consumer от повторного event.
- Проверить запрет рискованного действия: без равной версии и локального evidence оно не должно становиться доступным.
- Только затем соединить этот контракт с конкретными database, broker и внешними API и написать интеграционные tests на их реальные ошибки.
Историческая рамка и ограничения
CloudEvents 1.0.1 позволяет говорить о контексте события без привязки к одному transport. PostgreSQL 13 даёт точную модель локальной транзакции, а Kafka 2.7 показывает, что даже producer idempotence имеет отдельные предпосылки и границы session. Эти документы существовали к апрелю 2021 года, но не создают один общий стандарт «согласованности между сервисами». Transactional outbox, saga и compensating action часто используются как названия паттернов; здесь они не выданы за единый протокол и не приписаны учебному коду как готовая platform capability.
В пакете не запускались PostgreSQL, Kafka, queue, HTTP, внешний reserve API, два процесса, browser, CI, production build или deployment. Нет реального order, измеренной задержки и claims о стабильной distributed platform. Следующий шаг — выбрать одну реальную границу хранения и доказать в ней повторяемость owner update и ledger. После этого отдельно проверить, что consumer видит version gap как факт, а не как разрешение на произвольный rollback.
Проверяемые источники
- CloudEvents Specification v1.0.1: release record — официальная карточка выпуска, опубликованного в декабре 2020 года; поля id, source и type служат словарём учебного конверта, но fixture не является реализацией CloudEvents
- PostgreSQL 13: Transaction Isolation — версионная документация PostgreSQL 13, доступная к апрелю 2021 года; описывает границы локальной транзакции и необходимость retry при serialization failure
- PostgreSQL 13: INSERT и ON CONFLICT — официальный синтаксис и семантика локального уникального ограничения; пример ниже не утверждает атомарность между несколькими сервисами
- Apache Kafka 2.7.0: KafkaProducer API — versioned API линии 2.7, существовавшей в апреле 2021 года; ограничивает идемпотентность producer одной session и не устраняет повторную отправку прикладным кодом