somework / cqrs-bundle
CQRS utilities for Symfony Messenger integration.
Requires
- php: ^8.2
- psr/container: ^1.1 || ^2.0
- psr/log: ^3.0
- symfony/config: ^7.2 || ^8.0
- symfony/console: ^7.2 || ^8.0
- symfony/dependency-injection: ^7.2 || ^8.0
- symfony/filesystem: ^7.2 || ^8.0
- symfony/framework-bundle: ^7.2 || ^8.0
- symfony/http-kernel: ^7.2 || ^8.0
- symfony/messenger: ^7.2 || ^8.0
- symfony/service-contracts: ^3.5
Requires (Dev)
- doctrine/dbal: ^4.0
- doctrine/orm: ^3.6
- friendsofphp/php-cs-fixer: ^3.58
- open-telemetry/api: ^1.8
- phpstan/extension-installer: ^1.4
- phpstan/phpstan: ^2.1
- phpstan/phpstan-deprecation-rules: ^2.0
- phpstan/phpstan-phpunit: ^2.0
- phpstan/phpstan-strict-rules: ^2.0
- phpstan/phpstan-symfony: ^2.0
- phpunit/phpunit: ^11.5
- symfony/lock: ^7.2 || ^8.0
- symfony/rate-limiter: ^7.2 || ^8.0
Suggests
- doctrine/dbal: Required for DBAL-backed transactional outbox storage (DbalOutboxStorage)
- doctrine/doctrine-bundle: Provides the "doctrine.dbal.<name>_connection" services used by the transactional outbox
- doctrine/orm: Adds the outbox table to schema generation and Doctrine migrations (OutboxSchemaSubscriber)
- open-telemetry/api: Required for OpenTelemetry tracing bridge (^1.8)
- phpunit/phpunit: Required by the testing helpers in SomeWork\CqrsBundle\Testing (CqrsTestCase, CqrsAssertionsTrait)
- symfony/lock: Required for idempotency deduplication bridge (IdempotencyStamp to DeduplicateStamp)
- symfony/rate-limiter: Required for rate limiting message dispatch (RateLimitStampDecider)
Provides
None
Conflicts
- doctrine/dbal: <4.0 || >=5.0
- open-telemetry/api: <1.8
- symfony/lock: <7.2
- symfony/rate-limiter: <7.2
Replaces
None
- dev-main / 0.6.x-dev
- 0.5.x-dev
- v0.5.1
- v0.5.0
- v0.4.0
- v0.3.0
- v0.2.4
- v0.2.3
- v0.2.2
- v0.2.1
- v0.2.0
- v0.1.2
- v0.1.1
- v0.1.0
- dev-claude/confident-tesla-ildmkl
- dev-claude/determined-babbage-7hc30x
- dev-claude/upbeat-ride-cywy23
- dev-claude/trusting-lamport-6j7wgy
- dev-claude/jolly-mccarthy-j6ms92
- dev-claude/bold-cerf-ey7r4k
This package is auto-updated.
Last update: 2026-09-28 16:44:36 UTC
README
A Symfony bundle that wires Command, Query, and Event buses on top of Symfony Messenger. It auto-discovers handlers via PHP attributes, provides a configurable stamp pipeline, and ships with testing utilities and optional patterns such as a transactional outbox, idempotency, and rate limiting.
Why this bundle?
Symfony Messenger is a powerful transport layer, but it leaves CQRS wiring as an exercise for the developer. This bundle fills the gap:
- Auto-discovery -- Annotate handlers with
#[AsCommandHandler],#[AsQueryHandler], or#[AsEventHandler]and they are registered automatically, on the synchronous bus and (when configured) on the asynchronous bus of their type. No YAML tags, no manual wiring. - Stamp pipeline -- A composable
StampDeciderpipeline attaches retry policies, transport routing, serializer stamps, metadata, and dispatch-after-current-bus stamps per message type or per individual message class. - Dedicated buses --
CommandBus,QueryBus, andEventBuswith distinct semantics: commands support sync/async dispatch and can return a result, queries return the result of exactly one handler, events are fire-and-forget with zero-to-many handlers. - Testing utilities --
FakeCommandBus,FakeQueryBus, andFakeEventBusrecord dispatched messages;CqrsAssertionsTraitaddsassertDispatched()andassertNotDispatched()with callback-based property checks.
Architecture
flowchart LR
A[Your code] --> B[CommandBus / QueryBus / EventBus]
B --> C[DispatchModeDecider: sync, async or outbox]
C --> D[StampsDecider pipeline]
D --> E[Messenger bus]
E --> F[Handler]
E --> G[Transport]
G --> H[messenger:consume worker]
H --> F
Loading
The stamp pipeline runs the built-in deciders for rate limiting, retry policies, the #[Asynchronous] attribute, transport names, serializers, metadata, event sequence numbers, causation IDs, idempotency, and DispatchAfterCurrentBusStamp; deciders you register run by their priority (by default after the built-in ones and before the DispatchAfterCurrentBusStamp decider). Queries skip the dispatch-mode step; they are always handled synchronously.
How does it compare with plain Messenger?
| Capability | Plain Messenger | This bundle |
|---|---|---|
| Handler discovery | #[AsMessageHandler] or messenger.message_handler tags |
#[AsCommandHandler] / #[AsQueryHandler] / #[AsEventHandler], registered on the right sync and async buses |
| Bus API | MessageBusInterface::dispatch() for everything |
CommandBusInterface, QueryBusInterface, EventBusInterface with typed methods |
| Handler results | Read HandledStamp or use HandleTrait |
CommandBusInterface::dispatchSync() and QueryBusInterface::ask() return the result |
| Sync/async choice | routing per message class |
DispatchMode, #[Asynchronous], and per-message dispatch_modes configuration |
| Retry configuration | Per transport | Per message class or interface via RetryPolicy, bridged to the transport retry strategy |
| Stamps | Added by the caller | Composable StampDecider pipeline with priority ordering |
| Testing | InMemoryTransport or mocks |
Fake buses plus assertDispatched() / assertNotDispatched() |
| Event ordering | Not built-in | SequenceAware interface + AggregateSequenceStamp |
| Transactional outbox | Only with a Doctrine transport on the business connection | #[Outbox] / dispatch_modes / DispatchMode::OUTBOX on the buses (or OutboxWriter), DBAL storage and relay command, for any transport (AMQP, Redis, SQS, …) |
| OpenTelemetry | Not built-in | Middleware producing dispatch and consume spans |
Choose plain Messenger when your app has simple dispatch needs and you want no additional dependency. Choose this bundle when you want structured CQRS buses, per-message configuration, and testing utilities while staying close to Messenger. The bundle does not provide sagas, process managers, or event sourcing; use a full CQRS/ES framework if you need them.
Features
Core
CommandBuswith sync/async dispatch;dispatchSync()returns the handler resultQueryBus::ask()returns the result of the single handlerEventBuswith zero-to-many handlers and fire-and-forget semantics- Attribute-based handler discovery (
#[AsCommandHandler],#[AsQueryHandler],#[AsEventHandler]); implementing the handler marker interface (CommandHandler, …) with a typed__invoke()is the alternative - Compile-time check that every command and query has at most one handler per bus, counting handlers of its parent classes and interfaces
Stamp pipeline
- Composable
StampDecidersystem with priority ordering (@api-- extend it yourself) - Per-message retry policies via the
RetryPolicyinterface, with a transport-level retry strategy bridge - Per-message transport routing with Messenger's
TransportNamesStamp - Per-message serializer stamps
- Per-message metadata stamps with correlation and causation IDs
DispatchAfterCurrentBusStampcontrol per message
Patterns
- Causation ID propagation across nested dispatches
- Idempotency bridge (
IdempotencyStampto Messenger'sDeduplicateStamp) - Event ordering metadata with
SequenceAwareandAggregateSequenceStamp - Rate limiting via Symfony Rate Limiter
- Transactional outbox with DBAL storage and relay (retries with backoff), setup, failed-message and purge commands; messages reach it through the buses (
DispatchMode::OUTBOX,#[Outbox],dispatch_modes) orOutboxWriter
Developer experience
FakeCommandBus,FakeQueryBus,FakeEventBusfor unit testingassertDispatched()/assertNotDispatched()(inCqrsAssertionsTraitandCqrsTestCase) with callback-based property assertionssomework:cqrs:generatescaffolds a message and its handlersomework:cqrs:listhandler cataloguesomework:cqrs:debug-transportstransport configuration overviewsomework:cqrs:healthchecks handlers and transports (exit codes 0/1/2 for monitoring)
Observability
- OpenTelemetry middleware (spans for dispatching and for consuming messages in workers, trace context carried across transports)
- PSR-3 logging on a
cqrsMonolog channel, with a warning when an async dispatch ran synchronously
Integration
CommandBusInterface,QueryBusInterface,EventBusInterfacefor dependency injection and test doubles#[Asynchronous]attribute to make a message asynchronous without per-message YAML
Installation
Requirements
- PHP 8.2 or newer (0.7 will require PHP 8.3).
- Symfony 7.4 (7.4.9 or newer) or 8.1 and newer (FrameworkBundle and Messenger).
Optional packages enable additional features:
symfony/lock-- idempotency (IdempotencyStamp).symfony/rate-limiter-- rate limiting.doctrine/dbal4.3+ anddoctrine/doctrine-bundle-- transactional outbox.open-telemetry/api1.8+ -- tracing.phpunit/phpunit11.5, 12.5 or 13 -- the testing helpers (CqrsTestCase,CqrsAssertionsTrait).
Install the package
composer require somework/cqrs-bundle
With Symfony Flex the bundle is added to config/bundles.php automatically. Without Flex, register it manually:
// config/bundles.php return [ // ... SomeWork\CqrsBundle\SomeWorkCqrsBundle::class => ['all' => true], ];
No configuration file is required: by default every bus uses Messenger's default bus and all messages are handled synchronously. To customise the bundle, create config/packages/somework_cqrs.yaml; docs/flex-recipe/ contains a commented template with every option, and bin/console config:dump-reference somework_cqrs prints the full reference. The Flex recipe is not published to symfony/recipes-contrib yet, so Flex does not create this file for you.
Verify the installation
bin/console somework:cqrs:list
The command lists the registered command, query, and event handlers (or warns that none were found yet).
Quick start
The examples below assume a standard Symfony application where everything in src/ is registered as a service with autowire and autoconfigure enabled (the default config/services.yaml).
Step 1 -- Define a command
Commands are immutable DTOs that implement the Command marker interface:
<?php namespace App\Task; use SomeWork\CqrsBundle\Contract\Command; final class CreateTask implements Command { public function __construct( public readonly string $id, public readonly string $name, ) { } }
Step 2 -- Create the handler
Annotate the handler with #[AsCommandHandler] and type the __invoke() parameter with the command class:
<?php namespace App\Task; use SomeWork\CqrsBundle\Attribute\AsCommandHandler; use SomeWork\CqrsBundle\Contract\EventBusInterface; #[AsCommandHandler(CreateTask::class)] final class CreateTaskHandler { public function __construct( private readonly EventBusInterface $eventBus, ) { } public function __invoke(CreateTask $command): mixed { // Persist the task... $this->eventBus->dispatch(new TaskCreated($command->id, $command->name)); return $command->id; } }
An asynchronous event dispatched from a handler is sent once the handler has returned, after its
transaction committed; if the broker is down then, the event is lost and dispatchSync() throws
DeferredDispatchFailedException. For events that must not be lost, enable the
transactional outbox, mark the event class #[Outbox] (or map
it to outbox in dispatch_modes) and run the handler's database work in a transaction on the
outbox connection ($connection->transactional(), or Messenger's doctrine_transaction
middleware): the same dispatch() then stores the event in that transaction.
Step 3 -- Define a query and its handler
Queries implement Query; their handler returns the result:
<?php namespace App\Task; use SomeWork\CqrsBundle\Contract\Query; /** * @implements Query<array{id: string, name: string}> the result type of QueryBus::ask() for static analysis */ final class FindTask implements Query { public function __construct( public readonly string $id, ) { } }
<?php namespace App\Task; use SomeWork\CqrsBundle\Attribute\AsQueryHandler; #[AsQueryHandler(FindTask::class)] final class FindTaskHandler { /** * @return array{id: string, name: string} */ public function __invoke(FindTask $query): array { // Load the task from your storage... return ['id' => $query->id, 'name' => 'Write the docs']; } }
Step 4 -- Define an event and a listener
Events implement Event and may have any number of handlers, including none:
<?php namespace App\Task; use SomeWork\CqrsBundle\Contract\Event; final class TaskCreated implements Event { public function __construct( public readonly string $taskId, public readonly string $name, ) { } }
<?php namespace App\Task; use Psr\Log\LoggerInterface; use SomeWork\CqrsBundle\Attribute\AsEventHandler; #[AsEventHandler(TaskCreated::class)] final class LogTaskCreated { public function __construct( private readonly LoggerInterface $logger, ) { } public function __invoke(TaskCreated $event): void { $this->logger->info('Task {id} created', ['id' => $event->taskId]); } }
Step 5 -- Dispatch through the buses
Inject the bus interfaces (they are autowired) and dispatch:
<?php namespace App\Controller; use App\Task\CreateTask; use App\Task\FindTask; use SomeWork\CqrsBundle\Contract\CommandBusInterface; use SomeWork\CqrsBundle\Contract\QueryBusInterface; use Symfony\Bundle\FrameworkBundle\Controller\AbstractController; use Symfony\Component\HttpFoundation\JsonResponse; use Symfony\Component\HttpFoundation\Request; use Symfony\Component\Routing\Attribute\Route; final class TaskController extends AbstractController { public function __construct( private readonly CommandBusInterface $commandBus, private readonly QueryBusInterface $queryBus, ) { } #[Route('/tasks', methods: ['POST'])] public function create(Request $request): JsonResponse { $name = $request->getPayload()->getString('name'); // dispatchSync() handles the command right away and returns the handler result. $id = $this->commandBus->dispatchSync(new CreateTask(bin2hex(random_bytes(8)), $name)); return $this->json(['id' => $id], 201); } #[Route('/tasks/{id}', methods: ['GET'])] public function show(string $id): JsonResponse { return $this->json($this->queryBus->ask(new FindTask($id))); } }
CommandBusInterface::dispatch() returns the Messenger Envelope and handles the command synchronously or asynchronously depending on your configuration; dispatchSync() always handles it synchronously and returns the handler result. bin/console somework:cqrs:list now shows the three handlers.
Step 6 (optional) -- Handle commands asynchronously
Declare an asynchronous Messenger bus and a transport, then tell the bundle about them:
# config/packages/messenger.yaml framework: messenger: default_bus: command.bus buses: command.bus: ~ command.async_bus: ~ transports: async: '%env(MESSENGER_TRANSPORT_DSN)%'
# config/packages/somework_cqrs.yaml somework_cqrs: buses: command: command.bus command_async: command.async_bus transports: command_async: default: [async]
MESSENGER_TRANSPORT_DSN must point to a transport you have installed (Doctrine, AMQP, Redis, ...); for the Flex default doctrine://default?auto_setup=0, run composer require symfony/doctrine-messenger and bin/console messenger:setup-transports. Now $commandBus->dispatchAsync($command) sends the command to the async transport, and so does a plain dispatch() of a command class marked with #[Asynchronous] (SomeWork\CqrsBundle\Attribute\Asynchronous) or mapped to async under dispatch_modes.command.map. Handlers without an explicit bus are registered on the async bus automatically, so the worker finds them:
bin/console messenger:consume async
Events work the same way with buses.event_async and transports.event_async. See the Usage Guide for the dispatch-mode rules.
Documentation
Full documentation is available at somework.github.io/cqrs.
- Getting Started -- tutorial from installation to async dispatch and testing
- Usage Guide -- handler registration, dispatch modes, exceptions, console commands
- Configuration Reference -- every
somework_cqrsoption explained - Migrating from Symfony Messenger -- moving an existing Messenger application to the bundle
- Middleware & Stamp Pipeline -- built-in middleware and stamp deciders, custom deciders
- Testing Guide -- fake buses, assertions, integration testing
- Production Guide -- deployment, workers, monitoring
- Troubleshooting -- common issues and solutions
- Example application -- a runnable Symfony application with commands, queries, events and an async transport
- Upgrade Guide -- upgrading between versions of the bundle
- Changelog
Advanced topics
License
MIT. See LICENSE.