Skip to content
Closed
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
66 changes: 66 additions & 0 deletions docs/split-stream.md
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,72 @@ This impacts only the aggregate loaded by the repository. Subscriptions will sti
You can combine this feature with the snapshot feature to increase the performance even more.
:::

## Dynamic Consistency Boundary

Projections of the [dynamic consistency boundary](dynamic-consistency-boundary.md) also respect `#[SplitStream]`.
If a projection has an apply method for such an event,
the decision model is only built from the last of these events matching the projection's tags onwards.
No configuration is needed and nothing is archived, other projections and subscriptions still see all events.

```php
use Patchlevel\EventSourcing\Attribute\Apply;
use Patchlevel\EventSourcing\Projection\BasicProjection;

final class Balance extends BasicProjection
{
public function __construct(
private readonly BankAccountId $bankAccountId,
) {
}

public function initialState(): int
{
return 0;
}

/** @return list<string> */
protected function tagFilter(): array
{
return ["bank_account:{$this->bankAccountId->toString()}"];
}

#[Apply]
public function applyBalanceReported(int $state, BalanceReported $event): int
{
return $event->balanceInCents;
}

#[Apply]
public function applyMoneyDeposited(int $state, MoneyDeposited $event): int
{
return $state + $event->amountInCents;
}
}
```
:::note
Only the split events of the projection itself are considered.
A `BalanceReported` of another bank account does not cut off the history of this one,
and other projections in the same decision model still get all events they need.
:::

If a projection handles a split event but needs the whole history anyway, for example to count them,
you can opt out by overriding `splitEvents`.

```php
use Patchlevel\EventSourcing\Projection\BasicProjection;

final class NumberOfBalanceReports extends BasicProjection
{
// ...

/** @return list<class-string> */
public function splitEvents(): array
{
return [];
}
}
```

## Learn more

* [How to use message decorator](message-decorator.md)
Expand Down
7 changes: 5 additions & 2 deletions src/DecisionModel/StoreDecisionModelBuilder.php
Original file line number Diff line number Diff line change
Expand Up @@ -6,10 +6,12 @@

use Patchlevel\EventSourcing\Projection\CompositeProjection;
use Patchlevel\EventSourcing\Projection\Projection;
use Patchlevel\EventSourcing\Projection\SplitStreamPositions;
use Patchlevel\EventSourcing\Store\AppendCondition;
use Patchlevel\EventSourcing\Store\AppendStore;

use function max;
use function min;

