# Python Nats Jetstream Consumer Writing

> Используй при реализации или правке входящего NATS JetStream consumer на Python через nats-py: выбор владения и подготовки stream/subjects и durable consumer, повтор подписки при отсутствующей топологии, pull consumer, envelope и payload, mapping в application-вход, bounded concurrency, backpressure, ACK/NAK/TERM, логирование обработки, in-progress, redelivery, reconnect, идемпотентность и необязательный DLQ. Не применять для publisher, общего runtime процесса, исходящих адаптеров и application/domain-логики.

- Skill: `nemagu/python-nats-jetstream-consumer-writing` (Agent Skill, multi-file: 4 files)
- Install (CLI): `npx skillmds@latest add nemagu/python-nats-jetstream-consumer-writing`
- Raw SKILL.md: https://api.skillmd.com/api/skills/nemagu/python-nats-jetstream-consumer-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-consumer-writing

---


# NATS JetStream consumer на Python

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

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

1. Изучи контракт сообщения, публичный вход application-операции и требования доставки.
2. Зафиксируй consumer-конфигурацию, владельца stream/subjects, владельца durable
   consumer-а и политику повторной подписки.
3. Реализуй выбранный режим topology ensure до подписки; перед изменением уже
   существующей управляемой topology покажи таблицу текущих и требуемых значений.
4. Определи envelope, payload, mapping и политику результата сообщения.
5. Реализуй pull-loop с ограниченным параллелизмом и backpressure.
6. Добавь reconnect, progress acknowledgement и DLQ только по требованиям.
7. Добавь unit- и integration-тесты.

Если контракт полон, реализуй его без опросника. Не выводи право изменять stream,
subjects или durable consumer из самого факта подписки. Если ownership не указан,
сначала определи его из нормативной документации и существующего composition root;
для межсервисного входящего события считай stream/subjects принадлежащими publisher-у,
если нет явного обратного требования. Уточняй только отсутствующие решения, без которых
нельзя безопасно определить владельца, семантику доставки или разрешённое изменение
топологии.

## Consumer-контракт

Из требований должны быть известны:

- stream и filter subjects;
- режим владения stream/subjects и, для consumer-managed режима, управляемые
  параметры stream, достаточные для его создания;
- durable name;
- режим provisioning-а durable consumer-а;
- delivery policy и начальная позиция;
- explicit ack policy;
- `ack_wait`, `max_deliver` и backoff;
- batch size, fetch timeout и ограничения ожидающих pull;
- допустимый параллелизм и требования порядка;
- pending limits;
- положительная задержка или backoff повторной подписки с явными единицами;
- политика существующего consumer.

Не задавай универсальные значения. Отсутствующий параметр уточняй только тогда,
когда без него нельзя безопасно создать топологию, определить семантику доставки
или реализовать контролируемый retry без busy loop.

## Одно сообщение — одна операция

Для каждого сообщения выполняй последовательность:

1. извлечение JetStream metadata;
2. декодирование envelope и payload;
3. транспортная валидация;
4. явное преобразование в публичный application-вход;
5. один вызов operation;
6. выбор и отправка ACK-действия.

Не создавай в consumer Unit of Work и репозитории. Не вызывай несколько независимых application-операций. Composition root передаёт готовую operation со всеми портами.

## Транспортные модели

- Используй отдельные неизменяемые строгие Pydantic-модели envelope и payload.
- Запрещай неизвестные поля, если совместимость контракта не требует иного.
- Расположение metadata в headers или body определяется контрактом.
- Не используй application DTO как транспортную модель.
- Преобразуй поля явно и чистой функцией.
- Не передавай в application `Msg`, Pydantic-модель или JetStream metadata, если этого нет в публичном контракте.
- Проверяй поддерживаемую версию схемы.

## ACK, NAK и TERM

- Отправляй ACK только после успешного завершения application-операции.
- Используй server-confirmed ACK, когда это требуется контрактом доставки.
- NAK с задержкой применяй только для повторяемого результата.
- Задавай задержку NAK отдельным положительным целочисленным параметром
  конфигурации в миллисекундах с явным суффиксом `_milliseconds`. На границе
  `nats-py` явно преобразуй значение в секунды для `message.nak(delay=...)`.
  Не подменяй delayed NAK параметрами `ack_wait`, `max_deliver` или немедленным
  `message.nak()`; конкретное значение и default должны следовать требованиям.
