Skip to content

Commit b79d555

Browse files
committed
Add support for EventsCritereon and ToIndexCriteroen in InMemoryStore. Also add IndexHeader to the messages in the InMemoryStore
1 parent 0cfe743 commit b79d555

3 files changed

Lines changed: 135 additions & 6 deletions

File tree

src/Store/InMemoryStore.php

Lines changed: 47 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -8,19 +8,22 @@
88
use Patchlevel\EventSourcing\Aggregate\AggregateHeader;
99
use Patchlevel\EventSourcing\Message\HeaderNotFound;
1010
use Patchlevel\EventSourcing\Message\Message;
11+
use Patchlevel\EventSourcing\Metadata\Event\EventRegistry;
1112
use Patchlevel\EventSourcing\Store\Criteria\AggregateIdCriterion;
1213
use Patchlevel\EventSourcing\Store\Criteria\AggregateNameCriterion;
1314
use Patchlevel\EventSourcing\Store\Criteria\ArchivedCriterion;
1415
use Patchlevel\EventSourcing\Store\Criteria\Criteria;
16+
use Patchlevel\EventSourcing\Store\Criteria\EventsCriterion;
1517
use Patchlevel\EventSourcing\Store\Criteria\FromIndexCriterion;
1618
use Patchlevel\EventSourcing\Store\Criteria\FromPlayheadCriterion;
1719
use Patchlevel\EventSourcing\Store\Criteria\StreamCriterion;
20+
use Patchlevel\EventSourcing\Store\Criteria\ToIndexCriterion;
21+
use Patchlevel\EventSourcing\Store\Header\IndexHeader;
1822
use Patchlevel\EventSourcing\Store\Header\PlayheadHeader;
1923
use Patchlevel\EventSourcing\Store\Header\StreamNameHeader;
2024

2125
use function array_filter;
2226
use function array_map;
23-
use function array_push;
2427
use function array_reverse;
2528
use function array_slice;
2629
use function array_unique;
@@ -35,10 +38,14 @@
3538

3639
final class InMemoryStore implements StreamStore
3740
{
38-
/** @param array<positive-int|0, Message> $messages */
41+
private array $messages = [];
42+
43+
/** @param list<Message> $messages */
3944
public function __construct(
40-
private array $messages = [],
45+
array $messages = [],
46+
private readonly EventRegistry|null $eventRegistry = null,
4147
) {
48+
$this->save(...$messages);
4249
}
4350

4451
public function load(
@@ -71,7 +78,12 @@ public function count(Criteria|null $criteria = null): int
7178

7279
public function save(Message ...$messages): void
7380
{
74-
array_push($this->messages, ...$messages);
81+
$count = count($this->messages);
82+
83+
foreach ($messages as $message) {
84+
$count++;
85+
$this->messages[] = $message->withHeader(new IndexHeader($count));
86+
}
7587
}
7688

7789
/**
@@ -134,9 +146,11 @@ private function filter(Criteria|null $criteria): array
134146
return $this->messages;
135147
}
136148

149+
$eventRegistry = $this->eventRegistry;
150+
137151
return array_filter(
138152
$this->messages,
139-
static function (Message $message, int $index) use ($criteria): bool {
153+
static function (Message $message) use ($criteria, $eventRegistry): bool {
140154
foreach ($criteria->all() as $criterion) {
141155
switch ($criterion::class) {
142156
case AggregateIdCriterion::class:
@@ -222,10 +236,38 @@ static function (Message $message, int $index) use ($criteria): bool {
222236

223237
break;
224238
case FromIndexCriterion::class:
239+
try {
240+
$index = $message->header(IndexHeader::class)->index;
241+
} catch (HeaderNotFound) {
242+
return false;
243+
}
244+
225245
if ($index < $criterion->fromIndex) {
226246
return false;
227247
}
228248

249+
break;
250+
case ToIndexCriterion::class:
251+
try {
252+
$index = $message->header(IndexHeader::class)->index;
253+
} catch (HeaderNotFound) {
254+
return false;
255+
}
256+
257+
if ($index > $criterion->toIndex) {
258+
return false;
259+
}
260+
261+
break;
262+
case EventsCriterion::class:
263+
if ($eventRegistry === null) {
264+
throw new MissingEventRegistry($criterion::class);
265+
}
266+
267+
if (!in_array($eventRegistry->eventName($message->event()::class), $criterion->events)) {
268+
return false;
269+
}
270+
229271
break;
230272
default:
231273
throw new UnsupportedCriterion($criterion::class);

src/Store/MissingEventRegistry.php

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
<?php
2+
3+
declare(strict_types=1);
4+
5+
namespace Patchlevel\EventSourcing\Store;
6+
7+
use Patchlevel\EventSourcing\Metadata\Event\EventRegistry;
8+
use RuntimeException;
9+
10+
use function sprintf;
11+
12+
final class MissingEventRegistry extends RuntimeException
13+
{
14+
/** @param class-string $criterionClass */
15+
public function __construct(string $criterionClass)
16+
{
17+
parent::__construct(sprintf('criterion %s not supported without an %s given', $criterionClass, EventRegistry::class));
18+
}
19+
}

tests/Unit/Store/InMemoryStoreTest.php

Lines changed: 69 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,18 +7,24 @@
77
use DateTimeImmutable;
88
use Patchlevel\EventSourcing\Aggregate\AggregateHeader;
99
use Patchlevel\EventSourcing\Message\Message;
10+
use Patchlevel\EventSourcing\Metadata\Event\EventRegistry;
1011
use Patchlevel\EventSourcing\Store\ArchivedHeader;
1112
use Patchlevel\EventSourcing\Store\Criteria\AggregateIdCriterion;
1213
use Patchlevel\EventSourcing\Store\Criteria\AggregateNameCriterion;
1314
use Patchlevel\EventSourcing\Store\Criteria\ArchivedCriterion;
1415
use Patchlevel\EventSourcing\Store\Criteria\Criteria;
16+
use Patchlevel\EventSourcing\Store\Criteria\EventsCriterion;
1517
use Patchlevel\EventSourcing\Store\Criteria\FromIndexCriterion;
1618
use Patchlevel\EventSourcing\Store\Criteria\FromPlayheadCriterion;
1719
use Patchlevel\EventSourcing\Store\Criteria\StreamCriterion;
20+
use Patchlevel\EventSourcing\Store\Criteria\ToIndexCriterion;
1821
use Patchlevel\EventSourcing\Store\Header\PlayheadHeader;
1922
use Patchlevel\EventSourcing\Store\Header\StreamNameHeader;
2023
use Patchlevel\EventSourcing\Store\InMemoryStore;
24+
use Patchlevel\EventSourcing\Store\MissingEventRegistry;
2125
use Patchlevel\EventSourcing\Store\UnsupportedCriterion;
26+
use Patchlevel\EventSourcing\Tests\Unit\Fixture\Email;
27+
use Patchlevel\EventSourcing\Tests\Unit\Fixture\ProfileCreated;
2228
use Patchlevel\EventSourcing\Tests\Unit\Fixture\ProfileId;
2329
use Patchlevel\EventSourcing\Tests\Unit\Fixture\ProfileVisited;
2430
use PHPUnit\Framework\Attributes\CoversClass;
@@ -156,13 +162,33 @@ public function testLoadFromIndex(): void
156162

157163
$store = new InMemoryStore([$message1, $message2, $message3, $message4]);
158164

159-
$stream = $store->load(new Criteria(new FromIndexCriterion(2)));
165+
$stream = $store->load(new Criteria(new FromIndexCriterion(3)));
160166

161167
$messages = iterator_to_array($stream);
162168

163169
self::assertSame([$message3, $message4], $messages);
164170
}
165171

