Transactional outbox
Suppose a handler commits a database change and then sends a message to a broker. It can
fail between the two steps: the change is saved and the message is lost, or the message
goes out for a change that was rolled back. With the outbox, you store the message in a
database table in the same transaction as the business change. The
somework:cqrs:outbox:relay command later sends the stored messages to Messenger.
Stability
The public API (@api) of the outbox is:
OutboxWriterandContract\Outbox\OutboxWriterInterface,OutboxMessageandDbalOutboxStorage(public foraddTableToSchema()in migrations and for the setup and failed-message tools);- the storage contracts of
Contract\Outbox:OutboxStorageand the capabilitiesOutboxSchema,FailedOutboxMessages,OutboxMonitoringandTransactionalOutbox, with theOutboxStatusandFailedOutboxMessageDTOs; - the
#[Outbox]attribute,DispatchMode::OUTBOX, the stampsOutboxStoredStamp,RelayedFromOutboxStampandStoreInOutboxStamp, andTesting\FakeOutboxWriter.
Write with OutboxWriter (or through the buses), or type-hint the interfaces. The relay
(OutboxRelay, its RelayReporter and RelayResult) and RelayUnitOfWork are internal.
How it works
- Inside your database transaction, you write the business data and store the message: through
the buses (
DispatchMode::OUTBOX,#[Outbox]ordispatch_modes, see Through the buses), or withOutboxWriter::store(). - The transaction commits. The message row is saved only if the business change is.
somework:cqrs:outbox:relayreads due rows, transport by transport, decodes each one, and dispatches it through the Messenger bus of its type (see Relaying). The message goes to the transport stored with the row, or follows the Messenger routing when no transport was stored. Then the relay marks the row as published. A row that fails is retried later with an increasing delay, and given up aftermax_attemptsattempts (see Failures).
A message reaches the outbox only when you ask for it: with DispatchMode::OUTBOX, with the
#[Outbox] attribute or an outbox entry of dispatch_modes (then dispatch() stores it), or
with OutboxWriter. Everything else the buses dispatch is sent as usual.
When do I need the outbox?
- Your transport is not your database (AMQP, Redis, SQS, Kafka…): a message sent from a
handler is either sent before the commit (and goes out for a change that may still roll back)
or after it (and is lost when the send fails). The bundle's buses send asynchronous commands
and events after the handler returned (
dispatch_after_current_bus, on by default), so a failed send leaves the change committed and throwsDeferredDispatchFailedExceptionfromdispatchSync(); in a worker, the retried command skips the handler that already ran and the event is lost with a warning in the Messenger log. Store such messages withOutboxWriterinstead. - Your transport is a Doctrine transport on the same connection as your business data:
Messenger inserts the message in the current transaction, so it is already atomic, but only
when it is sent inside the transaction. Disable
dispatch_after_current_busfor those messages (somework_cqrs.dispatch_after_current_bus.event.map), or use the outbox anyway. - A message goes to several transports, some of which may fail: the outbox stores one row per transport and retries each on its own.
Requirements
doctrine/dbal4. Enabling the outbox without it fails container compilation with anInvalidConfigurationException.doctrine/doctrine-bundle, which provides thedoctrine.dbal.<name>_connectionservice the storage uses.- Optional:
symfony/lock, so only one relay runs at a time (see Relaying). - Optional:
doctrine/orm, which adds the outbox table to schema generation and Doctrine migrations (see Creating the table).
The Flex recipe of DoctrineBundle writes a doctrine.orm section. Without doctrine/orm the
container then fails with The doctrine/orm package is required when the doctrine.orm config
is set: remove that section, or install the ORM (composer require symfony/orm-pack).
Quick start
- Enable the outbox (
somework_cqrs.outbox.enabled: true, see Configuration). - Create the table:
bin/console somework:cqrs:outbox:setup, or a Doctrine migration. - Inside your transaction, dispatch the messages through the outbox (see
Through the buses), or store them with
OutboxWriter::store()(see Writing to the outbox). - Run
bin/console somework:cqrs:outbox:relayevery minute, or keep it running with--watch(see Relaying), and purge published rows every night (see Purging published rows). In development, see Development. - Watch
bin/console somework:cqrs:health(see Monitoring).
The rest of this page explains each step and the cases operations need to know.
Configuration
# config/packages/somework_cqrs.yaml
somework_cqrs:
outbox:
enabled: true
table_name: somework_cqrs_outbox
connection: default
serializer: messenger.default_serializer
auto_setup: true
max_attempts: 10
| Option | Default | Description |
|---|---|---|
enabled |
false |
Registers the outbox storage and the four console commands. It decides which services exist, so it must be a plain boolean, not an %env()% value. |
table_name |
somework_cqrs_outbox |
Name of the outbox table: letters, digits and underscores, optionally schema.table (on MySQL and MariaDB database.table: without a database selected on the connection, its dbname, the automatic setup is skipped, and the Doctrine schema listener only adds a table of the connection's database to generated migrations). A word reserved in MySQL, MariaDB, PostgreSQL or SQLite (such as order or user) is a configuration error, because the queries do not quote the name. |
connection |
default |
DBAL connection name. The storage uses the service doctrine.dbal.<name>_connection. Use the connection that holds your business data, otherwise store() is not part of the business transaction. |
serializer |
messenger.default_serializer |
Messenger serializer service id. It is exposed as the alias somework_cqrs.outbox.serializer for code that writes to the outbox, and the relay uses it to decode rows. |
relay_on_terminate |
false |
For development: runs the relay right after a request, a console command or a message a worker handled that stored messages in the outbox. See Development. A plain boolean (it decides which services exist). |
require_transaction |
true |
Refuses to store a message outside a transaction on the outbox connection (OutboxRequiresTransactionException): the message would not be part of the business change. Checked by storages that implement TransactionalOutbox (DbalOutboxStorage does). |
auto_setup |
true |
Creates the table on first use if it is missing, or adds the columns a table of an earlier version lacks, but never inside an open transaction. Indexes are left to somework:cqrs:outbox:setup. Set it to false when migrations manage the table. |
max_attempts |
10 |
Attempts after which the relay gives up on a row that cannot be decoded or sent (at least 1). A row whose transport fails gets three times as many. See Failures. |
Through the buses
The command and event buses store a message in the outbox instead of sending it when its dispatch
mode is outbox. The stamp pipeline runs as for an asynchronous dispatch (transports, retry,
serializer, metadata and causation, idempotency, rate limiting), and so does the middleware of
the bus the relay sends it on (buses.<type>_async, or the synchronous bus without one):
validation, the stamps of your own middleware and of router_context (tenant, user, request
context) and tracing run in the dispatching process, and a message they reject is not stored.
After all of it, right before Messenger's send_message, the bundle stores the message instead
of sending it, right away, in the current transaction: it is never deferred until the current handler has finished, also when the message
provides a DispatchAfterCurrentBusStamp among its default stamps. Three ways select the mode:
// 1. Explicitly, for one dispatch.
$this->eventBus->dispatch(new OrderPlaced($orderId), DispatchMode::OUTBOX);
// 2. On the message class: every dispatch() with the default mode stores it.
#[Outbox] // or #[Outbox(transport: 'orders')]
final class OrderPlaced implements Event { /* ... */ }
# 3. In the configuration, per message class or interface, or for a whole type.
somework_cqrs:
dispatch_modes:
event:
default: outbox # every event, unless mapped otherwise
map:
App\Event\AuditTrail: async
command:
map:
App\Application\Command\ChargeCard: outbox
The attribute and the configuration take part in the usual
resolution order: an exact map entry
wins over the attribute, which wins over entries for parents and interfaces and over the
default. A class with both #[Outbox] and #[Asynchronous], #[Outbox] on a query, and an
outbox mode in dispatch_modes while outbox.enabled is off fail the compilation, and so does
#[Outbox] with an undefined transport or without the outbox, for a message with a handler in
the application (the bundle knows messages through their handlers). For a message without one,
such as an integration event other services consume, a transport that is not a Messenger
transport is refused when it is stored (UnknownOutboxTransportException), so the business
transaction rolls back instead of committing a row the relay could never send.
The handler code stays the same: dispatch inside the transaction of the business change.
#[AsCommandHandler(command: PlaceOrder::class)]
final class PlaceOrderHandler
{
public function __construct(
private readonly Connection $connection,
private readonly EventBusInterface $eventBus,
) {
}
public function __invoke(PlaceOrder $command): mixed
{
$this->connection->transactional(function () use ($command): void {
$this->connection->insert('orders', ['id' => $command->orderId]);
$this->eventBus->dispatch(new OrderPlaced($command->orderId)); // OrderPlaced is #[Outbox]
});
return null;
}
}
With Messenger's doctrine_transaction middleware on the command bus, every handler already
runs in a transaction of the entity manager's connection, and the events it dispatches through
the outbox are stored in it without an explicit transactional().
- The transaction is checked before the stamp pipeline and the middleware run. With
require_transaction: true(the default), a message dispatched outside a transaction onoutbox.connection(none is open, or it is open on another connection) throwsOutboxRequiresTransactionExceptioninstead of being stored on its own. A connection withauto_commit: falseis always in a transaction. On it, the relay and the maintenance commands commit their own writes, the setup runs, and each message the relay handles in its own process (see The transports) is a unit of work of its own: a failed handler's writes are rolled back. - Middleware runs twice. Once when the message is stored and once when the relay dispatches
it on the same bus (then Messenger's deduplication, and the sending). When the relay dispatches
it, a stamp that middleware adds again for a class the stored message already carries (the
request context of
router_context, a tenant) is dropped: the caller's context wins. Middleware with side effects runs twice too; Doctrine'sdoctrine_transactionanddoctrine_open_transaction_logger(and DoctrineBridge 8.2'sDoctrineDbalTransactionMiddlewareandDoctrineDbalOpenTransactionLoggerMiddleware, under any service id) are skipped when the message is stored, wherever they are listed, and run in the relay (they would flush the caller's entity manager, open a transaction around the store, or report the caller's open transaction). Middleware must pass an outbox dispatch (an envelope carryingStoreInOutboxStamp) on to the next one: one that returns early makesdispatch()throw aLogicException. - The relay's run has no caller context. It happens in the relay's process, without the
caller's request, user or tenant, and without a
ReceivedStamp. Middleware that checks the dispatching context (authorization, for example) and skips received messages should also skip envelopes carryingRelayedFromOutboxStamp: a message stored through the buses was checked when it was stored. A message written withOutboxWriter::store()passed no bus middleware, so such a guard lets it through unchecked: check it before callingstore().
if (null !== $envelope->last(ReceivedStamp::class) || null !== $envelope->last(RelayedFromOutboxStamp::class)) {
return $stack->next()->handle($envelope, $stack);
}
dispatch() returns the envelope with an OutboxStoredStamp: the ids of the
stored rows and their transports. Nothing is sent until the relay runs.
- The transports are those of an asynchronous dispatch: transports.command_async /
transports.event_async, the transport of #[Outbox(transport: ...)], or
framework.messenger.routing when the relay sends the row. No async bus is needed, also with
transports.<type>_async set: the relay dispatches on buses.<type>_async, or on the
synchronous bus without one. A message with none of them is stored with a warning, and the
relay handles it synchronously in its own process when it relays it.
- Delays. A DelayStamp (passed by the caller, among the default stamps of the message,
or added by a middleware) is stored with the message untouched and applies when the relay
sends the row: the delay counts from that send, not from the dispatch. The row itself is due
at once, so the next relay run sends it; with a transport that supports delays the message is
then handled the delay later, and without a transport (handled in the relay's process) the
delay is ignored. The same holds for OutboxWriter::store(). A later version may start the
delay at the dispatch instead.
- What bypasses the outbox. dispatchSync() and ask() (they need the result) and
dispatchAsync() (an explicit asynchronous dispatch) ignore the outbox mode.
- Idempotency. An IdempotencyStamp becomes a DeduplicateStamp scoped to each row's
transport (<key>@<transport>) when the message is stored, and deduplicates when the relay
sends the row, not when it is stored: two dispatches with the same key in one transaction store
two rows, and the second is dropped when relayed. A DeduplicateStamp among the default stamps
of the message (DefaultStampsProviderInterface) is scoped and applied the same way. With a
lock store that ties its keys to the process (flock, semaphore, PostgreSQL advisory locks,
ZooKeeper), a message with a DeduplicateStamp (from an IdempotencyStamp, its default stamps
or the caller) is refused when it is stored (LogicException), since the relay could never
send it; OutboxWriter::store() refuses it too. The store is the one lock.factory uses at
runtime (a LOCK_DSN from the environment counts with its runtime value).
- Tests. assertStoredInOutbox() checks the fake buses for a dispatch with
DispatchMode::OUTBOX, or with the default mode of a class carrying #[Outbox], and the fakes
return an envelope with an OutboxStoredStamp for it (see Testing). A fake bus
does not know dispatch_modes.
Writing to the outbox
OutboxWriter stores a message without the stamp pipeline and the middleware of the buses (no
validation or authorization when it is stored): use it for messages that are not commands or
events, or to choose every stamp yourself. Inject it (type-hint OutboxWriterInterface, which
Testing\FakeOutboxWriter replaces in unit tests) and call store() inside your transaction:
<?php
declare(strict_types=1);
namespace App\Application\Command;
use App\Application\Event\OrderPlaced;
use Doctrine\DBAL\Connection;
use SomeWork\CqrsBundle\Attribute\AsCommandHandler;
use SomeWork\CqrsBundle\Contract\Outbox\OutboxWriterInterface;
#[AsCommandHandler(command: PlaceOrder::class)]
final class PlaceOrderHandler
{
public function __construct(
private readonly Connection $connection,
private readonly OutboxWriterInterface $outbox,
) {
}
public function __invoke(PlaceOrder $command): mixed
{
$this->connection->transactional(function () use ($command): void {
$this->connection->insert('orders', [
'id' => $command->orderId,
'customer_id' => $command->customerId,
]);
// Stored while this handler runs, the event continues the flow of PlaceOrder:
// same correlation id, PlaceOrder as its cause (pass a MessageMetadataStamp to override).
$this->outbox->store(new OrderPlaced($command->orderId));
});
return null;
}
}
OutboxWriter::store(object $message, ?string $transportName = null, StampInterface ...$stamps)
returns the stored rows. The main points:
- Use the same connection.
Connectionmust be the connection named inoutbox.connection; withrequire_transaction: true, a store outside a transaction on it throwsOutboxRequiresTransactionException. With Doctrine ORM, callstore()insideEntityManagerInterface::wrapInTransaction(). The entity manager of that connection shares the sameConnectioninstance. - The transport. Without a transport name, the writer sends the message where an
asynchronous dispatch through the CQRS buses would: the transports of
transports.command_asyncortransports.event_asyncfor the message, or the one of#[Asynchronous(transport: ...)]. It stores one row per transport, so a failing transport is retried alone. When none is configured, the row followsframework.messenger.routingwhen it is relayed. A transport name you pass wins. The relay sends the message there with Messenger'sTransportNamesStamp. - Only your stamps, plus the causation. The stamp pipeline does not run: the bundle adds
no retry, serializer or rate-limit stamps. Pass the stamps you need; they are serialized
with the message. The one exception is metadata: stored while a handler runs and without a
MessageMetadataStampof its own, the message gets one that keeps the correlation id of the handled message and names its message id as the cause. Outside a handler no metadata is added. - The bus is chosen for you. The relay dispatches commands on
buses.command_async(orbuses.command), events onbuses.event_async(orbuses.event), queries onbuses.query, and anything else on the default bus. That bus adds itsBusNameStamp, so the worker hands the message to the bus where its handlers are registered. ABusNameStampyou store yourself is kept.
Rows get a time-ordered UUIDv7 id: rows stored in the same millisecond by one process keep their order.
Without the writer, build the row yourself with
OutboxMessage::fromEnvelope(Envelope $envelope, SerializerInterface $serializer, ?string $transportName = null, ?DateTimeImmutable $createdAt = null)
and pass it to OutboxStorage::store(). Encode it with the somework_cqrs.outbox.serializer
service (the outbox.serializer option, by default messenger.default_serializer), which the
relay decodes rows with, not with a transport's serializer.
Creating the table
The table has to exist before the first store(). There are three ways to create it:
Setup command. Run it once per environment, for example in your deployment script. It creates the table if it is missing, adds the columns and indexes that a table created by an earlier version lacks (see Upgrading from 0.4), and does nothing otherwise:
Doctrine migrations. When doctrine/orm is installed, the bundle registers
OutboxSchemaSubscriber, a postGenerateSchema listener that is scoped to the configured
outbox.connection. It adds the outbox table to the schema of that connection, so
doctrine:migrations:diff and doctrine:schema:update include it. Set
auto_setup: false once a migration manages the table. For DBAL-only migrations without
the ORM, add the table in the migration yourself:
<?php
declare(strict_types=1);
namespace DoctrineMigrations;
use Doctrine\DBAL\Schema\Schema;
use Doctrine\Migrations\AbstractMigration;
use SomeWork\CqrsBundle\Outbox\DbalOutboxStorage;
final class Version20260101000000 extends AbstractMigration
{
public function up(Schema $schema): void
{
DbalOutboxStorage::addTableToSchema($schema, 'somework_cqrs_outbox');
}
public function down(Schema $schema): void
{
$schema->dropTable('somework_cqrs_outbox');
}
}
auto_setup: true (default). The first time a process uses the storage, it checks
whether the table exists and has the columns of this version (without locking anything). If
the table is missing, it creates it; if it lacks the columns of this version (a table of 0.4),
it adds them. This only happens outside a transaction: inside one, storing a message in a
table of 0.4 fails (and your transaction is rolled back) until the columns exist, which is why
the setup command runs before the deployment (the health check reports it as critical). The
relay and the failed command add the columns they need too; indexes are left to the setup
command.
Processes that start at the same time wait for each other (at most 30 seconds, less once
another one has added the columns). Adding the columns does not wait in the lock queue of the
table, where the writes would queue behind it: PostgreSQL tries LOCK TABLE … NOWAIT and
MariaDB ALTER TABLE … NOWAIT for up to 1 second. An autovacuum of the table gives way to a
waiting lock after deadlock_timeout, so while only autovacuum holds the table, PostgreSQL waits
that long (writes wait too), unless it prevents a wraparound, which does not give way; this needs
a role that sees the sessions of other roles (pg_read_all_stats), otherwise the upgrade fails
while autovacuum runs. MySQL and MariaDB only add the columns when that takes no time
(ALGORITHM=INSTANT; a compressed table, for example, would be rebuilt). MySQL cannot do
that: it does not try while any transaction of the server has been open for more than a
second (it does not tell which tables a transaction holds; this needs the PROCESS
privilege, without it each attempt holds up the writes to the table for up to 1 second).
Until the columns exist, every relay run fails with could not be changed: another session …
kept it locked (MySQL: is not changed while a transaction of the database server has been
open): run the setup command. On PostgreSQL this runs in one transaction, so it is safe behind a
pooler in transaction mode (PgBouncer). It never builds or drops an index, which can take long on
a big table: the relay and the health check warn until somework:cqrs:outbox:setup has done
it. Until then (the relay checks for the index every 10 seconds, also with auto_setup: false)
it fetches with one query along the index of 0.4, in the order the rows were stored: the
transports do not take turns, and rows that are not due (retries, given-up rows, paused
transports) are read past, which slows fetches when many of them are ahead. The health
check reports the outbox as critical while the table lacks the columns (storing a message inside
a transaction fails until then); once messages have waited there for more than 10 minutes, it also
says how many wait and for how long (the relay cannot send them until the setup command has run),
which reads all pending rows and takes seconds with a large backlog. It never creates the table inside an open
transaction: DDL would implicitly commit your transaction on MySQL or abort it on
PostgreSQL. store() normally runs inside your transaction, so a missing table then raises
a LogicException that tells you to run somework:cqrs:outbox:setup. Dates are stored in
UTC, so neither the time zone of the writing process nor daylight saving time changes the
relay order. In practice,
auto_setup only creates the table when the relay or purge command runs first. Do not rely
on it for the write path.
Table schema
| Column | Type | Nullable | Description |
|---|---|---|---|
id |
GUID | No | Primary key (UUIDv7) |
body |
TEXT | No | Encoded message body |
headers |
TEXT | No | JSON-encoded headers from the serializer |
transport_name |
VARCHAR(190) | Yes | Target transport; null follows the Messenger routing |
created_at |
DATETIME_IMMUTABLE | No | When the row was built |
published_at |
DATETIME_IMMUTABLE | Yes | When the relay sent it (null = unpublished) |
attempts |
INTEGER, default 0 |
No | Attempts to relay the row (counted when the relay claims it) |
available_at |
DATETIME_IMMUTABLE | Yes | Earliest time of the next attempt (null = never attempted, or requeued) |
failed_at |
DATETIME_IMMUTABLE | Yes | When the relay gave up on the row (null = still relayed) |
last_error |
TEXT | Yes | Exception class and message of the last failure |
claim_token |
VARCHAR(32) | Yes | The relay run that claimed the row for an attempt it has not finished |
claimed_at |
DATETIME_IMMUTABLE | Yes | When that claim was made (set on a due row: its last attempt was interrupted) |
signature |
VARCHAR(64) | Yes | Signature of the id, body and headers (see Security) |
An index on (published_at, failed_at, transport_name, available_at, created_at, id), named
idx_<table>_pending (idx_<hash>_pending for long table names), serves the relay, the
purge and the counts of the health check; idx_<table>_claimed on claimed_at finds
unfinished claims. attempts, available_at, failed_at, last_error, claim_token,
claimed_at, signature and these indexes were added in 0.5.0;
the index replaces idx_<table>_published_created of 0.4. Upgrade a table created
by an earlier version.
Upgrading from 0.4
Writes need the new columns (store() writes the signature), so upgrade the table before the
new version takes traffic. Outside a transaction, the automatic setup adds the columns before
the first write (it gives up after 1 second behind a transaction that holds the table); inside
one, store() fails with The outbox table "…" lacks columns this version of the bundle needs
(…). Upgrade it with "bin/console somework:cqrs:outbox:setup" or a Doctrine migration. Stop
the relays of the old version before the new ones start: an old relay ignores the new columns
(retry times, given-up rows, claims). The relay also needs the new index to stay fast. Add the
columns and the index with one of:
bin/console somework:cqrs:outbox:setup, over a direct database connection (setups that start at the same time wait for each other, except while one builds the index on PostgreSQL: then the others stop at once withAnother process is building the index, see below; the wait uses a database lock held by the session, which PgBouncer in transaction mode would hand to another client: on PostgreSQL the command usually notices it and refuses (releasing the lock), or fails saying that the lock stayed with another server connection; under light load it may not notice, so do not rely on it). While another setup holds the lock, it says so and waits for it. A signal (e.g. a deploy job that is terminated) stops it with the exit code128 + signal, once the running statement returns (on a network that drops the connection silently, only once libpq notices: set TCP keepalives, e.g.keepalives_idle, on the connection); the next setup continues. On PostgreSQL it builds the index withCREATE INDEX CONCURRENTLY(without the role'sstatement_timeout), so writes go on while it runs; it waits for transactions that started before (e.g. apg_dump). It rebuilds an index that an interrupted build left invalid, and stops withAnother process is building the indexwhile one is being built (e.g. by your migration, or by the server process of a setup that was killed). The old index is dropped only once the new one exists. Changing the table needs a moment without open transactions on it: adding the columns (and, on MySQL and MariaDB, the index) waits at most 5 seconds for them, then fails withcould not be changed: another session … kept it lockedinstead of blocking every write behind it; run the setup again when the table is less busy. On MySQL and MariaDB the index is built online, but the build needs that moment at its end too: if a transaction (e.g. a dump) holds the table then, the work of the build is lost. Withauto_setup: true, the first relay run adds the columns (it does not queue behind other sessions: see Creating the table for autovacuum, MySQL and tables that MySQL or MariaDB would rebuild), so the relay usually works before the setup command has run, only slower;doctrine:migrations:diffwhendoctrine/ormis installed (the schema listener includes the new columns and index);- a migration of your own:
-- PostgreSQL
ALTER TABLE somework_cqrs_outbox ADD attempts INT DEFAULT 0 NOT NULL;
ALTER TABLE somework_cqrs_outbox ADD available_at TIMESTAMP(0) WITHOUT TIME ZONE DEFAULT NULL;
ALTER TABLE somework_cqrs_outbox ADD failed_at TIMESTAMP(0) WITHOUT TIME ZONE DEFAULT NULL;
ALTER TABLE somework_cqrs_outbox ADD last_error TEXT DEFAULT NULL;
ALTER TABLE somework_cqrs_outbox ADD claim_token VARCHAR(32) DEFAULT NULL;
ALTER TABLE somework_cqrs_outbox ADD claimed_at TIMESTAMP(0) WITHOUT TIME ZONE DEFAULT NULL;
ALTER TABLE somework_cqrs_outbox ADD signature VARCHAR(64) DEFAULT NULL;
CREATE INDEX CONCURRENTLY idx_somework_cqrs_outbox_pending ON somework_cqrs_outbox (published_at, failed_at, transport_name, available_at, created_at, id);
CREATE INDEX CONCURRENTLY idx_somework_cqrs_outbox_claimed ON somework_cqrs_outbox (claimed_at);
DROP INDEX CONCURRENTLY idx_somework_cqrs_outbox_published_created;
-- MySQL / MariaDB
ALTER TABLE somework_cqrs_outbox ADD attempts INT DEFAULT 0 NOT NULL, ADD available_at DATETIME DEFAULT NULL,
ADD failed_at DATETIME DEFAULT NULL, ADD last_error LONGTEXT DEFAULT NULL, ADD claim_token VARCHAR(32) DEFAULT NULL,
ADD claimed_at DATETIME DEFAULT NULL, ADD signature VARCHAR(64) DEFAULT NULL;
CREATE INDEX idx_somework_cqrs_outbox_pending ON somework_cqrs_outbox (published_at, failed_at, transport_name, available_at, created_at, id);
CREATE INDEX idx_somework_cqrs_outbox_claimed ON somework_cqrs_outbox (claimed_at);
DROP INDEX idx_somework_cqrs_outbox_published_created ON somework_cqrs_outbox;
CREATE INDEX CONCURRENTLY cannot run inside a transaction: Doctrine migrations need
isTransactional() to return false for it. 0.4 had no purge command, so its table may hold
every message ever relayed. Building the indexes then takes a while, and a migration generated
by doctrine:migrations:diff uses a plain CREATE INDEX, which blocks writes on PostgreSQL
until it is done. It also drops the old index before it creates the new one, so the relay has
no index while the migration runs; put the DROP INDEX last. Purge the published rows first
(somework:cqrs:outbox:purge works on the 0.4 table and does not change it), or use the
statements above.
Relaying
bin/console somework:cqrs:outbox:relay # up to 100 messages
bin/console somework:cqrs:outbox:relay --limit=500
bin/console somework:cqrs:outbox:relay --watch # keeps relaying until stopped
| Option | Default | Description |
|---|---|---|
--limit, -l |
100 |
Maximum number of rows to process (relay or fail) in this run, or in each run with --watch (positive integer). |
--watch, -w |
off | Keeps relaying, like messenger:consume: runs again at once after a full run, and waits --sleep seconds when no row was due. Stops after the current row on SIGTERM or SIGINT (exit code 0), or at --time-limit. It holds the relay lock the whole time; when another relay holds it, it waits for it (a standby that takes over when the other relay stops), checking every --sleep seconds. |
--sleep |
1 |
Seconds to wait before looking again when no row was due, or for the lock (with --watch). |
--time-limit |
none | Stops watching after this many seconds (with --watch), e.g. to let a process manager restart it. |
--wait-for-lock |
0 |
Without --watch: seconds to wait for another relay to release the lock; if it keeps it, exits with 3. With 0, a run that finds the lock taken exits at once with 0. |
--no-reset |
off | Does not reset the application's services after the rows the relay handles itself. By default, after each row handled in the relay's process (a row without a transport, or a sync:// transport) or failing in a handler there, the relay resets the services (services_resetter, and Doctrine's entity managers: cleared, or reset when a failed flush closed one), as Messenger's workers do between messages. |
The relay fetches up to 50 due rows at a time and:
- claims them, in one statement per group of rows with the same number of attempts: each claim counts the attempt, marks the row with the token of the run and the time, and postpones the row until its next retry time, provided that no other relay claimed the row since it was read. A process that dies during the batch (a PHP fatal error, running out of memory, a killed worker) therefore does not start the next run with the same rows, and a relay that overlaps this one skips them;
- for each claimed row, decodes it with the outbox serializer;
- adds a
TransportNamesStampwith the stored transport name, if one was stored; - dispatches the envelope through the bus of the message type:
buses.command_async(elsebuses.command) for commands,buses.event_async(elsebuses.event) for events,buses.queryfor queries, the default bus for anything else; - marks the sent rows as published, every 2 seconds and at the end of the run;
- renews the claims of the batch every 20 seconds, so a slow batch keeps its rows, and skips a row whose claim another relay took over meanwhile;
- releases the claims of the rows it did not attempt (the run ended, its storage failed, or their transport was paused): their attempt is not counted.
A due row that is still claimed was being sent by a relay that died (its claim was not finished before the retry time). The relay claims and sends such a row on its own, before the others, so that if the row kills the process again, only that row is blamed.
The transports take turns, the one whose next row has waited longest first (rows stored without a transport name count as one transport), so the backlog of one transport, for example after an outage, does not hold up the others: each fetch orders the transports by the age of their next row and then takes one row of each in turn. (When more transports have a backlog than a fetch has rows, a transport whose next row is newer waits until the older ones are relayed.) Within a transport the relay takes the new rows first (never attempted, or requeued), in the order they were stored, then the rows that failed before and whose retry time has passed, in the order of their retry time. Rows that keep failing therefore do not hold up new rows. A relay that cannot keep up with the new rows of a transport retries its failed rows only once it catches up; the health check reports both (see Monitoring).
What happens in special cases:
- Failing rows are retried later. If a row cannot be decoded or dispatched, the relay
records the failure (
attempts,last_error) and postpones the row: the next attempt waits 1 minute, then 2, 4, 8 … minutes, at most 1 hour. It printsFailed to relay message "<id>" (attempt 1 of 10, next attempt after <time>): <reason>and moves on to the next row. When any row failed, a single run exits with code1(--watchgoes on, and its exit code does not change). See Failures for what happens after the last attempt. - A failing transport is paused, the other transports go on. A send that fails with
Messenger's
TransportException(the broker cannot be reached, or it rejects the message) counts against three timesmax_attempts: 30 attempts by default, about a day of retries, so a broker outage of some hours gives no row up. After 3 such failures in a row for one transport, the relay printsTransport "<name>" failed 3 times in a row; its other messages wait for the next run (30 seconds with --watch).and skips the rows of that transport for the rest of the run instead of walking its whole backlog. With--watch, the next runs skip them too, for 30 seconds, doubling with every pause up to 5 minutes until a row of that transport goes through. The rows of the other transports are relayed as usual. As new rows come first, 3 new rows are enough to detect an outage, and the attempts of older rows are not used up. A transport that accepted a message earlier in the run is up: single messages it rejects (e.g. too large) do not pause it before 10 failures in a row, or 3 failures in a row that took more than 10 seconds (a transport that went down during the run and makes every send wait for a timeout). A burst of new rows that the broker rejects before any row of the run went through is only told apart from an outage by trying them: each run tries 3 of them, so the burst delays the other rows of that transport, by one run per 3 rejected rows. When most of a transport's traffic is rejected (say 9 rows out of 10), the other rows keep waiting behind such bursts: fix what the broker rejects (the given-up rows show the error). Rows stored without a transport name share one such counter (Messages without a transport name failed to be sent 3 times in a row …): when one of the transports they are routed to is down, the others' rows may wait for the next run too. Store the transport name to keep transports apart. A handler that ran inline and failed counts againstmax_attempts, even when it failed to send another message: its side effects happened. - A storage failure stops the run. If the database is down, or sent rows cannot be
marked as published, the run stops right away with
Stopping: …and exit code1. Rows that were sent but not marked (at most those of the last 2 seconds) are sent again later (see Delivery guarantees). - Messages handled inline trigger a warning. If no transport received a message (no stored transport name and no routing), the bus of its type handles it synchronously inside the relay process. The relay prints and logs a warning and still marks the row as published. Store a transport name or add routing to avoid this.
- Messages that go nowhere trigger a warning. When a message is neither sent nor handled,
for example because Messenger's deduplication dropped it as a duplicate, or because it is
an event without handlers or routing, the relay prints and logs
Message "<id>" (<class>) was neither sent to a transport nor handled …and marks the row as published. - A retry dropped by the deduplication is a failed attempt. A row that carries a
DeduplicateStamptakes Messenger's deduplication lock when it is sent, and the lock stays held until a worker handled the message or its TTL (300 seconds by default) expires. An earlier attempt of the same row may hold it: one that sent the message but died before the row was marked as published, or one whose send failed without releasing it (the bundle's idempotency bridge releases it). The relay cannot tell these apart, and prefers a duplicate to a lost message: it marks only a first attempt that the deduplication dropped as published (a duplicate of another row); a dropped retry fails withMessenger's deduplication dropped this retry …and is retried after the backoff, when the lock has usually expired. A lock without a TTL never expires: release it, or the relay gives the row up aftermax_attempts. Release the lock (or wait for its TTL) before you requeue such a row: a requeued row starts again at its first attempt.OutboxWriterscopes the key to the transport of the row (<key>@<transport>), so the rows of a message stored for several transports, at once or one by one, do not drop each other; a row that follows the Messenger routing keeps the key. Rows you store withOutboxStorage::store()yourself need distinct keys per transport. A scoped key does not match the same key dispatched directly on a bus, so a message dispatched both ways is not deduplicated across the two. - SIGTERM and SIGINT stop the run after the current row. With the
pcntlextension, the relay finishes the row it is working on, starts no other, marks the sent rows as published, releases the claims of the others, printsStopped by signal <number> after <count> message(s); the remaining messages wait for the next run., releases the lock and exits with1(--watchprintsStopped by signal <number>.and exits with0). A second signal stops it at once. PHP handles signals between operations: a send blocked on the network is only interrupted by the transport's own timeout, so configure timeouts on your transports (and a grace period longer than them). While the relay waits for another process to add the columns to the table (at most 30 seconds, on upgrade day), the first signal takes effect after that wait. - Only one relay runs at a time. When
symfony/lockis installed, the command takes a lock named after the application, the connection and the table. The lock expires after 60 seconds, and the relay extends it every 10 seconds between rows; if the lock is lost, the run stops with exit code1. The application part isframework.cache.prefix_seedwhen you set it, the project directory otherwise (Symfony's default seed is not used: it differs between environments and debug modes). If every release is deployed to a new directory, setprefix_seedto a stable value (Symfony recommends this anyway), so the relays of the old and the new release share the lock. The lock comes from the application'slock.factoryservice. FrameworkBundle registers that service whenframework.lockis enabled, which is the default oncesymfony/lockis installed. Without that service, Symfony'sLockableTraitcreates a local semaphore or flock store. The default stores only guard one host, so configure a shared store inframework.lock(Redis, a database, …) when cron runs the relay on several servers. A second relay that finds the lock taken printsAnother outbox relay is already running.and exits with0. After a PHP fatal error the relay still releases the lock (in a shutdown function). A relay killed without cleanup (SIGKILL, the OOM killer, a container stopped after its grace period) cannot: the next runs exit with0without relaying until the lock expires, at most 60 seconds later. - Relays that overlap anyway skip each other's rows. Without
symfony/lock, or with a lock store that only guards one host, two relays can run at the same time. The claim keeps them from sending the same row twice: a relay that finds a row claimed by the other one skips it and printsSkipped <n> message(s) that another relay claimed first.The claims of a batch are renewed every 20 seconds, between sends. A row can therefore be sent twice when a single send takes longer than the claims last after the last renewal (at least 40 seconds on the first attempt: the claim holds 1 minute) and the other relay claims them meanwhile: the slow row itself, and the rows sent in the 2 seconds before it, which are not marked as published yet. Claims are computed with the clock of the relay host: keep the clocks of the relay hosts synchronised (NTP). The relay lock is also only extended between sends, so a single send that outlasts its TTL (60 seconds) lets another relay start; the claims keep it from sending the rows of the slow relay's batch until they run out.
| Exit code | Meaning |
|---|---|
0 |
All selected rows were relayed, no row was due, or another relay holds the lock (without --wait-for-lock); with --watch, it was stopped by a signal or --time-limit |
1 |
At least one row failed (which includes a paused transport), the storage failed (e.g. the database is down), a signal stopped the run, or the lock could not be acquired or was lost |
2 |
Invalid options |
3 |
Another relay kept the lock for --wait-for-lock seconds |
Throughput. Each fetch lists the transports with pending rows once per run (again when a
fetch comes back short, or after 10 seconds), reads as many rows of each transport as the batch
of 50 needs, and reads up to 50 transports in one statement (UNION ALL). Relaying 20 000
rows took (PHP 8.4, local PostgreSQL 16 and MariaDB 10.11, a bus that only records the
messages) about 2.5 s with 1 transport, 6 s with 10 and 20 s with 100; with a simulated network
round trip of 0.5 ms per statement, reading the transports together is 1.5 times (10
transports) to 2 times (100 transports) faster than reading them one by one. Many transports
cost more statements per row: prefer a few transports with many rows each.
The relay handles at most --limit rows per run and then exits. Run it on a schedule, for
example from cron:
or keep it running under a process manager, like a Messenger worker; messages then leave the
outbox within --sleep seconds instead of up to a minute:
; supervisor
[program:outbox-relay]
command=php /path/to/project/bin/console somework:cqrs:outbox:relay --watch --time-limit=3600
autorestart=true
; exits with 1 when the database or the lock store fails: restarted at once, about every
; second while it is down (startsecs=0: never FATAL); it resumes when the database is back,
; or when the relay lock expires (up to 60 s) if the lock store is that database
startsecs=0
stopsignal=TERM
stopwaitsecs=30
A second relay from cron finds the lock taken and exits at once with 0. A second --watch
(another server, or a process that restarts while the old one finishes its row) waits for the
lock and takes over when the first one stops: run one or two, not one per server for
throughput, as only the relay holding the lock works. See
Production for running it under a process manager.
Development
In production the relay runs on a schedule or with --watch, so a stored message waits until
the next run. In development there usually is no such process: a message dispatched through the
outbox then only sits in the table, and with a sync:// transport its handlers do not run
until someone runs the relay. They run in the relay's process, not in the request.
Two ways to see the messages handled while developing:
- Relay on terminate (recommended). The relay runs right after each request, console command
or worker message that stored messages in the outbox, once its transaction is committed
(
kernel.terminate,console.terminate, and Messenger's worker events). The messages take the real path through the outbox table; withsync://they are handled at once, in the same PHP process, after the response was sent. Messages that their handlers store in turn are relayed in the same way.
What to expect:
- Where the handlers run. In the same PHP process, after the response: before the
profiler stores the profile of the request, so their log lines show in its Logs panel.
The client gets the response first only where Symfony can finish it early (PHP-FPM,
FrankenPHP); with
php -S,symfony servewithout PHP-FPM or mod_php, the request takes as long as the relayed handlers. Exceptions of the relayed handlers do not reach the error page: the relay records the failure on the row and logs it (cqrschannel), like a worker would. - Failed rows. A row whose handler failed is retried after the backoff (1 minute, then 2,
4 … minutes): by the relay that runs after the next request that stores a message, or by
bin/console somework:cqrs:outbox:relay(or--watch) at any time. Aftermax_attempts, the relay gives it up:somework:cqrs:outbox:failedlists it, and--requeueretries it after you fixed the handler.failedalso lists a row whose message class you renamed (it does not decode the body), but after--requeuethe relay fails to decode it again (and gives it up aftermax_attempts), and--signrefuses it (the class in the body cannot be found): restore the old class (or aclass_alias()of it) until the row is relayed, or delete the row with--delete. - When it does not run. While a transaction is still open on the outbox connection (the
rows are not committed yet): with Doctrine's
auto_commit: falsethat is always the case, so use--watchthere. After a command that a signal interrupted, and after the relay command itself (its handlers' messages wait for its next run). A notice in the log tells when the relay was skipped for an open transaction. - How many. Up to 10 runs of 100 messages after each request, command or worker
message: the relay goes on while a run stops at its limit, or while the relayed handlers
store more. A bulk load (fixtures, an import) of more than 1 000 messages leaves the rest
in the table, with a notice in the log (also when exactly 1 000 were left): run
bin/console somework:cqrs:outbox:relay, or keep--watchrunning while you load. - Several requests at once. Only one relay runs at a time: the next one waits up to
2 seconds for the lock, then leaves its messages to the relay holding it, which fetches
until no message is due. Use either
relay_on_terminateor a relay with--watch, not both: the watcher keeps the lock, so each request that stores messages would wait those 2 seconds before leaving them to it. - Services are not reset between the rows it handles, as with
sync://in a request. - A watching relay. Keep
bin/console somework:cqrs:outbox:relay --watchrunning in a terminal (or indocker compose, next tomessenger:consume).
Keep require_transaction on in development too: it catches the dispatches outside a
transaction that production would refuse. The outbox cannot be switched off per environment for
messages with #[Outbox] (the build fails without it); use relay_on_terminate instead. In
tests, the fake buses and Testing\FakeOutboxWriter replace the outbox (see
Testing).
Failures
A row that still fails after max_attempts attempts (default 10, about 4 hours of retries),
or after three times as many when its transport keeps failing (about a day), is given up: the relay prints Gave up on message "<id>" after 10 attempt(s): <reason>, sets
failed_at and never selects the row again. It stays in the table for inspection. List the
given-up rows and hand them back to the relay once the cause is fixed:
bin/console somework:cqrs:outbox:failed # id, transport, dates, attempts, last error
bin/console somework:cqrs:outbox:failed --limit=200
bin/console somework:cqrs:outbox:failed --requeue # every given-up row
bin/console somework:cqrs:outbox:failed --requeue <id> <id>
bin/console somework:cqrs:outbox:failed --requeue --transport=async <id> # and send it to another transport
Requeued rows start again with attempts = 0 and no last_error.
A row stored for a transport that does not exist (a typo, a renamed transport) is given up on
its first run, without an attempt: The transport "<name>" does not exist (is the row from another application sharing this table?). Fix the code that
stores it, then run "somework:cqrs:outbox:failed --requeue --transport=<name> <id>".
When the relay process dies during an attempt (a PHP fatal error, running out of memory, a
killed process, a lost database connection), the row keeps its claim (claimed_at) and the
error of the attempt before. It is tried again, on its own, after its retry delay; the rows
the process had claimed with it but not attempted yet count that attempt too. As the cause of
an interrupted attempt is unknown (a hanging broker as well as a message that crashes the
process), the larger budget applies: after three times max_attempts attempts, the next run
gives the row up without another attempt, because the row may be what kills the process, with
the error The last attempt did not finish (…); the message may have been sent. Previous
error: <error of the attempt before>. The message may have reached its transport: check the
consumer before you requeue it. When you lower max_attempts, a row
that already had more failed attempts gets one more attempt and, if it fails, is given up with
its real error.
Publishing a row clears its failed_at, last_error and claim. purge never deletes given-up rows; delete those you do not want to relay with
somework:cqrs:outbox:failed --delete <id>….
Monitoring
somework:cqrs:healthincludes an outbox check. It warns when the relay gave up on rows; when rows failed (or their attempt was interrupted more than once: a relay that keeps dying on them) and wait for another attempt while the oldest of them was stored more than 10 minutes ago (a transport outage, or rows that cannot be sent); and when the oldest due row has waited more than 10 minutes (the relay does not run, does not keep up, or pauses their failing transport); when a claim ran out (at the retry time of its attempt) more than 10 minutes ago without a relay taking it over (a relay died, or hangs on a send without a timeout, and no relay runs since); and when the table needssomework:cqrs:outbox:setup(e.g. its index is missing or invalid). It is critical when the table cannot be read, and when the table lacks the columns of this version (storing a message inside a transaction fails until the setup has run). The check reads at most 10 000 rows per count (it reportsmore than 10000) and finds the oldest due row with one probe per transport along the index, so it stays cheap on a large backlog.- The relay logs failed attempts, paused transports, messages handled inline or dropped, a
table that needs the setup command, and
runs stopped by a signal (warning), and given-up rows and stopped runs (error), besides
printing them: on the
cqrschannel with MonologBundle (monolog.logger.cqrs), otherwise to the application'sloggerservice. - Alert on
somework:cqrs:health(exit code1for a warning,2for critical) and on given-up rows, not on every exit code1of the relay: the relay also exits with1when a single row failed and will be retried, and when a signal stopped it (a deployment). Runs that keep exiting with1, or errors in its log, need a look.
Purging published rows
Published rows stay in the table until you purge them:
bin/console somework:cqrs:outbox:purge # published more than 7 days ago
bin/console somework:cqrs:outbox:purge --older-than="12 hours"
Rows are deleted in batches of 1000, so purging a large backlog does not hold one long lock.
--older-than takes a number and a unit (second, minute, hour, day, week, month or
year, singular or plural), such as 7 days, 12 hours or 1 month, with at most 6 digits.
The default is 7 days. Only rows whose published_at is older than that are deleted.
Unpublished rows, including given-up ones, are never deleted. An invalid value exits with code 2. Schedule the purge, for example
daily.
Delivery guarantees
- Atomic write. The row exists only if your transaction commits.
- At-least-once delivery. The relay marks rows as published after dispatching them,
every 2 seconds. A crash in between (it affects the rows sent in the last 2 seconds), a
failed
markPublished(), or a single send that outlasts the claims (at least 40 seconds) while another relay runs (it affects that row and the rows sent in the 2 seconds before it) can send the same message twice; so does the retry of a row carrying aDeduplicateStampafter such a crash, once the deduplication lock expired. A message routed to several transports is sent to all of them again when one of them fails: store one row per transport to avoid that. Consumers must be idempotent, for example by recording processed message ids under a unique constraint. The bundle's idempotency bridge does not prevent these duplicates either: theDeduplicateStampof a row stored with anIdempotencyStampis applied when the relay sends it, but its lock is released once a worker handled the message (or expires afteridempotency.ttl), and the relay retries the unmarked row until a send goes through. - Order. New rows of a transport are relayed in the order of their
created_at, the time (to the second) at whichOutboxMessage::fromEnvelope()encoded them in PHP, then of their time-ordered id. That is not the order in which they were committed: a row stored early in a long transaction is relayed after newer rows that were already published. There is no order across transports, which take turns. A row that fails is postponed, so later rows overtake it, and so do the rows of other transports while its transport is paused. Several workers consuming the transport can also process messages out of order. - Latency. Messages leave the outbox only when the relay runs. Your schedule sets the delay.
Custom storage
The DBAL storage covers relational databases that DoctrineBundle can connect to. To keep
the outbox somewhere else, implement SomeWork\CqrsBundle\Contract\Outbox\OutboxStorage:
<?php
declare(strict_types=1);
namespace SomeWork\CqrsBundle\Contract\Outbox;
use DateTimeImmutable;
use SomeWork\CqrsBundle\Outbox\OutboxMessage;
/**
* Persists messages in a transactional outbox for reliable async dispatch.
*
* The relay works in claims: it fetches due messages, claims them with a token of its run
* (counting the attempt before anything is sent, so an attempt that kills the process still
* counts), sends them, and then marks them published, records their failure, or releases the
* ones it did not get to. A claim that is never finished leaves claimedAt set: when the message
* is due again, the next relay knows that the attempt was interrupted.
*
* Every message a storage returns must carry the id, body, headers and signature exactly as
* they were stored.
*/
interface OutboxStorage
{
/**
* Stores a message; call it inside the database transaction of the business change.
*/
public function store(OutboxMessage $message): void;
/**
* Returns the messages that are due: neither published nor given up, and either never
* attempted (or requeued) or past their retry time. The transports take turns, the one whose
* next message has waited longest first (the messages without a transport name count as one
* transport); within a transport, the messages never attempted come first, in the order they
* were stored, then the others, in the order of their retry time.
*
* @param list<string|null> $excludedTransports Transports whose messages are skipped; null
* stands for messages stored without a transport name
*
* @return list<OutboxMessage>
*/
public function fetchUnpublished(int $limit, array $excludedTransports = []): array;
/**
* Claims fetched messages for an attempt, atomically per message: a message is claimed only
* while it is neither published nor given up, and its attempts and transport name are still
* the fetched ones (so two relays cannot claim the same attempt). A claimed message gets the
* token, claimedAt (now), one more attempt, and its retry time: $retryAt[<fetched attempts>].
* It is not due before that time, so an attempt that never finishes is retried after it. A
* run claims each message at most once (a new run uses a new token).
*
* @param list<OutboxMessage> $messages As fetchUnpublished() returned them
* @param array<int, DateTimeImmutable> $retryAt Retry times keyed by the fetched number of attempts
* @param string $token Identifies the claims of one relay run (not empty)
*
* @return list<string> The ids of the claimed messages, in the order of $messages
*/
public function claim(array $messages, array $retryAt, string $token): array;
/**
* Renews the claims of fetched messages that are still claimed with $token: claimedAt becomes
* now and the retry time $retryAt[<fetched attempts>], so the claims of a long batch do not run
* out before the relay gets to them.
*
* @param list<OutboxMessage> $messages As fetchUnpublished() returned them
* @param array<int, DateTimeImmutable> $retryAt Retry times keyed by the fetched number of attempts
*
* @return list<string> The ids still claimed with $token, in the order of $messages
*/
public function renew(array $messages, array $retryAt, string $token): array;
/**
* Undoes the claims of messages the relay did not attempt (it stopped, or paused their
* transport): their attempts, retry time and claimedAt go back to the fetched values. Messages
* no longer claimed with $token are left alone.
*
* @param list<OutboxMessage> $messages As fetchUnpublished() returned them
*/
public function release(array $messages, string $token): void;
/**
* Marks sent messages as published and clears their claim and failure. Ids that do not exist
* or are already published are ignored.
*
* @param list<string> $ids
*/
public function markPublished(array $ids): void;
/**
* Records the failure of a claimed message and ends its claim: the number of attempts, the
* error, and the retry time; with $retryAt null the message is given up and never returned by
* fetchUnpublished() again.
*
* @return bool false when the message is no longer claimed with $token (e.g. published, or
* claimed by another relay); nothing is recorded then
*/
public function recordFailure(string $id, string $token, int $attempts, string $error, ?DateTimeImmutable $retryAt): bool;
/**
* Deletes messages published before the given date and returns how many were deleted.
*/
public function purgePublished(DateTimeImmutable $publishedBefore): int;
}
An implementation must meet these rules:
store()must write through the same transaction as your business data. Otherwise the outbox guarantees nothing.fetchUnpublished()returns due messages only (unpublished, not given up, retry time passed). The transports take turns, the one whose next message has waited longest first; within a transport, first the messages never attempted, in the order they were stored (createdAt, then id), then the others, in the order of their retry time. It skips the excluded transports. The relay excludes a transport after 3 send failures in a row (after 10, or 3 over at least 10 seconds, once it accepted a message in the run); a storage that ignores the exclusion makes it stop early instead of reaching the other transports.claim()updates each message in one atomic step (e.g. a conditionalUPDATE): only while it is neither published nor given up, and its attempts and transport name are still the fetched ones. It returns the ids it claimed; the relay skips the others. A claimed message is not due before the given retry time.renew()movesclaimedAtto now and the retry time forward for the messages still claimed with the token, and returns their ids: the relay renews the claims of a batch every 20 seconds, so they do not run out before it gets to them.release()restores the attempts, the retry time andclaimedAtof the fetched messages that are still claimed with the token.recordFailure()only changes a message that is still claimed with the token, and ends the claim.markPublished()also ends the claim.- Rebuild each message with
new OutboxMessage(string $id, string $body, string $headers, DateTimeImmutable $createdAt, ?string $transportName = null, int $attempts = 0, ?string $lastError = null, ?DateTimeImmutable $claimedAt = null, ?DateTimeImmutable $availableAt = null, ?string $signature = null). Keep the values exactly as stored, the id, body, headers and signature byte for byte (the relay verifies the signature). The relay decides when to give up fromattempts, and recognises an interrupted attempt byclaimedAt.$headersis the JSON string produced byfromEnvelope(), and$idand$bodymust not be empty.
Name your implementation under outbox.storage (a service id, or a class name, which the
bundle registers as an autowired service):
# config/packages/somework_cqrs.yaml
somework_cqrs:
outbox:
enabled: true
storage: App\Outbox\MongoOutboxStorage
The relay and purge commands, OutboxWriter and the OutboxStorage alias then use it.
doctrine/dbal is not needed, and table_name, connection and auto_setup are ignored:
they only configure the DBAL storage. The relay lock is named after the service id.
The other features need more than OutboxStorage. Implement the interfaces of
SomeWork\CqrsBundle\Contract\Outbox your storage can support:
| Interface | Methods | Used by |
|---|---|---|
OutboxSchema |
setup(?\Closure $onWait = null): void, pendingChanges(): list<string> |
somework:cqrs:outbox:setup; the relay and the health check report what pendingChanges() returns |
FailedOutboxMessages |
fetchFailed(int $limit, array $ids = []): list<FailedOutboxMessage>, requeueFailed(array $ids = [], ?string $transportName = null, ?\Closure $sign = null): int, deleteFailed(array $ids): int |
somework:cqrs:outbox:failed |
OutboxMonitoring |
status(): OutboxStatus |
the outbox check of somework:cqrs:health |
TransactionalOutbox |
isInTransaction(): bool |
outbox.require_transaction (OutboxWriter and DispatchMode::OUTBOX refuse to store outside a transaction); outbox.relay_on_terminate waits for the transaction to be committed |
fetchFailed() returns only the given ids when there are any. When requeueFailed() gets
$sign, it calls $sign($message) with each requeued row as stored (id, body, headers) and
stores the returned string as the row's signature; that is how --requeue --sign works.
FailedOutboxMessage may carry messageType (the serializer's type header), bodyClass (the
message class of a PHP-serialized body), bodyClasses (every class the body instantiates; both
read without unserializing it, with SerializedBody::inspect()) and digest
(FailedOutboxMessage::digest($body, $headers)), which --sign shows to the operator and
checks before it signs a row. Fill them only for messages asked for by id: a listing of all
given-up messages should not read every body.
Without them, setup and failed exit with 1 and say which interface is missing, and the
health check reports the outbox as not checked. DbalOutboxStorage implements all four.
Without
TransactionalOutbox,outbox.require_transactionis not enforced. The bundle cannot tell whether a transaction is open, soOutboxWriterandDispatchMode::OUTBOXalso store messages outside one, which are then not part of the business change, andrelay_on_terminaterelays without waiting for the transaction to be committed. The same holds when the application replaces thesomework_cqrs.outbox.storageservice instead of configuring its storage underoutbox.storageor decorating it, and when the class of the storage service is not known when the container is built (a service created by a factory without a class: declare its class). In debug mode, the container compilation log warns about it (grep require_transaction var/cache/<env>/*Compiler.log): implementTransactionalOutbox(isInTransaction()returns whether a store would join an open transaction of your business data), or setoutbox.require_transaction: falseto acknowledge it. 0.5.x only warns; a later minor version may refuse to build.
DbalOutboxStorage also runs each message the relay handles in its own process (no transport,
sync://) as a unit of work of its own on a connection with auto_commit: false, and rolls
back a transaction such a handler left open. That interface (RelayUnitOfWork) is internal: a
custom storage relays without it.
To add behaviour to the storage instead (logging, metrics), decorate it:
#[AsDecorator('somework_cqrs.outbox.storage')] on a class that implements OutboxStorage
and takes the inner storage. The relay, the purge command and your code then go through the
decorator, while setup, failed, the health check and the relay's report of pending
changes keep working on the configured storage behind it (somework_cqrs.outbox.base_storage,
also somework_cqrs.outbox.dbal_storage for the DBAL storage), so the decorator does not have
to implement the capabilities. Your own services get them by type-hinting OutboxSchema,
FailedOutboxMessages or OutboxMonitoring, which autowire to that storage when it implements
them.
Security
The relay decodes due rows with the outbox serializer and dispatches the message it finds.
With Messenger's default PHP serializer, decoding runs unserialize(), which can execute code
through classes your application loads. Whoever can write rows (e.g. through an SQL injection
in your application) could therefore run code in the relay, so the rows are signed.
Signed rows
With outbox.signing.enabled (the default), the storage signs every row it stores with
HMAC-SHA256 (a key derived from framework.secret, or from outbox.signing.secret) over the
id, the body and the headers, and the relay verifies the signature before it decodes the
row. A row without a valid signature is given up at once, without being decoded:
The message is not signed, so it was not decoded … or The signature of the message does not
match ….
What signing protects, and what it does not:
- It keeps rows the application did not store from reaching
unserialize()and the handlers: an attacker who can write to the table but does not know the secret cannot forge a row. - It does not stop someone with write access from changing the state of rows: deleting them,
marking them published or given up, changing
transport_name(not signed, so thatoutbox:failed --requeue --transportworks), or copying a signed row under the same id to send a message again. Consumers must be idempotent anyway (see Delivery guarantees). - It does not help when the secret leaks: rotate it.
Operating it:
- Store through the
OutboxStorageservice orOutboxWriter. The signature is added by a decorator ofsomework_cqrs.outbox.storage; rows stored by SQL, or through aDbalOutboxStorageyou create yourself, are not signed. (The bundle offers no autowiring alias ofDbalOutboxStoragefor that reason.) - Rotate the secret by moving the old one to
outbox.signing.previous_secretsuntil the rows signed with it are relayed. Each entry may be an environment variable (previous_secrets: ['%env(OUTBOX_PREVIOUS_SECRET)%'], with an empty default for the variable when there is none: empty entries are ignored); the whole list cannot come from a single variable. Rotatingframework.secret(e.g.APP_SECRET) rotates the outbox secret too, unlessoutbox.signing.secretis set. - Rows of an earlier version are not signed. During a rolling deployment, instances of 0.4
keep storing unsigned rows until the last one is replaced, so set
outbox.signing.accept_unsigned: truefor the upgrade and remove it once those rows are relayed. Meanwhile the relay decodes any unsigned row, forged ones included: keep the window short. Rows with a wrong signature are always given up. - The secret must not be empty. With an empty
framework.secret(e.g. an unsetAPP_SECRET), every service that stores outbox rows fails to start: set a secret, oroutbox.signing.secret. - A row you checked (e.g. one stored while signing was disabled) is signed with the current
secret and handed back to the relay with
somework:cqrs:outbox:failed --requeue --sign <id> …. It shows the rows first: the class thetypeheader names, the message class of a PHP-serialized body, every class the body would instantiate (read as text, never unserialized) and a SHA-256 prefix of the body. A serialized type name (#[AsMessage(serializedTypeName: 'shop.place_order')]with the Symfony serializer) is shown after its class,App\Message\PlaceOrder (shop.place_order): the outbox serializer tells the class (MessageTypeAwareSerializerInterface), or the type map of Messenger's Symfony serializer when it is the outbox serializer; any other header names the class itself. It refuses a row whosetypeheader does not match the class in its body; a body that instantiates a class that is neither the envelope, a stamp, a command, query or event (as the message) nor a type declared by the properties of those (recursively); and a body with custom serialization (a class implementing onlySerializable, whose data the review cannot read). A forged row therefore cannot smuggle an object of any other class (anunserialize()gadget) past the review. A message that is not a command, query or event, or an object your application stores in an untyped property (mixed,object, arrays), is refused too; allow its class or interface with--allow-class=App\Money(the classes you allow, and the types their properties declare, are trusted). In an interactive terminal it asks for confirmation. It signs the bodies it showed: a row whose body changed in the meantime stops the command. The review cannot tell a forged row from a genuine one when it only carries the application's own classes (with data chosen by whoever wrote it): only sign rows your application stored, and delete the others. - A storage of your own must return the id, body, headers and signature exactly as stored.
signing.enabled: false restores the trust model of Messenger's Doctrine transport: the
table is trusted.
Other measures
- Keep write access narrow. Let the application connect with a role that can only read and
write rows (
SELECT,INSERT,UPDATE,DELETE); every command of the outbox works with it once the table is set up. Runsomework:cqrs:outbox:setupor your migrations with a role that may change the schema, and setauto_setup: false. - Prefer a serializer that does not unserialize PHP objects, e.g.
serializer: messenger.transport.symfony_serializer(JSON). It still instantiates the class itstypeheader names (with the body as constructor arguments), so signing matters as much;--signchecks that header, but cannot list other classes in a body that is not PHP-serialized. Messages made of primitives, as the bundle recommends, encode without extra normalizers. Rows written with another serializer cannot be decoded after the switch: relay them first. - Symfony refuses unsigned
RunProcessMessageandRunCommandMessage, so a forged row cannot start a process or a console command through Messenger's own handlers. - The relay drops the stamps that only describe a dispatch in progress (
ReceivedStamp,SentStamp,HandledStampand the other non-sendable stamps). Serializers never write them, so only a forged row holds them; aReceivedStampwould make the relay handle the message itself instead of sending it.
Error texts. The relay stores the message of the exception of a failed attempt in
last_error, prints it and logs it, and the OpenTelemetry middleware records exceptions on
spans. Exception messages can contain personal data. Rows the relay gave up on keep their
error until you requeue them or delete them (somework:cqrs:outbox:failed --delete <id>…), so
include them in your retention policy (see Personal data). Control characters are replaced by spaces in stored errors and in the
output of somework:cqrs:outbox:failed.
Message size. A fetch reads the bodies of at most 8 MiB of messages (a larger message is read on its own); the rest waits for the next fetch of the same run. The relay needs a few times the size of the largest message in memory: decoding and sending copy it.
Limitations
- Polling only. Messages leave the outbox when the relay runs. Change data capture is not supported.
- One database. The outbox table must be reachable through the same connection and transaction as your business data. Distributed transactions are not supported.
- One table per application. The relay of an application sends every row of its table:
a second application (or kernel) sharing the table would have its rows given up as
unsigned or badly signed (another secret) or for an unknown transport, or sent by the wrong
relay. Give each application its own
table_name, schema or database. - Long transactions slow the relay. On PostgreSQL, while any transaction holds an old
snapshot (a long report,
pg_dumpon the primary, an idle-in-transaction session), the rows relayed since stay in the index as dead entries, and every fetch walks them: relaying gets slower the longer the snapshot is held, until it ends (then autovacuum cleans up). InnoDB (MySQL, MariaDB) behaves the same way while an old read view is open (amysqldump --single-transaction, a longREPEATABLE READtransaction): purge cannot remove the marked rows' old versions, and the history list length grows. Keep long transactions off the primary (take dumps from a replica), and watchpg_stat_activityforidle in transactionsessions, orinformation_schema.innodb_trxon MySQL and MariaDB. Larger--limitruns suffer less. - Reads go to the primary. With a
PrimaryReadReplicaConnection, the storage switches to the primary before the relay, the health check and the outbox commands read, so they never see a lagging replica (a row already published would be sent again). - Given-up rows wait for you. Rows the relay gave up on stay in the table until you requeue or delete them. Once a row is relayed, Messenger's retry and failure transports take over.