NATS JetStream consumer на Python
Реализуй входящую presentation-границу JetStream, которая надёжно получает сообщение, проверяет транспортный контракт, вызывает одну application-операцию и подтверждает результат брокеру.
Порядок работы
- Изучи контракт сообщения, публичный вход application-операции и требования доставки.
- Зафиксируй consumer-конфигурацию, владельца stream/subjects, владельца durable
consumer-а и политику повторной подписки.
- Реализуй выбранный режим topology ensure до подписки; перед изменением уже
существующей управляемой topology покажи таблицу текущих и требуемых значений.
- Определи envelope, payload, mapping и политику результата сообщения.
- Реализуй pull-loop с ограниченным параллелизмом и backpressure.
- Добавь reconnect, progress acknowledgement и DLQ только по требованиям.
- Добавь 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.
Одно сообщение — одна операция
Для каждого сообщения выполняй последовательность:
- извлечение JetStream metadata;
- декодирование envelope и payload;
- транспортная валидация;
- явное преобразование в публичный application-вход;
- один вызов operation;
- выбор и отправка 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.
Параллелизм и 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.
Повтор подписки и 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.
1---2name: python-nats-jetstream-consumer-writing3description: Используй при реализации или правке входящего 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-логики.4---56# NATS JetStream consumer на Python78Реализуй входящую presentation-границу JetStream, которая надёжно получает сообщение, проверяет транспортный контракт, вызывает одну application-операцию и подтверждает результат брокеру.910## Порядок работы11121. Изучи контракт сообщения, публичный вход application-операции и требования доставки.132. Зафиксируй consumer-конфигурацию, владельца stream/subjects, владельца durable14 consumer-а и политику повторной подписки.153. Реализуй выбранный режим topology ensure до подписки; перед изменением уже16 существующей управляемой topology покажи таблицу текущих и требуемых значений.174. Определи envelope, payload, mapping и политику результата сообщения.185. Реализуй pull-loop с ограниченным параллелизмом и backpressure.196. Добавь reconnect, progress acknowledgement и DLQ только по требованиям.207. Добавь unit- и integration-тесты.2122Если контракт полон, реализуй его без опросника. Не выводи право изменять stream,23subjects или durable consumer из самого факта подписки. Если ownership не указан,24сначала определи его из нормативной документации и существующего composition root;25для межсервисного входящего события считай stream/subjects принадлежащими publisher-у,26если нет явного обратного требования. Уточняй только отсутствующие решения, без которых27нельзя безопасно определить владельца, семантику доставки или разрешённое изменение28топологии.2930## Consumer-контракт3132Из требований должны быть известны:3334- stream и filter subjects;35- режим владения stream/subjects и, для consumer-managed режима, управляемые36 параметры stream, достаточные для его создания;37- durable name;38- режим provisioning-а durable consumer-а;39- delivery policy и начальная позиция;40- explicit ack policy;41- `ack_wait`, `max_deliver` и backoff;42- batch size, fetch timeout и ограничения ожидающих pull;43- допустимый параллелизм и требования порядка;44- pending limits;45- положительная задержка или backoff повторной подписки с явными единицами;46- политика существующего consumer.4748Не задавай универсальные значения. Отсутствующий параметр уточняй только тогда,49когда без него нельзя безопасно создать топологию, определить семантику доставки50или реализовать контролируемый retry без busy loop.5152## Одно сообщение — одна операция5354Для каждого сообщения выполняй последовательность:55561. извлечение JetStream metadata;572. декодирование envelope и payload;583. транспортная валидация;594. явное преобразование в публичный application-вход;605. один вызов operation;616. выбор и отправка ACK-действия.6263Не создавай в consumer Unit of Work и репозитории. Не вызывай несколько независимых application-операций. Composition root передаёт готовую operation со всеми портами.6465## Транспортные модели6667- Используй отдельные неизменяемые строгие Pydantic-модели envelope и payload.68- Запрещай неизвестные поля, если совместимость контракта не требует иного.69- Расположение metadata в headers или body определяется контрактом.70- Не используй application DTO как транспортную модель.71- Преобразуй поля явно и чистой функцией.72- Не передавай в application `Msg`, Pydantic-модель или JetStream metadata, если этого нет в публичном контракте.73- Проверяй поддерживаемую версию схемы.7475## ACK, NAK и TERM7677- Отправляй ACK только после успешного завершения application-операции.78- Используй server-confirmed ACK, когда это требуется контрактом доставки.79- NAK с задержкой применяй только для повторяемого результата.80- Задавай задержку NAK отдельным положительным целочисленным параметром81 конфигурации в миллисекундах с явным суффиксом `_milliseconds`. На границе82 `nats-py` явно преобразуй значение в секунды для `message.nak(delay=...)`.83 Не подменяй delayed NAK параметрами `ack_wait`, `max_deliver` или немедленным84 `message.nak()`; конкретное значение и default должны следовать требованиям.85- TERM применяй только для заведомо постоянной ошибки сообщения.86- Не считай TERM помещением в DLQ.87- Для долгой обработки отправляй `in_progress` по требованиям.88- При отмене не отправляй ACK или TERM.89- Ошибка отправки ACK означает, что сообщение может быть доставлено повторно.9091Решение принимает явная чистая политика уровня presentation, основанная на публичных результатах и ошибках application-операции. Не импортируй domain errors и внутренние инфраструктурные исключения. Подробности — в [delivery-policy.md](references/delivery-policy.md).9293## Параллелизм и backpressure9495- Ограничивай число сообщений в обработке конфигурируемой ёмкостью.96- Не выполняй pull, когда локальная ёмкость исчерпана.97- Размер fetch не должен превышать свободную ёмкость.98- Управляй обработчиками через ограниченный `TaskGroup`; не создавай неограниченные задачи.99- Ошибка одного сообщения разрешается его политикой и не отменяет успешно обработанные сообщения.100- Если для отсутствующей topology согласован retry, отсутствие stream, subject101 или durable переводит consumer в not-ready и цикл повторной подготовки, а не102 завершает обязательную runtime-задачу.103- Непредусмотренная ошибка fetch/subscription loop остаётся ошибкой обязательной104 runtime-задачи.105- При строгом порядке используй concurrency `1`.106107## Reconnect108109- Используй встроенный reconnect `nats-py` с заданными ограничениями и задержками.110- Не создавай параллельный ручной reconnect-loop.111- Callbacks могут менять readiness и пробуждать ожидающие задачи.112- После восстановления пересоздавай pull subscription только при необходимости.113- Исчерпание reconnect завершает задачу с ошибкой.114- Ожидания reconnect не должны создавать busy loop и должны прерываться отменой.115116## Топология consumer117118Выбирай независимо режим владения stream/subjects и режим provisioning-а durable.119Не смешивай их: право создать собственный durable не даёт права менять stream.120121### Stream и subjects122123- **Publisher-owned / external** — основной вариант для межсервисных событий.124 Publisher обеспечивает stream/subjects до публикации по правилам профильного125 publisher-адаптера. Consumer проверяет, что subject существует и принадлежит126 ожидаемому stream, но никогда не создаёт, не расширяет, не обновляет и не127 удаляет stream/subjects. Отсутствие внешней topology обрабатывает как128 зависимость, которая ещё не готова, согласно политике retry/readiness.129- **Consumer-managed additive** — допустим только когда контракт явно назначает130 consumer владельцем этой stream topology. Отсутствующий stream создавай из131 полной типизированной конфигурации; в существующий stream добавляй только132 отсутствующие требуемые subjects, сохраняя остальные subjects и все133 неуправляемые параметры.134- **Platform-provisioned** — stream/subjects создаёт deployment, init job или135 оператор. Runtime consumer выполняет только проверку; реакцию на отсутствие136 topology — retry или startup failure — закрепи эксплуатационным контрактом.137138Для любого режима не удаляй и не пересоздавай stream автоматически. Не меняй139retention, storage, replicas, limits и другие существующие параметры без явного140владения, требования и разрешения. Конфликт принадлежности subject и несовместимая141существующая topology являются фатальными ошибками конфигурации, а не поводом для142бесконечного retry.143144### Durable consumer145146- **Consumer-owned ensure** — сервис владеет уникальным durable и идемпотентно147 создаёт его при отсутствии с полным явным `ConsumerConfig`, затем перечитывает148 и валидирует итог. Этот режим не разрешает менять stream/subjects.149- **Externally provisioned bind-only** — deployment или оператор создаёт durable;150 runtime только получает `consumer_info`, валидирует контракт и привязывается.151 Отсутствующий durable не создаётся неявно.152153В обоих режимах сравнивай существующую и требуемую конфигурацию. Несовместимый154durable не изменяй, не удаляй и не пересоздавай автоматически; заверши подготовку155явной topology-ошибкой с несовпадающими полями. Имена durable должны быть156стабильными и уникальными для логического подписчика, чтобы независимые сервисы157не делили одну очередь сообщений случайно.158159Не используй `pull_subscribe` как скрытый provisioning. После ensure/validation160привязывайся через bind-only API (`pull_subscribe_bind` либо точный эквивалент161используемой версии клиента). При consumer-owned ensure обработай гонку нескольких162реплик по схеме create-or-observe: после конфликта перечитай durable и проверь его.163164Отделяй topology ensure от цикла обработки сообщений, но вызывай его снова после165фактической ошибки отсутствующей topology. Подробный выбор режима, алгоритмы и166классификация исходов приведены в167[topology-and-subscription-retry.md](references/topology-and-subscription-retry.md).168169## Повтор подписки и readiness170171- Ошибку отсутствующего stream, покрытия subject или durable считай временным172 состоянием подписки только в режиме с согласованным retry. Если контракт требует173 fail-fast при ошибке platform provisioning, заверши startup явной ошибкой.174- Снимай readiness до topology ensure и возвращай её только после успешной привязки175 subscription; liveness heartbeat процесса при этом продолжает работать.176- После неуспешной попытки ожидай настроенную задержку или backoff без busy loop.177 Ожидание должно немедленно прерываться stop event или отменой.178- Повторяй последовательность `stream check/ensure -> durable check/ensure -> bind`179 до успеха или остановки, если выбран retry и требования не задают предел попыток.180- Не маскируй retry-loop-ом неверную конфигурацию, несовместимую топологию,181 исчерпание reconnect или неизвестную ошибку `nats-py`.182- Логируй переходы состояния и итог попытки согласно logging-контракту; не создавай183 одинаковую error-запись на каждой попытке при длительном отсутствии топологии.184185## At-least-once и идемпотентность186187- Явно учитывай возможность повторной доставки, включая потерю ACK после успешной операции.188- Не обещай exactly-once.189- Передавай message/event ID в application-вход, если это предусмотрено контрактом.190- Не создавай универсальный dedup cache в consumer.191- Атомарная дедупликация и бизнес-изменение принадлежат application-операции.192- Одинаковый ID и одинаковое содержимое должны давать идемпотентный результат.193- Одинаковый ID с другим содержимым должен считаться конфликтом.194- Результат «уже обработано» подтверждай ACK.195- Не генерируй отсутствующий идентификатор.196197## DLQ198199DLQ добавляй только по требованиям:200201- `max_deliver` сам по себе не создаёт DLQ;202- TERM не является DLQ;203- сохраняй исходный payload и только безопасные metadata и причину;204- окончательно подтверждай исходное сообщение только после успешной записи в DLQ;205- ошибка DLQ не должна приводить к потере исходного сообщения;206- replay сохраняет исходные идентификаторы корреляции и идемпотентности;207- JetStream advisories используй для наблюдаемости, а не как хранилище DLQ.208209Если DLQ отсутствует, предусмотренные требованиями метрики и оповещения должны позволять обнаружить исчерпание доставок.210211## Логирование212213Применяй `python-service-logging-writing` и logging-профиль consumer-а. Consumer214владеет operation context сообщения и одной итоговой записью после определения215и выполнения транспортного исхода.216217- Связывай согласованные metadata сообщения до валидации payload, чтобы218 окончательный отказ имел диагностический контекст.219- Сохраняй стабильный message ID при redelivery; не генерируй его вопреки220 контракту владельца.221- Измеряй полную длительность от принятия сообщения до итогового ACK-действия222 монотонными часами.223- Записывай ACK, NAK или TERM только после фактического завершения действия. Если224 транспортный сбой нельзя выразить согласованным outcome, сначала дополни225 logging-контракт, а не выдавай намерение за результат.226- Для outcome и других классификаторов используй перечисления logging-слоя, а не227 перечисления application DTO. Даже при совпадении значений один к одному228 выполняй явное исчерпывающее преобразование.229- Очищай контекст сообщения в `finally`, включая validation error, timeout,230 cancellation и ошибку отправки подтверждения.231- Не логируй payload, чувствительные headers и произвольный текст исключения.232- Application, delivery policy и NATS-адаптер не повторяют итоговую запись или233 stack trace consumer-а.234235## Тестирование236237Unit-тестами покрой:238239- выбранные режимы ownership/provisioning и запрет недопустимых мутаций;240- для consumer-managed режима — план создания отсутствующего stream, добавление241 только отсутствующих subjects и сохранение посторонних subjects и242 неуправляемых параметров;243- для external/platform stream — отсутствие вызовов create/update при любой244 реакции на отсутствующую topology;245- создание собственного durable либо validate/bind externally provisioned durable246 согласно выбранному режиму;247- валидацию совместимого durable и отказ от мутации несовместимого;248- классификацию retryable и fatal ошибок topology/subscription;249- stop-aware паузу, повтор ensure/subscription и переходы readiness;250- envelope, payload и mapping;251- таблицу публичный результат → ACK/NAK/TERM;252- передачу настроенной задержки в `message.nak(delay=...)` для повторяемого253 результата;254- версию схемы и неизвестные поля;255- расчёт свободной ёмкости;256- отмену и отсутствие подтверждения при cancellation.257258Интеграционными тестами с настоящим JetStream покрой:259260- consumer-managed stream: создание отсутствующего stream и добавление261 отсутствующего subject без удаления существующих;262- external/platform stream: отсутствие его создания или изменения consumer-ом;263- consumer-owned durable: создание с полным контрактом, повторный и конкурентный264 ensure и последующий bind;265- externally provisioned durable: validate/bind без неявного создания;266- при выбранном retry — запуск до появления topology и успешную подписку после267 повторной попытки;268- конкурентную подготовку topology несколькими экземплярами;269- конфликт или несовместимость топологии как фатальную ошибку;270- ACK после успеха;271- delayed NAK и redelivery;272- TERM;273- reconnect;274- повтор после потерянного ACK;275- `in_progress`;276- concurrency и backpressure;277- mismatch consumer-конфигурации;278- DLQ, если он предусмотрен.279280Тестовая инфраструктура должна подниматься автоматически с уникальными именами и освобождаться после прогона.281282## Границы283284В область скила входят выбор ownership-модели, разрешённый additive topology ensure285для stream/subjects, ensure или bind-only durable consumer-а, повтор подписки,286pull consumer, транспортная валидация и mapping, bounded287concurrency, backpressure, ACK policy, reconnect, redelivery, `in_progress`, metadata288идемпотентности, необязательный DLQ, логирование и тесты.289290Не входят общий runtime и сигналы процесса, publisher, исходящий адаптер, Unit of Work, репозитории, реализация application/domain-логики и HTTP.291292## Проверка результата293294- Consumer-конфигурация соответствует требованиям.295- Режимы владения stream/subjects и provisioning-а durable заданы явно и296 реализованы независимо.297- В consumer-managed режиме отсутствующий stream создаётся из полной конфигурации,298 а subjects добавляются без разрушительных изменений.299- В publisher-owned и platform-provisioned режимах consumer не создаёт и не300 изменяет stream/subjects.301- Собственный durable явно создаётся и валидируется либо только валидируется и302 привязывается — согласно выбранному режиму; скрытого provisioning-а нет.303- При выбранном retry отсутствующая topology не завершает процесс: readiness304 снята, ожидание stop-aware, после паузы выполняется новая попытка подготовки и305 подписки; при согласованном fail-fast возвращается явная startup error.306- Конфликт и несовместимость топологии не скрываются бесконечными повторами.307- Одно сообщение вызывает одну готовую application-операцию.308- Транспортные модели не пересекают application-границу.309- ACK отправляется только после успеха.310- Параллелизм ограничен и pull учитывает свободную ёмкость.311- Reconnect не дублируется ручным циклом.312- At-least-once и идемпотентность отражены явно.313- Топология не изменяется разрушительно.314- Поля сообщения, длительность и итог ACK/NAK/TERM соответствуют logging-профилю,315 контекст не протекает между конкурентными обработчиками.316- Unit- и integration-тесты проверяют собственную логику, а не внутренности `nats-py`.