Why integration events

ПроблемаДоменный факт уехал на шину как есть, со своими полями и типами. Сосед разобрал его и построился на нём: теперь переименование поля внутри домена ломает чужой сервис, и узнаём мы об этом из его алертов.

Тот же сервис, что в главах про DDD: счёт, его модули, порты и адаптеры. Здесь появляется граница контекста — та, за которой начинаются чужие сервисы. У факта на ней другое имя, другая форма и другие правила изменения.

  1. Шаг 1 из 9

    Домен уехал на шину

    Справа публикация из главы про инфраструктуру. Факт уходит в outbox так, как он написан в домене: json.Marshal берёт структуру Issued целиком, а типом сообщения становится внутреннее имя invoice.issued.

    Пока это читаем только мы, всё честно. Но сообщение прочитал сосед — доставка подписалась и начала собирать посылку. С этой минуты внутренняя структура домена стала контрактом: переименовал поле Total — сломал чужой парсер; добавил поле для себя — рассказал соседям то, чего им знать не нужно.

    Работает это и в обратную сторону. Доставке понадобился адрес клиента, которого в факте нет, — и просьба добавить поле приходит в домен.

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

    Шаг 2 из 9

    Интеграционное событие

    Интеграционное событие — отдельный тип в отдельном пакете, и версия стоит прямо в его пути: infrastructure/broker/contracts/v1.

    Лежит он в инфраструктуре шины, а не в корне сервиса, и это не мелочь. Форма сообщения — деталь того транспорта, которым оно едет: у шины это JSON в конверте с ключом партиции, у gRPC был бы protobuf со своими правилами совместимости, у HTTP — тело ответа и статус. Транспортов у сервиса бывает несколько, и у каждого контракт свой; общий каталог contracts делает вид, что он один.

    Типы внутри примитивные: строки и int64. Ни vo.Money, ни invoice.Number за границу не уезжают — там JSON, и поля названы так, как их прочитает чужой сервис. Это DTO границы, и ничего больше.

    Имя сообщения — константа с версией на конце: billing.invoice.issued.v1. По нему читателю видно, откуда сообщение, про что оно и какой схемы. Внутреннее FactName() осталось своим: его мы переименуем свободно, это имя — нет.

    Набор полей здесь — решение, а не отражение факта. В контракт попадает то, что соседям действительно нужно: добавить поле потом можно, убрать — уже нет.

    Шаг 3 из 9

    Конверт

    Данные факта — одно, метаданные сообщения — другое, поэтому они и лежат отдельно: тело в Payload, вокруг него идентификатор, тип, время и ключ.

    ID нужен получателю: доставка at-least-once обещает, что сообщение придёт, но не обещает, что один раз. По этому идентификатору повтор узнают — к этому вернёмся на входящей стороне.

    Key — про порядок. Брокер держит порядок внутри партиции, а не во всём топике; ключ, равный номеру счёта, кладёт все сообщения одного счёта в одну партицию: «выставлен» придёт раньше «оплачен». Между разными счетами порядка нет, и он не нужен.

    Type стоит в конверте, а не только в теле: подписчик решает, разбирать ли тело, ещё не разобрав его.

    Шаг 4 из 9

    Кто переводит

    Перевод живёт в инфраструктуре: домен про контракт не знает и знать не должен. Contract берёт доменный факт и отдаёт готовый конверт.

    switch по типу факта выглядит скучно — в этом и смысл. Он явный: чтобы новое сообщение уехало наружу, кто-то должен добавить его здесь. Пока не добавил, факт остаётся внутренним, и публикация его пропускает — ErrInternalFact.

    Именно это чаще всего и забывают. Публикуешь всё подряд — и каждый новый доменный факт становится публичным в день, когда его написали, молча и без чьего-либо решения.

    Рядом — тот, кто Contract зовёт: broker.Publisher. Это он реализует порт invoice.Publisher, который домен объявил ещё во второй главе: перевести факт, сериализовать и отдать строку на запись. Отправлять он ничего не отправляет: шина в транзакцию не входит — это было в главе про инфраструктуру.

    Записывает строку не он, а Appender — узкий порт из одного метода, объявленный здесь же, у того, кто им пользуется. За ним стоит postgres.Outbox, и после этой главы он занимается ровно одним делом: кладёт строку в таблицу. Ни домена, ни контракта в его импортах больше нет — раньше он знал и то, и другое.

    Направление зависимостей от этого выпрямилось: знает про контракт тот, кто с шиной и работает, а хранилище — только про строки. Связывает их сборка: wire.Bind(new(broker.Appender), new(*postgres.Outbox)).

    Шаг 5 из 9

    Нам тоже присылают

    Обратная сторона границы. Платёжный шлюз публикует payments.payment.received.v1, и по этому сообщению счёт становится оплаченным — значит, кто-то должен его принять и позвать сценарий справа.

    Сделать это плохо можно двумя способами.

    Первый: разобрать чужой JSON прямо в сценарии. Тогда application начинает знать формат соседа, и переименование поля в чужом сервисе доезжает до нашего домена.

    Второй: позвать сценарий как есть. То же сообщение приедет второй раз — после ретрая брокера, после перезапуска консьюмера, после переезда партиции. Второй раз MarkPaid вернёт ErrAlreadyPaid, консьюмер сочтёт это ошибкой обработки и вернёт сообщение в очередь. И так по кругу, пока кто-нибудь не посмотрит в логи.

    Шаг 6 из 9

    Приём: перевести, вспомнить, простить

    Хендлер лежит в слайсе платежа, рядом со своим сценарием, и делает три вещи.

    Переводит. PaymentReceived — наша копия чужой схемы, и она стоит на границе. Изменится схема у соседа — правим этот файл; домен не заметит.

    Помнит. Remember отвечает «впервые» или «уже было», и вызов стоит внутри той же единицы работы, что и сценарий. Иначе между пометкой и изменением остаётся щель, ровно в которую и попадает повтор. Поэтому реализация Do получила одну строку: транзакция в ctx уже есть — значит, выполняем работу внутри неё, а не начинаем вторую.

    Прощает. ErrAlreadyPaid — не ошибка обработки: счёт уже в нужном состоянии, сообщение подтверждаем и идём дальше. Это и есть идемпотентность на входе — два одинаковых сообщения дают один результат.

    Шаг 7 из 9

    Память о сообщениях

    Реализация простая: таблица с уникальным идентификатором и INSERT ... ON CONFLICT DO NOTHING. Вставилась строка — сообщение новое; ноль строк — это сообщение мы уже обрабатывали.

    Соединение берётся тем же conn(ctx, db), что и у остальных адаптеров, — значит, пометка ложится в транзакцию операции. Либо счёт оплачен и сообщение помечено, либо ни того, ни другого.

    О чём не надо забывать: таблица растёт. Идентификаторы старше срока хранения сообщений в брокере не нужны — их чистят по времени, партициями, как обычные логи.

    Шаг 8 из 9

    Как менять контракт

    Правило одно, и оно не про код: схему читают чужие, значит, меняем её так, чтобы читатель не заметил.

    • Добавить необязательное поле — можно. Старый читатель его проигнорирует.
    • Переименовать, удалить, сменить смысл или тип — нельзя. Это v2: новый тип сообщения, который живёт рядом с v1, пока последний читатель не переехал.

    Значит, надо знать, кто читает. Без этого v1 не выключить никогда — и версии будут копиться, а не сменяться.

    И граница, за которую лучше не выходить: событие говорит «что случилось», а не «сделай». Как только в контракте появляется ChargePenalty, это уже не факт, а команда — и сосед начинает управлять нашим доменом через шину.

  2. Шаг 9 из 9

    Что получилось

    Щёлкни по файлу, чтобы прочитать.

    • infrastructure/broker/contracts/v1/events.go — что видят соседи.
    • infrastructure/broker/contracts/v1/envelope.go — конверт: id, тип, время, ключ.
    • infrastructure/broker/outbound.go — перевод факта в сообщение.
    • infrastructure/broker/publisher.go — публикация: переводит и отдаёт строку на запись.
    • infrastructure/postgres/outbox.go — одна строка в таблицу, и больше ничего.
    • applications/payment/transport/consumer/received.go — приём чужого сообщения.
    • infrastructure/postgres/inbox.go — память об обработанных.
    • cmd/api/wire.go — сборка: порт публикации теперь за broker.Publisher.

    Домен не изменился ни на строку: у него по-прежнему свои факты и своё имя для каждого. Всё новое легло на границу — по обе её стороны.

Листать шаги можно стрелками ← и →.