Search by

gian-tiaga / spiral-outbox

gian_tiaga

Transactional outbox package for Spiral applications

Package info

github.com/falur/spiral-outbox

pkg:composer/gian-tiaga/spiral-outbox

Statistics

Installs: 16

Dependents: 0

Suggesters: 0

Stars: 0

Open Issues: 0

v0.1.0 2026-09-14 14:37 UTC

This package is auto-updated.

Last update: 2026-09-14 14:48:12 UTC


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