From 11322523710984107dd1acc0669790b911cb3d2c Mon Sep 17 00:00:00 2001 From: David Badura Date: Sat, 3 Oct 2026 13:23:29 +0200 Subject: [PATCH] Replace catch_up, throw_on_error and run_after_aggregate_save with sync option Subscriptions can now be processed sync after an aggregate has been saved, filtered by ids and groups, while the rest is still handled by the worker. Catch up and throw on error only apply to this sync run and no longer decorate the global subscription engine used by the worker and commands. --- docs/UPGRADE-4.0.md | 52 +++++++++ docs/configuration.md | 73 +++++++++---- docs/installation.md | 10 +- src/DependencyInjection/Configuration.php | 24 ++--- .../PatchlevelEventSourcingExtension.php | 60 ++++++----- .../PatchlevelEventSourcingBundleTest.php | 101 +++++++++++++----- 6 files changed, 219 insertions(+), 101 deletions(-) diff --git a/docs/UPGRADE-4.0.md b/docs/UPGRADE-4.0.md index 1db502ab..15dc432c 100644 --- a/docs/UPGRADE-4.0.md +++ b/docs/UPGRADE-4.0.md @@ -14,6 +14,58 @@ This guide only covers the changes of the bundle itself. ## Subscription +### Sync Subscriptions + +The options `catch_up`, `throw_on_error` and `run_after_aggregate_save` have been removed +in favor of the new `sync` option. +Before, `catch_up` and `throw_on_error` decorated the global subscription engine, +so they also affected the worker and the console commands. +Now they only apply to the sync run after an aggregate has been saved. + +before: + +```yaml +when@dev: + patchlevel_event_sourcing: + subscription: + catch_up: true + throw_on_error: true + run_after_aggregate_save: true +``` +after: + +```yaml +when@dev: + patchlevel_event_sourcing: + subscription: + sync: + throw_on_error: true +``` +Catching up is now always active for sync runs. +The `limit` option of `catch_up` is now called `catch_up_limit`. +The `limit` option of `run_after_aggregate_save` has been removed without replacement. + +before: + +```yaml +patchlevel_event_sourcing: + subscription: + catch_up: + limit: 10 + run_after_aggregate_save: + groups: ['sync'] +``` +after: + +```yaml +patchlevel_event_sourcing: + subscription: + sync: + groups: ['sync'] + catch_up_limit: 10 +``` +If you run your tests with a worker-less setup, replace the options in `when@test` the same way. + ### Retry Strategy The deprecated `retry_strategy` option has been removed. Use `retry_strategies` instead. diff --git a/docs/configuration.md b/docs/configuration.md index 8beb9871..2adbfd1f 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -352,43 +352,76 @@ If you are using the [doctrine-test-bundle](https://github.com/dmaicher/doctrine you can use the `static_in_memory` store for testing. ::: -### Catch Up +### Sync Subscriptions -If aggregates are used in the processors and new events are generated there, -then they are not part of the current subscription engine `run` and will only be processed during the next run or boot. -This is usually not a problem in dev or prod environment because a worker is used -and these events will be processed at some point. But in testing it is not so easy. -For this reason, you can activate the `catch_up` option. +By default, all subscriptions are processed asynchronously by a worker +that runs the `event-sourcing:subscription:run` command. +If you want subscriptions to be processed directly after an aggregate has been saved, +you can activate the `sync` option. +Without further configuration, all subscriptions are processed synchronously. +This is useful for development and testing, so you don't have to run a worker. ```yaml -patchlevel_event_sourcing: - subscription: - catch_up: true +when@dev: + patchlevel_event_sourcing: + subscription: + sync: true ``` -### Throw on Error - -You can activate the `throw_on_error` option to throw an exception if a subscription engine run has an error. -This is useful for testing or development to get directly feedback if something is wrong. +In production you usually want only some subscriptions to be processed synchronously, +e.g. projections that must be up to date in the same request. +For this you can filter the subscriptions by `ids` and `groups`. ```yaml patchlevel_event_sourcing: subscription: - throw_on_error: true + sync: + groups: ['sync'] ``` -:::warning -This option should not be used in production. The normal behavior is to log the error and continue. +The group can be defined directly on the subscriber. + +```php +use Patchlevel\EventSourcing\Attribute\Projector; + +#[Projector('profile', group: 'sync')] +final class ProfileProjector +{ + // ... +} +``` +:::note +If you define both `ids` and `groups`, a subscription must match both to be processed synchronously. ::: -### Run After Aggregate Save +:::note +Sync subscriptions are processed in the same process as the request. +Keep them fast, otherwise your requests will be slowed down. +The worker still processes sync subscriptions as well, for example to retry failed ones. +::: -If you want to run the subscription engine after an aggregate is saved, you can activate this option. -This is useful for testing or development, so you don't have run a worker to process the events. +If a sync subscriber saves an aggregate itself, the new events are processed in the same run as well. +You can limit how often the subscription engine catches up with the `catch_up_limit` option. ```yaml patchlevel_event_sourcing: subscription: - run_after_aggregate_save: true + sync: + catch_up_limit: 10 ``` +You can also activate the `throw_on_error` option to throw an exception if a sync subscription has an error. +This is useful for testing or development to get directly feedback if something is wrong. + +```yaml +when@dev: + patchlevel_event_sourcing: + subscription: + sync: + throw_on_error: true +``` +:::warning +This option should not be used in production. The normal behavior is to log the error and continue. +This option only affects the sync run after an aggregate has been saved. The worker and the console commands are not affected. +::: + ### Auto Setup If you want to automatically setup the subscription engine, you can activate this option. diff --git a/docs/installation.md b/docs/installation.md index f023176d..ae1f471a 100644 --- a/docs/installation.md +++ b/docs/installation.md @@ -55,9 +55,8 @@ patchlevel_event_sourcing: when@dev: patchlevel_event_sourcing: subscription: - catch_up: true - throw_on_error: true - run_after_aggregate_save: true + sync: + throw_on_error: true rebuild_after_file_change: true auto_setup: true @@ -66,9 +65,8 @@ when@test: subscription: store: type: 'static_in_memory' - catch_up: true - throw_on_error: true - run_after_aggregate_save: true + sync: + throw_on_error: true ``` ## Dotenv diff --git a/src/DependencyInjection/Configuration.php b/src/DependencyInjection/Configuration.php index 85096bff..b75982ce 100644 --- a/src/DependencyInjection/Configuration.php +++ b/src/DependencyInjection/Configuration.php @@ -30,13 +30,12 @@ * }, * retry_strategies: array}>, * default_retry_strategy: string, - * catch_up: array{enabled: bool, limit: positive-int|null}, - * throw_on_error: array{enabled: bool}, - * run_after_aggregate_save: array{ + * sync: array{ * enabled: bool, * ids: list, * groups: list, - * limit: positive-int|null + * catch_up_limit: positive-int|null, + * throw_on_error: bool * }, * auto_setup: array{ * enabled: bool, @@ -245,25 +244,14 @@ public function getConfigTreeBuilder(): TreeBuilder ->scalarNode('default_retry_strategy')->defaultValue('default')->end() - ->arrayNode('catch_up') - ->canBeEnabled() - ->addDefaultsIfNotSet() - ->children() - ->integerNode('limit')->defaultNull()->end() - ->end() - ->end() - - ->arrayNode('throw_on_error') - ->canBeEnabled() - ->end() - - ->arrayNode('run_after_aggregate_save') + ->arrayNode('sync') ->canBeEnabled() ->addDefaultsIfNotSet() ->children() ->arrayNode('ids')->scalarPrototype()->end()->end() ->arrayNode('groups')->scalarPrototype()->end()->end() - ->integerNode('limit')->defaultNull()->end() + ->integerNode('catch_up_limit')->defaultNull()->min(1)->end() + ->booleanNode('throw_on_error')->defaultFalse()->end() ->end() ->end() diff --git a/src/DependencyInjection/PatchlevelEventSourcingExtension.php b/src/DependencyInjection/PatchlevelEventSourcingExtension.php index 22e730b7..412e6313 100644 --- a/src/DependencyInjection/PatchlevelEventSourcingExtension.php +++ b/src/DependencyInjection/PatchlevelEventSourcingExtension.php @@ -519,34 +519,7 @@ static function (ChildDefinition $definition): void { 'method' => 'onWorkerRunningEvent', ]); - if ($config['subscription']['throw_on_error']['enabled']) { - $container->register(ThrowOnErrorSubscriptionEngine::class) - ->setDecoratedService(SubscriptionEngine::class) - ->setArguments([ - new Reference('.inner'), - ]); - } - - if ($config['subscription']['catch_up']['enabled']) { - $container->register(CatchUpSubscriptionEngine::class) - ->setDecoratedService(SubscriptionEngine::class) - ->setArguments([ - new Reference('.inner'), - $config['subscription']['catch_up']['limit'], - ]); - } - - if ($config['subscription']['run_after_aggregate_save']['enabled']) { - $container->register(RunSubscriptionEngineRepositoryManager::class) - ->setDecoratedService(RepositoryManager::class) - ->setArguments([ - new Reference('.inner'), - new Reference(SubscriptionEngine::class), - $config['subscription']['run_after_aggregate_save']['ids'] ?: null, - $config['subscription']['run_after_aggregate_save']['groups'] ?: null, - $config['subscription']['run_after_aggregate_save']['limit'], - ]); - } + $this->configureSyncSubscription($config, $container); if ($config['subscription']['auto_setup']['enabled']) { $container->register(AutoSetupListener::class) @@ -582,6 +555,37 @@ static function (ChildDefinition $definition): void { ]); } + /** @param Config $config */ + private function configureSyncSubscription(array $config, ContainerBuilder $container): void + { + if (!$config['subscription']['sync']['enabled']) { + return; + } + + $container->register('event_sourcing.subscription.sync_engine', CatchUpSubscriptionEngine::class) + ->setArguments([ + new Reference(DefaultSubscriptionEngine::class), + $config['subscription']['sync']['catch_up_limit'], + ]); + + if ($config['subscription']['sync']['throw_on_error']) { + $container->register('event_sourcing.subscription.sync_engine.throw_on_error', ThrowOnErrorSubscriptionEngine::class) + ->setDecoratedService('event_sourcing.subscription.sync_engine') + ->setArguments([ + new Reference('.inner'), + ]); + } + + $container->register(RunSubscriptionEngineRepositoryManager::class) + ->setDecoratedService(RepositoryManager::class) + ->setArguments([ + new Reference('.inner'), + new Reference('event_sourcing.subscription.sync_engine'), + $config['subscription']['sync']['ids'] ?: null, + $config['subscription']['sync']['groups'] ?: null, + ]); + } + /** @param Config $config */ private function configureHydrator(array $config, ContainerBuilder $container): void { diff --git a/tests/Unit/PatchlevelEventSourcingBundleTest.php b/tests/Unit/PatchlevelEventSourcingBundleTest.php index 48944bf6..96f6e20e 100644 --- a/tests/Unit/PatchlevelEventSourcingBundleTest.php +++ b/tests/Unit/PatchlevelEventSourcingBundleTest.php @@ -83,6 +83,7 @@ use Patchlevel\EventSourcing\Subscription\Engine\MessageLoader; use Patchlevel\EventSourcing\Subscription\Engine\StoreMessageLoader; use Patchlevel\EventSourcing\Subscription\Engine\SubscriptionEngine; +use Patchlevel\EventSourcing\Subscription\Engine\ThrowOnErrorSubscriptionEngine; use Patchlevel\EventSourcing\Subscription\Repository\RunSubscriptionEngineRepositoryManager; use Patchlevel\EventSourcing\Subscription\RetryStrategy\ClockBasedRetryStrategy; use Patchlevel\EventSourcing\Subscription\RetryStrategy\NoRetryStrategy; @@ -1083,7 +1084,58 @@ public function testDecorator(): void self::assertInstanceOf(SplitStreamDecorator::class, $container->get(SplitStreamDecorator::class)); } - public function testRunSubscriptionEngineRepositoryManager(): void + public function testSyncDisabledByDefault(): void + { + $container = new ContainerBuilder(); + + $this->compileContainer( + $container, + [ + 'patchlevel_event_sourcing' => [ + 'connection' => ['service' => 'doctrine.dbal.eventstore_connection'], + ], + ], + ); + + self::assertFalse($container->hasDefinition(RunSubscriptionEngineRepositoryManager::class)); + self::assertFalse($container->hasDefinition('event_sourcing.subscription.sync_engine')); + self::assertNotInstanceOf( + RunSubscriptionEngineRepositoryManager::class, + $container->get(RepositoryManager::class), + ); + } + + public function testSyncAllSubscriptions(): void + { + $container = new ContainerBuilder(); + + $this->compileContainer( + $container, + [ + 'patchlevel_event_sourcing' => [ + 'connection' => ['service' => 'doctrine.dbal.eventstore_connection'], + 'subscription' => ['sync' => true], + ], + ], + ); + + self::assertInstanceOf( + RunSubscriptionEngineRepositoryManager::class, + $container->get(RepositoryManager::class), + ); + + $definition = $container->getDefinition(RunSubscriptionEngineRepositoryManager::class); + self::assertNull($definition->getArgument(2)); + self::assertNull($definition->getArgument(3)); + + self::assertInstanceOf( + CatchUpSubscriptionEngine::class, + $container->get('event_sourcing.subscription.sync_engine'), + ); + self::assertInstanceOf(DefaultSubscriptionEngine::class, $container->get(SubscriptionEngine::class)); + } + + public function testSyncFilteredSubscriptions(): void { $container = new ContainerBuilder(); @@ -1093,10 +1145,11 @@ public function testRunSubscriptionEngineRepositoryManager(): void 'patchlevel_event_sourcing' => [ 'connection' => ['service' => 'doctrine.dbal.eventstore_connection'], 'subscription' => [ - 'run_after_aggregate_save' => [ - 'ids' => ['a'], - 'groups' => ['b'], - 'limit' => 10, + 'sync' => [ + 'ids' => ['profile'], + 'groups' => ['sync'], + 'catch_up_limit' => 10, + 'throw_on_error' => true, ], ], ], @@ -1107,6 +1160,16 @@ public function testRunSubscriptionEngineRepositoryManager(): void RunSubscriptionEngineRepositoryManager::class, $container->get(RepositoryManager::class), ); + + $definition = $container->getDefinition(RunSubscriptionEngineRepositoryManager::class); + self::assertSame(['profile'], $definition->getArgument(2)); + self::assertSame(['sync'], $definition->getArgument(3)); + + self::assertInstanceOf( + ThrowOnErrorSubscriptionEngine::class, + $container->get('event_sourcing.subscription.sync_engine'), + ); + self::assertInstanceOf(DefaultSubscriptionEngine::class, $container->get(SubscriptionEngine::class)); } public function testSubscriptionEngineInMemoryStore(): void @@ -1163,28 +1226,6 @@ public function testSubscriptionEngineStaticInMemoryStore(): void ); } - public function testCatchUpSubscriptionEngine(): void - { - $container = new ContainerBuilder(); - - $this->compileContainer( - $container, - [ - 'patchlevel_event_sourcing' => [ - 'connection' => ['service' => 'doctrine.dbal.eventstore_connection'], - 'subscription' => [ - 'catch_up' => ['limit' => 10], - ], - ], - ], - ); - - self::assertInstanceOf( - CatchUpSubscriptionEngine::class, - $container->get(SubscriptionEngine::class), - ); - } - public function testSubscriberSameConnectionError(): void { $container = new ContainerBuilder(); @@ -1446,8 +1487,10 @@ public function testFullBuild(): void ], ], 'subscription' => [ - 'catch_up' => ['limit' => 10], - 'throw_on_error' => true, + 'sync' => [ + 'catch_up_limit' => 10, + 'throw_on_error' => true, + ], ], ], ],