diff --git a/docs/dynamic-consistency-boundary.md b/docs/dynamic-consistency-boundary.md index 0605d0695..b2fade6dd 100644 --- a/docs/dynamic-consistency-boundary.md +++ b/docs/dynamic-consistency-boundary.md @@ -304,6 +304,8 @@ final class CreateHotelHandler :::note Handlers build a Decision Model from the projections and then append events with an optimistic AppendCondition. If any relevant event arrives between read and write, the append fails and you can retry. +The AppendCondition contains the query of the Decision Model and the index of the last event it has seen as `after`. +The append fails if an event matching the query exists with a higher index than `after`. ::: :::tip diff --git a/src/DecisionModel/StoreDecisionModelBuilder.php b/src/DecisionModel/StoreDecisionModelBuilder.php index 2b48aac85..7eaef7b6b 100644 --- a/src/DecisionModel/StoreDecisionModelBuilder.php +++ b/src/DecisionModel/StoreDecisionModelBuilder.php @@ -9,6 +9,8 @@ use Patchlevel\EventSourcing\Store\AppendCondition; use Patchlevel\EventSourcing\Store\AppendStore; +use function max; + /** @experimental */ final class StoreDecisionModelBuilder implements DecisionModelBuilder { @@ -31,7 +33,7 @@ public function build( $highestId = 0; foreach ($stream as $message) { - $highestId = $stream->index() ?? 0; + $highestId = max($highestId, $stream->index() ?? 0); $state = $projection->apply($state, $message); } diff --git a/src/Store/AppendCondition.php b/src/Store/AppendCondition.php index 2697c84d7..019f0efa9 100644 --- a/src/Store/AppendCondition.php +++ b/src/Store/AppendCondition.php @@ -4,20 +4,19 @@ namespace Patchlevel\EventSourcing\Store; -use InvalidArgumentException; - -/** @experimental */ +/** + * The append fails if the store contains an event matching the query with an index higher than `after`. + * `after` is the highest index the caller was aware of while building the decision model. + * With `after` 0, no event may match the query at all. + * + * @experimental + */ final class AppendCondition { + /** @param positive-int|0 $after */ public function __construct( - public readonly Query $query = new Query(), - public readonly int|null $highestSequenceNumber = null, + public readonly Query $query, + public readonly int $after = 0, ) { - if ($query->subQueries !== [] && $highestSequenceNumber === null) { - throw new InvalidArgumentException( - 'An AppendCondition with a non-empty query needs a highestSequenceNumber. ' - . 'Pass 0 to require that no matching event exists yet.', - ); - } } } diff --git a/src/Store/InMemoryStore.php b/src/Store/InMemoryStore.php index 601b36072..510804df7 100644 --- a/src/Store/InMemoryStore.php +++ b/src/Store/InMemoryStore.php @@ -124,11 +124,10 @@ public function append(iterable $messages, AppendCondition|null $appendCondition : array_values($messages); $this->transactional(function () use ($messages, $appendCondition): void { - if ($appendCondition instanceof AppendCondition && $appendCondition->highestSequenceNumber !== null) { + if ($appendCondition instanceof AppendCondition) { $matched = $this->matchQuery($appendCondition->query); - $highestSequenceNumber = $matched === [] ? 0 : array_key_last($matched); - if ($highestSequenceNumber !== $appendCondition->highestSequenceNumber) { + if ($matched !== [] && array_key_last($matched) > $appendCondition->after) { throw new AppendConditionNotMet($appendCondition); } } diff --git a/src/Store/TaggableDoctrineDbalStore.php b/src/Store/TaggableDoctrineDbalStore.php index 13c935b3c..124c4e714 100644 --- a/src/Store/TaggableDoctrineDbalStore.php +++ b/src/Store/TaggableDoctrineDbalStore.php @@ -482,21 +482,16 @@ public function append(iterable $messages, AppendCondition|null $appendCondition implode(' UNION ALL ', $selects), ); - if ($appendCondition instanceof AppendCondition && $appendCondition->highestSequenceNumber !== null) { + if ($appendCondition instanceof AppendCondition) { $queryBuilder = $this->connection->createQueryBuilder() ->select('events.id') ->from($this->config['table_name'], 'events') - ->orderBy('events.id', 'DESC') - ->setMaxResults(1); + ->where('events.id > :appendConditionAfter'); $this->queryCondition($queryBuilder, $appendCondition->query); - if ($appendCondition->highestSequenceNumber === 0) { - $query .= ' WHERE NOT EXISTS (' . $queryBuilder->getSQL() . ')'; - } else { - $query .= ' WHERE (' . $queryBuilder->getSQL() . ') = :highestId'; - $parameters['highestId'] = $appendCondition->highestSequenceNumber; - } + $query .= ' WHERE NOT EXISTS (' . $queryBuilder->getSQL() . ')'; + $parameters['appendConditionAfter'] = $appendCondition->after; $parameters = array_merge( $parameters, @@ -515,7 +510,7 @@ public function append(iterable $messages, AppendCondition|null $appendCondition throw new UniqueConstraintViolation($e); } - if ($affectedRows === 0 && $appendCondition && $appendCondition->highestSequenceNumber !== null) { + if ($affectedRows === 0 && $appendCondition instanceof AppendCondition) { throw new AppendConditionNotMet($appendCondition); } diff --git a/tests/Integration/Store/TaggableDoctrineDbalStoreTest.php b/tests/Integration/Store/TaggableDoctrineDbalStoreTest.php index e800a04c3..401fce2f4 100644 --- a/tests/Integration/Store/TaggableDoctrineDbalStoreTest.php +++ b/tests/Integration/Store/TaggableDoctrineDbalStoreTest.php @@ -847,6 +847,70 @@ public function testAppendRaceCondition(): void ); } + public function testAppendConditionWithAfterHigherThanLastMatchingEvent(): void + { + $profileId1 = ProfileId::generate(); + $profileId2 = ProfileId::generate(); + + $messages = [ + 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()])), + 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($messages); + + $message = Message::create(new ExternEvent('test message')) + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new TagsHeader(['profile:' . $profileId1->toString()])); + + $this->store->append( + [$message], + new AppendCondition( + new Query(new SubQuery(['profile:' . $profileId1->toString()])), + 2, + ), + ); + + self::assertStreamEquals([...$messages, $message], $this->store->load()); + } + + public function testAppendConditionWithMatchingEventAfter(): void + { + $profileId = ProfileId::generate(); + + $messages = [ + 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()])), + Message::create(new ExternEvent('test message')) + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new TagsHeader(['profile:' . $profileId->toString()])), + ]; + + $this->store->append($messages); + + $this->expectException(AppendConditionNotMet::class); + + $this->store->append( + [ + Message::create(new ExternEvent('test message')) + ->withHeader(new StreamNameHeader('foo')) + ->withHeader(new TagsHeader(['profile:' . $profileId->toString()])), + ], + new AppendCondition( + new Query(new SubQuery(['profile:' . $profileId->toString()])), + 1, + ), + ); + } + public function testAppendConditionRejectsEntireBatch(): void { $profileId1 = ProfileId::generate(); diff --git a/tests/Unit/Store/AppendConditionNotMetTest.php b/tests/Unit/Store/AppendConditionNotMetTest.php index c3f666da9..ea62ddec4 100644 --- a/tests/Unit/Store/AppendConditionNotMetTest.php +++ b/tests/Unit/Store/AppendConditionNotMetTest.php @@ -6,6 +6,7 @@ use Patchlevel\EventSourcing\Store\AppendCondition; use Patchlevel\EventSourcing\Store\AppendConditionNotMet; +use Patchlevel\EventSourcing\Store\Query; use PHPUnit\Framework\Attributes\CoversClass; use PHPUnit\Framework\TestCase; @@ -14,7 +15,7 @@ final class AppendConditionNotMetTest extends TestCase { public function testCreate(): void { - $exception = new AppendConditionNotMet($condition = new AppendCondition()); + $exception = new AppendConditionNotMet($condition = new AppendCondition(new Query())); self::assertSame( 'Append condition not met', diff --git a/tests/Unit/Store/AppendConditionTest.php b/tests/Unit/Store/AppendConditionTest.php index 1dd7ba43c..08f7f50ba 100644 --- a/tests/Unit/Store/AppendConditionTest.php +++ b/tests/Unit/Store/AppendConditionTest.php @@ -4,7 +4,6 @@ namespace Patchlevel\EventSourcing\Tests\Unit\Store; -use InvalidArgumentException; use Patchlevel\EventSourcing\Store\AppendCondition; use Patchlevel\EventSourcing\Store\Query; use Patchlevel\EventSourcing\Store\SubQuery; @@ -21,31 +20,16 @@ public function testInstantiate(): void $condition = new AppendCondition($query, 42); self::assertSame($query, $condition->query); - self::assertSame(42, $condition->highestSequenceNumber); + self::assertSame(42, $condition->after); } - public function testInstantiateWithDefaults(): void - { - $condition = new AppendCondition(); - - self::assertEquals(new Query(), $condition->query); - self::assertNull($condition->highestSequenceNumber); - } - - public function testInstantiateWithZeroSequenceAndQuery(): void + public function testAfterDefaultsToZero(): void { $query = new Query(new SubQuery(['foo'])); - $condition = new AppendCondition($query, 0); + $condition = new AppendCondition($query); self::assertSame($query, $condition->query); - self::assertSame(0, $condition->highestSequenceNumber); - } - - public function testNonEmptyQueryWithoutSequenceNumberIsRejected(): void - { - $this->expectException(InvalidArgumentException::class); - - new AppendCondition(new Query(new SubQuery(['foo']))); + self::assertSame(0, $condition->after); } } diff --git a/tests/Unit/Store/InMemoryStoreTest.php b/tests/Unit/Store/InMemoryStoreTest.php index 40700bf12..6ac2004dc 100644 --- a/tests/Unit/Store/InMemoryStoreTest.php +++ b/tests/Unit/Store/InMemoryStoreTest.php @@ -809,6 +809,42 @@ public function testAppendWithConditionRejectsEntireBatch(): void self::assertCount(1, iterator_to_array($store->load())); } + public function testAppendWithAfterHigherThanLastMatchingEvent(): 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]); + + $store->append( + [ + (new Message(new ProfileVisited(ProfileId::fromString('3')))) + ->withHeader(new TagsHeader(['profile-1'])), + ], + new AppendCondition(new Query(new SubQuery(['profile-1'])), 2), + ); + + self::assertCount(3, iterator_to_array($store->load())); + } + + public function testAppendWithMatchingEventAfter(): void + { + $message1 = $this->message(new ProfileVisited(ProfileId::fromString('1')), 1, ['profile-1']); + $message2 = $this->message(new ProfileVisited(ProfileId::fromString('2')), 2, ['profile-1']); + + $store = new InMemoryStore([$message1, $message2]); + + $this->expectException(AppendConditionNotMet::class); + + $store->append( + [ + (new Message(new ProfileVisited(ProfileId::fromString('3')))) + ->withHeader(new TagsHeader(['profile-1'])), + ], + new AppendCondition(new Query(new SubQuery(['profile-1'])), 1), + ); + } + public function testAppendWithZeroSequenceConditionMet(): void { $store = new InMemoryStore(); diff --git a/tests/Unit/Store/TaggableDoctrineDbalStoreTest.php b/tests/Unit/Store/TaggableDoctrineDbalStoreTest.php index 06077ff66..b75802e97 100644 --- a/tests/Unit/Store/TaggableDoctrineDbalStoreTest.php +++ b/tests/Unit/Store/TaggableDoctrineDbalStoreTest.php @@ -2134,7 +2134,7 @@ public function testAppendWithAppendCondition(): void ->expects($this->once()) ->method('executeStatement') ->with( - 'INSERT INTO event_store (stream, playhead, event_id, event_name, event_payload, tags, recorded_on, archived, custom_headers) SELECT stream, playhead, event_id, event_name, event_payload, tags, recorded_on, archived, custom_headers FROM (SELECT :stream0 AS stream, :playhead0 AS playhead, :event_id0 AS event_id, :event_name0 AS event_name, :event_payload0 AS event_payload, :tags0 AS tags, :recorded_on0 AS recorded_on, :archived0 AS archived, :custom_headers0 AS custom_headers) AS data WHERE (SELECT events.id 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 ORDER BY events.id DESC LIMIT 1) = :highestId', + 'INSERT INTO event_store (stream, playhead, event_id, event_name, event_payload, tags, recorded_on, archived, custom_headers) SELECT stream, playhead, event_id, event_name, event_payload, tags, recorded_on, archived, custom_headers FROM (SELECT :stream0 AS stream, :playhead0 AS playhead, :event_id0 AS event_id, :event_name0 AS event_name, :event_payload0 AS event_payload, :tags0 AS tags, :recorded_on0 AS recorded_on, :archived0 AS archived, :custom_headers0 AS custom_headers) AS data WHERE NOT EXISTS (SELECT events.id 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 > :appendConditionAfter)', [ 'stream0' => 'profile-1', 'playhead0' => 1, @@ -2145,7 +2145,7 @@ public function testAppendWithAppendCondition(): void 'recorded_on0' => $recordedOn, 'archived0' => false, 'custom_headers0' => '[]', - 'highestId' => 5, + 'appendConditionAfter' => 5, 'param1' => 'profile-1', ], [ @@ -2220,7 +2220,7 @@ public function testAppendWithAppendConditionZeroSequence(): void ->expects($this->once()) ->method('executeStatement') ->with( - 'INSERT INTO event_store (stream, playhead, event_id, event_name, event_payload, tags, recorded_on, archived, custom_headers) SELECT stream, playhead, event_id, event_name, event_payload, tags, recorded_on, archived, custom_headers FROM (SELECT :stream0 AS stream, :playhead0 AS playhead, :event_id0 AS event_id, :event_name0 AS event_name, :event_payload0 AS event_payload, :tags0 AS tags, :recorded_on0 AS recorded_on, :archived0 AS archived, :custom_headers0 AS custom_headers) AS data WHERE NOT EXISTS (SELECT events.id FROM event_store events ORDER BY events.id DESC LIMIT 1)', + 'INSERT INTO event_store (stream, playhead, event_id, event_name, event_payload, tags, recorded_on, archived, custom_headers) SELECT stream, playhead, event_id, event_name, event_payload, tags, recorded_on, archived, custom_headers FROM (SELECT :stream0 AS stream, :playhead0 AS playhead, :event_id0 AS event_id, :event_name0 AS event_name, :event_payload0 AS event_payload, :tags0 AS tags, :recorded_on0 AS recorded_on, :archived0 AS archived, :custom_headers0 AS custom_headers) AS data WHERE NOT EXISTS (SELECT events.id FROM event_store events WHERE events.id > :appendConditionAfter)', [ 'stream0' => 'profile-1', 'playhead0' => 1, @@ -2231,6 +2231,7 @@ public function testAppendWithAppendConditionZeroSequence(): void 'recorded_on0' => $recordedOn, 'archived0' => false, 'custom_headers0' => '[]', + 'appendConditionAfter' => 0, ], [ 'tags0' => Type::getType(Types::JSON), @@ -2303,7 +2304,7 @@ public function testAppendConditionNotMet(): void ->expects($this->once()) ->method('executeStatement') ->with( - 'INSERT INTO event_store (stream, playhead, event_id, event_name, event_payload, tags, recorded_on, archived, custom_headers) SELECT stream, playhead, event_id, event_name, event_payload, tags, recorded_on, archived, custom_headers FROM (SELECT :stream0 AS stream, :playhead0 AS playhead, :event_id0 AS event_id, :event_name0 AS event_name, :event_payload0 AS event_payload, :tags0 AS tags, :recorded_on0 AS recorded_on, :archived0 AS archived, :custom_headers0 AS custom_headers) AS data WHERE (SELECT events.id FROM event_store events ORDER BY events.id DESC LIMIT 1) = :highestId', + 'INSERT INTO event_store (stream, playhead, event_id, event_name, event_payload, tags, recorded_on, archived, custom_headers) SELECT stream, playhead, event_id, event_name, event_payload, tags, recorded_on, archived, custom_headers FROM (SELECT :stream0 AS stream, :playhead0 AS playhead, :event_id0 AS event_id, :event_name0 AS event_name, :event_payload0 AS event_payload, :tags0 AS tags, :recorded_on0 AS recorded_on, :archived0 AS archived, :custom_headers0 AS custom_headers) AS data WHERE NOT EXISTS (SELECT events.id FROM event_store events WHERE events.id > :appendConditionAfter)', [ 'stream0' => 'profile-1', 'playhead0' => 1, @@ -2314,7 +2315,7 @@ public function testAppendConditionNotMet(): void 'recorded_on0' => $recordedOn, 'archived0' => false, 'custom_headers0' => '[]', - 'highestId' => 5, + 'appendConditionAfter' => 5, ], [ 'tags0' => Type::getType(Types::JSON),