From 46588fb06cdeee2180fe5d5bc6c94360e32f7ac6 Mon Sep 17 00:00:00 2001 From: David Badura Date: Sat, 3 Oct 2026 19:08:38 +0200 Subject: [PATCH] Align empty sub query and empty append handling across stores The dbal store skipped empty sub queries, while the in memory store and SubQuery::match() treat them as "match all". A query like `new Query(new SubQuery(), new SubQuery(['a']))` returned only the events tagged `a` on the database and every event in memory. The dbal store now returns all events in that case and resolves an empty only-last-event sub query to the last event of the store, the same as the in memory store. Appending no messages produced invalid sql on the dbal store. Both stores now treat it as a no-op, like StoreEventAppender already does. --- src/Store/InMemoryStore.php | 4 + src/Store/TaggableDoctrineDbalStore.php | 18 ++-- .../Store/TaggableDoctrineDbalStoreTest.php | 96 +++++++++++++++++++ tests/Unit/Store/InMemoryStoreTest.php | 37 +++++++ 4 files changed, 147 insertions(+), 8 deletions(-) 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();