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