/** @experimental */
final class StoreDecisionModelBuilder implements DecisionModelBuilder
Expand All @@ -23,10 +25,11 @@
public function build(
array $projections,
): DecisionModel {
$projection = new CompositeProjection($projections);
$from = SplitStreamPositions::resolve($this->store, $projections);
$projection = new CompositeProjection($projections, $from);

$query = $projection->query();
$stream = $this->store->query($query);
$stream = $this->store->query($query, $from === [] ? 0 : min($from));

Check warning on line 32 in src/DecisionModel/StoreDecisionModelBuilder.php

View workflow job for this annotation

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

Escaped Mutant for Mutator "IncrementInteger": @@ @@ $projection = new CompositeProjection($projections, $from); $query = $projection->query(); - $stream = $this->store->query($query, $from === [] ? 0 : min($from)); + $stream = $this->store->query($query, $from === [] ? 1 : min($from)); $state = $projection->initialState();

$state = $projection->initialState();

Expand Down
14 changes: 13 additions & 1 deletion src/Projection/BasicProjection.php
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
namespace Patchlevel\EventSourcing\Projection;

use Patchlevel\EventSourcing\Attribute\Apply;
use Patchlevel\EventSourcing\Attribute\SplitStream;
use Patchlevel\EventSourcing\Message\Message;
use Patchlevel\EventSourcing\Store\SubQuery;
use ReflectionClass;
Expand All @@ -13,18 +14,20 @@
use ReflectionNamedType;
use ReflectionUnionType;

use function array_filter;
use function array_key_exists;
use function array_keys;
use function array_map;
use function array_merge;
use function array_values;
use function class_exists;

/**
* @experimental
* @template S = mixed
* @implements Projection<S>
*/
abstract class BasicProjection implements Projection, SubQueryProvider
abstract class BasicProjection implements Projection, SubQueryProvider, SplitStreamProvider
{
/** @var array<class-string, string>|null $applyMethods */
private array|null $applyMethods = null;
Expand Down Expand Up @@ -67,6 +70,15 @@ public function subQuery(): SubQuery
return $this->subQuery;
}

/** @return list<class-string> */
public function splitEvents(): array
{
return array_values(array_filter(
$this->eventTypeFilter(),
static fn (string $eventClass): bool => (new ReflectionClass($eventClass))->getAttributes(SplitStream::class) !== [],
));
}

/** @return list<class-string> */
protected function eventTypeFilter(): array
{
Expand Down
14 changes: 13 additions & 1 deletion src/Projection/CompositeProjection.php
Original file line number Diff line number Diff line change
Expand Up @@ -5,16 +5,21 @@
namespace Patchlevel\EventSourcing\Projection;

use Patchlevel\EventSourcing\Message\Message;
use Patchlevel\EventSourcing\Store\Header\IndexHeader;
use Patchlevel\EventSourcing\Store\Query;

use function array_map;

/** @experimental */
final class CompositeProjection
{
/** @param array<string, Projection> $projections */
/**
* @param array<string, Projection> $projections
* @param array<string, positive-int|0> $from messages before this index are not applied to the projection
*/
public function __construct(
private readonly array $projections,
private readonly array $from = [],
) {
}

Expand Down Expand Up @@ -51,6 +56,13 @@
public function apply(mixed $state, Message $message): mixed
{
foreach ($this->projections as $name => $projection) {
if (
($this->from[$name] ?? 0) > 0

Check warning on line 60 in src/Projection/CompositeProjection.php

View workflow job for this annotation

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

Escaped Mutant for Mutator "DecrementInteger": @@ @@ { foreach ($this->projections as $name => $projection) { if ( - ($this->from[$name] ?? 0) > 0 + ($this->from[$name] ?? -1) > 0 && $message->header(IndexHeader::class)->index < $this->from[$name] ) { continue;
&& $message->header(IndexHeader::class)->index < $this->from[$name]

Check warning on line 61 in src/Projection/CompositeProjection.php

View workflow job for this annotation

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

Escaped Mutant for Mutator "LessThanNegotiation": @@ @@ foreach ($this->projections as $name => $projection) { if ( ($this->from[$name] ?? 0) > 0 - && $message->header(IndexHeader::class)->index < $this->from[$name] + && $message->header(IndexHeader::class)->index >= $this->from[$name] ) { continue; }
) {
continue;

Check warning on line 63 in src/Projection/CompositeProjection.php

View workflow job for this annotation

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

Escaped Mutant for Mutator "Continue_": @@ @@ ($this->from[$name] ?? 0) > 0 && $message->header(IndexHeader::class)->index < $this->from[$name] ) { - continue; + break; } $state[$name] = $projection->apply($state[$name], $message);
}

$state[$name] = $projection->apply($state[$name], $message);
}

Expand Down
68 changes: 68 additions & 0 deletions src/Projection/SplitStreamPositions.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
<?php

declare(strict_types=1);

namespace Patchlevel\EventSourcing\Projection;

use Patchlevel\EventSourcing\Store\AppendStore;
use Patchlevel\EventSourcing\Store\Query;
use Patchlevel\EventSourcing\Store\SubQuery;

use function array_map;
use function array_values;
use function max;

/** @internal */
final class SplitStreamPositions
{
/**
* Returns for each projection the index from which on it needs events.
* That is the index of its last split event, or 0 if it has none.
*
* @param array<string, Projection> $projections
*
* @return array<string, positive-int|0>
*/
public static function resolve(AppendStore $store, array $projections): array
{
$positions = array_map(static fn (): int => 0, $projections);
$lookups = [];

foreach ($projections as $name => $projection) {
if (!$projection instanceof SplitStreamProvider || !$projection instanceof SubQueryProvider) {
continue;
}

$splitEvents = $projection->splitEvents();

if ($splitEvents === []) {
continue;
}

$subQuery = $projection->subQuery();

$lookups[$name] = new SubQuery(
$subQuery->tags,
$splitEvents,
$subQuery->streamName,
true,
);
}

if ($lookups === []) {
return $positions;
}

foreach ($store->query(new Query(...array_values($lookups))) as $index => $message) {
foreach ($lookups as $name => $lookup) {
if (!$lookup->match($message)) {
continue;
}

$positions[$name] = max($positions[$name], $index);
}
}

return $positions;
}
}
17 changes: 17 additions & 0 deletions src/Projection/SplitStreamProvider.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
<?php

declare(strict_types=1);

namespace Patchlevel\EventSourcing\Projection;

/** @experimental */
interface SplitStreamProvider
{
/**
* Events which contain the whole state the projection needs.
* Everything before the last of these events is not loaded anymore.
*
* @return list<class-string>
*/
public function splitEvents(): array;
}
7 changes: 5 additions & 2 deletions src/Projection/StoreProjectionBuilder.php
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@

use Patchlevel\EventSourcing\Store\AppendStore;

use function min;

/** @experimental */
final class StoreProjectionBuilder implements ProjectionBuilder
{
Expand All @@ -22,10 +24,11 @@
public function build(
array $projections,
): array {
$projection = new CompositeProjection($projections);
$from = SplitStreamPositions::resolve($this->store, $projections);
$projection = new CompositeProjection($projections, $from);

$query = $projection->query();
$stream = $this->store->query($query);
$stream = $this->store->query($query, $from === [] ? 0 : min($from));

Check warning on line 31 in src/Projection/StoreProjectionBuilder.php

View workflow job for this annotation

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

Escaped Mutant for Mutator "IncrementInteger": @@ @@ $projection = new CompositeProjection($projections, $from); $query = $projection->query(); - $stream = $this->store->query($query, $from === [] ? 0 : min($from)); + $stream = $this->store->query($query, $from === [] ? 1 : min($from)); $state = $projection->initialState();

$state = $projection->initialState();

Expand Down
Loading
Loading