gian-tiaga / spiral-outbox
Transactional outbox package for Spiral applications
Requires
- php: >=8.5
- cuyz/valinor: ^2.6
- cycle/database: ^2.16
- cycle/migrations: ^4.3
- php-http/discovery: ^1.20
- psr/clock: ^1.0
- psr/log: ^3.0
- ramsey/uuid: ^4.9
- spiral/framework: ^3.16
Requires (Dev)
- gian-tiaga/phpstan-strict-rules: ^0.1.0
- php-cs-fixer/shim: ^3.95.23
- phpstan/phpstan: ^2.1.54
- phpunit/phpunit: ^13.1
Suggests
None
Provides
None
Conflicts
None
Replaces
None
README
gian-tiaga/spiral-outbox — transactional outbox для Spiral. Пакет сохраняет интеграционное
событие в бизнес-транзакции, а после commit создаёт и независимо выполняет одну или несколько
доставок через Spiral Queue.
Установка
composer require gian-tiaga/spiral-outbox:^0.1.0
Приложение подключает Spiral\RoadRunnerBridge\Bootloader\QueueBootloader и
GianTiaga\SpiralOutbox\Bootloader\OutboxBootloader после bootloader миграций Cycle. Пакет
добавляет свой каталог миграций, команды outbox:relay и outbox:status и перехватчик
Worker.
Целевой алгоритм
Что хранит Outbox
outbox_events → что произошло
outbox_deliveries → какую Job, куда и по каким правилам нужно доставить
Событие неизменяемо и общего статуса не имеет. Каждая доставка хранит свой класс Job, очередь, список задержек, таймаут, статус и число неудачных попыток. Успех одной доставки не закрывает соседние доставки события.
Запись события
OutboxEventStoreContract::add() вызывается внутри уже открытой бизнес-транзакции. Хранилище
просит у сериализатора строку полезной нагрузки, добавляет одну строку outbox_events с
routed_at = null и возвращает её идентификатор. Карта маршрутов и очередь в этой транзакции
не читаются. Без транзакции запись запрещена; ошибка записи откатывает доменные изменения
вместе с событием.
Сериализация события
Полезную нагрузку собирает не хранилище, а реализация OutboxEventSerializerContract: событие
превращается в строку, а строка колонки payload вместе с классом события — обратно в объект.
Реализация по умолчанию — JSON на cuyz/valinor, который собирает событие по типам его
конструктора. Приложение подменяет её своей привязкой контракта в контейнере: ни хранилище, ни
загрузчик о формате не знают.
Обычный проход relay: этап 1 — создать доставки
Relay выбирает события с routed_at = null через FOR UPDATE SKIP LOCKED. Для каждого
OutboxRoute он создаёт строку outbox_deliveries и ставит событию routed_at. Доставки и
отметка маршрутизации фиксируются одним commit.
Если маршрут отсутствует или неверен, частичные доставки не создаются, а routed_at остаётся
null. Новый маршрут не применяется к уже маршрутизированной истории. Unique по
event_id, job_class и queue_name защищает от одинаковой доставки.
Обычный проход relay: этап 2 — отправить доставки
В новой транзакции relay выбирает готовые pending, блокирует их через
FOR UPDATE SKIP LOCKED, ставит dispatched, dispatched_at и delivery_timeout_at, затем
делает commit. Только после commit он отправляет Job в очередь доставки.
Полное имя класса Job служит типом задачи Spiral Queue. Отдельного имени Job и записей в
queue.registry.handlers нет. Полезная нагрузка содержит outboxEventId, outboxDeliveryId,
внутренний outboxDispatchId, уникальный для каждой отправки, и признак последней попытки
outboxLastAttempt. Идентификатор события relay уже читает из доставки, поэтому отдельного
запроса он не стоит.
Маршрут проверяет публичный invoke(): обязательные скалярные параметры должны приходить
из этой полезной нагрузки. Объектные зависимости из контейнера и необязательные параметры
разрешены.
Значение статусов доставки
| Статус | Значение |
|---|---|
pending |
Доставку нужно отправить впервые или повторить после available_at |
dispatched |
Relay зафиксировал отправку и ждёт результат Worker |
completed |
Worker успешно выполнил Job; статус финальный |
failed |
Доставка окончательно не выполнена; статус финальный |
Работа Worker и Job
Перехватчик пакета находит доставку по outboxDeliveryId, требует статус dispatched,
совпадение внутреннего outboxDispatchId с текущей отправкой, сверяет полный класс Job и
загружает событие. Поздний результат старой Job не меняет новую отправку.
final class PublishStockJob extends JobHandler { public function invoke( string $outboxDeliveryId, OutboxMessageLoaderContract $outboxMessageLoader, ): void { $event = $outboxMessageLoader->load( outboxDeliveryId: $outboxDeliveryId, expectedMessageClass: AvailableStockChangedEvent::class, ); // Вызов прикладного сценария с данными $event. } }
Успех переводит только эту доставку в completed. RetryableOutboxException оформляет
повторяемую ошибку. Любое другое исключение сразу переводит только эту доставку в failed.
Все переходы условно проверяют текущий статус.
Расписание повторов
Число повторов равно длине retryDelaysSeconds; отдельного maxAttempts нет. Список
[60, 120] даёт первую попытку и два повтора. Пустой список [] означает финальный отказ
после первой ошибки.
При повторе доставка возвращается в pending, получает новое available_at, а
dispatched_at и delivery_timeout_at очищаются. Другие доставки события продолжают работу.
Отправка сообщает Job, последняя ли это попытка: bool $outboxLastAttempt равен true, когда
следующей паузы у доставки уже нет и её неудача станет окончательной. Признак считается из
снимка пауз и числа попыток самой доставки, поэтому своего счётчика попыток Job не заводит и
повторами по-прежнему владеет только Outbox. Job принимает признак необязательно — он нужен
только тем обработчикам, которым на последней неудаче нужно сделать что-то отдельное.
final class PublishStockJob extends JobHandler { public function invoke(string $outboxDeliveryId, bool $outboxLastAttempt): void {} }
Потерянная Job и таймаут результата
deliveryTimeoutSeconds покрывает ожидание Worker, выполнение и запас. Пока
delivery_timeout_at не наступил, вторую Job relay не создаёт. После таймаута попытка
считается потерянной и проходит обычную ветку повтора или финального отказа.
Таймаут не останавливает старую Job. Обработчик обязан безопасно переносить повторный запуск; гарантия пакета остаётся at-least-once.
Ошибка отправки в RabbitMQ
Если очередь отвергла Job, relay оформляет повторяемую ошибку только этой доставки и
продолжает остальные строки пачки. Он не ждёт delivery_timeout_at, потому что отказ известен.
Кто делает повторы
Повторами владеет только Outbox. RabbitMQ хранит и отдаёт текущую Job, Worker выполняет её и
сообщает результат. RetryPolicyInterceptor, retry-атрибуты Job и возврат упавшей Job в
RabbitMQ не используются. Для AMQP pipeline задаётся requeue_on_fail: false.
Очереди и Worker
Имя очереди — строковое значение enum приложения. То же значение используется в
OutboxRoute, конфигурации Spiral Queue и RoadRunner. Один relay обслуживает все очереди, а
каждую физическую очередь читает отдельная группа Worker.
Постоянный relay
Без параметров команда выполняет один проход. Постоянный процесс запускается отдельно от HTTP-приложения и после миграций:
php app.php outbox:relay php app.php outbox:relay --loop
При пустом проходе он ждёт секунду. --sleep=1 принимает целое значение от 1 до 3600.
SIGTERM и SIGINT запрещают новый проход и завершают процесс после текущего. Для цикла
нужно расширение PHP pcntl. Фатальная ошибка базы или конфигурации завершает процесс, а
отказ отправки одной доставки остаётся ошибкой этой доставки.
Состояние Outbox
Команда outbox:status печатает таблицей текущее состояние обмена по базе окружения:
php app.php outbox:status
| Показатель | Значение |
|---|---|
| Немаршрутизированные события | Число строк outbox_events с routed_at = null и возраст самой старой из них |
| Доставки по статусам | Число строк outbox_deliveries в каждом из pending, dispatched, completed, failed |
| Возраст старейшей незакрытой доставки | Время с создания старейшей строки в статусе pending или dispatched |
Команда только читает и всегда завершается успехом: пустые таблицы — это нули, а не ошибка. Данных события, полезной нагрузки и реквизитов в выводе нет. Длину самой очереди в брокере команда не спрашивает — её показывает панель брокера.
Граница будущего микросервиса
Будущий микросервис не читает таблицы приложения. Он получает стабильное integration event через свою RabbitMQ-очередь и владеет своими inbox и Outbox. PHP-класс локальной Job не является межсервисным контрактом.
Контракты реализации
Событие
Событие реализует IntegrationEventContract. Поле события — string, int, float, bool,
null, backed enum, DateTimeImmutable, вложенный объект по тем же правилам, list<T> и
array<string, T> с указанным T. Видимость свойства не ограничена: непубличные свойства тоже
попадают в полезную нагрузку, — но одноимённый параметр конструктора обязателен для свойства
любой видимости, потому что именно им событие собирается обратно. Статическое свойство полем
события не является.
Нетипизированный array, iterable, mixed, объект без класса и свойство без одноимённого
параметра конструктора ловит правило статического анализа
gian-tiaga/phpstan-strict-rules: параметр integrationEventInterface в phpstan.neon
называет ему интерфейс события. В рантайме остаётся отказ значения, которое нельзя закодировать,
— NAN, INF, строка с невалидным UTF-8, замыкание или ресурс: пакет отвергает его до вставки
строки. Секреты и персональные данные в событие не кладут.
Маршрут
Для класса события задаётся непустой список объектов OutboxRoute:
return [ 'eventsTableName' => 'outbox_events', 'deliveriesTableName' => 'outbox_deliveries', 'batchSize' => 100, 'routes' => [ AvailableStockChangedEvent::class => [ new OutboxRoute( job: PublishStockJob::class, queue: StockQueue::Publish->value, retryDelaysSeconds: [60, 120], deliveryTimeoutSeconds: 3600, ), ], ], ];
job — конкретный наследник JobHandler с публичным invoke() и обязательным параметром
string $outboxDeliveryId; кроме него из полезной нагрузки разрешены только
string $outboxEventId, string $outboxDispatchId и bool $outboxLastAttempt; queue — непустая строка;
retryDelaysSeconds — список положительных целых секунд; deliveryTimeoutSeconds —
положительное целое число. Настройки копируются в доставку и потом не меняются вместе с
конфигурацией.
Таблицы и миграции
Имена таблиц задаются до первой миграции. После неё переименование выполняется отдельной
миграцией приложения. Пакет добавляет свой каталог в migration.vendorDirectories.
outbox_events хранит id, event_type, payload, created_at, routed_at.
outbox_deliveries хранит id, event_id, job_class, queue_name, снимок повторов и
таймаута, статус, попытки, dispatch_id, времена переходов и безопасную причину отказа.
Логирование
Записи связываются полями outboxEventId и outboxDeliveryId. В контекст входят класс
события, класс Job, очередь, статус, попытки и времена. Полезная нагрузка, исходные сообщения
внешних исключений, секреты и персональные данные в журнал не попадают.
Проверки пакета
composer install
composer cs
composer phpstan
composer test
Тестам нужен настоящий PostgreSQL: поведение транзакций и SKIP LOCKED живёт в самой СУБД
и подделкой не проверяется. Адрес базы приходит переменными окружения — своей базы у пакета
нет, он создаёт в чужой только собственные таблицы и удаляет их после прогона:
OUTBOX_TEST_DB_HOST=127.0.0.1 \
OUTBOX_TEST_DB_PORT=5432 \
OUTBOX_TEST_DB_DATABASE=outbox_test \
OUTBOX_TEST_DB_USERNAME=outbox \
OUTBOX_TEST_DB_PASSWORD=outbox \
composer test