From 209ff593d3481accf73c1e2505beca159cb3a764 Mon Sep 17 00:00:00 2001 From: "renovate[bot]" <29139614+renovate[bot]@users.noreply.github.com> Date: Tue, 30 Sep 2025 19:49:47 +0000 Subject: [PATCH 01/11] Update dependency mkdocs-material to v9.6.21 | datasource | package | from | to | | ---------- | --------------- | ------ | ------ | | pypi | mkdocs-material | 9.6.20 | 9.6.21 | Signed-off-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com> --- docs/requirements.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/requirements.txt b/docs/requirements.txt index 1044a91a1..b381ee4d3 100644 --- a/docs/requirements.txt +++ b/docs/requirements.txt @@ -1,7 +1,7 @@ mkdocs==1.6.1 mike==2.1.3 markdown==3.9 -mkdocs-material==9.6.20 +mkdocs-material==9.6.21 # Markdown extensions Pygments==2.19.2 From 171d70ae15a280297b7b451d102d998ec5223625 Mon Sep 17 00:00:00 2001 From: Daniel Badura Date: Tue, 7 Oct 2025 12:23:39 +0200 Subject: [PATCH 02/11] The `RecordedOnArgumentResolver` can now also handle message with `RecordedOnHeader` --- .../RecordedOnArgumentResolver.php | 5 +++++ .../RecordedOnArgumentResolverTest.php | 19 ++++++++++++++++++- 2 files changed, 23 insertions(+), 1 deletion(-) diff --git a/src/Subscription/Subscriber/ArgumentResolver/RecordedOnArgumentResolver.php b/src/Subscription/Subscriber/ArgumentResolver/RecordedOnArgumentResolver.php index 2a2e7dc45..2320a7ab6 100644 --- a/src/Subscription/Subscriber/ArgumentResolver/RecordedOnArgumentResolver.php +++ b/src/Subscription/Subscriber/ArgumentResolver/RecordedOnArgumentResolver.php @@ -8,11 +8,16 @@ use Patchlevel\EventSourcing\Aggregate\AggregateHeader; use Patchlevel\EventSourcing\Message\Message; use Patchlevel\EventSourcing\Metadata\Subscriber\ArgumentMetadata; +use Patchlevel\EventSourcing\Store\Header\RecordedOnHeader; final class RecordedOnArgumentResolver implements ArgumentResolver { public function resolve(ArgumentMetadata $argument, Message $message): DateTimeImmutable { + if ($message->hasHeader(RecordedOnHeader::class)) { + return $message->header(RecordedOnHeader::class)->recordedOn; + } + return $message->header(AggregateHeader::class)->recordedOn; } diff --git a/tests/Unit/Subscription/Subscriber/ArgumentResolver/RecordedOnArgumentResolverTest.php b/tests/Unit/Subscription/Subscriber/ArgumentResolver/RecordedOnArgumentResolverTest.php index 7e284a0a4..3ef1a8e8c 100644 --- a/tests/Unit/Subscription/Subscriber/ArgumentResolver/RecordedOnArgumentResolverTest.php +++ b/tests/Unit/Subscription/Subscriber/ArgumentResolver/RecordedOnArgumentResolverTest.php @@ -8,6 +8,7 @@ use Patchlevel\EventSourcing\Aggregate\AggregateHeader; use Patchlevel\EventSourcing\Message\Message; use Patchlevel\EventSourcing\Metadata\Subscriber\ArgumentMetadata; +use Patchlevel\EventSourcing\Store\Header\RecordedOnHeader; use Patchlevel\EventSourcing\Subscription\Subscriber\ArgumentResolver\RecordedOnArgumentResolver; use PHPUnit\Framework\Attributes\CoversClass; use PHPUnit\Framework\TestCase; @@ -35,7 +36,7 @@ public function testSupport(): void ); } - public function testResolve(): void + public function testResolveFromAggregateHeader(): void { $date = new DateTimeImmutable(); @@ -57,4 +58,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, + ), + ); + } } From 2fb266fe9f86a76f523426e900bbcc2f71fbddf9 Mon Sep 17 00:00:00 2001 From: "renovate[bot]" <29139614+renovate[bot]@users.noreply.github.com> Date: Wed, 15 Oct 2025 10:28:34 +0000 Subject: [PATCH 03/11] Update dependency mkdocs-material to v9.6.22 | datasource | package | from | to | | ---------- | --------------- | ------ | ------ | | pypi | mkdocs-material | 9.6.21 | 9.6.22 | Signed-off-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com> --- docs/requirements.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/requirements.txt b/docs/requirements.txt index b381ee4d3..f005d7f96 100644 --- a/docs/requirements.txt +++ b/docs/requirements.txt @@ -1,7 +1,7 @@ mkdocs==1.6.1 mike==2.1.3 markdown==3.9 -mkdocs-material==9.6.21 +mkdocs-material==9.6.22 # Markdown extensions Pygments==2.19.2 From 729744a5ac66b47b015ce55701aa9db4dbf1ed23 Mon Sep 17 00:00:00 2001 From: "renovate[bot]" <29139614+renovate[bot]@users.noreply.github.com> Date: Sat, 1 Nov 2025 16:40:08 +0000 Subject: [PATCH 04/11] Update dependency mkdocs-material to v9.6.23 | datasource | package | from | to | | ---------- | --------------- | ------ | ------ | | pypi | mkdocs-material | 9.6.22 | 9.6.23 | Signed-off-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com> --- docs/requirements.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/requirements.txt b/docs/requirements.txt index f005d7f96..2224640d3 100644 --- a/docs/requirements.txt +++ b/docs/requirements.txt @@ -1,7 +1,7 @@ mkdocs==1.6.1 mike==2.1.3 markdown==3.9 -mkdocs-material==9.6.22 +mkdocs-material==9.6.23 # Markdown extensions Pygments==2.19.2 From 31811ec08388a935cb82b5690b1d6b436569f410 Mon Sep 17 00:00:00 2001 From: "renovate[bot]" <29139614+renovate[bot]@users.noreply.github.com> Date: Mon, 3 Nov 2025 22:10:30 +0000 Subject: [PATCH 05/11] Update dependency markdown to v3.10 | datasource | package | from | to | | ---------- | -------- | ---- | ---- | | pypi | markdown | 3.9 | 3.10 | Signed-off-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com> --- docs/requirements.txt | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/requirements.txt b/docs/requirements.txt index 2224640d3..72b521da7 100644 --- a/docs/requirements.txt +++ b/docs/requirements.txt @@ -1,6 +1,6 @@ mkdocs==1.6.1 mike==2.1.3 -markdown==3.9 +markdown==3.10 mkdocs-material==9.6.23 # Markdown extensions From b638d9ddf987e7910f85b5244321ccf986ea8169 Mon Sep 17 00:00:00 2001 From: Daniel Badura Date: Mon, 10 Nov 2025 22:38:05 +0100 Subject: [PATCH 06/11] Fix typo in docs: RoomBocked -> RoomBooked --- docs/pages/aggregate.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/pages/aggregate.md b/docs/pages/aggregate.md index a04efcde4..8439d427c 100644 --- a/docs/pages/aggregate.md +++ b/docs/pages/aggregate.md @@ -565,7 +565,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; @@ -575,7 +575,7 @@ final class Hotel extends BasicAggregateRoot } #[Apply] - protected function applyRoomBocked(RoomBocked $event): void + protected function applyRoomBooked(RoomBooked $event): void { $this->people++; } From 613f426cfcb8e5fbf9035c9ac8e30dbdcaddee17 Mon Sep 17 00:00:00 2001 From: "renovate[bot]" <29139614+renovate[bot]@users.noreply.github.com> Date: Wed, 12 Nov 2025 01:44:12 +0000 Subject: [PATCH 07/11] Update Docs dependencies | datasource | package | from | to | | ---------- | ------------------ | ------- | ------- | | pypi | mkdocs-material | 9.6.23 | 9.7.0 | | pypi | pymdown-extensions | 10.16.1 | 10.17.1 | Signed-off-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com> --- docs/requirements.txt | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/docs/requirements.txt b/docs/requirements.txt index 72b521da7..5d276069b 100644 --- a/docs/requirements.txt +++ b/docs/requirements.txt @@ -1,11 +1,11 @@ mkdocs==1.6.1 mike==2.1.3 markdown==3.10 -mkdocs-material==9.6.23 +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 From 62f3e11536acb71ad2c9b7a7532cd77dfab6724f Mon Sep 17 00:00:00 2001 From: "renovate[bot]" <29139614+renovate[bot]@users.noreply.github.com> Date: Fri, 14 Nov 2025 05:40:10 +0000 Subject: [PATCH 08/11] Update postgres Docker tag to v17.7 | datasource | package | from | to | | ---------- | -------- | ---- | ---- | | docker | postgres | 17.6 | 17.7 | Signed-off-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com> --- .github/workflows/benchmark.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/benchmark.yml b/.github/workflows/benchmark.yml index eced960d4..e4b5a28f8 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 From 0cfe743e1de514f8256e664757bd45a5504d7734 Mon Sep 17 00:00:00 2001 From: "renovate[bot]" <29139614+renovate[bot]@users.noreply.github.com> Date: Fri, 14 Nov 2025 22:14:44 +0000 Subject: [PATCH 09/11] Lock file maintenance (#778) Signed-off-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com> Co-authored-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com> --- composer.lock | 3 ++- tools/composer.lock | 2 +- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/composer.lock b/composer.lock index dca6406eb..faed8f802 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" }, { @@ -9654,5 +9655,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/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" } From 297c88c8afbf4726100f3e3e68297ddc718224dc Mon Sep 17 00:00:00 2001 From: "renovate[bot]" <29139614+renovate[bot]@users.noreply.github.com> Date: Thu, 20 Nov 2025 18:06:28 +0000 Subject: [PATCH 10/11] Update actions/checkout action to v6 | datasource | package | from | to | | ----------- | ---------------- | ---- | -- | | github-tags | actions/checkout | v5 | v6 | Signed-off-by: renovate[bot] <29139614+renovate[bot]@users.noreply.github.com> --- .github/workflows/backward-compatibility-check.yml | 2 +- .github/workflows/benchmark.yml | 4 ++-- .github/workflows/coding-standard.yml | 2 +- .github/workflows/deptrac.yml | 2 +- .github/workflows/docs-build-try.yml | 2 +- .github/workflows/docs-build.yml | 2 +- .github/workflows/docs-check.yml | 2 +- .github/workflows/integration.yml | 8 ++++---- .github/workflows/mutation-tests-diff.yml | 2 +- .github/workflows/mutation-tests.yml | 2 +- .github/workflows/phpstan.yml | 2 +- .github/workflows/psalm.yml | 2 +- ...lease-on-milestone-closed-triggering-release-event.yml | 2 +- .github/workflows/unit.yml | 2 +- 14 files changed, 18 insertions(+), 18 deletions(-) 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 e4b5a28f8..2fdb63a55 100644 --- a/.github/workflows/benchmark.yml +++ b/.github/workflows/benchmark.yml @@ -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" From efb2da877bfbba758e4028eac805eb163962f544 Mon Sep 17 00:00:00 2001 From: Daniel Badura Date: Tue, 18 Nov 2025 15:49:51 +0100 Subject: [PATCH 11/11] Add support for EventsCritereon and ToIndexCriteroen in InMemoryStore. Also add IndexHeader, EventIdHeader and RecordedOnHeader to the messages in the InMemoryStore and fixed FromIndexCriteroen and transactional behaviour. --- baseline.xml | 5 - deptrac-baseline.yaml | 9 + deptrac.yaml | 2 - src/Store/Header/IndexHeader.php | 5 +- src/Store/InMemoryStore.php | 86 ++++- src/Store/MissingEventRegistry.php | 19 ++ tests/Unit/Store/InMemoryStoreTest.php | 456 ++++++++++++++++++++++--- 7 files changed, 509 insertions(+), 73 deletions(-) create mode 100644 src/Store/MissingEventRegistry.php diff --git a/baseline.xml b/baseline.xml index 2e51abc62..3e15b6e95 100644 --- a/baseline.xml +++ b/baseline.xml @@ -159,11 +159,6 @@ convertToPHPValue($data['recorded_on'], $platform)]]> - - - messages]]> - - diff --git a/deptrac-baseline.yaml b/deptrac-baseline.yaml index fe42e99fd..83ffb4da8 100644 --- a/deptrac-baseline.yaml +++ b/deptrac-baseline.yaml @@ -14,3 +14,12 @@ deptrac: - Patchlevel\EventSourcing\Aggregate\AggregateRoot Patchlevel\EventSourcing\Attribute\Subscriber: - Patchlevel\EventSourcing\Subscription\RunMode + Patchlevel\EventSourcing\Store\DoctrineDbalStore: + - Patchlevel\EventSourcing\Aggregate\AggregateHeader + Patchlevel\EventSourcing\Store\DoctrineDbalStoreStream: + - Patchlevel\EventSourcing\Aggregate\AggregateHeader + Patchlevel\EventSourcing\Store\InMemoryStore: + - Patchlevel\EventSourcing\Aggregate\AggregateHeader + - Patchlevel\EventSourcing\Metadata\Event\EventRegistry + Patchlevel\EventSourcing\Store\MissingEventRegistry: + - Patchlevel\EventSourcing\Metadata\Event\EventRegistry diff --git a/deptrac.yaml b/deptrac.yaml index ded380b15..9d7c64eb8 100644 --- a/deptrac.yaml +++ b/deptrac.yaml @@ -186,8 +186,6 @@ deptrac: - Cryptography - MetadataAggregate Store: - - Aggregate - - Attribute - Clock - Message - Metadata 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 e2833430c..55b88ccb1 100644 --- a/src/Store/InMemoryStore.php +++ b/src/Store/InMemoryStore.php @@ -6,21 +6,30 @@ use Closure; use Patchlevel\EventSourcing\Aggregate\AggregateHeader; +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\AggregateIdCriterion; use Patchlevel\EventSourcing\Store\Criteria\AggregateNameCriterion; 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; @@ -35,10 +44,16 @@ final class InMemoryStore implements StreamStore { - /** @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( @@ -71,7 +86,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; + } + }); } /** @@ -81,7 +116,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 */ @@ -134,9 +176,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 AggregateIdCriterion::class: @@ -222,7 +266,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); @@ -57,10 +73,19 @@ public function testLoadMessages(): void public function testLoadByAggregateId(): void { $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) - ->withHeader(new AggregateHeader('profile', '1', 1, new DateTimeImmutable())); + ->withHeader(new AggregateHeader('profile', '1', 1, new DateTimeImmutable())) + ->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 AggregateHeader('profile', '2', 1, new DateTimeImmutable())); - $message3 = new Message(new ProfileVisited(ProfileId::fromString('3'))); + ->withHeader(new AggregateHeader('profile', '2', 1, new DateTimeImmutable())) + ->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]); @@ -74,10 +99,19 @@ public function testLoadByAggregateId(): void public function testLoadByAggregateName(): void { $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) - ->withHeader(new AggregateHeader('foo', '1', 1, new DateTimeImmutable())); + ->withHeader(new AggregateHeader('foo', '1', 1, new DateTimeImmutable())) + ->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 AggregateHeader('bar', '2', 1, new DateTimeImmutable())); - $message3 = new Message(new ProfileVisited(ProfileId::fromString('3'))); + ->withHeader(new AggregateHeader('bar', '2', 1, new DateTimeImmutable())) + ->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]); @@ -91,10 +125,19 @@ public function testLoadByAggregateName(): 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]); @@ -108,11 +151,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]); @@ -126,13 +178,25 @@ public function testLoadByStreamNameWithLike(): void public function testLoadFromPlayhead(): void { $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) - ->withHeader(new AggregateHeader('foo', '1', 1, new DateTimeImmutable())); + ->withHeader(new AggregateHeader('foo', '1', 1, new DateTimeImmutable())) + ->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 AggregateHeader('foo', '1', 2, new DateTimeImmutable())); + ->withHeader(new AggregateHeader('foo', '1', 2, new DateTimeImmutable())) + ->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]); @@ -146,13 +210,25 @@ public function testLoadFromPlayhead(): void public function testLoadFromIndex(): void { $message1 = (new Message(new ProfileVisited(ProfileId::fromString('1')))) - ->withHeader(new AggregateHeader('foo', '1', 1, new DateTimeImmutable())); + ->withHeader(new AggregateHeader('foo', '1', 1, new DateTimeImmutable())) + ->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 AggregateHeader('foo', '1', 2, new DateTimeImmutable())); + ->withHeader(new AggregateHeader('foo', '1', 2, new DateTimeImmutable())) + ->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]); @@ -160,17 +236,87 @@ 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 AggregateHeader('foo', '1', 1, new DateTimeImmutable())) + ->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 AggregateHeader('foo', '1', 2, new DateTimeImmutable())) + ->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(AggregateHeader::class)->playhead, + $messages[0]->header(AggregateHeader::class)->playhead, + ); + self::assertSame( + 1, + $messages[0]->header(IndexHeader::class)->index, + ); + self::assertSame( + $message2->header(AggregateHeader::class)->playhead, + $messages[1]->header(AggregateHeader::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]); @@ -184,8 +330,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]); @@ -196,11 +348,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); @@ -210,8 +428,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]); @@ -224,8 +448,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]); @@ -238,8 +468,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]); @@ -253,8 +489,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]); @@ -264,8 +506,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([]); @@ -282,11 +530,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); @@ -299,15 +556,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]); @@ -317,12 +603,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]); @@ -349,15 +647,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]);