diff --git a/src/Store/InMemoryStore.php b/src/Store/InMemoryStore.php index 601b36072..71d7c9b26 100644 --- a/src/Store/InMemoryStore.php +++ b/src/Store/InMemoryStore.php @@ -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); diff --git a/src/Store/TaggableDoctrineDbalStore.php b/src/Store/TaggableDoctrineDbalStore.php index 13c935b3c..4905d8b48 100644 --- a/src/Store/TaggableDoctrineDbalStore.php +++ b/src/Store/TaggableDoctrineDbalStore.php @@ -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 @@ -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']); @@ -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') diff --git a/tests/Integration/Store/TaggableDoctrineDbalStoreTest.php b/tests/Integration/Store/TaggableDoctrineDbalStoreTest.php index e800a04c3..9dfef80a8 100644 --- a/tests/Integration/Store/TaggableDoctrineDbalStoreTest.php +++ b/tests/Integration/Store/TaggableDoctrineDbalStoreTest.php @@ -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(); @@ -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'); diff --git a/tests/Unit/Store/InMemoryStoreTest.php b/tests/Unit/Store/InMemoryStoreTest.php index 40700bf12..62c74956b 100644 --- a/tests/Unit/Store/InMemoryStoreTest.php +++ b/tests/Unit/Store/InMemoryStoreTest.php @@ -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();