- TERM применяй только для заведомо постоянной ошибки сообщения.
- Не считай TERM помещением в DLQ.
- Для долгой обработки отправляй `in_progress` по требованиям.
- При отмене не отправляй ACK или TERM.
- Ошибка отправки ACK означает, что сообщение может быть доставлено повторно.

Решение принимает явная чистая политика уровня presentation, основанная на публичных результатах и ошибках application-операции. Не импортируй domain errors и внутренние инфраструктурные исключения. Подробности — в [delivery-policy.md](references/delivery-policy.md).

## Параллелизм и backpressure

- Ограничивай число сообщений в обработке конфигурируемой ёмкостью.
- Не выполняй pull, когда локальная ёмкость исчерпана.
- Размер fetch не должен превышать свободную ёмкость.
- Управляй обработчиками через ограниченный `TaskGroup`; не создавай неограниченные задачи.
- Ошибка одного сообщения разрешается его политикой и не отменяет успешно обработанные сообщения.
- Если для отсутствующей topology согласован retry, отсутствие stream, subject
  или durable переводит consumer в not-ready и цикл повторной подготовки, а не
  завершает обязательную runtime-задачу.
- Непредусмотренная ошибка fetch/subscription loop остаётся ошибкой обязательной
  runtime-задачи.
- При строгом порядке используй concurrency `1`.

## Reconnect

- Используй встроенный reconnect `nats-py` с заданными ограничениями и задержками.
- Не создавай параллельный ручной reconnect-loop.
- Callbacks могут менять readiness и пробуждать ожидающие задачи.
- После восстановления пересоздавай pull subscription только при необходимости.
- Исчерпание reconnect завершает задачу с ошибкой.
- Ожидания reconnect не должны создавать busy loop и должны прерываться отменой.

## Топология consumer

Выбирай независимо режим владения stream/subjects и режим provisioning-а durable.
Не смешивай их: право создать собственный durable не даёт права менять stream.

### Stream и subjects

- **Publisher-owned / external** — основной вариант для межсервисных событий.
  Publisher обеспечивает stream/subjects до публикации по правилам профильного
  publisher-адаптера. Consumer проверяет, что subject существует и принадлежит
  ожидаемому stream, но никогда не создаёт, не расширяет, не обновляет и не
  удаляет stream/subjects. Отсутствие внешней topology обрабатывает как
  зависимость, которая ещё не готова, согласно политике retry/readiness.
- **Consumer-managed additive** — допустим только когда контракт явно назначает
  consumer владельцем этой stream topology. Отсутствующий stream создавай из
  полной типизированной конфигурации; в существующий stream добавляй только
  отсутствующие требуемые subjects, сохраняя остальные subjects и все
  неуправляемые параметры.
- **Platform-provisioned** — stream/subjects создаёт deployment, init job или
  оператор. Runtime consumer выполняет только проверку; реакцию на отсутствие
  topology — retry или startup failure — закрепи эксплуатационным контрактом.

Для любого режима не удаляй и не пересоздавай stream автоматически. Не меняй
retention, storage, replicas, limits и другие существующие параметры без явного
владения, требования и разрешения. Конфликт принадлежности subject и несовместимая
существующая topology являются фатальными ошибками конфигурации, а не поводом для
бесконечного retry.

### Durable consumer

- **Consumer-owned ensure** — сервис владеет уникальным durable и идемпотентно
  создаёт его при отсутствии с полным явным `ConsumerConfig`, затем перечитывает
  и валидирует итог. Этот режим не разрешает менять stream/subjects.
- **Externally provisioned bind-only** — deployment или оператор создаёт durable;
  runtime только получает `consumer_info`, валидирует контракт и привязывается.
  Отсутствующий durable не создаётся неявно.

