diff --git a/.github/workflows/backward-compatibility-check.yml b/.github/workflows/backward-compatibility-check.yml index c9b094cbe..36b41a550 100644 --- a/.github/workflows/backward-compatibility-check.yml +++ b/.github/workflows/backward-compatibility-check.yml @@ -22,7 +22,7 @@ jobs: steps: - name: "Checkout" - uses: actions/checkout@v5 + uses: actions/checkout@v6 with: fetch-depth: 0 diff --git a/.github/workflows/benchmark.yml b/.github/workflows/benchmark.yml index eced960d4..2fdb63a55 100644 --- a/.github/workflows/benchmark.yml +++ b/.github/workflows/benchmark.yml @@ -14,7 +14,7 @@ jobs: services: postgres: # Docker Hub image - image: "postgres:17.6" + image: "postgres:17.7" # Provide the password for postgres env: POSTGRES_PASSWORD: postgres @@ -46,7 +46,7 @@ jobs: extensions: pdo_sqlite - name: "Checkout base" - uses: actions/checkout@v5 + uses: actions/checkout@v6 with: ref: ${{ github.base_ref }} @@ -58,7 +58,7 @@ jobs: run: "vendor/bin/phpbench run tests/Benchmark --progress=none --report=default --tag=base" - name: "Checkout" - uses: actions/checkout@v5 + uses: actions/checkout@v6 with: clean: false diff --git a/.github/workflows/coding-standard.yml b/.github/workflows/coding-standard.yml index 4d4ac1727..c56fa1c9c 100644 --- a/.github/workflows/coding-standard.yml +++ b/.github/workflows/coding-standard.yml @@ -26,7 +26,7 @@ jobs: steps: - name: "Checkout" - uses: actions/checkout@v5 + uses: actions/checkout@v6 - name: "Install PHP" uses: "shivammathur/setup-php@2.35.5" diff --git a/.github/workflows/deptrac.yml b/.github/workflows/deptrac.yml index 5f0a777ae..5053e45b8 100644 --- a/.github/workflows/deptrac.yml +++ b/.github/workflows/deptrac.yml @@ -26,7 +26,7 @@ jobs: steps: - name: "Checkout" - uses: actions/checkout@v5 + uses: actions/checkout@v6 - name: "Install PHP" uses: "shivammathur/setup-php@2.35.5" diff --git a/.github/workflows/docs-build-try.yml b/.github/workflows/docs-build-try.yml index aa1b5acb3..bd943e1c8 100644 --- a/.github/workflows/docs-build-try.yml +++ b/.github/workflows/docs-build-try.yml @@ -11,7 +11,7 @@ jobs: name: Deploy docs runs-on: ubuntu-latest steps: - - uses: actions/checkout@v5 + - uses: actions/checkout@v6 with: fetch-depth: 0 diff --git a/.github/workflows/docs-build.yml b/.github/workflows/docs-build.yml index 608f62436..e84f21482 100644 --- a/.github/workflows/docs-build.yml +++ b/.github/workflows/docs-build.yml @@ -13,7 +13,7 @@ jobs: name: Deploy docs runs-on: ubuntu-latest steps: - - uses: actions/checkout@v5 + - uses: actions/checkout@v6 with: fetch-depth: 0 diff --git a/.github/workflows/docs-check.yml b/.github/workflows/docs-check.yml index c583b9bd3..34a06e230 100644 --- a/.github/workflows/docs-check.yml +++ b/.github/workflows/docs-check.yml @@ -25,7 +25,7 @@ jobs: steps: - name: "Checkout" - uses: actions/checkout@v5 + uses: actions/checkout@v6 - name: "Install PHP" uses: "shivammathur/setup-php@2.35.5" diff --git a/.github/workflows/integration.yml b/.github/workflows/integration.yml index b24bd1eeb..70e3794a5 100644 --- a/.github/workflows/integration.yml +++ b/.github/workflows/integration.yml @@ -49,7 +49,7 @@ jobs: steps: - name: "Checkout" - uses: actions/checkout@v5 + uses: actions/checkout@v6 - name: "Install PHP" uses: "shivammathur/setup-php@2.35.5" @@ -104,7 +104,7 @@ jobs: steps: - name: "Checkout" - uses: actions/checkout@v5 + uses: actions/checkout@v6 - name: "Install PHP" uses: "shivammathur/setup-php@2.35.5" @@ -159,7 +159,7 @@ jobs: steps: - name: "Checkout" - uses: actions/checkout@v5 + uses: actions/checkout@v6 - name: "Install PHP" uses: "shivammathur/setup-php@2.35.5" @@ -193,7 +193,7 @@ jobs: steps: - name: "Checkout" - uses: actions/checkout@v5 + uses: actions/checkout@v6 - name: "Install PHP" uses: "shivammathur/setup-php@2.35.5" diff --git a/.github/workflows/mutation-tests-diff.yml b/.github/workflows/mutation-tests-diff.yml index 6a841fb61..2f212e305 100644 --- a/.github/workflows/mutation-tests-diff.yml +++ b/.github/workflows/mutation-tests-diff.yml @@ -22,7 +22,7 @@ jobs: steps: - name: "Checkout" - uses: actions/checkout@v5 + uses: actions/checkout@v6 with: fetch-depth: 0 diff --git a/.github/workflows/mutation-tests.yml b/.github/workflows/mutation-tests.yml index 4d85bd37d..927ce758b 100644 --- a/.github/workflows/mutation-tests.yml +++ b/.github/workflows/mutation-tests.yml @@ -26,7 +26,7 @@ jobs: steps: - name: "Checkout" - uses: actions/checkout@v5 + uses: actions/checkout@v6 - name: "Install PHP" uses: "shivammathur/setup-php@2.35.5" diff --git a/.github/workflows/phpstan.yml b/.github/workflows/phpstan.yml index 459815cd5..abb9baa21 100644 --- a/.github/workflows/phpstan.yml +++ b/.github/workflows/phpstan.yml @@ -26,7 +26,7 @@ jobs: steps: - name: "Checkout" - uses: actions/checkout@v5 + uses: actions/checkout@v6 - name: "Install PHP" uses: "shivammathur/setup-php@2.35.5" diff --git a/.github/workflows/psalm.yml b/.github/workflows/psalm.yml index 0127d9473..980447146 100644 --- a/.github/workflows/psalm.yml +++ b/.github/workflows/psalm.yml @@ -26,7 +26,7 @@ jobs: steps: - name: "Checkout" - uses: actions/checkout@v5 + uses: actions/checkout@v6 - name: "Install PHP" uses: "shivammathur/setup-php@2.35.5" diff --git a/.github/workflows/release-on-milestone-closed-triggering-release-event.yml b/.github/workflows/release-on-milestone-closed-triggering-release-event.yml index 67fae3527..8af28e8f5 100644 --- a/.github/workflows/release-on-milestone-closed-triggering-release-event.yml +++ b/.github/workflows/release-on-milestone-closed-triggering-release-event.yml @@ -18,7 +18,7 @@ jobs: steps: - name: "Checkout" - uses: actions/checkout@v5 + uses: actions/checkout@v6 - name: "Release" uses: "laminas/automatic-releases@v1" diff --git a/.github/workflows/unit.yml b/.github/workflows/unit.yml index f1931de20..cd6cfeb60 100644 --- a/.github/workflows/unit.yml +++ b/.github/workflows/unit.yml @@ -37,7 +37,7 @@ jobs: steps: - name: "Checkout" - uses: actions/checkout@v5 + uses: actions/checkout@v6 - name: "Install PHP" uses: "shivammathur/setup-php@2.35.5" diff --git a/baseline.xml b/baseline.xml index e93e5b25c..aa5103eb2 100644 --- a/baseline.xml +++ b/baseline.xml @@ -161,11 +161,6 @@ - - - messages]]> - - diff --git a/composer.lock b/composer.lock index 56ce046b3..d7c6d1c23 100644 --- a/composer.lock +++ b/composer.lock @@ -3622,6 +3622,7 @@ "issues": "https://github.com/doctrine/annotations/issues", "source": "https://github.com/doctrine/annotations/tree/2.0.2" }, + "abandoned": true, "time": "2024-09-05T10:17:24+00:00" }, { @@ -9713,5 +9714,5 @@ "platform-dev": { "ext-pdo_sqlite": "~8.2.0 || ~8.3.0 || ~8.4.0" }, - "plugin-api-version": "2.6.0" + "plugin-api-version": "2.9.0" } diff --git a/deptrac.yaml b/deptrac.yaml index 46ded22fb..651e6de52 100644 --- a/deptrac.yaml +++ b/deptrac.yaml @@ -195,8 +195,6 @@ deptrac: - Identifier - MetadataAggregate Store: - - Aggregate - - Attribute - Clock - Message - Metadata diff --git a/docs/pages/aggregate.md b/docs/pages/aggregate.md index 43132c272..6b00e1e8a 100644 --- a/docs/pages/aggregate.md +++ b/docs/pages/aggregate.md @@ -561,7 +561,7 @@ final class Hotel extends BasicAggregateRoot throw new NoPlaceException($name); } - $this->recordThat(new RoomBocked($name)); + $this->recordThat(new RoomBooked($name)); if ($this->people !== self::SIZE) { return; @@ -571,7 +571,7 @@ final class Hotel extends BasicAggregateRoot } #[Apply] - protected function applyRoomBocked(RoomBocked $event): void + protected function applyRoomBooked(RoomBooked $event): void { $this->people++; } diff --git a/docs/requirements.txt b/docs/requirements.txt index 1044a91a1..5d276069b 100644 --- a/docs/requirements.txt +++ b/docs/requirements.txt @@ -1,11 +1,11 @@ mkdocs==1.6.1 mike==2.1.3 -markdown==3.9 -mkdocs-material==9.6.20 +markdown==3.10 +mkdocs-material==9.7.0 # Markdown extensions Pygments==2.19.2 -pymdown-extensions==10.16.1 +pymdown-extensions==10.17.1 # MkDocs plugins mkdocs-material-extensions==1.3.1 diff --git a/src/Store/Header/IndexHeader.php b/src/Store/Header/IndexHeader.php index e5baa9816..8e659b6ec 100644 --- a/src/Store/Header/IndexHeader.php +++ b/src/Store/Header/IndexHeader.php @@ -4,10 +4,7 @@ namespace Patchlevel\EventSourcing\Store\Header; -/** - * @psalm-immutable - * @experimental - */ +/** @psalm-immutable */ final class IndexHeader { /** @param positive-int $index */ diff --git a/src/Store/InMemoryStore.php b/src/Store/InMemoryStore.php index 5113a773a..ea12e01c1 100644 --- a/src/Store/InMemoryStore.php +++ b/src/Store/InMemoryStore.php @@ -5,19 +5,28 @@ namespace Patchlevel\EventSourcing\Store; use Closure; +use Patchlevel\EventSourcing\Clock\SystemClock; use Patchlevel\EventSourcing\Message\HeaderNotFound; use Patchlevel\EventSourcing\Message\Message; +use Patchlevel\EventSourcing\Metadata\Event\EventRegistry; use Patchlevel\EventSourcing\Store\Criteria\ArchivedCriterion; use Patchlevel\EventSourcing\Store\Criteria\Criteria; +use Patchlevel\EventSourcing\Store\Criteria\EventsCriterion; use Patchlevel\EventSourcing\Store\Criteria\FromIndexCriterion; use Patchlevel\EventSourcing\Store\Criteria\FromPlayheadCriterion; use Patchlevel\EventSourcing\Store\Criteria\StreamCriterion; +use Patchlevel\EventSourcing\Store\Criteria\ToIndexCriterion; +use Patchlevel\EventSourcing\Store\Header\EventIdHeader; +use Patchlevel\EventSourcing\Store\Header\IndexHeader; use Patchlevel\EventSourcing\Store\Header\PlayheadHeader; +use Patchlevel\EventSourcing\Store\Header\RecordedOnHeader; use Patchlevel\EventSourcing\Store\Header\StreamNameHeader; +use Psr\Clock\ClockInterface; +use Ramsey\Uuid\Uuid; +use Throwable; use function array_filter; use function array_map; -use function array_push; use function array_reverse; use function array_slice; use function array_unique; @@ -32,10 +41,16 @@ final class InMemoryStore implements Store { - /** @param array $messages */ + /** @var array<0|positive-int, Message> */ + private array $messages = []; + + /** @param list $messages */ public function __construct( - private array $messages = [], + array $messages = [], + private readonly EventRegistry|null $eventRegistry = null, + private readonly ClockInterface $clock = new SystemClock(), ) { + $this->save(...$messages); } public function load( @@ -68,7 +83,27 @@ public function count(Criteria|null $criteria = null): int public function save(Message ...$messages): void { - array_push($this->messages, ...$messages); + $this->transactional(function () use ($messages): void { + $count = count($this->messages); + + foreach ($messages as $message) { + $count++; + + if (!$message->hasHeader(IndexHeader::class)) { + $message = $message->withHeader(new IndexHeader($count)); + } + + if (!$message->hasHeader(EventIdHeader::class)) { + $message = $message->withHeader(new EventIdHeader(Uuid::uuid7()->toString())); + } + + if (!$message->hasHeader(RecordedOnHeader::class)) { + $message = $message->withHeader(new RecordedOnHeader($this->clock->now())); + } + + $this->messages[] = $message; + } + }); } /** @@ -78,7 +113,14 @@ public function save(Message ...$messages): void */ public function transactional(Closure $function): void { - $function(); + $messages = $this->messages; + try { + $function(); + } catch (Throwable $e) { + $this->messages = $messages; + + throw $e; + } } /** @return list */ @@ -127,9 +169,11 @@ private function filter(Criteria|null $criteria): array return $this->messages; } + $eventRegistry = $this->eventRegistry; + return array_filter( $this->messages, - static function (Message $message, int $index) use ($criteria): bool { + static function (Message $message) use ($criteria, $eventRegistry): bool { foreach ($criteria->all() as $criterion) { switch ($criterion::class) { case StreamCriterion::class: @@ -187,7 +231,35 @@ static function (Message $message, int $index) use ($criteria): bool { break; case FromIndexCriterion::class: - if ($index < $criterion->fromIndex) { + try { + $index = $message->header(IndexHeader::class)->index; + } catch (HeaderNotFound) { + return false; + } + + if ($index <= $criterion->fromIndex) { + return false; + } + + break; + case ToIndexCriterion::class: + try { + $index = $message->header(IndexHeader::class)->index; + } catch (HeaderNotFound) { + return false; + } + + if ($index >= $criterion->toIndex) { + return false; + } + + break; + case EventsCriterion::class: + if ($eventRegistry === null) { + throw new MissingEventRegistry($criterion::class); + } + + if (!in_array($eventRegistry->eventName($message->event()::class), $criterion->events)) { return false; } diff --git a/src/Store/MissingEventRegistry.php b/src/Store/MissingEventRegistry.php new file mode 100644 index 000000000..0987425f0 --- /dev/null +++ b/src/Store/MissingEventRegistry.php @@ -0,0 +1,19 @@ +withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)), + (new Message(new ProfileVisited(ProfileId::fromString('2')))) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)), ]; $store = new InMemoryStore($expected); @@ -53,10 +70,19 @@ public function testLoadMessages(): void public function testLoadByStreamName(): void { $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) - ->withHeader(new StreamNameHeader('foo')); + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) - ->withHeader(new StreamNameHeader('bar')); - $message3 = new Message(new ProfileVisited(ProfileId::fromString('3'))); + ->withHeader(new StreamNameHeader('bar')) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); + $message3 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) + ->withHeader(new EventIdHeader('019aa604-94b8-7182-b1dc-f5d4aa9652ca')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(3)); $store = new InMemoryStore([$message1, $message2, $message3]); @@ -70,11 +96,20 @@ public function testLoadByStreamName(): void public function testLoadByStreamNameWithLike(): void { $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) - ->withHeader(new StreamNameHeader('foo-3')); + ->withHeader(new StreamNameHeader('foo-3')) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) - ->withHeader(new StreamNameHeader('bar-1')); + ->withHeader(new StreamNameHeader('bar-1')) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); $message3 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) - ->withHeader(new StreamNameHeader('bar-2')); + ->withHeader(new StreamNameHeader('bar-2')) + ->withHeader(new EventIdHeader('019aa604-94b8-7182-b1dc-f5d4aa9652ca')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(3)); $store = new InMemoryStore([$message1, $message2, $message3]); @@ -88,13 +123,25 @@ public function testLoadByStreamNameWithLike(): void public function testLoadFromPlayhead(): void { $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) - ->withHeader(new PlayheadHeader(1)); + ->withHeader(new PlayheadHeader(1)) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) - ->withHeader(new PlayheadHeader(2)); + ->withHeader(new PlayheadHeader(2)) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); $message3 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) ->withHeader(new StreamNameHeader('foo-1')) - ->withHeader(new PlayheadHeader(3)); - $message4 = new Message(new ProfileVisited(ProfileId::fromString('3'))); + ->withHeader(new PlayheadHeader(3)) + ->withHeader(new EventIdHeader('019aa604-94b8-7182-b1dc-f5d4aa9652ca')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(3)); + $message4 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) + ->withHeader(new EventIdHeader('019aa607-3a33-7f47-bb66-4223f2390a30')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(4)); $store = new InMemoryStore([$message1, $message2, $message3, $message4]); @@ -108,13 +155,25 @@ public function testLoadFromPlayhead(): void public function testLoadFromIndex(): void { $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) - ->withHeader(new PlayheadHeader(1)); + ->withHeader(new PlayheadHeader(1)) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) - ->withHeader(new PlayheadHeader(2)); + ->withHeader(new PlayheadHeader(2)) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); $message3 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) ->withHeader(new StreamNameHeader('foo-1')) - ->withHeader(new PlayheadHeader(3)); - $message4 = new Message(new ProfileVisited(ProfileId::fromString('3'))); + ->withHeader(new PlayheadHeader(3)) + ->withHeader(new EventIdHeader('019aa604-94b8-7182-b1dc-f5d4aa9652ca')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(3)); + $message4 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) + ->withHeader(new EventIdHeader('019aa607-3a33-7f47-bb66-4223f2390a30')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(4)); $store = new InMemoryStore([$message1, $message2, $message3, $message4]); @@ -122,17 +181,89 @@ public function testLoadFromIndex(): void $messages = iterator_to_array($stream); - self::assertSame([$message3, $message4], $messages); + self::assertCount(2, $messages); + self::assertSame( + $message3->header(PlayheadHeader::class)->playhead, + $messages[0]->header(PlayheadHeader::class)->playhead, + ); + self::assertSame( + 3, + $messages[0]->header(IndexHeader::class)->index, + ); + self::assertFalse($message4->hasHeader(PlayheadHeader::class)); + self::assertSame( + 4, + $messages[1]->header(IndexHeader::class)->index, + ); + } + + public function testLoadToIndex(): void + { + $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) + ->withHeader(new StreamNameHeader('foo-1')) + ->withHeader(new PlayheadHeader(1)) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); + $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) + ->withHeader(new StreamNameHeader('foo-1')) + ->withHeader(new PlayheadHeader(2)) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); + $message3 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) + ->withHeader(new StreamNameHeader('foo-1')) + ->withHeader(new PlayheadHeader(3)) + ->withHeader(new EventIdHeader('019aa604-94b8-7182-b1dc-f5d4aa9652ca')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(3)); + $message4 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) + ->withHeader(new EventIdHeader('019aa607-3a33-7f47-bb66-4223f2390a30')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(4)); + + $store = new InMemoryStore([$message1, $message2, $message3, $message4]); + + $stream = $store->load(new Criteria(new ToIndexCriterion(3))); + + $messages = iterator_to_array($stream); + + self::assertCount(2, $messages); + self::assertSame( + $message1->header(PlayheadHeader::class)->playhead, + $messages[0]->header(PlayheadHeader::class)->playhead, + ); + self::assertSame( + 1, + $messages[0]->header(IndexHeader::class)->index, + ); + self::assertSame( + $message2->header(PlayheadHeader::class)->playhead, + $messages[1]->header(PlayheadHeader::class)->playhead, + ); + self::assertSame( + 2, + $messages[1]->header(IndexHeader::class)->index, + ); } public function testLoadByStreamNameWithLikeAll(): void { $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) - ->withHeader(new StreamNameHeader('foo-3')); + ->withHeader(new StreamNameHeader('foo-3')) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) - ->withHeader(new StreamNameHeader('bar-1')); + ->withHeader(new StreamNameHeader('bar-1')) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); $message3 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) - ->withHeader(new StreamNameHeader('bar-2')); + ->withHeader(new StreamNameHeader('bar-2')) + ->withHeader(new EventIdHeader('019aa604-94b8-7182-b1dc-f5d4aa9652ca')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(3)); $store = new InMemoryStore([$message1, $message2, $message3]); @@ -146,8 +277,14 @@ public function testLoadByStreamNameWithLikeAll(): void public function testLoadArchived(): void { $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) - ->withHeader(new ArchivedHeader()); - $message2 = new Message(new ProfileVisited(ProfileId::fromString('2'))); + ->withHeader(new ArchivedHeader()) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); + $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); $store = new InMemoryStore([$message1, $message2]); @@ -158,11 +295,77 @@ public function testLoadArchived(): void self::assertSame([$message1], $messages); } + public function testLoadByEventName(): void + { + $message1 = (new Message(new ProfileCreated(ProfileId::fromString('1'), Email::fromString('s@b.de')))) + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); + $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) + ->withHeader(new StreamNameHeader('bar')) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); + $message3 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) + ->withHeader(new EventIdHeader('019aa604-94b8-7182-b1dc-f5d4aa9652ca')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(3)); + + $store = new InMemoryStore( + [$message1, $message2, $message3], + new EventRegistry([ + 'profile_created' => ProfileCreated::class, + 'profile_visited' => ProfileVisited::class, + ]), + ); + + $stream = $store->load(new Criteria(new EventsCriterion(['profile_created']))); + $messages = iterator_to_array($stream); + + self::assertSame([$message1], $messages); + + $stream = $store->load(new Criteria(new EventsCriterion(['profile_visited']))); + $messages = iterator_to_array($stream); + + self::assertSame([$message2, $message3], $messages); + self::assertSame([], iterator_to_array($store->load(new Criteria(new EventsCriterion(['profile_deleted']))))); + } + + public function testLoadByEventNameWithoutRegistry(): void + { + $message1 = (new Message(new ProfileCreated(ProfileId::fromString('1'), Email::fromString('s@b.de')))) + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); + $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) + ->withHeader(new StreamNameHeader('bar')) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); + $message3 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) + ->withHeader(new EventIdHeader('019aa604-94b8-7182-b1dc-f5d4aa9652ca')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(3)); + + $store = new InMemoryStore([$message1, $message2, $message3]); + + $this->expectException(MissingEventRegistry::class); + $store->load(new Criteria(new EventsCriterion(['profile_created']))); + } + public function testLoadUnsupportedCriterion(): void { $store = new InMemoryStore([ - new Message(new ProfileVisited(ProfileId::fromString('1'))), - new Message(new ProfileVisited(ProfileId::fromString('2'))), + (new Message(new ProfileVisited(ProfileId::fromString('1')))) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)), + (new Message(new ProfileVisited(ProfileId::fromString('2')))) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)), ]); $this->expectException(UnsupportedCriterion::class); @@ -172,8 +375,14 @@ public function testLoadUnsupportedCriterion(): void public function testLoadLimit(): void { - $message1 = new Message(new ProfileVisited(ProfileId::fromString('1'))); - $message2 = new Message(new ProfileVisited(ProfileId::fromString('2'))); + $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); + $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); $store = new InMemoryStore([$message1, $message2]); @@ -186,8 +395,14 @@ public function testLoadLimit(): void public function testLoadOffset(): void { - $message1 = new Message(new ProfileVisited(ProfileId::fromString('1'))); - $message2 = new Message(new ProfileVisited(ProfileId::fromString('2'))); + $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); + $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); $store = new InMemoryStore([$message1, $message2]); @@ -200,8 +415,14 @@ public function testLoadOffset(): void public function testLoadBackwards(): void { - $message1 = new Message(new ProfileVisited(ProfileId::fromString('1'))); - $message2 = new Message(new ProfileVisited(ProfileId::fromString('2'))); + $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); + $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); $store = new InMemoryStore([$message1, $message2]); @@ -215,8 +436,14 @@ public function testLoadBackwards(): void public function testCount(): void { $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) - ->withHeader(new ArchivedHeader()); - $message2 = new Message(new ProfileVisited(ProfileId::fromString('2'))); + ->withHeader(new ArchivedHeader()) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); + $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); $store = new InMemoryStore([$message1, $message2]); @@ -226,8 +453,14 @@ public function testCount(): void public function testSaveEmpty(): void { $expected = [ - new Message(new ProfileVisited(ProfileId::fromString('1'))), - new Message(new ProfileVisited(ProfileId::fromString('2'))), + (new Message(new ProfileVisited(ProfileId::fromString('1')))) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)), + (new Message(new ProfileVisited(ProfileId::fromString('2')))) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)), ]; $store = new InMemoryStore([]); @@ -244,11 +477,20 @@ public function testSaveEmpty(): void public function testSaveWithExistingMessages(): void { $startMessages = [ - new Message(new ProfileVisited(ProfileId::fromString('1'))), - new Message(new ProfileVisited(ProfileId::fromString('2'))), + (new Message(new ProfileVisited(ProfileId::fromString('1')))) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)), + (new Message(new ProfileVisited(ProfileId::fromString('2')))) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)), ]; - $message1 = new Message(new ProfileVisited(ProfileId::fromString('3'))); + $message1 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) + ->withHeader(new EventIdHeader('019aa604-94b8-7182-b1dc-f5d4aa9652ca')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(3)); $store = new InMemoryStore($startMessages); @@ -261,15 +503,44 @@ public function testSaveWithExistingMessages(): void self::assertSame([...$startMessages, $message1], $messages); } + public function testSaveWithoutHeaders(): void + { + $store = new InMemoryStore([new Message(new ProfileVisited(ProfileId::fromString('3')))]); + $store->save(new Message(new ProfileVisited(ProfileId::fromString('1')))); + + $stream = $store->load(); + $messages = iterator_to_array($stream); + + self::assertCount(2, $messages); + + foreach ($messages as $message) { + self::assertTrue($message->hasHeader(EventIdHeader::class)); + self::assertTrue($message->hasHeader(RecordedOnHeader::class)); + self::assertTrue($message->hasHeader(IndexHeader::class)); + } + } + public function testStreams(): void { $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) - ->withHeader(new StreamNameHeader('foo')); + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) - ->withHeader(new StreamNameHeader('bar')); + ->withHeader(new StreamNameHeader('bar')) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); $message3 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) - ->withHeader(new StreamNameHeader('bar')); - $message4 = new Message(new ProfileVisited(ProfileId::fromString('3'))); + ->withHeader(new StreamNameHeader('bar')) + ->withHeader(new EventIdHeader('019aa604-94b8-7182-b1dc-f5d4aa9652ca')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(3)); + $message4 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) + ->withHeader(new EventIdHeader('019aa607-3a33-7f47-bb66-4223f2390a30')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(4)); $store = new InMemoryStore([$message1, $message2, $message3, $message4]); @@ -279,12 +550,24 @@ public function testStreams(): void public function testRemove(): void { $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) - ->withHeader(new StreamNameHeader('foo')); + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) - ->withHeader(new StreamNameHeader('bar')); + ->withHeader(new StreamNameHeader('bar')) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); $message3 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) - ->withHeader(new StreamNameHeader('bar')); - $message4 = new Message(new ProfileVisited(ProfileId::fromString('3'))); + ->withHeader(new StreamNameHeader('bar')) + ->withHeader(new EventIdHeader('019aa604-94b8-7182-b1dc-f5d4aa9652ca')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(3)); + $message4 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) + ->withHeader(new EventIdHeader('019aa607-3a33-7f47-bb66-4223f2390a30')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(4)); $store = new InMemoryStore([$message1, $message2, $message3, $message4]); @@ -311,15 +594,63 @@ static function () use (&$called): void { self::assertTrue($called); } + public function testTransactionalThrowAndResetting(): void + { + $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); + $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) + ->withHeader(new StreamNameHeader('bar')) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); + + $called = false; + $catched = false; + + $store = new InMemoryStore([$message1]); + + try { + $store->transactional( + static function () use (&$called, $message2, $store): void { + $called = true; + $store->save($message2); + + throw new RuntimeException('test'); + }, + ); + } catch (RuntimeException) { + $catched = true; + } + + self::assertTrue($catched); + self::assertTrue($called); + self::assertCount(1, iterator_to_array($store->load())); + } + public function testClear(): void { $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) - ->withHeader(new StreamNameHeader('foo')); + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new EventIdHeader('019aa600-56ef-7ca3-b92a-37c53851e2c2')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(1)); $message2 = (new Message(new ProfileVisited(ProfileId::fromString('2')))) - ->withHeader(new StreamNameHeader('bar')); + ->withHeader(new StreamNameHeader('bar')) + ->withHeader(new EventIdHeader('019aa600-8834-752a-ae2e-d8650e84f403')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(2)); $message3 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) - ->withHeader(new StreamNameHeader('bar')); - $message4 = new Message(new ProfileVisited(ProfileId::fromString('3'))); + ->withHeader(new StreamNameHeader('bar')) + ->withHeader(new EventIdHeader('019aa604-94b8-7182-b1dc-f5d4aa9652ca')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(3)); + $message4 = (new Message(new ProfileVisited(ProfileId::fromString('3')))) + ->withHeader(new EventIdHeader('019aa607-3a33-7f47-bb66-4223f2390a30')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable())) + ->withHeader(new IndexHeader(4)); $store = new InMemoryStore([$message1, $message2, $message3, $message4]); diff --git a/tests/Unit/Subscription/Subscriber/ArgumentResolver/RecordedOnArgumentResolverTest.php b/tests/Unit/Subscription/Subscriber/ArgumentResolver/RecordedOnArgumentResolverTest.php index 259165272..28857a6ed 100644 --- a/tests/Unit/Subscription/Subscriber/ArgumentResolver/RecordedOnArgumentResolverTest.php +++ b/tests/Unit/Subscription/Subscriber/ArgumentResolver/RecordedOnArgumentResolverTest.php @@ -35,7 +35,7 @@ public function testSupport(): void ); } - public function testResolve(): void + public function testResolveFromAggregateHeader(): void { $date = new DateTimeImmutable(); @@ -52,4 +52,20 @@ public function testResolve(): void ), ); } + + public function testResolveFromRecordedOnHeader(): void + { + $date = new DateTimeImmutable(); + + $resolver = new RecordedOnArgumentResolver(); + $message = (new Message(new stdClass()))->withHeader(new RecordedOnHeader($date)); + + self::assertSame( + $date, + $resolver->resolve( + new ArgumentMetadata('foo', DateTimeImmutable::class), + $message, + ), + ); + } } diff --git a/tools/composer.lock b/tools/composer.lock index d22ea2598..69821488d 100644 --- a/tools/composer.lock +++ b/tools/composer.lock @@ -2889,5 +2889,5 @@ "php": "~8.4.0" }, "platform-dev": {}, - "plugin-api-version": "2.6.0" + "plugin-api-version": "2.9.0" }