Publisher-адаптер NATS JetStream на Python
Реализуй исходящий технологический адаптер, который преобразует DTO порта в согласованное сообщение JetStream и считает публикацию успешной только после подтверждения брокера.
Порядок работы
- Изучи исходящий порт, DTO, формат сообщений и маршруты.
- Проверь, достаточно ли DTO для публикации без обращения к другим источникам.
- Зафиксируй типизированную конфигурацию маршрутов и управляемой топологии.
- Реализуй явные transport-модели, mapping, сериализацию и headers.
- Реализуй подготовку топологии до первой публикации.
- Реализуй publish и преобразование технологических ошибок в ошибки порта.
- Добавь unit- и integration-тесты.
Если требования полны, следуй им без дополнительного проектирования. Уточняй только отсутствующие решения, без которых нельзя безопасно реализовать контракт или изменить существующую топологию.
Исходящий порт и DTO
- Адаптер реализует уже определённый application-порт.
- Принимай только DTO этого порта.
- Если порт принимает закрытое объединение конкретных DTO, используй объявленный
application
TypeAlias в сигнатуре и выполняй явный исчерпывающий dispatch по
runtime-типу. Не заменяй alias Any, общим transport DTO или перечислением
типов ресурсов внутри application.
- DTO содержит семантические данные и необходимые metadata, но не имена stream, subject и параметры JetStream.
- Не задавай универсальный набор полей DTO.
- Не извлекай недостающие данные из репозиториев, доменных объектов или других application-операций.
- Недостаточный DTO приводит к типизированной невосстановимой ошибке порта до обращения к NATS.
- Не передавай наружу
PubAck, Msg и типы nats-py, если публичный порт этого не требует.
Маршрутизация
- Используй явный реестр «тип DTO → маршрут и codec», переданный при сборке приложения.
- Бери stream и subject из типизированной конфигурации.
- Не зашивай универсальные
created, updated, deleted, restored.
- Не выбирай маршрут по неявному анализу payload.
- Неизвестный тип DTO приводит к типизированной невосстановимой ошибке до сетевого вызова.
- Проверяй полноту и непротиворечивость реестра при подготовке адаптера.
Transport-модель и mapping
- Для каждого внешнего контракта используй отдельную строгую неизменяемую Pydantic-модель.
- Не используй DTO порта как transport-модель.
- Преобразуй поля явно чистой функцией.
- Не сериализуй DTO целиком через
asdict, model_dump или аналогичный универсальный механизм.
- Версию схемы, тип сообщения, время, correlation и causation добавляй только при наличии в контракте и DTO.
- Не генерируй отсутствующие семантические metadata.
Пример разделения ролей приведён в publisher-design.md.
Сериализация и headers
- Формат определяется контрактом: JSON, bytes или отдельный codec.
- Для JSON явно определяй представление UUID, datetime, enum, decimal и bytes.
- Сериализуй один раз непосредственно перед публикацией.
- Serializer должен быть чистым и не выполнять I/O.
- Формируй только предусмотренные контрактом headers.
- Указывай content type, schema version и encoding только когда они являются частью контракта.
- Ошибка mapping, валидации или сериализации является невосстановимой ошибкой порта и возникает до сетевого вызова.
Идентификатор и дедупликация
- Получай идентификатор сообщения из DTO порта.
- Не генерируй его в адаптере.
- Включай идентификатор в body или headers только в месте, заданном внешним
контрактом.
- Если требования отдельно используют дедупликацию JetStream, передавай его как
Nats-Msg-Id; наличие идентификатора в payload само по себе этого не требует.
- Не обещай exactly-once: потеря подтверждения может привести к повторной попытке и дубликату за пределами окна дедупликации.
Публикация
- Используй готовый
JetStreamContext, переданный composition root.
- Не создавай и не закрывай общее NATS-соединение внутри адаптера.
- Считай вызов успешным только после получения publish acknowledgement.
- Если результат порту не нужен, метод может возвращать
None.
- Не выполняй скрытые повторы внутри
publish, если они прямо не предусмотрены требованиями.
- Не создавай собственный reconnect-loop; используй настроенный reconnect клиента
nats-py.
- Не выполняй topology ensure перед каждой публикацией.
- Вызов неподготовленного адаптера заверши явной ошибкой состояния.
Ошибки порта
- Преобразуй исключения
nats-py в типизированные ошибки, предусмотренные исходящим портом.
- Не выбрасывай из адаптера application internal errors и не раскрывай технологические исключения в публичном контракте.
- Различай как минимум невосстановимые ошибки контракта или конфигурации и временные ошибки публикации, если это допускает порт.
- При неопределённом результате публикации не утверждай, что сообщение не было принято брокером.
- Сохраняй исходное исключение через
raise ... from error.
- Включай только минимальный безопасный контекст: тип сообщения, идентификатор и маршрут, если это разрешено контрактом.
- Не включай payload, чувствительные headers, credentials и адреса подключения.
Подготовка топологии
- Выполняй topology ensure в lifecycle адаптера до его передачи application-операциям.
- Имена stream и subject обязательны и приходят из конфигурации.
- Отсутствующий stream создавай автоматически из согласованной типизированной конфигурации.
- В существующий stream автоматически добавляй отсутствующие требуемые subject.
- Не удаляй лишние subject.
- При обновлении получай актуальный
StreamConfig и изменяй только subjects, сохраняя остальные параметры.
- Не меняй retention, storage, replicas, limits и другие параметры без явного требования и разрешения пользователя.
- Сравнивай только параметры, которыми конфигурация явно управляет.
- Конфликт принадлежности subject или несовместимая топология приводит к startup error.
- Не удаляй и не пересоздавай stream автоматически.
- Обеспечь корректный итог при конкурентном запуске нескольких экземпляров: после ограниченного конфликта перечитай состояние и проверь его.
Перед изменением управляемых параметров существующего stream покажи пользователю таблицу текущих и требуемых значений. Подробности — в topology-ensure.md.
Lifecycle и readiness
- Оформляй подготовку адаптера асинхронной фабрикой контекста или явной lifecycle-функцией.
- Не публикуй до успешной проверки доступности JetStream и топологии.
- Readiness хранит composition root или runtime, а не интерфейс исходящего порта.
- Повторную проверку после reconnect выполняй только по требованиям lifecycle или после фактической ошибки отсутствующей топологии.
Пакетная публикация
- Реализуй batch-метод только при его наличии в порте.
- Обычный
publish не обязан делегировать пакетному методу.
- Явно следуй требованиям частичного успеха: JetStream не делает набор отдельных publish атомарным.
- Возвращай поэлементный результат только через типы порта, без
PubAck и stream sequence.
- Ограничивай параллелизм; не создавай неограниченные задачи.
- Сохраняй порядок только по требованиям.
- Ключ упорядочивания получай из DTO порта, а не выводи из payload или subject.
Логирование
Адаптер не логирует успешные публикации, обработанные попытки и ошибки. Он
возвращает результат или типизированную ошибку порта с безопасным контекстом.
Тестирование
Unit-тестами покрой:
- mapping каждого DTO;
- выбор маршрута;
- сериализацию и headers;
- неизвестный тип и недостаточные metadata;
- классификацию ошибок;
- отсутствие сетевого вызова при локальной ошибке;
- отсутствие логирования.
Интеграционными тестами с настоящим JetStream покрой:
- создание отсутствующего stream;
- добавление отсутствующих subject;
- сохранение посторонних subject и неуправляемых параметров;
- повторную и конкурентную подготовку;
- конфликт топологии;
- реальную публикацию payload и headers;
- получение publish acknowledgement;
- преобразование ошибки недоступного JetStream.
Инфраструктура поднимается автоматически, использует уникальные имена и освобождается после тестов. Проверяй собственное поведение адаптера, а не внутренности nats-py.
Границы
В область скила входят реализация исходящего порта, маршрутизация, transport-модели, mapping, сериализация, headers, publish acknowledgement, ошибки порта, topology ensure и тесты.
Не входят периодический или NATS-специализированный publisher worker, consumer, общий runtime, outbox, Unit of Work, репозитории, application/domain-логика и HTTP.
Проверка результата
- DTO порта достаточен и не содержит NATS-маршрутизацию.
- Mapping и сериализация явные.
- Маршруты конфигурируемы и проверены до первой публикации.
- Отсутствующие stream и subject обеспечиваются безопасно.
- Остальная топология не изменяется без разрешения.
- Успех подтверждён JetStream.
- Технологические ошибки не пересекают порт.
- Скрытые повторы отсутствуют.
- Адаптер ничего не логирует.
- Unit- и integration-тесты покрывают контракт и топологию.
1---2name: python-nats-jetstream-publisher-adapter-writing3description: Используй при реализации или правке исходящего NATS JetStream publisher-адаптера на Python через nats-py и Pydantic: реализация application-порта, явный mapping DTO в transport-модель, маршрутизация, сериализация, headers, PubAck, типизированные ошибки и безопасное создание stream или добавление subject. Не применять для воркеров, consumer-ов, outbox, Unit of Work и application/domain-логики.4---56# Publisher-адаптер NATS JetStream на Python78Реализуй исходящий технологический адаптер, который преобразует DTO порта в согласованное сообщение JetStream и считает публикацию успешной только после подтверждения брокера.910## Порядок работы11121. Изучи исходящий порт, DTO, формат сообщений и маршруты.132. Проверь, достаточно ли DTO для публикации без обращения к другим источникам.143. Зафиксируй типизированную конфигурацию маршрутов и управляемой топологии.154. Реализуй явные transport-модели, mapping, сериализацию и headers.165. Реализуй подготовку топологии до первой публикации.176. Реализуй publish и преобразование технологических ошибок в ошибки порта.187. Добавь unit- и integration-тесты.1920Если требования полны, следуй им без дополнительного проектирования. Уточняй только отсутствующие решения, без которых нельзя безопасно реализовать контракт или изменить существующую топологию.2122## Исходящий порт и DTO2324- Адаптер реализует уже определённый application-порт.25- Принимай только DTO этого порта.26- Если порт принимает закрытое объединение конкретных DTO, используй объявленный27 application `TypeAlias` в сигнатуре и выполняй явный исчерпывающий dispatch по28 runtime-типу. Не заменяй alias `Any`, общим transport DTO или перечислением29 типов ресурсов внутри application.30- DTO содержит семантические данные и необходимые metadata, но не имена stream, subject и параметры JetStream.31- Не задавай универсальный набор полей DTO.32- Не извлекай недостающие данные из репозиториев, доменных объектов или других application-операций.33- Недостаточный DTO приводит к типизированной невосстановимой ошибке порта до обращения к NATS.34- Не передавай наружу `PubAck`, `Msg` и типы `nats-py`, если публичный порт этого не требует.3536## Маршрутизация3738- Используй явный реестр «тип DTO → маршрут и codec», переданный при сборке приложения.39- Бери stream и subject из типизированной конфигурации.40- Не зашивай универсальные `created`, `updated`, `deleted`, `restored`.41- Не выбирай маршрут по неявному анализу payload.42- Неизвестный тип DTO приводит к типизированной невосстановимой ошибке до сетевого вызова.43- Проверяй полноту и непротиворечивость реестра при подготовке адаптера.4445## Transport-модель и mapping4647- Для каждого внешнего контракта используй отдельную строгую неизменяемую Pydantic-модель.48- Не используй DTO порта как transport-модель.49- Преобразуй поля явно чистой функцией.50- Не сериализуй DTO целиком через `asdict`, `model_dump` или аналогичный универсальный механизм.51- Версию схемы, тип сообщения, время, correlation и causation добавляй только при наличии в контракте и DTO.52- Не генерируй отсутствующие семантические metadata.5354Пример разделения ролей приведён в [publisher-design.md](references/publisher-design.md).5556## Сериализация и headers5758- Формат определяется контрактом: JSON, bytes или отдельный codec.59- Для JSON явно определяй представление UUID, datetime, enum, decimal и bytes.60- Сериализуй один раз непосредственно перед публикацией.61- Serializer должен быть чистым и не выполнять I/O.62- Формируй только предусмотренные контрактом headers.63- Указывай content type, schema version и encoding только когда они являются частью контракта.64- Ошибка mapping, валидации или сериализации является невосстановимой ошибкой порта и возникает до сетевого вызова.6566## Идентификатор и дедупликация6768- Получай идентификатор сообщения из DTO порта.69- Не генерируй его в адаптере.70- Включай идентификатор в body или headers только в месте, заданном внешним71 контрактом.72- Если требования отдельно используют дедупликацию JetStream, передавай его как73 `Nats-Msg-Id`; наличие идентификатора в payload само по себе этого не требует.74- Не обещай exactly-once: потеря подтверждения может привести к повторной попытке и дубликату за пределами окна дедупликации.7576## Публикация7778- Используй готовый `JetStreamContext`, переданный composition root.79- Не создавай и не закрывай общее NATS-соединение внутри адаптера.80- Считай вызов успешным только после получения publish acknowledgement.81- Если результат порту не нужен, метод может возвращать `None`.82- Не выполняй скрытые повторы внутри `publish`, если они прямо не предусмотрены требованиями.83- Не создавай собственный reconnect-loop; используй настроенный reconnect клиента `nats-py`.84- Не выполняй topology ensure перед каждой публикацией.85- Вызов неподготовленного адаптера заверши явной ошибкой состояния.8687## Ошибки порта8889- Преобразуй исключения `nats-py` в типизированные ошибки, предусмотренные исходящим портом.90- Не выбрасывай из адаптера application internal errors и не раскрывай технологические исключения в публичном контракте.91- Различай как минимум невосстановимые ошибки контракта или конфигурации и временные ошибки публикации, если это допускает порт.92- При неопределённом результате публикации не утверждай, что сообщение не было принято брокером.93- Сохраняй исходное исключение через `raise ... from error`.94- Включай только минимальный безопасный контекст: тип сообщения, идентификатор и маршрут, если это разрешено контрактом.95- Не включай payload, чувствительные headers, credentials и адреса подключения.9697## Подготовка топологии9899- Выполняй topology ensure в lifecycle адаптера до его передачи application-операциям.100- Имена stream и subject обязательны и приходят из конфигурации.101- Отсутствующий stream создавай автоматически из согласованной типизированной конфигурации.102- В существующий stream автоматически добавляй отсутствующие требуемые subject.103- Не удаляй лишние subject.104- При обновлении получай актуальный `StreamConfig` и изменяй только subjects, сохраняя остальные параметры.105- Не меняй retention, storage, replicas, limits и другие параметры без явного требования и разрешения пользователя.106- Сравнивай только параметры, которыми конфигурация явно управляет.107- Конфликт принадлежности subject или несовместимая топология приводит к startup error.108- Не удаляй и не пересоздавай stream автоматически.109- Обеспечь корректный итог при конкурентном запуске нескольких экземпляров: после ограниченного конфликта перечитай состояние и проверь его.110111Перед изменением управляемых параметров существующего stream покажи пользователю таблицу текущих и требуемых значений. Подробности — в [topology-ensure.md](references/topology-ensure.md).112113## Lifecycle и readiness114115- Оформляй подготовку адаптера асинхронной фабрикой контекста или явной lifecycle-функцией.116- Не публикуй до успешной проверки доступности JetStream и топологии.117- Readiness хранит composition root или runtime, а не интерфейс исходящего порта.118- Повторную проверку после reconnect выполняй только по требованиям lifecycle или после фактической ошибки отсутствующей топологии.119120## Пакетная публикация121122- Реализуй batch-метод только при его наличии в порте.123- Обычный `publish` не обязан делегировать пакетному методу.124- Явно следуй требованиям частичного успеха: JetStream не делает набор отдельных publish атомарным.125- Возвращай поэлементный результат только через типы порта, без `PubAck` и stream sequence.126- Ограничивай параллелизм; не создавай неограниченные задачи.127- Сохраняй порядок только по требованиям.128- Ключ упорядочивания получай из DTO порта, а не выводи из payload или subject.129130## Логирование131132Адаптер не логирует успешные публикации, обработанные попытки и ошибки. Он133возвращает результат или типизированную ошибку порта с безопасным контекстом.134135## Тестирование136137Unit-тестами покрой:138139- mapping каждого DTO;140- выбор маршрута;141- сериализацию и headers;142- неизвестный тип и недостаточные metadata;143- классификацию ошибок;144- отсутствие сетевого вызова при локальной ошибке;145- отсутствие логирования.146147Интеграционными тестами с настоящим JetStream покрой:148149- создание отсутствующего stream;150- добавление отсутствующих subject;151- сохранение посторонних subject и неуправляемых параметров;152- повторную и конкурентную подготовку;153- конфликт топологии;154- реальную публикацию payload и headers;155- получение publish acknowledgement;156- преобразование ошибки недоступного JetStream.157158Инфраструктура поднимается автоматически, использует уникальные имена и освобождается после тестов. Проверяй собственное поведение адаптера, а не внутренности `nats-py`.159160## Границы161162В область скила входят реализация исходящего порта, маршрутизация, transport-модели, mapping, сериализация, headers, publish acknowledgement, ошибки порта, topology ensure и тесты.163164Не входят периодический или NATS-специализированный publisher worker, consumer, общий runtime, outbox, Unit of Work, репозитории, application/domain-логика и HTTP.165166## Проверка результата167168- DTO порта достаточен и не содержит NATS-маршрутизацию.169- Mapping и сериализация явные.170- Маршруты конфигурируемы и проверены до первой публикации.171- Отсутствующие stream и subject обеспечиваются безопасно.172- Остальная топология не изменяется без разрешения.173- Успех подтверждён JetStream.174- Технологические ошибки не пересекают порт.175- Скрытые повторы отсутствуют.176- Адаптер ничего не логирует.177- Unit- и integration-тесты покрывают контракт и топологию.