В обоих режимах сравнивай существующую и требуемую конфигурацию. Несовместимый
durable не изменяй, не удаляй и не пересоздавай автоматически; заверши подготовку
явной topology-ошибкой с несовпадающими полями. Имена durable должны быть
стабильными и уникальными для логического подписчика, чтобы независимые сервисы
не делили одну очередь сообщений случайно.

Не используй `pull_subscribe` как скрытый provisioning. После ensure/validation
привязывайся через bind-only API (`pull_subscribe_bind` либо точный эквивалент
используемой версии клиента). При consumer-owned ensure обработай гонку нескольких
реплик по схеме create-or-observe: после конфликта перечитай durable и проверь его.

Отделяй topology ensure от цикла обработки сообщений, но вызывай его снова после
фактической ошибки отсутствующей topology. Подробный выбор режима, алгоритмы и
классификация исходов приведены в
[topology-and-subscription-retry.md](references/topology-and-subscription-retry.md).

## Повтор подписки и readiness

- Ошибку отсутствующего stream, покрытия subject или durable считай временным
  состоянием подписки только в режиме с согласованным retry. Если контракт требует
  fail-fast при ошибке platform provisioning, заверши startup явной ошибкой.
- Снимай readiness до topology ensure и возвращай её только после успешной привязки
  subscription; liveness heartbeat процесса при этом продолжает работать.
- После неуспешной попытки ожидай настроенную задержку или backoff без busy loop.
  Ожидание должно немедленно прерываться stop event или отменой.
- Повторяй последовательность `stream check/ensure -> durable check/ensure -> bind`
  до успеха или остановки, если выбран retry и требования не задают предел попыток.
- Не маскируй retry-loop-ом неверную конфигурацию, несовместимую топологию,
  исчерпание reconnect или неизвестную ошибку `nats-py`.
- Логируй переходы состояния и итог попытки согласно logging-контракту; не создавай
  одинаковую error-запись на каждой попытке при длительном отсутствии топологии.

## At-least-once и идемпотентность

- Явно учитывай возможность повторной доставки, включая потерю ACK после успешной операции.
- Не обещай exactly-once.
- Передавай message/event ID в application-вход, если это предусмотрено контрактом.
- Не создавай универсальный dedup cache в consumer.
- Атомарная дедупликация и бизнес-изменение принадлежат application-операции.
- Одинаковый ID и одинаковое содержимое должны давать идемпотентный результат.
- Одинаковый ID с другим содержимым должен считаться конфликтом.
- Результат «уже обработано» подтверждай ACK.
- Не генерируй отсутствующий идентификатор.

## DLQ

DLQ добавляй только по требованиям:

- `max_deliver` сам по себе не создаёт DLQ;
- TERM не является DLQ;
- сохраняй исходный payload и только безопасные metadata и причину;
- окончательно подтверждай исходное сообщение только после успешной записи в DLQ;
- ошибка DLQ не должна приводить к потере исходного сообщения;
- replay сохраняет исходные идентификаторы корреляции и идемпотентности;
- JetStream advisories используй для наблюдаемости, а не как хранилище DLQ.

Если DLQ отсутствует, предусмотренные требованиями метрики и оповещения должны позволять обнаружить исчерпание доставок.

## Логирование

Применяй `python-service-logging-writing` и logging-профиль consumer-а. Consumer
владеет operation context сообщения и одной итоговой записью после определения
и выполнения транспортного исхода.

- Связывай согласованные metadata сообщения до валидации payload, чтобы
  окончательный отказ имел диагностический контекст.
- Сохраняй стабильный message ID при redelivery; не генерируй его вопреки
  контракту владельца.
- Измеряй полную длительность от принятия сообщения до итогового ACK-действия
  монотонными часами.
- Записывай ACK, NAK или TERM только после фактического завершения действия. Если
  транспортный сбой нельзя выразить согласованным outcome, сначала дополни
  logging-контракт, а не выдавай намерение за результат.
- Для outcome и других классификаторов используй перечисления logging-слоя, а не
  перечисления application DTO. Даже при совпадении значений один к одному
  выполняй явное исчерпывающее преобразование.
- Очищай контекст сообщения в `finally`, включая validation error, timeout,
  cancellation и ошибку отправки подтверждения.
