From 24288fa3b877136cb6acc77e4786a0aac8ad9c8b Mon Sep 17 00:00:00 2001 From: David Badura Date: Mon, 5 Oct 2026 12:17:50 +0200 Subject: [PATCH] Add event emitter, GIN index for tags and DCB with any append store The event emitter is opt-in via subscription.event_emitter, because removing a subscription also removes its subscription stream from the event store. With a taggable store the PostgreSQLPlatformMiddleware is registered on the event store connection, so the GIN index for tags gets created. DCB now uses the configured store, so the in memory store can be used in tests as well. --- composer.lock | 8 +- docs/configuration.md | 50 ++++++- src/DependencyInjection/Configuration.php | 7 +- .../PatchlevelEventSourcingExtension.php | 69 ++++++++- src/Doctrine/DbalConnectionFactory.php | 6 +- .../PatchlevelEventSourcingBundleTest.php | 139 ++++++++++++++++++ 6 files changed, 267 insertions(+), 12 deletions(-) diff --git a/composer.lock b/composer.lock index 3de774f9..502e98ac 100644 --- a/composer.lock +++ b/composer.lock @@ -420,12 +420,12 @@ "source": { "type": "git", "url": "https://github.com/patchlevel/event-sourcing.git", - "reference": "51c07650652266b11c641a3a25d9955757bd3b31" + "reference": "542b287b76c17450c3c049989dcb7a4864c99877" }, "dist": { "type": "zip", - "url": "https://api.github.com/repos/patchlevel/event-sourcing/zipball/51c07650652266b11c641a3a25d9955757bd3b31", - "reference": "51c07650652266b11c641a3a25d9955757bd3b31", + "url": "https://api.github.com/repos/patchlevel/event-sourcing/zipball/542b287b76c17450c3c049989dcb7a4864c99877", + "reference": "542b287b76c17450c3c049989dcb7a4864c99877", "shasum": "" }, "require": { @@ -507,7 +507,7 @@ "issues": "https://github.com/patchlevel/event-sourcing/issues", "source": "https://github.com/patchlevel/event-sourcing/tree/4.0.x" }, - "time": "2026-10-03T10:26:29+00:00" + "time": "2026-10-05T07:55:57+00:00" }, { "name": "patchlevel/hydrator", diff --git a/docs/configuration.md b/docs/configuration.md index ee4f25ae..4b65db59 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -178,7 +178,7 @@ patchlevel_event_sourcing: Following store types are available: - `dbal_stream` *default* -- `dbal_taggable` *required for DCB* +- `dbal_taggable` - `in_memory` - `custom` @@ -186,6 +186,14 @@ Following store types are available: If you use `custom` store type, you need to set the service id under `patchlevel_event_sourcing.store.service`. ::: +:::tip +With `dbal_taggable` the bundle registers the `PostgreSQLPlatformMiddleware` on the event store connection, +so a GIN index on the `tags` column is created on PostgreSQL. +This works automatically if you configure the connection via `url` or use a doctrine bundle connection +like `doctrine.dbal.eventstore_connection`. +For any other connection service you have to register the middleware yourself. +::: + ### Change table Name You can change the table name of the event store. @@ -550,6 +558,23 @@ final class SubscriptionErrorListener You can find all available events in the [library documentation](/docs/event-sourcing/latest/subscription). ::: +### Event Emitter + +Subscribers can emit new events into their own `subscription_` stream with the `EventEmitter`. +This is disabled by default and can be enabled like this: + +```yaml +patchlevel_event_sourcing: + subscription: + event_emitter: true +``` +:::warning +If a subscription is removed, its `subscription_` stream is removed from the event store as well. +::: + +:::note +The event emitter does not work with a read only store. +::: ## Command Bus You can enable the command bus integration to use your aggregates as command handlers. @@ -869,4 +894,25 @@ You can then specify this service here: patchlevel_event_sourcing: clock: service: 'my_own_clock_service' -``` \ No newline at end of file +``` + +## Dynamic Consistency Boundary + +You can enable the experimental dynamic consistency boundary (DCB). +This registers the `DecisionModelBuilder` and the `EventAppender` services. + +```yaml +patchlevel_event_sourcing: + store: + type: 'dbal_taggable' + dcb: true +``` +:::note +DCB needs a store that supports appending, like `dbal_taggable` or `in_memory` for tests. +::: + +:::tip +We recommend PostgreSQL for DCB, because only there the tag queries can use an index. +::: + +If you want to learn more about DCB, read the [library documentation](/docs/event-sourcing/latest/dynamic-consistency-boundary). diff --git a/src/DependencyInjection/Configuration.php b/src/DependencyInjection/Configuration.php index c772bc8b..c4c3e7d2 100644 --- a/src/DependencyInjection/Configuration.php +++ b/src/DependencyInjection/Configuration.php @@ -52,7 +52,8 @@ * enabled: bool, * retries_in_ms: list, * detection_window: string - * } + * }, + * event_emitter: array{enabled: bool} * }, * connection: ?array{ * service: ?string, @@ -264,6 +265,10 @@ public function getConfigTreeBuilder(): TreeBuilder ->scalarNode('detection_window')->defaultValue('PT5M')->end() ->end() ->end() + + ->arrayNode('event_emitter') + ->canBeEnabled() + ->end() ->end() ->end() diff --git a/src/DependencyInjection/PatchlevelEventSourcingExtension.php b/src/DependencyInjection/PatchlevelEventSourcingExtension.php index 354df807..41abd823 100644 --- a/src/DependencyInjection/PatchlevelEventSourcingExtension.php +++ b/src/DependencyInjection/PatchlevelEventSourcingExtension.php @@ -92,6 +92,7 @@ use Patchlevel\EventSourcing\Snapshot\Adapter\Psr6SnapshotAdapter; use Patchlevel\EventSourcing\Snapshot\DefaultSnapshotStore; use Patchlevel\EventSourcing\Snapshot\SnapshotStore; +use Patchlevel\EventSourcing\Store\Dbal\PostgreSQLPlatformMiddleware; use Patchlevel\EventSourcing\Store\InMemoryStore; use Patchlevel\EventSourcing\Store\ReadOnlyStore; use Patchlevel\EventSourcing\Store\Store; @@ -102,7 +103,9 @@ use Patchlevel\EventSourcing\Subscription\Cleanup\DefaultCleaner; use Patchlevel\EventSourcing\Subscription\Engine\CatchUpSubscriptionEngine; use Patchlevel\EventSourcing\Subscription\Engine\DefaultSubscriptionEngine; +use Patchlevel\EventSourcing\Subscription\Engine\Event\OnSubscriptionRemoved; use Patchlevel\EventSourcing\Subscription\Engine\GapResolverStoreMessageLoader; +use Patchlevel\EventSourcing\Subscription\Engine\Listener\RemoveSubscriptionStreamListener; use Patchlevel\EventSourcing\Subscription\Engine\MessageLoader; use Patchlevel\EventSourcing\Subscription\Engine\StoreMessageLoader; use Patchlevel\EventSourcing\Subscription\Engine\SubscriptionEngine; @@ -115,6 +118,7 @@ use Patchlevel\EventSourcing\Subscription\Store\InMemorySubscriptionStore; use Patchlevel\EventSourcing\Subscription\Store\SubscriptionStore; use Patchlevel\EventSourcing\Subscription\Subscriber\ArgumentResolver\ArgumentResolver; +use Patchlevel\EventSourcing\Subscription\Subscriber\ArgumentResolver\EventEmitterResolver; use Patchlevel\EventSourcing\Subscription\Subscriber\ArgumentResolver\LookupResolver; use Patchlevel\EventSourcing\Subscription\Subscriber\MetadataSubscriberAccessorRepository; use Patchlevel\EventSourcing\Subscription\Subscriber\SubscriberAccessorRepository; @@ -156,6 +160,7 @@ use Symfony\Component\EventDispatcher\EventDispatcher; use function class_exists; +use function preg_match; use function sprintf; /** @psalm-import-type Config from Configuration */ @@ -533,6 +538,7 @@ static function (ChildDefinition $definition): void { ]); $this->configureSyncSubscription($config, $container); + $this->configureEventEmitter($config, $container); if ($config['subscription']['auto_setup']['enabled']) { $container->register(AutoSetupListener::class) @@ -599,6 +605,33 @@ private function configureSyncSubscription(array $config, ContainerBuilder $cont ]); } + /** @param Config $config */ + private function configureEventEmitter(array $config, ContainerBuilder $container): void + { + if (!$config['subscription']['event_emitter']['enabled']) { + return; + } + + if ($config['store']['read_only']) { + throw new InvalidArgumentException('Event emitter does not support a read only store'); + } + + $container->register(EventEmitterResolver::class) + ->setArguments([new Reference(Store::class)]) + ->addTag('event_sourcing.argument_resolver'); + + $container->register(RemoveSubscriptionStreamListener::class) + ->setArguments([ + new Reference(Store::class), + new Reference('logger', ContainerInterface::NULL_ON_INVALID_REFERENCE), + ]) + ->addTag('kernel.event_listener', [ + 'event' => OnSubscriptionRemoved::class, + 'dispatcher' => 'event_sourcing.subscription.event_dispatcher', + ]) + ->addTag('monolog.logger', ['channel' => 'event_sourcing']); + } + /** @param Config $config */ private function configureHydrator(array $config, ContainerBuilder $container): void { @@ -681,11 +714,19 @@ private function configureConnection(array $config, ContainerBuilder $container) return; } + $middlewares = []; + + if ($this->usesTaggableStore($config)) { + $container->register(PostgreSQLPlatformMiddleware::class); + $middlewares[] = new Reference(PostgreSQLPlatformMiddleware::class); + } + if ($config['connection']['url'] !== null) { $container->register('event_sourcing.dbal_connection', Connection::class) ->setFactory([DbalConnectionFactory::class, 'createConnection']) ->setArguments([ $config['connection']['url'], + $middlewares, ]); if ($config['connection']['provide_dedicated_connection']) { @@ -710,6 +751,26 @@ private function configureConnection(array $config, ContainerBuilder $container) } $container->setAlias('event_sourcing.dbal_connection', $config['connection']['service']); + + if ( + $middlewares === [] + || !preg_match('/^doctrine\.dbal\.(.+)_connection$/', $config['connection']['service'], $matches) + ) { + return; + } + + $container->getDefinition(PostgreSQLPlatformMiddleware::class) + ->addTag('doctrine.middleware', ['connection' => $matches[1]]); + } + + /** @param Config $config */ + private function usesTaggableStore(array $config): bool + { + return $config['store']['type'] === 'dbal_taggable' + || ( + $config['store']['migrate_to_new_store']['enabled'] + && $config['store']['migrate_to_new_store']['type'] === 'dbal_taggable' + ); } /** @param Config $config */ @@ -1196,19 +1257,19 @@ private function configureDCB(array $config, ContainerBuilder $container): void return; } - if ($config['store']['type'] !== 'dbal_taggable') { + if ($config['store']['type'] === 'dbal_stream') { throw new InvalidArgumentException( - 'DCB requires a taggable store, please use "dbal_taggable" as store type.', + 'DCB requires a store that supports appending, please use "dbal_taggable", "in_memory" or a custom store.', ); } $container->register(StoreDecisionModelBuilder::class) - ->setArguments([new Reference(TaggableDoctrineDbalStore::class)]); + ->setArguments([new Reference(Store::class)]); $container->setAlias(DecisionModelBuilder::class, StoreDecisionModelBuilder::class); $container->register(StoreEventAppender::class) ->setArguments([ - new Reference(TaggableDoctrineDbalStore::class), + new Reference(Store::class), new Reference(EventTagExtractor::class), ]); $container->setAlias(EventAppender::class, StoreEventAppender::class); diff --git a/src/Doctrine/DbalConnectionFactory.php b/src/Doctrine/DbalConnectionFactory.php index babf31bb..89dbbc33 100644 --- a/src/Doctrine/DbalConnectionFactory.php +++ b/src/Doctrine/DbalConnectionFactory.php @@ -4,7 +4,9 @@ namespace Patchlevel\EventSourcingBundle\Doctrine; +use Doctrine\DBAL\Configuration; use Doctrine\DBAL\Connection; +use Doctrine\DBAL\Driver\Middleware; use Doctrine\DBAL\DriverManager; use Doctrine\DBAL\Tools\DsnParser; @@ -27,10 +29,12 @@ final class DbalConnectionFactory 'sqlite3' => 'pdo_sqlite', ]; - public static function createConnection(string $url): Connection + /** @param list $middlewares */ + public static function createConnection(string $url, array $middlewares = []): Connection { return DriverManager::getConnection( (new DsnParser(self::DEFAULT_SCHEME_MAP))->parse($url), + (new Configuration())->setMiddlewares($middlewares), ); } } diff --git a/tests/Unit/PatchlevelEventSourcingBundleTest.php b/tests/Unit/PatchlevelEventSourcingBundleTest.php index 6d46b1ae..9f47af52 100644 --- a/tests/Unit/PatchlevelEventSourcingBundleTest.php +++ b/tests/Unit/PatchlevelEventSourcingBundleTest.php @@ -72,6 +72,8 @@ use Patchlevel\EventSourcing\Snapshot\DefaultSnapshotStore; use Patchlevel\EventSourcing\Snapshot\SnapshotStore; use Patchlevel\EventSourcing\Store\ArchivedHeader; +use Patchlevel\EventSourcing\Store\Dbal\PostgreSQLPlatform as EventSourcingPostgreSQLPlatform; +use Patchlevel\EventSourcing\Store\Dbal\PostgreSQLPlatformMiddleware; use Patchlevel\EventSourcing\Store\InMemoryStore; use Patchlevel\EventSourcing\Store\ReadOnlyStore; use Patchlevel\EventSourcing\Store\Store; @@ -82,7 +84,9 @@ use Patchlevel\EventSourcing\Subscription\Cleanup\DefaultCleaner; use Patchlevel\EventSourcing\Subscription\Engine\CatchUpSubscriptionEngine; use Patchlevel\EventSourcing\Subscription\Engine\DefaultSubscriptionEngine; +use Patchlevel\EventSourcing\Subscription\Engine\Event\OnSubscriptionRemoved; use Patchlevel\EventSourcing\Subscription\Engine\GapResolverStoreMessageLoader; +use Patchlevel\EventSourcing\Subscription\Engine\Listener\RemoveSubscriptionStreamListener; use Patchlevel\EventSourcing\Subscription\Engine\MessageLoader; use Patchlevel\EventSourcing\Subscription\Engine\StoreMessageLoader; use Patchlevel\EventSourcing\Subscription\Engine\SubscriptionEngine; @@ -95,6 +99,7 @@ use Patchlevel\EventSourcing\Subscription\Store\DoctrineSubscriptionStore; use Patchlevel\EventSourcing\Subscription\Store\InMemorySubscriptionStore; use Patchlevel\EventSourcing\Subscription\Store\SubscriptionStore; +use Patchlevel\EventSourcing\Subscription\Subscriber\ArgumentResolver\EventEmitterResolver; use Patchlevel\EventSourcingBundle\DependencyInjection\PatchlevelEventSourcingExtension; use Patchlevel\EventSourcingBundle\EventBus\SymfonyEventBus; use Patchlevel\EventSourcingBundle\Normalizer\SymfonyExtension; @@ -623,6 +628,96 @@ public function testDCB(): void self::assertInstanceOf(StoreDecisionModelBuilder::class, $container->get(DecisionModelBuilder::class)); } + public function testDCBWithInMemoryStore(): void + { + $container = new ContainerBuilder(); + + $this->compileContainer( + $container, + [ + 'patchlevel_event_sourcing' => [ + 'connection' => ['service' => 'doctrine.dbal.eventstore_connection'], + 'store' => ['type' => 'in_memory'], + 'dcb' => true, + ], + ], + ); + + self::assertInstanceOf(StoreEventAppender::class, $container->get(EventAppender::class)); + self::assertInstanceOf(StoreDecisionModelBuilder::class, $container->get(DecisionModelBuilder::class)); + } + + public function testDCBWithStreamStore(): void + { + $this->expectException(InvalidArgumentException::class); + + $this->compileContainer( + new ContainerBuilder(), + [ + 'patchlevel_event_sourcing' => [ + 'connection' => ['service' => 'doctrine.dbal.eventstore_connection'], + 'dcb' => true, + ], + ], + ); + } + + public function testPostgreSQLPlatformMiddlewareForDoctrineConnection(): void + { + $container = new ContainerBuilder(); + + $this->compileContainer( + $container, + [ + 'patchlevel_event_sourcing' => [ + 'connection' => ['service' => 'doctrine.dbal.eventstore_connection'], + 'store' => ['type' => 'dbal_taggable'], + ], + ], + ); + + self::assertEquals( + [['connection' => 'eventstore']], + $container->getDefinition(PostgreSQLPlatformMiddleware::class)->getTag('doctrine.middleware'), + ); + } + + public function testPostgreSQLPlatformMiddlewareForUrlConnection(): void + { + $container = new ContainerBuilder(); + + $this->compileContainer( + $container, + [ + 'patchlevel_event_sourcing' => [ + 'connection' => ['url' => 'pdo-pgsql://user:secret@localhost/app?serverVersion=16'], + 'store' => ['type' => 'dbal_taggable'], + ], + ], + ); + + $connection = $container->get('event_sourcing.dbal_connection'); + + self::assertInstanceOf(Connection::class, $connection); + self::assertInstanceOf(EventSourcingPostgreSQLPlatform::class, $connection->getDatabasePlatform()); + } + + public function testNoPostgreSQLPlatformMiddlewareForStreamStore(): void + { + $container = new ContainerBuilder(); + + $this->compileContainer( + $container, + [ + 'patchlevel_event_sourcing' => [ + 'connection' => ['service' => 'doctrine.dbal.eventstore_connection'], + ], + ], + ); + + self::assertFalse($container->has(PostgreSQLPlatformMiddleware::class)); + } + public function testMessageLoader(): void { $container = new ContainerBuilder(); @@ -1431,6 +1526,50 @@ public function testSubscriptionEventDispatcher(): void ); } + public function testEventEmitter(): void + { + $container = new ContainerBuilder(); + + $this->compileContainer( + $container, + [ + 'patchlevel_event_sourcing' => [ + 'connection' => ['service' => 'doctrine.dbal.eventstore_connection'], + 'subscription' => ['event_emitter' => true], + ], + ], + ); + + self::assertTrue( + $container->getDefinition(EventEmitterResolver::class)->hasTag('event_sourcing.argument_resolver'), + ); + self::assertEquals( + [ + [ + 'event' => OnSubscriptionRemoved::class, + 'dispatcher' => 'event_sourcing.subscription.event_dispatcher', + ], + ], + $container->getDefinition(RemoveSubscriptionStreamListener::class)->getTag('kernel.event_listener'), + ); + } + + public function testEventEmitterWithReadOnlyStore(): void + { + $this->expectException(InvalidArgumentException::class); + + $this->compileContainer( + new ContainerBuilder(), + [ + 'patchlevel_event_sourcing' => [ + 'connection' => ['service' => 'doctrine.dbal.eventstore_connection'], + 'store' => ['read_only' => true], + 'subscription' => ['event_emitter' => true], + ], + ], + ); + } + public function testRetryStrategy(): void { $container = new ContainerBuilder();