172+
public function testLoadToIndex(): void
173+
{
174+
$message1 = (new Message(new ProfileVisited(ProfileId::fromString('1'))))
175+
->withHeader(new AggregateHeader('foo', '1', 1, new DateTimeImmutable()));
176+
$message2 = (new Message(new ProfileVisited(ProfileId::fromString('2'))))
177+
->withHeader(new AggregateHeader('foo', '1', 2, new DateTimeImmutable()));
178+
$message3 = (new Message(new ProfileVisited(ProfileId::fromString('3'))))
179+
->withHeader(new StreamNameHeader('foo-1'))
180+
->withHeader(new PlayheadHeader(3));
181+
$message4 = new Message(new ProfileVisited(ProfileId::fromString('3')));
182+
183+
$store = new InMemoryStore([$message1, $message2, $message3, $message4]);
184+
185+
$stream = $store->load(new Criteria(new ToIndexCriterion(2)));
186+
187+
$messages = iterator_to_array($stream);
188+
189+
self::assertSame([$message1, $message2], $messages);
190+
}
191+
166192
public function testLoadByStreamNameWithLikeAll(): void
167193
{
168194
$message1 = (new Message(new ProfileVisited(ProfileId::fromString('1'))))
@@ -196,6 +222,48 @@ public function testLoadArchived(): void
196222
self::assertSame([$message1], $messages);
197223
}
198224

225+
public function testLoadByEventName(): void
226+
{
227+
$message1 = (new Message(new ProfileCreated(ProfileId::fromString('1'), Email::fromString('s@b.de'))))
228+
->withHeader(new StreamNameHeader('foo'));
229+
$message2 = (new Message(new ProfileVisited(ProfileId::fromString('2'))))
230+
->withHeader(new StreamNameHeader('bar'));
231+
$message3 = new Message(new ProfileVisited(ProfileId::fromString('3')));
232+
233+
$store = new InMemoryStore(
234+
[$message1, $message2, $message3],
235+
new EventRegistry([
236+
'profile_created' => ProfileCreated::class,
237+
'profile_visited' => ProfileVisited::class,
238+
]),
239+
);
240+
241+
$stream = $store->load(new Criteria(new EventsCriterion(['profile_created'])));
242+
$messages = iterator_to_array($stream);
243+
244+
self::assertSame([$message1], $messages);
245+
246+
$stream = $store->load(new Criteria(new EventsCriterion(['profile_visited'])));
247+
$messages = iterator_to_array($stream);
248+
249+
self::assertSame([$message2, $message3], $messages);
250+
self::assertSame([], iterator_to_array($store->load(new Criteria(new EventsCriterion(['profile_deleted'])))));
251+
}
252+
253+
public function testLoadByEventNameWithoutRegistry(): void
254+
{
255+
$message1 = (new Message(new ProfileCreated(ProfileId::fromString('1'), Email::fromString('s@b.de'))))
256+
->withHeader(new StreamNameHeader('foo'));
257+
$message2 = (new Message(new ProfileVisited(ProfileId::fromString('2'))))
258+
->withHeader(new StreamNameHeader('bar'));
259+
$message3 = new Message(new ProfileVisited(ProfileId::fromString('3')));
260+
261+
$store = new InMemoryStore([$message1, $message2, $message3]);
262+
263+
$this->expectException(MissingEventRegistry::class);
264+
$store->load(new Criteria(new EventsCriterion(['profile_created'])));
265+
}
266+
199267
public function testLoadUnsupportedCriterion(): void
200268
{
201269
$store = new InMemoryStore([

0 commit comments

Comments
 (0)