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 2e51abc62..3e15b6e95 100644
--- a/baseline.xml
+++ b/baseline.xml
@@ -159,11 +159,6 @@
convertToPHPValue($data['recorded_on'], $platform)]]>
-
-
- messages]]>
-
-
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/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/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++;
}
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 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 @@
+hasHeader(RecordedOnHeader::class)) {
+ return $message->header(RecordedOnHeader::class)->recordedOn;
+ }
+
return $message->header(AggregateHeader::class)->recordedOn;
}
diff --git a/tests/Unit/Store/InMemoryStoreTest.php b/tests/Unit/Store/InMemoryStoreTest.php
index 9d66afba2..4c606305f 100644
--- a/tests/Unit/Store/InMemoryStoreTest.php
+++ b/tests/Unit/Store/InMemoryStoreTest.php
@@ -7,22 +7,32 @@
use DateTimeImmutable;
use Patchlevel\EventSourcing\Aggregate\AggregateHeader;
use Patchlevel\EventSourcing\Message\Message;
+use Patchlevel\EventSourcing\Metadata\Event\EventRegistry;
use Patchlevel\EventSourcing\Store\ArchivedHeader;
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 Patchlevel\EventSourcing\Store\InMemoryStore;
+use Patchlevel\EventSourcing\Store\MissingEventRegistry;
use Patchlevel\EventSourcing\Store\UnsupportedCriterion;
+use Patchlevel\EventSourcing\Tests\Unit\Fixture\Email;
+use Patchlevel\EventSourcing\Tests\Unit\Fixture\ProfileCreated;
use Patchlevel\EventSourcing\Tests\Unit\Fixture\ProfileId;
use Patchlevel\EventSourcing\Tests\Unit\Fixture\ProfileVisited;
use PHPUnit\Framework\Attributes\CoversClass;
use PHPUnit\Framework\TestCase;
+use RuntimeException;
use stdClass;
use function iterator_to_array;
@@ -41,8 +51,14 @@ public function testLoadEmpty(): void
public function testLoadMessages(): 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($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]);
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,
+ ),
+ );
+ }
}
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"
}