Middleware & Stamp Pipeline
The bundle extends Symfony Messenger in two places:
- Messenger middleware that runs inside the Messenger buses, both when a message is dispatched and when a worker handles a received message.
- A stamp decider pipeline that runs in the CQRS facades (
CommandBus,QueryBus,EventBus) before the message is handed to Messenger, and adds stamps to the envelope.
How the stamp pipeline works
When a facade dispatches a message, it first resolves the dispatch mode, then
runs the StampsDecider aggregator. The aggregator calls every registered
StampDecider in priority order (highest first). Each decider receives the
message, the resolved DispatchMode and the current stamp array, and returns
the (possibly modified) array. The result is passed to Messenger's
MessageBusInterface::dispatch() on the sync or async bus. A message dispatched through the
outbox gets the stamps of an asynchronous dispatch and goes
through the async bus (or the sync bus without one), which stores it instead of sending it.
sequenceDiagram
participant C as Caller
participant B as CommandBus
participant D as DispatchModeDecider
participant S as StampsDecider
participant SD as Stamp deciders (by priority)
participant M as Messenger bus (sync or async; async for the outbox)
C->>B: dispatch(command, mode, ...stamps)
B->>D: resolve(command, mode)
D-->>B: SYNC, ASYNC or OUTBOX
B->>S: decide(command, resolved mode, caller stamps)
loop highest priority first
S->>SD: decide(command, mode, stamps)
SD-->>S: stamps
end
S-->>B: final stamps
B->>M: dispatch(command, stamps)
Things to know about the pipeline:
- The initial stamps are the ones the caller passed.
dispatchSync()andask()remove aDispatchAfterCurrentBusStampfirst, because they need the result immediately. - Deciders receive the resolved mode,
SYNCorASYNC, neverDEFAULT(ASYNCfor a dispatch through the outbox).QueryBus::ask()always passesSYNC. - Each decider sees the stamps added by the deciders before it.
- The pipeline only runs for dispatches through the CQRS facades. It does not run
when a worker handles a received message, when you dispatch on a
MessageBusInterfacedirectly, or when the outbox relay sends stored messages. - Caller stamps win. The built-in deciders do not replace or duplicate a
MessageMetadataStamp,SerializerStamp,TransportNamesStamp,AggregateSequenceStamp,DeduplicateStamporDispatchAfterCurrentBusStamppassed by the caller, and an explicit causation id is kept. Retry policy stamps are added only for stamp classes the caller did not pass.
Middleware classes
The bundle registers its Messenger middleware automatically; you do not list it
under framework.messenger.buses.*.middleware.
Position in the bus
Most bundle middleware is inserted right after Messenger's
dispatch_after_current_bus middleware. Messages deferred with
DispatchAfterCurrentBusStamp are released later from that point of the stack,
so middleware placed before it would be skipped for them. The others need
another place: TraceContextCaptureMiddleware goes first, to record the trace
context before a message is deferred; OutboxPrepareMiddleware right after
add_default_stamps_middleware; DeduplicationLockReleaseMiddleware right
after deduplicate_middleware; and OutboxStoreMiddleware right before
send_message. With FrameworkBundle's default middleware, a bus handled by the
bundle looks like this (abridged; bundle middleware in brackets):
[TraceContextCaptureMiddleware] when a tracer provider is registered
add_default_stamps_middleware Symfony 7.4+
[OutboxPrepareMiddleware] when the outbox is enabled
add_bus_name_stamp_middleware
reject_redelivered_message_middleware
dispatch_after_current_bus
[OpenTelemetryMiddleware] when a tracer provider is registered
[CausationIdMiddleware] when causation_id.enabled is true
[AllowNoHandlerMiddleware] event buses only
failed_message_processing_middleware
deduplicate_middleware with framework.lock
[DeduplicationLockReleaseMiddleware] when the idempotency bridge is active
... your own middleware ... (Doctrine's transaction middleware is skipped by outbox stores)
[OutboxStoreMiddleware] when the outbox is enabled
send_message
handle_message
On a bus without dispatch_after_current_bus (for example with
default_middleware: false), the bundle middleware is placed first (after
TraceContextCaptureMiddleware), except OutboxStoreMiddleware, which goes before
send_message or handle_message, or last. DeduplicationLockReleaseMiddleware is only added to buses that contain
Messenger's deduplicate_middleware.
The bundle inserts its middleware in compiler passes that run after Messenger's
MessengerPass has built the middleware lists of the buses (priority -24 of the
beforeOptimization phase; MessengerPass runs at 0 up to Symfony 8.1 and at -16
from Symfony 8.2 on). A compiler pass of your own that must see the bundle's
middleware, or reorder it, needs a priority below -25.
The "CQRS buses" below are the bus ids the facades dispatch on, with aliases
resolved: every configured buses.* entry, plus default_bus when one of
buses.command, buses.query or buses.event is not set (the facade then falls
back to it). With all three set, default_bus gets none of the bundle's
middleware, so an unrelated default bus (mailer, notifier) stays untouched.
AllowNoHandlerMiddleware
Catches Messenger's NoHandlerForMessageException for messages implementing
SomeWork\CqrsBundle\Contract\Event, so events can be published before anyone
listens to them. Any other message rethrows the exception.
It is added to the event bus (buses.event, or default_bus when that is not
set) and to buses.event_async. Because it only silences events, a bus shared
by commands and events still reports commands without a handler.
CausationIdMiddleware
While a message is handled, pushes its MessageMetadataStamp onto the
CausationIdContext stack and pops it afterwards (also when the handler
throws). The metadata deciders read this stack: a message dispatched from
inside the handler inherits the parent's correlation id and gets the parent's
message id as its causation id. A message without a MessageMetadataStamp is
pushed as "no parent", so the messages its handlers dispatch start a new flow.
With buses, the other CQRS buses get a variant
(somework_cqrs.messenger.middleware.causation_id_isolation) that only pushes
"no parent".
Configure it under somework_cqrs.causation_id:
somework_cqrs:
causation_id:
enabled: true # default
buses: # default []: all CQRS buses
- messenger.bus.commands
enabled: false removes both the middleware and CausationIdStampDecider.
Entries in buses must be Messenger bus service ids; an unknown id fails the
container compilation. CausationIdContext is tagged kernel.reset, so workers
start every message with an empty stack.
OpenTelemetryMiddleware
Creates one span each time a message passes through a CQRS bus:
| Situation | Span name | Span kind |
|---|---|---|
| Dispatch (synchronous, sending to a transport, or storing in the outbox) | cqrs.dispatch <ShortClassName> |
PRODUCER |
| The outbox relay sends (or handles) a stored message | cqrs.dispatch <ShortClassName> |
PRODUCER |
| A worker handles a received message | cqrs.consume <ShortClassName> |
CONSUMER |
A message dispatched through the outbox (DispatchMode::OUTBOX,
#[Outbox]) therefore has two cqrs.dispatch spans: one when it is stored, although
nothing is sent then, and one when the relay dispatches the row, as a child of the first
(the stored TraceContextStamp). The worker's cqrs.consume span is a child of the
relay's span. OutboxWriter::store() opens no span: it stores the current trace context.
- The tracer is named
somework.cqrs. Spans carry the attributescqrs.message.class(the FQCN) andcqrs.message.type(command,query,event, orunknownfor other messages on the same bus). - A synchronous dispatch produces a single
cqrs.dispatchspan that also covers the handlers. There is no separate handler span. - The status is
OK, orERRORwith the exception recorded when the rest of the stack throws.
Trace propagation. On dispatch, the middleware sets a TraceContextStamp
holding the W3C traceparent/tracestate headers of the dispatch span; a
TraceContextStamp already on the envelope (for example one passed by the
caller or captured for a deferred dispatch) becomes the parent of that span and
is replaced. The stamp travels with the message through the
transport, and the worker's cqrs.consume span uses it as its parent, so the
consumer continues the producer's trace.
Activation. The middleware is registered on all CQRS buses when
open-telemetry/api (1.8 or newer) is installed and the container has a service
named OpenTelemetry\API\Trace\TracerProviderInterface. Without that service no
middleware is added. See Production: OpenTelemetry
for wiring a tracer provider.
DeduplicationLockReleaseMiddleware
Part of the idempotency bridge. Messenger's deduplicate_middleware acquires a
lock for each DeduplicateStamp and, for messages handled synchronously, keeps
it until the TTL expires. When the dispatch throws (a failing handler, or a
transport that cannot send), this middleware releases the lock, so the caller
can retry with the same idempotency key. It is registered when the idempotency
bridge is active and a lock.factory service exists. See
Idempotency.
OutboxPrepareMiddleware and OutboxStoreMiddleware
Registered when the outbox is enabled; other messages pass through them untouched. For a message dispatched through the outbox:
OutboxPrepareMiddleware, right after Messenger'sadd_default_stamps_middleware, hides itsDeduplicateStamps (including default stamps) fromdeduplicate_middleware, so the lock is taken when the relay sends the message, and drops itsDispatchAfterCurrentBusStamps: the message is stored now.OutboxStoreMiddleware, after your own middleware (validation, context stamps), stores it in the current transaction instead of sending or handling it. When the relay dispatches the stored message on the bus, it drops the stamps that middleware adds again for a class the stored message already carries, so the caller's context (e.g.router_context) wins over the relay's.- Doctrine's
doctrine_transactionanddoctrine_open_transaction_logger(and DoctrineBridge 8.2'sDoctrineDbalTransactionMiddlewareandDoctrineDbalOpenTransactionLoggerMiddleware, whatever service id you register them under), wherever they are listed on a CQRS bus (also more than once), are wrapped so that a message being stored skips them (they would flush the caller's entity manager, open a transaction around the store, or report the caller's open transaction); they run when the relay dispatches it.
The relay's dispatch carries RelayedFromOutboxStamp and runs without the caller's
context and without a ReceivedStamp: middleware that checks the dispatching context
(authorization) should skip it as it skips received messages.
Middleware between them must call the next middleware for an outbox dispatch:
otherwise the bus throws a LogicException.
Built-in stamp deciders
All built-in deciders are registered with fixed priorities. Deciders marked "per type" are registered once for commands, once for queries and once for events.
| Priority | Decider | Applies to | Registered when | Leaves alone |
|---|---|---|---|---|
| 225 | RateLimitStampDecider |
per type | a limiter is mapped under rate_limiting |
- |
| 200 | RetryPolicyStampDecider |
per type | always | policy stamps of a class the caller passed |
| 175 | MessageTransportStampDecider |
commands, queries, events | always | an existing TransportNamesStamp |
| 150 | MessageSerializerStampDecider |
per type | always | an existing SerializerStamp |
| 125 | MessageMetadataStampDecider |
per type | always | an existing MessageMetadataStamp |
| 110 | SequenceStampDecider |
events | sequence.enabled (default true) |
an existing AggregateSequenceStamp |
| 100 | CausationIdStampDecider |
all messages | causation_id.enabled (default true) |
an explicit causation id |
| 50 | IdempotencyStampDecider |
all messages | idempotency.enabled (default true); without symfony/lock it is a no-op (a warning at the first IdempotencyStamp) |
an existing DeduplicateStamp |
| -10 | DispatchAfterCurrentBusStampDecider |
commands and events | always | an existing DispatchAfterCurrentBusStamp |
RateLimitStampDecider (225)
Consumes one token from the Symfony rate limiter mapped to the message (see
rate_limiting); the limiter key is the message
class. When the limit is exceeded it logs a warning and throws
RateLimitExceededException, so nothing is dispatched. See
Rate limiting.
RetryPolicyStampDecider (200)
Appends the stamps returned by the RetryPolicy resolved for the message (exact
class, parent classes, interfaces, type default). The built-in policies return
no stamps; transport-level retries are configured with retry_strategy.
MessageTransportStampDecider (175)
Adds a TransportNamesStamp with the transports chosen for the message and mode.
A TransportNamesStamp passed by the caller wins; otherwise, in this order:
- the transports configured for exactly the message class under
transports.<command|command_async|query|event|event_async>.map; - on asynchronous and outbox dispatches, the transport named by
#[Outbox(transport: '...')]or#[Asynchronous(transport: '...')](#[Outbox]is read first; a class carrying both fails the build); - the transports configured for a parent class or interface, then the section's
default; - on asynchronous and outbox dispatches of a class with a bare
#[Asynchronous]or#[Outbox], theasynctransport, unlessframework.messenger.routingor#[AsMessage(transport: ...)]routes the message.
When nothing applies it adds nothing and Messenger's routing decides. See
transports and
Async routing with the #[Asynchronous] attribute.
MessageSerializerStampDecider (150)
Adds the SerializerStamp returned by the MessageSerializer resolved for the
message (exact class, parent classes, interfaces, type default, global default),
if any.
MessageMetadataStampDecider (125)
Adds the MessageMetadataStamp returned by the MessageMetadataProvider
resolved for the message. The default provider generates a random message id,
which is also the correlation id of the first message of a flow. While another
message is handled, the provider's stamp takes the correlation id of that
message and its message id as causation id (unless the provider set a
causation id).
SequenceStampDecider (110)
For events implementing SequenceAware, adds an AggregateSequenceStamp with
the aggregate type, the aggregate id and the sequence number. See
Event ordering.
CausationIdStampDecider (100)
When a message is dispatched while another one is being handled with a
MessageMetadataStamp passed by the caller, sets its causation id to the
parent's message id and keeps the caller's correlation id. A copy of the
handled message's own stamp (forwarded, e.g. $received->withExtra(...))
becomes a new stamp: new message id, same correlation id and extras, the
handled message as cause. It runs after the
metadata deciders, so the stamp already exists.
IdempotencyStampDecider (50)
Turns an IdempotencyStamp into Messenger's DeduplicateStamp, with the key
<message class>::<idempotency key> and the TTL from idempotency.ttl. The
IdempotencyStamp stays on the envelope. See Idempotency.
DispatchAfterCurrentBusStampDecider (-10)
For asynchronous dispatches, adds DispatchAfterCurrentBusStamp unless
dispatch_after_current_bus disables it for the message. A stamp passed by the
caller is kept by the decider; dispatchSync() and ask() drop it, and so does
OutboxPrepareMiddleware for an outbox dispatch (the message is stored at once). See
dispatch_after_current_bus.
Creating custom stamp deciders
Implement SomeWork\CqrsBundle\Contract\StampDecider (@api) to add your own
stamps to every dispatch through the facades.
1. Implement the interface
<?php
declare(strict_types=1);
namespace App\Infrastructure\Cqrs;
use SomeWork\CqrsBundle\Bus\DispatchMode;
use SomeWork\CqrsBundle\Contract\StampDecider;
use Symfony\Component\Messenger\Stamp\StampInterface;
use Symfony\Component\Security\Core\Authentication\Token\Storage\TokenStorageInterface;
final class AuditTrailStampDecider implements StampDecider
{
public function __construct(
private readonly TokenStorageInterface $tokenStorage,
) {
}
/**
* @param array<int, StampInterface> $stamps
*
* @return array<int, StampInterface>
*/
public function decide(object $message, DispatchMode $mode, array $stamps): array
{
foreach ($stamps as $stamp) {
if ($stamp instanceof AuditTrailStamp) {
return $stamps; // a stamp passed by the caller wins
}
}
$stamps[] = new AuditTrailStamp(
userId: $this->tokenStorage->getToken()?->getUserIdentifier(),
dispatchedAt: new \DateTimeImmutable(),
);
return $stamps;
}
}
AuditTrailStamp is your own class implementing Messenger's StampInterface.
decide() receives:
$message: the message being dispatched;$mode: the resolved mode,DispatchMode::SYNCorDispatchMode::ASYNC(an outbox dispatch decides its stamps asASYNC, with aStoreInOutboxStampamong them);$stamps: the current stamps (caller stamps plus those added by higher-priority deciders).
Return the array with your changes. You can also remove or replace stamps, but keep stamps passed by the caller unless you have a reason not to.
2. Register it
With autoconfiguration, implementing StampDecider adds the
somework_cqrs.dispatch_stamp_decider tag automatically, with priority 0: the
decider runs after every built-in decider except
DispatchAfterCurrentBusStampDecider (-10), so it can add its own
DispatchAfterCurrentBusStamp. Use a priority below -10 to see the final
stamps. Set the priority explicitly to control where the decider runs, either in
the service definition:
services:
App\Infrastructure\Cqrs\AuditTrailStampDecider:
tags:
- { name: 'somework_cqrs.dispatch_stamp_decider', priority: 130 }
or with Symfony's attribute on the class:
<?php
use SomeWork\CqrsBundle\Contract\StampDecider;
use Symfony\Component\DependencyInjection\Attribute\AsTaggedItem;
#[AsTaggedItem(priority: 130)]
final class AuditTrailStampDecider implements StampDecider
{
// ...
}
Without autoconfiguration, add the tag (with its priority) yourself.
3. Restrict it to message types (optional)
Implement SomeWork\CqrsBundle\Contract\MessageTypeAwareStampDecider (@api) to
run the decider only for some messages. messageTypes() returns classes or
interfaces; the decider is called only for messages that are an instance of one
of them:
<?php
declare(strict_types=1);
namespace App\Infrastructure\Cqrs;
use SomeWork\CqrsBundle\Bus\DispatchMode;
use SomeWork\CqrsBundle\Contract\Command;
use SomeWork\CqrsBundle\Contract\MessageTypeAwareStampDecider;
use Symfony\Component\Messenger\Stamp\StampInterface;
final class CommandAuditStampDecider implements MessageTypeAwareStampDecider
{
/**
* @return list<class-string>
*/
public function messageTypes(): array
{
return [Command::class];
}
/**
* @param array<int, StampInterface> $stamps
*
* @return array<int, StampInterface>
*/
public function decide(object $message, DispatchMode $mode, array $stamps): array
{
// Only called for Command instances.
return $stamps;
}
}
Without MessageTypeAwareStampDecider, the decider runs for every message
dispatched through the facades.
4. Choose a priority
Higher priorities run first. Pick a value relative to the built-in deciders:
- above 225: before rate limiting, for example to reject a dispatch early;
- between 175 and 200: after retry stamps, before transport routing (a
TransportNamesStampadded here takes precedence over thetransportsconfiguration and#[Asynchronous]); - between 125 and 150: after serialization, before metadata;
- between 100 and 125: after the metadata stamp exists, before the causation id is added;
- between -10 and 50 (the default
0is here): after almost everything, beforeDispatchAfterCurrentBusStampDecider(-10), which then sees yourDispatchAfterCurrentBusStamp. The causation id (100) and the idempotency bridge (50) have already run: anIdempotencyStampyou add here is not turned into aDeduplicateStamp(use a priority above 50), and aMessageMetadataStampyou add or replace here gets no causation id (use a priority above 100, or set it yourself); - below -10: after every built-in decider, to see the final stamps.
5. Test it
A decider is a plain class, so a unit test can call decide() directly:
<?php
declare(strict_types=1);
namespace App\Tests\Infrastructure\Cqrs;
use App\Infrastructure\Cqrs\AuditTrailStamp;
use App\Infrastructure\Cqrs\AuditTrailStampDecider;
use PHPUnit\Framework\TestCase;
use SomeWork\CqrsBundle\Bus\DispatchMode;
use Symfony\Component\Security\Core\Authentication\Token\Storage\TokenStorageInterface;
final class AuditTrailStampDeciderTest extends TestCase
{
public function testAddsAuditTrailStamp(): void
{
$decider = new AuditTrailStampDecider($this->createStub(TokenStorageInterface::class));
$stamps = $decider->decide(new \stdClass(), DispatchMode::SYNC, []);
self::assertCount(1, $stamps);
self::assertInstanceOf(AuditTrailStamp::class, $stamps[0]);
}
}