Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 34 additions & 0 deletions docs/store.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
7 changes: 6 additions & 1 deletion src/Store/AppendStore.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
10 changes: 8 additions & 2 deletions src/Store/InMemoryStore.php
Original file line number Diff line number Diff line change
Expand Up @@ -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
{
Expand Down Expand Up @@ -111,9 +112,14 @@
});
}

public function query(Query $query): Stream
/** @param positive-int|0 $from */
public function query(Query $query, int $from = 0): Stream

Check warning on line 116 in src/Store/InMemoryStore.php

View workflow job for this annotation

GitHub Actions / Mutation tests on diff (locked, 8.5, ubuntu-latest)

Escaped Mutant for Mutator "IncrementInteger": @@ @@ } /** @PARAM positive-int|0 $from */ - public function query(Query $query, int $from = 0): Stream + public function query(Query $query, int $from = 1): Stream { return new Stream(array_filter( $this->matchQuery($query),
{
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<Message> $messages */
Expand Down
9 changes: 8 additions & 1 deletion src/Store/TaggableDoctrineDbalStore.php
Original file line number Diff line number Diff line change
Expand Up @@ -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('*')
Expand All @@ -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(
Expand Down
70 changes: 70 additions & 0 deletions tests/Integration/Store/TaggableDoctrineDbalStoreTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
37 changes: 37 additions & 0 deletions tests/Unit/Store/InMemoryStoreTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
47 changes: 47 additions & 0 deletions tests/Unit/Store/TaggableDoctrineDbalStoreTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Loading