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
2 changes: 2 additions & 0 deletions docs/dynamic-consistency-boundary.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 3 additions & 1 deletion src/DecisionModel/StoreDecisionModelBuilder.php
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@
use Patchlevel\EventSourcing\Store\AppendCondition;
use Patchlevel\EventSourcing\Store\AppendStore;

use function max;

/** @experimental */
final class StoreDecisionModelBuilder implements DecisionModelBuilder
{
Expand All @@ -31,7 +33,7 @@
$highestId = 0;

foreach ($stream as $message) {
$highestId = $stream->index() ?? 0;
$highestId = max($highestId, $stream->index() ?? 0);

Check warning on line 36 in src/DecisionModel/StoreDecisionModelBuilder.php

View workflow job for this annotation

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

Escaped Mutant for Mutator "DecrementInteger": @@ @@ $highestId = 0; foreach ($stream as $message) { - $highestId = max($highestId, $stream->index() ?? 0); + $highestId = max($highestId, $stream->index() ?? -1); $state = $projection->apply($state, $message); }

Check warning on line 36 in src/DecisionModel/StoreDecisionModelBuilder.php

View workflow job for this annotation

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

Escaped Mutant for Mutator "IncrementInteger": @@ @@ $highestId = 0; foreach ($stream as $message) { - $highestId = max($highestId, $stream->index() ?? 0); + $highestId = max($highestId, $stream->index() ?? 1); $state = $projection->apply($state, $message); }
$state = $projection->apply($state, $message);
}

Expand Down
21 changes: 10 additions & 11 deletions src/Store/AppendCondition.php
Original file line number Diff line number Diff line change
Expand Up @@ -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.',
);
}
}
}
5 changes: 2 additions & 3 deletions src/Store/InMemoryStore.php
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
}
Expand Down
15 changes: 5 additions & 10 deletions src/Store/TaggableDoctrineDbalStore.php
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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);
}

Expand Down
64 changes: 64 additions & 0 deletions tests/Integration/Store/TaggableDoctrineDbalStoreTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
3 changes: 2 additions & 1 deletion tests/Unit/Store/AppendConditionNotMetTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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',
Expand Down
24 changes: 4 additions & 20 deletions tests/Unit/Store/AppendConditionTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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);
}
}
36 changes: 36 additions & 0 deletions tests/Unit/Store/InMemoryStoreTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
11 changes: 6 additions & 5 deletions tests/Unit/Store/TaggableDoctrineDbalStoreTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -2145,7 +2145,7 @@ public function testAppendWithAppendCondition(): void
'recorded_on0' => $recordedOn,
'archived0' => false,
'custom_headers0' => '[]',
'highestId' => 5,
'appendConditionAfter' => 5,
'param1' => 'profile-1',
],
[
Expand Down Expand Up @@ -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,
Expand All @@ -2231,6 +2231,7 @@ public function testAppendWithAppendConditionZeroSequence(): void
'recorded_on0' => $recordedOn,
'archived0' => false,
'custom_headers0' => '[]',
'appendConditionAfter' => 0,
],
[
'tags0' => Type::getType(Types::JSON),
Expand Down Expand Up @@ -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,
Expand All @@ -2314,7 +2315,7 @@ public function testAppendConditionNotMet(): void
'recorded_on0' => $recordedOn,
'archived0' => false,
'custom_headers0' => '[]',
'highestId' => 5,
'appendConditionAfter' => 5,
],
[
'tags0' => Type::getType(Types::JSON),
Expand Down
Loading