diff --git a/docs/store.md b/docs/store.md index c9a899c85..86a88ff44 100644 --- a/docs/store.md +++ b/docs/store.md @@ -401,6 +401,40 @@ The stream cannot rewind, so you can only iterate over it once. If you want to iterate over it again, you have to call the `load` method again. ::: +### Query + +Stores that implement `AppendStore`, like the `TaggableDoctrineDbalStore` and the `InMemoryStore`, +can load events by tags and event types with a `Query`. +This is what the [dynamic consistency boundary](dynamic-consistency-boundary.md) uses to build a decision model. + +```php +use Patchlevel\EventSourcing\Store\AppendStore; +use Patchlevel\EventSourcing\Store\Query; +use Patchlevel\EventSourcing\Store\SubQuery; + +/** @var AppendStore $store */ +$stream = $store->query( + new Query( + new SubQuery(['hotel:1'], [GuestIsCheckedIn::class, GuestIsCheckedOut::class]), + ), +); +``` +A message matches a sub query if it has all of its tags and one of its event types. +The query returns all messages that match at least one sub query, ordered by their index. + +You can pass an index as second argument to only load events from this index onwards (inclusive). + +```php +use Patchlevel\EventSourcing\Store\AppendStore; +use Patchlevel\EventSourcing\Store\Query; +use Patchlevel\EventSourcing\Store\SubQuery; + +/** @var AppendStore $store */ +$stream = $store->query( + new Query(new SubQuery(['hotel:1'])), + 1000, +); +``` ### Count You can count the number of events in the store with the `count` method. diff --git a/src/Store/AppendStore.php b/src/Store/AppendStore.php index b54b466bf..3bc987ed2 100644 --- a/src/Store/AppendStore.php +++ b/src/Store/AppendStore.php @@ -16,5 +16,10 @@ public function append( AppendCondition|null $appendCondition = null, ): void; - public function query(Query $query): Stream; + /** + * Returns all events matching the query, starting at the given index (inclusive). + * + * @param positive-int|0 $from + */ + public function query(Query $query, int $from = 0): Stream; } diff --git a/src/Store/InMemoryStore.php b/src/Store/InMemoryStore.php index 36cb17e18..ddad4022d 100644 --- a/src/Store/InMemoryStore.php +++ b/src/Store/InMemoryStore.php @@ -43,6 +43,7 @@ use function str_starts_with; use const ARRAY_FILTER_USE_BOTH; +use const ARRAY_FILTER_USE_KEY; final class InMemoryStore implements Store, AppendStore { @@ -111,9 +112,14 @@ public function save(Message ...$messages): void }); } - public function query(Query $query): Stream + /** @param positive-int|0 $from */ + public function query(Query $query, int $from = 0): Stream { - return new Stream($this->matchQuery($query)); + return new Stream(array_filter( + $this->matchQuery($query), + static fn (int $index): bool => $index >= $from, + ARRAY_FILTER_USE_KEY, + )); } /** @param iterable $messages */ diff --git a/src/Store/TaggableDoctrineDbalStore.php b/src/Store/TaggableDoctrineDbalStore.php index 1d6bab3d6..0b2676516 100644 --- a/src/Store/TaggableDoctrineDbalStore.php +++ b/src/Store/TaggableDoctrineDbalStore.php @@ -522,7 +522,8 @@ public function append(iterable $messages, AppendCondition|null $appendCondition }); } - public function query(Query $query): Stream + /** @param positive-int|0 $from */ + public function query(Query $query, int $from = 0): Stream { $builder = $this->connection->createQueryBuilder() ->select('*') @@ -531,6 +532,12 @@ public function query(Query $query): Stream $this->queryCondition($builder, $query); + if ($from > 0) { + $builder + ->andWhere('events.id >= :queryFrom') + ->setParameter('queryFrom', $from); + } + return new Stream( $this->buildGenerator( $this->connection->executeQuery( diff --git a/tests/Integration/Store/TaggableDoctrineDbalStoreTest.php b/tests/Integration/Store/TaggableDoctrineDbalStoreTest.php index c0d189297..4aa81256f 100644 --- a/tests/Integration/Store/TaggableDoctrineDbalStoreTest.php +++ b/tests/Integration/Store/TaggableDoctrineDbalStoreTest.php @@ -825,6 +825,76 @@ public function testQueryWithEmptyOnlyLastEventSubQuery(): void self::assertStreamEquals([$message1, $message3], $stream); } + public function testQueryFrom(): void + { + $profileId1 = ProfileId::generate(); + $profileId2 = ProfileId::generate(); + + $message1 = Message::create(new ProfileCreated($profileId1, 'test')) + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable('2020-01-01 00:00:00'))) + ->withHeader(new TagsHeader(['profile:' . $profileId1->toString()])); + $message2 = Message::create(new ProfileCreated($profileId2, 'test')) + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable('2020-01-01 00:00:00'))) + ->withHeader(new TagsHeader(['profile:' . $profileId2->toString()])); + $message3 = Message::create(new ExternEvent('test message')) + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new TagsHeader(['profile:' . $profileId1->toString()])); + + $this->store->append([$message1, $message2, $message3]); + + $stream = $this->store->query( + new Query(new SubQuery(['profile:' . $profileId1->toString()])), + 2, + ); + + self::assertStreamEquals([$message3], $stream); + } + + public function testQueryFromIsInclusive(): void + { + $profileId1 = ProfileId::generate(); + $profileId2 = ProfileId::generate(); + + $message1 = Message::create(new ProfileCreated($profileId1, 'test')) + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable('2020-01-01 00:00:00'))) + ->withHeader(new TagsHeader(['profile:' . $profileId1->toString()])); + $message2 = Message::create(new ProfileCreated($profileId2, 'test')) + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable('2020-01-01 00:00:00'))) + ->withHeader(new TagsHeader(['profile:' . $profileId2->toString()])); + + $this->store->append([$message1, $message2]); + + self::assertStreamEquals([$message2], $this->store->query(new Query(), 2)); + } + + public function testQueryFromWithOnlyLastEventBeforeFrom(): void + { + $profileId1 = ProfileId::generate(); + $profileId2 = ProfileId::generate(); + + $message1 = Message::create(new ProfileCreated($profileId1, 'test')) + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable('2020-01-01 00:00:00'))) + ->withHeader(new TagsHeader(['profile:' . $profileId1->toString()])); + $message2 = Message::create(new ProfileCreated($profileId2, 'test')) + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new RecordedOnHeader(new DateTimeImmutable('2020-01-01 00:00:00'))) + ->withHeader(new TagsHeader(['profile:' . $profileId2->toString()])); + + $this->store->append([$message1, $message2]); + + $stream = $this->store->query( + new Query(new SubQuery(['profile:' . $profileId1->toString()], onlyLastEvent: true)), + 2, + ); + + self::assertStreamEquals([], $stream); + } + public function testComplexQuery(): void { $profileId1 = ProfileId::generate(); diff --git a/tests/Unit/Store/InMemoryStoreTest.php b/tests/Unit/Store/InMemoryStoreTest.php index 8c1033e96..a6ffd17c8 100644 --- a/tests/Unit/Store/InMemoryStoreTest.php +++ b/tests/Unit/Store/InMemoryStoreTest.php @@ -770,6 +770,43 @@ public function testAppendWithoutMessages(): void self::assertSame([$message], $store->load()->toList()); } + public function testQueryFrom(): void + { + $message1 = $this->message(new ProfileVisited(ProfileId::fromString('1')), 1, ['profile-1']); + $message2 = $this->message(new ProfileVisited(ProfileId::fromString('2')), 2, ['profile-2']); + $message3 = $this->message(new ProfileVisited(ProfileId::fromString('3')), 3, ['profile-1']); + + $store = new InMemoryStore([$message1, $message2, $message3]); + + $stream = $store->query(new Query(new SubQuery(['profile-1'])), 2); + + self::assertSame([$message3], $stream->toList()); + } + + public function testQueryFromIsInclusive(): void + { + $message1 = $this->message(new ProfileVisited(ProfileId::fromString('1')), 1); + $message2 = $this->message(new ProfileVisited(ProfileId::fromString('2')), 2); + + $store = new InMemoryStore([$message1, $message2]); + + $stream = $store->query(new Query(), 2); + + self::assertSame([$message2], $stream->toList()); + } + + public function testQueryFromWithOnlyLastEventBeforeFrom(): void + { + $message1 = $this->message(new ProfileVisited(ProfileId::fromString('1')), 1, ['profile-1']); + $message2 = $this->message(new ProfileVisited(ProfileId::fromString('2')), 2, ['profile-2']); + + $store = new InMemoryStore([$message1, $message2]); + + $stream = $store->query(new Query(new SubQuery(['profile-1'], onlyLastEvent: true)), 2); + + self::assertSame([], $stream->toList()); + } + public function testAppendAssignsIndex(): void { $store = new InMemoryStore(); diff --git a/tests/Unit/Store/TaggableDoctrineDbalStoreTest.php b/tests/Unit/Store/TaggableDoctrineDbalStoreTest.php index b75802e97..90d162829 100644 --- a/tests/Unit/Store/TaggableDoctrineDbalStoreTest.php +++ b/tests/Unit/Store/TaggableDoctrineDbalStoreTest.php @@ -2825,6 +2825,53 @@ public function testQueryWithOnlyLastEvent(): void self::assertSame(null, $stream->position()); } + public function testQueryWithFrom(): void + { + $connection = $this->createMock(Connection::class); + $result = $this->createMock(Result::class); + $result + ->expects($this->once()) + ->method('iterateAssociative') + ->willReturn(new EmptyIterator()); + + $connection + ->expects($this->once()) + ->method('executeQuery') + ->with( + 'SELECT * FROM event_store events INNER JOIN (SELECT id FROM (SELECT id FROM event_store WHERE stream = :param1) j GROUP BY j.id) ej ON ej.id = events.id WHERE events.id >= :queryFrom ORDER BY events.id ASC', + ['param1' => 'profile-1', 'queryFrom' => 5], + $this->isArray(), + ) + ->willReturn($result); + + $connection + ->expects($this->exactly(5)) + ->method('getDatabasePlatform') + ->willReturn(new SQLitePlatform()); + $connection + ->expects($this->exactly(3)) + ->method('createQueryBuilder') + ->willReturnCallback( + static fn (): QueryBuilder => new QueryBuilder($connection), + ); + + $eventSerializer = $this->createMock(EventSerializer::class); + $eventRegistry = new EventRegistry([]); + $headersSerializer = $this->createMock(HeadersSerializer::class); + + $doctrineDbalStore = new TaggableDoctrineDbalStore( + $connection, + $eventSerializer, + $eventRegistry, + $headersSerializer, + ); + + $stream = $doctrineDbalStore->query(new Query(new SubQuery(streamName: 'profile-1')), 5); + + self::assertSame(null, $stream->index()); + self::assertSame(null, $stream->position()); + } + public function testQueryWithMultipleSubQueries(): void { $connection = $this->createMock(Connection::class);