- Не логируй payload, чувствительные headers и произвольный текст исключения.
- Application, delivery policy и NATS-адаптер не повторяют итоговую запись или
  stack trace consumer-а.

## Тестирование

Unit-тестами покрой:

- выбранные режимы ownership/provisioning и запрет недопустимых мутаций;
- для consumer-managed режима — план создания отсутствующего stream, добавление
  только отсутствующих subjects и сохранение посторонних subjects и
  неуправляемых параметров;
- для external/platform stream — отсутствие вызовов create/update при любой
  реакции на отсутствующую topology;
- создание собственного durable либо validate/bind externally provisioned durable
  согласно выбранному режиму;
- валидацию совместимого durable и отказ от мутации несовместимого;
- классификацию retryable и fatal ошибок topology/subscription;
- stop-aware паузу, повтор ensure/subscription и переходы readiness;
- envelope, payload и mapping;
- таблицу публичный результат → ACK/NAK/TERM;
- передачу настроенной задержки в `message.nak(delay=...)` для повторяемого
  результата;
- версию схемы и неизвестные поля;
- расчёт свободной ёмкости;
- отмену и отсутствие подтверждения при cancellation.

Интеграционными тестами с настоящим JetStream покрой:

- consumer-managed stream: создание отсутствующего stream и добавление
  отсутствующего subject без удаления существующих;
- external/platform stream: отсутствие его создания или изменения consumer-ом;
- consumer-owned durable: создание с полным контрактом, повторный и конкурентный
  ensure и последующий bind;
- externally provisioned durable: validate/bind без неявного создания;
- при выбранном retry — запуск до появления topology и успешную подписку после
  повторной попытки;
- конкурентную подготовку topology несколькими экземплярами;
- конфликт или несовместимость топологии как фатальную ошибку;
- ACK после успеха;
- delayed NAK и redelivery;
- TERM;
- reconnect;
- повтор после потерянного ACK;
- `in_progress`;
- concurrency и backpressure;
- mismatch consumer-конфигурации;
- DLQ, если он предусмотрен.

Тестовая инфраструктура должна подниматься автоматически с уникальными именами и освобождаться после прогона.

## Границы

В область скила входят выбор ownership-модели, разрешённый additive topology ensure
для stream/subjects, ensure или bind-only durable consumer-а, повтор подписки,
pull consumer, транспортная валидация и mapping, bounded
concurrency, backpressure, ACK policy, reconnect, redelivery, `in_progress`, metadata
идемпотентности, необязательный DLQ, логирование и тесты.

Не входят общий runtime и сигналы процесса, publisher, исходящий адаптер, Unit of Work, репозитории, реализация application/domain-логики и HTTP.

## Проверка результата

- Consumer-конфигурация соответствует требованиям.
- Режимы владения stream/subjects и provisioning-а durable заданы явно и
  реализованы независимо.
- В consumer-managed режиме отсутствующий stream создаётся из полной конфигурации,
  а subjects добавляются без разрушительных изменений.
- В publisher-owned и platform-provisioned режимах consumer не создаёт и не
  изменяет stream/subjects.
- Собственный durable явно создаётся и валидируется либо только валидируется и
  привязывается — согласно выбранному режиму; скрытого provisioning-а нет.
- При выбранном retry отсутствующая topology не завершает процесс: readiness
  снята, ожидание stop-aware, после паузы выполняется новая попытка подготовки и
  подписки; при согласованном fail-fast возвращается явная startup error.
- Конфликт и несовместимость топологии не скрываются бесконечными повторами.
- Одно сообщение вызывает одну готовую application-операцию.
- Транспортные модели не пересекают application-границу.
- ACK отправляется только после успеха.
- Параллелизм ограничен и pull учитывает свободную ёмкость.
- Reconnect не дублируется ручным циклом.
- At-least-once и идемпотентность отражены явно.
- Топология не изменяется разрушительно.
- Поля сообщения, длительность и итог ACK/NAK/TERM соответствуют logging-профилю,
  контекст не протекает между конкурентными обработчиками.
- Unit- и integration-тесты проверяют собственную логику, а не внутренности `nats-py`.

