# Python Nats Jetstream Publisher Adapter Writing

> Используй при реализации или правке исходящего NATS JetStream publisher-адаптера на Python через nats-py и Pydantic: реализация application-порта, явный mapping DTO в transport-модель, маршрутизация, сериализация, headers, PubAck, типизированные ошибки и безопасное создание stream или добавление subject. Не применять для воркеров, consumer-ов, outbox, Unit of Work и application/domain-логики.

- Skill: `nemagu/python-nats-jetstream-publisher-adapter-writing` (Agent Skill, multi-file: 4 files)
- Install (CLI): `npx skillmds@latest add nemagu/python-nats-jetstream-publisher-adapter-writing`
- Raw SKILL.md: https://api.skillmd.com/api/skills/nemagu/python-nats-jetstream-publisher-adapter-writing/raw
- Safety review: pending
- Works with: Claude Code, Claude.ai, OpenAI Codex
- Category: Coding & Dev Tools
- Author: Nemagu (https://skillmd.com/u/nemagu)
- Updated: 2026-09-21
- Page: https://skillmd.com/skills/nemagu/python-nats-jetstream-publisher-adapter-writing

---


# Publisher-адаптер NATS JetStream на Python

Реализуй исходящий технологический адаптер, который преобразует DTO порта в согласованное сообщение JetStream и считает публикацию успешной только после подтверждения брокера.

## Порядок работы

1. Изучи исходящий порт, DTO, формат сообщений и маршруты.
2. Проверь, достаточно ли DTO для публикации без обращения к другим источникам.
3. Зафиксируй типизированную конфигурацию маршрутов и управляемой топологии.
4. Реализуй явные transport-модели, mapping, сериализацию и headers.
5. Реализуй подготовку топологии до первой публикации.
6. Реализуй publish и преобразование технологических ошибок в ошибки порта.
7. Добавь 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](references/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](references/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-тесты покрывают контракт и топологию.

