К основному содержимому
Программа курса

Лабораторная 03: Событие придёт дважды

Добавьте очередь и обработчики задач. Реализуйте outbox и дедупликацию, проверьте сбой между записью результата и подтверждением сообщения.

Код стенда, установка и шаблон отчёта

Лабораторная 3. Outbox, повторная доставка и один локальный эффект

Добавьте асинхронную обработку к тому же сервису. Цель — связать две локальные транзакции без общей транзакции etcd+NATS. Требуются Lab1–2 и темы транзакций, сообщений и идемпотентности. Ориентир длительности — 90 минут плюс самостоятельное завершение.

Запуск

Подготовьте зависимости по README. Из корня репозитория:

PYTHONDONTWRITEBYTECODE=1 /tmp/ds-course-venv/bin/python labs/distributed-systems/runner.py --lab 3 --solution starter --report /tmp/lab03.json

В starter/protocol.py завершите TODO для create и complete. Транспорт relay/worker уже дан в service.py; изучите порядок его действий и fault hooks. Эталон выбирается --solution reference, assertions общие.

Задание

  1. Убедитесь, что создание task, idempotency и outbox атомарно. Relay читает outbox, публикует событие и только после publish ACK отмечает запись опубликованной. Сбой между этими действиями допускает повторную публикацию.
  2. Реализуйте атомарную фиксацию dedup + result + completed task + counter. Проверка текущих ревизий защищает счётчик от потерянных обновлений при конкурентных изменениях. Результат по уникальному ключу сам по себе не защищает соседний эффект от повтора.
  3. ACK consumer разрешён только после commit. Сохраните стабильный event_id. Дедупликация в брокере намеренно не используется, чтобы она не скрывала дефект приложения.

Эффект — увеличение локального etcd counter на amount. Внешний HTTP-вызов, письмо или платёж не входят в транзакцию и потребовали бы другого протокола. Вычисление перед commit может выполняться повторно.

Три сценария

  • Сохранённый outbox, публикации ещё нет. Relay аварийно завершается непосредственно перед публикацией. В etcd есть pending intent, stream пуст. После нового запуска задача должна завершиться.
  • Потерян publish ACK. Relay публикует в настоящий JetStream с reply inbox без подписчика. Брокер сохраняет сообщение, но приложение не получает ACK и не отмечает outbox. Повторная публикация создаёт вторую копию; worker обязан применить эффект один раз. Это проверяется реальным stream count, dedup и счётчиком.
  • Результат сохранён, consumer ACK потерян. Worker фиксирует результат и завершается до ACK. Durable consumer повторяет доставку после AckWait. Новый worker распознаёт event_id, не увеличивает counter повторно и подтверждает сообщение.

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

Что сдавать

  • Diff протокола, рисунок двух транзакционных границ и отдельного ACK брокеру.
  • JSON с тремя сценариями, stream count, redelivery и counter. В штатном наборе эффекты 2, 3 и 5 дают итог 10 при трёх результатах, даже если сообщений больше трёх.
  • Выдержку событий одного task_id: публикация/повтор → delivery → commit → повтор delivery → suppression → ACK.
  • Отчёт и ADR: retention, размер атомарной области, пределы общего счётчика как учебной модели.

Чтение и вопросы

Kafka помогает отличить журнал и consumer от транзакции приложения; Helland — определить локальную атомарность; Sagas — обсудить компенсацию, если эффект нельзя включить в одну транзакцию. Соответствующие исследовательские деки предназначены для подготовки преподавателя. Механика используется из JetStream consumers.

Почему ACK до commit может потерять задачу? Почему commit до ACK требует dedup? Почему «broker supports exactly once» не доказывает уникальность внешнего платежа? Что произойдёт при очистке dedup раньше повторной доставки? Дополнение: спроектируйте внешний эффект с idempotency API получателя; явно укажите новую границу ответственности.

Восстановление и ограничения

Один NATS file-store broker переживает остановку процесса с сохранением volume; уничтожение диска не покрыто. Runner восстанавливает свои процессы/сети и удаляет только собственный проект, как описано в README. Не переносите служебные fault hooks в публичный сервис. Ни один сценарий не подменяет etcd или JetStream заглушкой в памяти.

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