Skip to content

Repository files navigation

SomeWork CQRS Bundle

CI codecov PHPStan Level 8 License: MIT Latest Version Downloads

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 StampDecider pipeline 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, and EventBus with 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, and FakeEventBus record dispatched messages; CqrsAssertionsTrait adds assertDispatched() and assertNotDispatched() 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

  • CommandBus with sync/async dispatch; dispatchSync() returns the handler result
  • QueryBus::ask() returns the result of the single handler
  • EventBus with 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 StampDecider system with priority ordering (@api -- extend it yourself)
  • Per-message retry policies via the RetryPolicy interface, 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
  • DispatchAfterCurrentBusStamp control per message

Patterns

  • Causation ID propagation across nested dispatches
  • Idempotency bridge (IdempotencyStamp to Messenger's DeduplicateStamp)
  • Event ordering metadata with SequenceAware and AggregateSequenceStamp
  • 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) or OutboxWriter

Developer experience

  • FakeCommandBus, FakeQueryBus, FakeEventBus for unit testing
  • assertDispatched() / assertNotDispatched() (in CqrsAssertionsTrait and CqrsTestCase) with callback-based property assertions
  • somework:cqrs:generate scaffolds a message and its handler
  • somework:cqrs:list handler catalogue
  • somework:cqrs:debug-transports transport configuration overview
  • somework:cqrs:health checks 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 cqrs Monolog channel, with a warning when an async dispatch ran synchronously

Integration

  • CommandBusInterface, QueryBusInterface, EventBusInterface for 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/dbal 4.3+ and doctrine/doctrine-bundle -- transactional outbox.
  • open-telemetry/api 1.8+ -- tracing.
  • phpunit/phpunit 11.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.

Advanced topics

License

MIT. See LICENSE.

About

Symfony bundle wiring Command, Query and Event buses on top of Symfony Messenger. Attribute-based handler discovery, configurable stamp pipeline, PHPStan level 8.

Resources

Contributing

Security policy

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages