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
4 changes: 4 additions & 0 deletions src/Store/InMemoryStore.php
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,10 @@ public function append(iterable $messages, AppendCondition|null $appendCondition
? iterator_to_array($messages, false)
: array_values($messages);

if ($messages === []) {
return;
}

$this->transactional(function () use ($messages, $appendCondition): void {
if ($appendCondition instanceof AppendCondition && $appendCondition->highestSequenceNumber !== null) {
$matched = $this->matchQuery($appendCondition->query);
Expand Down
18 changes: 10 additions & 8 deletions src/Store/TaggableDoctrineDbalStore.php
Original file line number Diff line number Diff line change
Expand Up @@ -468,6 +468,10 @@ public function append(iterable $messages, AppendCondition|null $appendCondition
$position++;
}

if ($selects === []) {
return;
}

// The rows are wrapped in a derived table so that the append condition
// below applies to the whole batch. Appending it straight after a
// `UNION ALL` chain would bind it to the last SELECT only, letting every
Expand Down Expand Up @@ -882,15 +886,17 @@ private function queryCondition(QueryBuilder $builder, Query $query): void
return;
}

foreach ($query->subQueries as $subQuery) {
if ($subQuery->empty() && !$subQuery->onlyLastEvent) {
return; // an empty sub query matches all events, same as SubQuery::match()
}
}

$subqueries = [];

$uniqueParameterGenerator = $this->uniqueParameterGenerator();

foreach ($query->subQueries as $subQuery) {
if ($subQuery->empty()) {
continue;
}

$subQueryBuilder = $this->connection->createQueryBuilder()
->select('id')
->from($this->config['table_name']);
Expand Down Expand Up @@ -940,10 +946,6 @@ private function queryCondition(QueryBuilder $builder, Query $query): void
$subqueries[] = $subQueryBuilder->getSQL();
}

if ($subqueries === []) {
return;
}

$joinQueryBuilder = $this->connection->createQueryBuilder()
->select('id')
->from('(' . implode(' UNION ALL ', $subqueries) . ')', 'j')
Expand Down
96 changes: 96 additions & 0 deletions tests/Integration/Store/TaggableDoctrineDbalStoreTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -774,6 +774,57 @@ public function testOnlyLastEvent(): void
self::assertStreamEquals([$message2], $stream);
}

public function testQueryWithEmptySubQueryMatchesAll(): 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(),
new SubQuery(['profile:' . $profileId1->toString()]),
));

self::assertStreamEquals([$message1, $message2], $stream);
}

public function testQueryWithEmptyOnlyLastEventSubQuery(): 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([]));

$this->store->append([$message1, $message2, $message3]);

$stream = $this->store->query(new Query(
new SubQuery(onlyLastEvent: true),
new SubQuery(['profile:' . $profileId1->toString()]),
));

self::assertStreamEquals([$message1, $message3], $stream);
}

public function testComplexQuery(): void
{
$profileId1 = ProfileId::generate();
Expand Down Expand Up @@ -984,6 +1035,51 @@ public function testAppendWithEmptyQuery(): void
self::assertStreamEquals([...$messages, $message], $this->store->load());
}

public function testAppendWithoutMessages(): void
{
$profileId = ProfileId::generate();

$message = Message::create(new ProfileCreated($profileId, 'test'))
->withHeader(new StreamNameHeader('foo'))
->withHeader(new RecordedOnHeader(new DateTimeImmutable('2020-01-01 00:00:00')))
->withHeader(new TagsHeader(['profile:' . $profileId->toString()]));

$this->store->append([$message]);
$this->store->append([]);
$this->store->append([], new AppendCondition(new Query(), 0));

self::assertStreamEquals([$message], $this->store->load());
}

public function testAppendConditionWithEmptySubQuery(): void
{
$profileId = ProfileId::generate();

$message1 = Message::create(new ProfileCreated($profileId, 'test'))
->withHeader(new StreamNameHeader('foo'))
->withHeader(new RecordedOnHeader(new DateTimeImmutable('2020-01-01 00:00:00')))
->withHeader(new TagsHeader(['profile:' . $profileId->toString()]));

$this->store->append([$message1]);

$message2 = Message::create(new ExternEvent('test message'))
->withHeader(new StreamNameHeader('foo'))
->withHeader(new TagsHeader([]));

$this->expectException(AppendConditionNotMet::class);

$this->store->append(
[$message2],
new AppendCondition(
new Query(
new SubQuery(),
new SubQuery(['unknown']),
),
0,
),
);
}

public function testStreams(): void
{
$profileId = ProfileId::fromString('0190e47e-77e9-7b90-bf62-08bbf0ab9b4b');
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 @@ -733,6 +733,43 @@ public function testQueryOnlyLastEvent(): void
self::assertSame([$message2], $stream->toList());
}

public function testQueryWithEmptySubQueryMatchesAll(): 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(), new SubQuery(['profile-1'])));

self::assertSame([$message1, $message2], $stream->toList());
}

public function testQueryWithEmptyOnlyLastEventSubQuery(): 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-3']);

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

$stream = $store->query(new Query(new SubQuery(onlyLastEvent: true), new SubQuery(['profile-1'])));

self::assertSame([$message1, $message3], $stream->toList());
}

public function testAppendWithoutMessages(): void
{
$message = $this->message(new ProfileVisited(ProfileId::fromString('1')), 1, ['profile-1']);

$store = new InMemoryStore([$message]);

$store->append([]);
$store->append([], new AppendCondition(new Query(), 0));

self::assertSame([$message], $store->load()->toList());
}

public function testAppendAssignsIndex(): void
{
$store = new InMemoryStore();
Expand Down
Loading