Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
52 changes: 52 additions & 0 deletions docs/UPGRADE-4.0.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
73 changes: 53 additions & 20 deletions docs/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
10 changes: 4 additions & 6 deletions docs/installation.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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

Expand Down
24 changes: 6 additions & 18 deletions src/DependencyInjection/Configuration.php
Original file line number Diff line number Diff line change
Expand Up @@ -30,13 +30,12 @@
* },
* retry_strategies: array<string, array{type: string, service: string, options: array<string, mixed>}>,
* 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<string>,
* groups: list<string>,
* limit: positive-int|null
* catch_up_limit: positive-int|null,
* throw_on_error: bool
* },
* auto_setup: array{
* enabled: bool,
Expand Down Expand Up @@ -206,7 +205,7 @@
->addDefaultsIfNotSet()
->children()
->enumNode('type')
->values(['dbal', 'in_memory', 'static_in_memory', 'custom'])

Check warning on line 208 in src/DependencyInjection/Configuration.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.5, ubuntu-latest)

Escaped Mutant for Mutator "ArrayItemRemoval": @@ @@ ->addDefaultsIfNotSet() ->children() ->enumNode('type') - ->values(['dbal', 'in_memory', 'static_in_memory', 'custom']) + ->values(['in_memory', 'static_in_memory', 'custom']) ->defaultValue('dbal') ->end() ->scalarNode('service')->defaultNull()->end()
->defaultValue('dbal')
->end()
->scalarNode('service')->defaultNull()->end()
Expand All @@ -228,13 +227,13 @@
->arrayNode('options')->variablePrototype()->end()->end()
->end()
->end()
->defaultValue([

Check warning on line 230 in src/DependencyInjection/Configuration.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.5, ubuntu-latest)

Escaped Mutant for Mutator "ArrayItemRemoval": @@ @@ ->end() ->end() ->defaultValue([ - 'default' => [ - 'type' => 'clock_based', - 'options' => [ - 'base_delay' => 5, - 'delay_factor' => 2, - 'max_attempts' => 5, - ], - ], 'no_retry' => [ 'type' => 'no_retry', ],
'default' => [
'type' => 'clock_based',
'options' => [

Check warning on line 233 in src/DependencyInjection/Configuration.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.5, ubuntu-latest)

Escaped Mutant for Mutator "ArrayItemRemoval": @@ @@ 'default' => [ 'type' => 'clock_based', 'options' => [ - 'base_delay' => 5, 'delay_factor' => 2, 'max_attempts' => 5, ],
'base_delay' => 5,

Check warning on line 234 in src/DependencyInjection/Configuration.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.5, ubuntu-latest)

Escaped Mutant for Mutator "IncrementInteger": @@ @@ 'default' => [ 'type' => 'clock_based', 'options' => [ - 'base_delay' => 5, + 'base_delay' => 6, 'delay_factor' => 2, 'max_attempts' => 5, ],

Check warning on line 234 in src/DependencyInjection/Configuration.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.5, ubuntu-latest)

Escaped Mutant for Mutator "DecrementInteger": @@ @@ 'default' => [ 'type' => 'clock_based', 'options' => [ - 'base_delay' => 5, + 'base_delay' => 4, 'delay_factor' => 2, 'max_attempts' => 5, ],
'delay_factor' => 2,

Check warning on line 235 in src/DependencyInjection/Configuration.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.5, ubuntu-latest)

Escaped Mutant for Mutator "IncrementInteger": @@ @@ 'type' => 'clock_based', 'options' => [ 'base_delay' => 5, - 'delay_factor' => 2, + 'delay_factor' => 3, 'max_attempts' => 5, ], ],

Check warning on line 235 in src/DependencyInjection/Configuration.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.5, ubuntu-latest)

Escaped Mutant for Mutator "DecrementInteger": @@ @@ 'type' => 'clock_based', 'options' => [ 'base_delay' => 5, - 'delay_factor' => 2, + 'delay_factor' => 1, 'max_attempts' => 5, ], ],
'max_attempts' => 5,

Check warning on line 236 in src/DependencyInjection/Configuration.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.5, ubuntu-latest)

Escaped Mutant for Mutator "IncrementInteger": @@ @@ 'options' => [ 'base_delay' => 5, 'delay_factor' => 2, - 'max_attempts' => 5, + 'max_attempts' => 6, ], ], 'no_retry' => [

Check warning on line 236 in src/DependencyInjection/Configuration.php

View workflow job for this annotation

GitHub Actions / Mutation tests (locked, 8.5, ubuntu-latest)

Escaped Mutant for Mutator "DecrementInteger": @@ @@ 'options' => [ 'base_delay' => 5, 'delay_factor' => 2, - 'max_attempts' => 5, + 'max_attempts' => 4, ], ], 'no_retry' => [
],
],
'no_retry' => [
Expand All @@ -245,25 +244,14 @@

->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()

Check warning on line 253 in src/DependencyInjection/Configuration.php

View workflow job for this annotation

GitHub Actions / Mutation tests on diff (locked, 8.5, ubuntu-latest)

Escaped Mutant for Mutator "DecrementInteger": @@ @@ ->children() ->arrayNode('ids')->scalarPrototype()->end()->end() ->arrayNode('groups')->scalarPrototype()->end()->end() - ->integerNode('catch_up_limit')->defaultNull()->min(1)->end() + ->integerNode('catch_up_limit')->defaultNull()->min(0)->end() ->booleanNode('throw_on_error')->defaultFalse()->end() ->end() ->end()

Check warning on line 253 in src/DependencyInjection/Configuration.php

View workflow job for this annotation

GitHub Actions / Mutation tests on diff (locked, 8.5, ubuntu-latest)

Escaped Mutant for Mutator "IncrementInteger": @@ @@ ->children() ->arrayNode('ids')->scalarPrototype()->end()->end() ->arrayNode('groups')->scalarPrototype()->end()->end() - ->integerNode('catch_up_limit')->defaultNull()->min(1)->end() + ->integerNode('catch_up_limit')->defaultNull()->min(2)->end() ->booleanNode('throw_on_error')->defaultFalse()->end() ->end() ->end()
->booleanNode('throw_on_error')->defaultFalse()->end()
->end()
->end()

Expand Down
60 changes: 32 additions & 28 deletions src/DependencyInjection/PatchlevelEventSourcingExtension.php
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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
{
Expand Down
Loading
Loading