diff --git a/docs/configuration.md b/docs/configuration.md index 43b9d18b..e9c136be 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -284,7 +284,7 @@ patchlevel_event_sourcing: :::tip You can find out more about subscriptions in the library -[documentation](https://event-sourcing.patchlevel.io/latest/subscription/). +[documentation](/docs/event-sourcing/latest/subscription). ::: ### Store @@ -456,6 +456,29 @@ patchlevel_event_sourcing: gap_detection: detection_window: 'PT5M' ``` +### Engine Events + +The subscription engine dispatches events while processing, e.g. `OnHandleMessageError` or `OnResult`. +It uses its own event dispatcher `event_sourcing.subscription.event_dispatcher`, +so you have to pass it to the `AsEventListener` attribute. + +```php +use Patchlevel\EventSourcing\Subscription\Engine\Event\OnHandleMessageError; +use Symfony\Component\EventDispatcher\Attribute\AsEventListener; + +#[AsEventListener(dispatcher: 'event_sourcing.subscription.event_dispatcher')] +final class SubscriptionErrorListener +{ + public function __invoke(OnHandleMessageError $event): void + { + // logging, metrics, ... + } +} +``` +:::note +You can find all available events in the [library documentation](/docs/event-sourcing/latest/subscription). +::: + ## Command Bus You can enable the command bus integration to use your aggregates as command handlers. diff --git a/src/DependencyInjection/PatchlevelEventSourcingExtension.php b/src/DependencyInjection/PatchlevelEventSourcingExtension.php index e7e4b488..354df807 100644 --- a/src/DependencyInjection/PatchlevelEventSourcingExtension.php +++ b/src/DependencyInjection/PatchlevelEventSourcingExtension.php @@ -153,6 +153,7 @@ use Symfony\Component\DependencyInjection\Extension\Extension; use Symfony\Component\DependencyInjection\Parameter; use Symfony\Component\DependencyInjection\Reference; +use Symfony\Component\EventDispatcher\EventDispatcher; use function class_exists; use function sprintf; @@ -491,7 +492,6 @@ static function (ChildDefinition $definition): void { ->setArguments([ new TaggedIteratorArgument('event_sourcing.subscriber'), new Reference(SubscriberMetadataFactory::class), - new TaggedIteratorArgument('event_sourcing.argument_resolver'), ]); $container->setAlias(SubscriberAccessorRepository::class, MetadataSubscriberAccessorRepository::class); @@ -506,6 +506,8 @@ static function (ChildDefinition $definition): void { $container->setAlias(Cleaner::class, DefaultCleaner::class); + $container->register('event_sourcing.subscription.event_dispatcher', EventDispatcher::class); + $container->register(DefaultSubscriptionEngine::class) ->setArguments([ new Reference(MessageLoader::class), @@ -514,6 +516,8 @@ static function (ChildDefinition $definition): void { new Reference(RetryStrategyRepository::class), new Reference('logger', ContainerInterface::NULL_ON_INVALID_REFERENCE), new Reference(Cleaner::class), + new Reference('event_sourcing.subscription.event_dispatcher'), + new TaggedIteratorArgument('event_sourcing.argument_resolver'), ]) ->addTag('monolog.logger', ['channel' => 'event_sourcing']); @@ -849,6 +853,7 @@ private function configureStoreMigration(array $config, ContainerBuilder $contai ->setArguments([ new Reference('event_sourcing.dbal_connection'), new Reference(EventSerializer::class), + new Reference(EventRegistry::class), new Reference(HeadersSerializer::class), new Reference('event_sourcing.clock'), $config['store']['migrate_to_new_store']['options'], diff --git a/tests/Unit/PatchlevelEventSourcingBundleTest.php b/tests/Unit/PatchlevelEventSourcingBundleTest.php index 3cf46a06..6d46b1ae 100644 --- a/tests/Unit/PatchlevelEventSourcingBundleTest.php +++ b/tests/Unit/PatchlevelEventSourcingBundleTest.php @@ -76,6 +76,7 @@ use Patchlevel\EventSourcing\Store\ReadOnlyStore; use Patchlevel\EventSourcing\Store\Store; use Patchlevel\EventSourcing\Store\StreamDoctrineDbalStore; +use Patchlevel\EventSourcing\Store\TaggableDoctrineDbalStore; use Patchlevel\EventSourcing\Subscription\Cleanup\Cleaner; use Patchlevel\EventSourcing\Subscription\Cleanup\Dbal\DbalCleanupTaskHandler; use Patchlevel\EventSourcing\Subscription\Cleanup\DefaultCleaner; @@ -94,7 +95,6 @@ use Patchlevel\EventSourcing\Subscription\Store\DoctrineSubscriptionStore; use Patchlevel\EventSourcing\Subscription\Store\InMemorySubscriptionStore; use Patchlevel\EventSourcing\Subscription\Store\SubscriptionStore; -use Patchlevel\EventSourcing\Subscription\Subscriber\MetadataSubscriberAccessorRepository; use Patchlevel\EventSourcingBundle\DependencyInjection\PatchlevelEventSourcingExtension; use Patchlevel\EventSourcingBundle\EventBus\SymfonyEventBus; use Patchlevel\EventSourcingBundle\Normalizer\SymfonyExtension; @@ -135,6 +135,7 @@ use Symfony\Component\DependencyInjection\Definition; use Symfony\Component\DependencyInjection\Dumper\XmlDumper; use Symfony\Component\DependencyInjection\Reference; +use Symfony\Component\EventDispatcher\EventDispatcher; use Symfony\Component\HttpKernel\DependencyInjection\ServicesResetter; use Symfony\Component\Messenger\MessageBusInterface; @@ -354,6 +355,26 @@ public function testMigrateStore(): void ); } + public function testMigrateToTaggableStore(): void + { + $container = new ContainerBuilder(); + + $this->compileContainer( + $container, + [ + 'patchlevel_event_sourcing' => [ + 'connection' => ['service' => 'doctrine.dbal.eventstore_connection'], + 'store' => [ + 'migrate_to_new_store' => ['type' => 'dbal_taggable'], + ], + ], + ], + ); + + self::assertInstanceOf(StreamDoctrineDbalStore::class, $container->get(Store::class)); + self::assertInstanceOf(TaggableDoctrineDbalStore::class, $container->get('event_sourcing.store.new_store')); + } + public function testSymfonyEventBus(): void { $eventBus = $this->createMock(MessageBusInterface::class); @@ -1380,13 +1401,33 @@ public function testAutoconfigureArgumentResolver(): void ); self::assertTrue($container->getDefinition(DummyArgumentResolver::class)->hasTag('event_sourcing.argument_resolver')); - self::assertInstanceOf( - TaggedIteratorArgument::class, - $container->getDefinition(MetadataSubscriberAccessorRepository::class)->getArgument(2), + + $argument = $container->getDefinition(DefaultSubscriptionEngine::class)->getArgument(7); + + self::assertInstanceOf(TaggedIteratorArgument::class, $argument); + self::assertEquals('event_sourcing.argument_resolver', $argument->getTag()); + } + + public function testSubscriptionEventDispatcher(): void + { + $container = new ContainerBuilder(); + + $this->compileContainer( + $container, + [ + 'patchlevel_event_sourcing' => [ + 'connection' => ['service' => 'doctrine.dbal.eventstore_connection'], + ], + ], ); + self::assertEquals( - 'event_sourcing.argument_resolver', - $container->getDefinition(MetadataSubscriberAccessorRepository::class)->getArgument(2)->getTag(), + new Reference('event_sourcing.subscription.event_dispatcher'), + $container->getDefinition(DefaultSubscriptionEngine::class)->getArgument(6), + ); + self::assertInstanceOf( + EventDispatcher::class, + $container->get('event_sourcing.subscription.event_dispatcher'), ); }