From a07f42bb3eaed7b4522189d5022ecf5f760eba99 Mon Sep 17 00:00:00 2001 From: Gustavo Freze Date: Sun, 27 Sep 2026 12:55:51 -0300 Subject: [PATCH] feat: Carry the correlation id of the unit of work in an optional column. --- README.md | 51 +++++- src/DoctrineOutboxRepository.php | 16 +- src/Internal/OutboxInsert.php | 61 ++++--- src/Internal/OutboxWriter.php | 23 ++- src/Schema/Columns.php | 11 +- src/Schema/ColumnsBuilder.php | 19 ++- .../DoctrineOutboxRepositoryTest.php | 157 ++++++++++++++++++ tests/Integration/OutboxTableFactory.php | 11 ++ 8 files changed, 308 insertions(+), 41 deletions(-) diff --git a/README.md b/README.md index 30b3b2d..fd321d6 100644 --- a/README.md +++ b/README.md @@ -75,6 +75,14 @@ CREATE TABLE outbox_events The library writes to `id`, `aggregate_id`, `aggregate_type`, `event_type`, `revision`, `aggregate_version`, `payload`, and `occurred_at`. It never writes to `created_at`. The database fills it automatically. +A table may also carry the correlation id of the unit of work that emitted each event, so the consumer of the event +can keep the trace of the request that caused it. The column is optional and nullable, and the library writes it only +when the layout enables it (see [Carrying the correlation id](#carrying-the-correlation-id)): + +```sql +correlation_id VARCHAR(255) NULL COMMENT 'The correlation id of the unit of work that emitted the event, absent when none was in flight (e.g. 018f8e94-1c2a-7c3d-9b4e-5f6a7b8c9d0e).' +``` + For aggregates whose identities are not UUID strings, use VARCHAR columns and configure `IdentityColumnType::STRING` (see [Customizing the table layout](#customizing-the-table-layout)): @@ -121,12 +129,13 @@ pure enums, and date-times, and unwrapping single-property wrappers to their inn serializers before it for integration events that need custom shaping (see [Writing a custom payload serializer](#writing-a-custom-payload-serializer)). -| Parameter | Type | Required | Description | -|---------------|-------------------------------|:--------:|----------------------------------------------------------------------------------------------------------------------------------------------------------------------| -| `connection` | `Connection` | Yes | Doctrine DBAL connection used for all INSERT statements. | -| `serializers` | `PayloadSerializers` | Yes | Ordered collection of payload serializers operating on integration event records, first match wins. | -| `translators` | `IntegrationEventTranslators` | Yes | Ordered collection of translators mapping domain events to integration events. Records without a matching translator are persisted carrying the domain event itself. | -| `tableLayout` | `TableLayout` | No | Table and column configuration, defaults to `outbox_events` with BINARY(16) ids. | +| Parameter | Type | Required | Description | +|-----------------|-------------------------------|:--------:|----------------------------------------------------------------------------------------------------------------------------------------------------------------------| +| `connection` | `Connection` | Yes | Doctrine DBAL connection used for all INSERT statements. | +| `serializers` | `PayloadSerializers` | Yes | Ordered collection of payload serializers operating on integration event records, first match wins. | +| `translators` | `IntegrationEventTranslators` | Yes | Ordered collection of translators mapping domain events to integration events. Records without a matching translator are persisted carrying the domain event itself. | +| `tableLayout` | `TableLayout` | No | Table and column configuration, defaults to `outbox_events` with BINARY(16) ids. | +| `correlationId` | `Closure` | No | Reads the correlation id of the unit of work in flight at each write, an empty string when none. Written only when the layout enables the column. | ### Producing events from an aggregate @@ -348,6 +357,7 @@ require both `name:` and `type:`, all other methods require only `name:`. | `withAggregateType(name:)` | `aggregate_type` | | Renames the aggregate type column. | | `withAggregateVersion(name:)` | `aggregate_version` | | Renames the aggregate version column. | | `withCreatedAt(name:)` | `created_at` | | Renames the record creation timestamp column. | +| `withCorrelationId(name:)` | none, disabled | | Enables the correlation id column under the given name. | `TableLayout::builder()` controls the table name, columns, and unique constraint name. @@ -365,6 +375,35 @@ Constraint violation detection works with MySQL, MariaDB, PostgreSQL, and SQL Se constraint name in their violation messages. SQLite is not supported because it omits the constraint name. All unique violations with SQLite fall under `DuplicateOutboxEvent`. +### Carrying the correlation id + +Enable the column in the layout and give the repository a reader of the correlation id in flight. The reader runs at +each write, so one repository serves every request of a long-lived container. An empty string stores NULL, which is +what a worker outside any request writes. + +```php +withColumns(columns: Columns::builder()->withCorrelationId(name: 'correlation_id')->build()) + ->build(), + correlationId: static fn(): string => $correlationId->toString() +); +``` + ### Writing a custom payload serializer `PayloadSerializerReflection` covers integration events that `tiny-blocks/mapper` maps by reflection: scalars, diff --git a/src/DoctrineOutboxRepository.php b/src/DoctrineOutboxRepository.php index edde723..51c3ad3 100644 --- a/src/DoctrineOutboxRepository.php +++ b/src/DoctrineOutboxRepository.php @@ -4,6 +4,7 @@ namespace TinyBlocks\Outbox; +use Closure; use Doctrine\DBAL\Connection; use TinyBlocks\BuildingBlocks\Event\EventRecord; use TinyBlocks\BuildingBlocks\Event\EventRecords; @@ -18,18 +19,29 @@ private TableLayout $tableLayout; private OutboxWriter $writer; + /** + * @param Connection $connection The connection whose open transaction the rows join. + * @param PayloadSerializers $serializers The ordered payload serializers, first match wins. + * @param IntegrationEventTranslators $translators The ordered translators from domain to integration events. + * @param TableLayout|null $tableLayout The table and column configuration, the default layout when null. + * @param (Closure(): string)|null $correlationId Reads the correlation id of the unit of work in flight at + * each write, returning an empty string when there is none. It is + * written only when the layout enables the correlation id column. + */ public function __construct( private Connection $connection, PayloadSerializers $serializers, IntegrationEventTranslators $translators, - ?TableLayout $tableLayout = null + ?TableLayout $tableLayout = null, + ?Closure $correlationId = null ) { $this->tableLayout = ($tableLayout ?? TableLayout::default()); $this->writer = new OutboxWriter( connection: $connection, serializers: $serializers, tableLayout: $this->tableLayout, - translators: $translators + translators: $translators, + correlationId: $correlationId ); } diff --git a/src/Internal/OutboxInsert.php b/src/Internal/OutboxInsert.php index f45bed5..c1010e4 100644 --- a/src/Internal/OutboxInsert.php +++ b/src/Internal/OutboxInsert.php @@ -18,41 +18,50 @@ private function __construct(public string $sql, public array $parameters) public static function from( EventRecord|IntegrationEventRecord $record, SerializedPayload $payload, - TableLayout $tableLayout + TableLayout $tableLayout, + string $correlationId = '' ): OutboxInsert { - $template = <<columns; $idValue = $columns->id->convert(identityValue: $record->id->toString()); $aggregateIdValue = $columns->aggregateId->convert(identityValue: $record->aggregateId->identityValue()); + $names = [ + $columns->id->name(), + $columns->aggregateId->name(), + $columns->aggregateType, + $columns->eventType, + $columns->revision, + $columns->aggregateVersion, + $columns->payload, + $columns->occurredAt + ]; + + $parameters = [ + 'id' => $idValue, + 'aggregateId' => $aggregateIdValue, + 'aggregateType' => $record->aggregateType, + 'eventType' => $record->eventType->value, + 'revision' => $record->revision->value, + 'aggregateVersion' => $record->aggregateVersion->value, + 'payload' => $payload->toJson(), + 'occurredAt' => $record->occurredAt->toIso8601() + ]; + + if (!is_null($columns->correlationId)) { + $names[] = $columns->correlationId; + $parameters['correlationId'] = $correlationId === '' ? null : $correlationId; + } + + $placeholders = array_map(static fn(string $key): string => sprintf(':%s', $key), array_keys($parameters)); + return new OutboxInsert( sql: sprintf( - $template, + 'INSERT INTO %s (%s) VALUES (%s)', $tableLayout->tableName, - $columns->id->name(), - $columns->aggregateId->name(), - $columns->aggregateType, - $columns->eventType, - $columns->revision, - $columns->aggregateVersion, - $columns->payload, - $columns->occurredAt + implode(', ', $names), + implode(', ', $placeholders) ), - parameters: [ - 'id' => $idValue, - 'aggregateId' => $aggregateIdValue, - 'aggregateType' => $record->aggregateType, - 'eventType' => $record->eventType->value, - 'revision' => $record->revision->value, - 'aggregateVersion' => $record->aggregateVersion->value, - 'payload' => $payload->toJson(), - 'occurredAt' => $record->occurredAt->toIso8601() - ] + parameters: $parameters ); } } diff --git a/src/Internal/OutboxWriter.php b/src/Internal/OutboxWriter.php index b72302d..f2eeec5 100644 --- a/src/Internal/OutboxWriter.php +++ b/src/Internal/OutboxWriter.php @@ -4,6 +4,7 @@ namespace TinyBlocks\Outbox\Internal; +use Closure; use Doctrine\DBAL\Connection; use Doctrine\DBAL\Exception\UniqueConstraintViolationException; use TinyBlocks\BuildingBlocks\Event\EventRecord; @@ -18,24 +19,34 @@ final readonly class OutboxWriter { + /** + * @param Connection $connection The connection whose open transaction the rows join. + * @param PayloadSerializers $serializers The ordered payload serializers, first match wins. + * @param TableLayout $tableLayout The table and column configuration. + * @param IntegrationEventTranslators $translators The ordered translators from domain to integration events. + * @param (Closure(): string)|null $correlationId Reads the correlation id of the unit of work in flight. + */ public function __construct( private Connection $connection, private PayloadSerializers $serializers, private TableLayout $tableLayout, - private IntegrationEventTranslators $translators + private IntegrationEventTranslators $translators, + private ?Closure $correlationId = null ) { } public function write(EventRecord $eventRecord): void { $translator = $this->translators->findFor(record: $eventRecord); + $correlationId = $this->correlationId(); try { if (is_null($translator)) { $insert = OutboxInsert::from( record: $eventRecord, payload: SerializedPayload::fromEvent(event: $eventRecord->event), - tableLayout: $this->tableLayout + tableLayout: $this->tableLayout, + correlationId: $correlationId ); $this->connection->executeStatement(sql: $insert->sql, params: $insert->parameters); @@ -57,7 +68,8 @@ public function write(EventRecord $eventRecord): void $insert = OutboxInsert::from( record: $record, payload: $payloadSerializer->serialize(record: $record), - tableLayout: $this->tableLayout + tableLayout: $this->tableLayout, + correlationId: $correlationId ); $this->connection->executeStatement(sql: $insert->sql, params: $insert->parameters); @@ -77,4 +89,9 @@ public function write(EventRecord $eventRecord): void ); } } + + private function correlationId(): string + { + return is_null($this->correlationId) ? '' : ($this->correlationId)(); + } } diff --git a/src/Schema/Columns.php b/src/Schema/Columns.php index 48d14d0..6596dec 100644 --- a/src/Schema/Columns.php +++ b/src/Schema/Columns.php @@ -15,7 +15,8 @@ private function __construct( public string $occurredAt, public IdentityColumn $aggregateId, public string $aggregateType, - public string $aggregateVersion + public string $aggregateVersion, + public ?string $correlationId = null ) { } @@ -31,6 +32,8 @@ private function __construct( * @param IdentityColumn $aggregateId The identity column for the owning aggregate. * @param string $aggregateType The column name for the aggregate type identifier. * @param string $aggregateVersion The column name for the per-aggregate version counter. + * @param string|null $correlationId The column name for the correlation id of the unit of work that emitted the + * event, or null when the table carries none. * @return Columns The built column configuration. */ public static function from( @@ -42,7 +45,8 @@ public static function from( string $occurredAt, IdentityColumn $aggregateId, string $aggregateType, - string $aggregateVersion + string $aggregateVersion, + ?string $correlationId = null ): Columns { return new Columns( id: $id, @@ -53,7 +57,8 @@ public static function from( occurredAt: $occurredAt, aggregateId: $aggregateId, aggregateType: $aggregateType, - aggregateVersion: $aggregateVersion + aggregateVersion: $aggregateVersion, + correlationId: $correlationId ); } diff --git a/src/Schema/ColumnsBuilder.php b/src/Schema/ColumnsBuilder.php index b56edf6..05cc925 100644 --- a/src/Schema/ColumnsBuilder.php +++ b/src/Schema/ColumnsBuilder.php @@ -21,6 +21,7 @@ final class ColumnsBuilder private IdentityColumnType $aggregateIdType = IdentityColumnType::BINARY; private string $aggregateType = 'aggregate_type'; private string $aggregateVersion = 'aggregate_version'; + private ?string $correlationId = null; private function __construct() { @@ -52,7 +53,8 @@ public function build(): Columns occurredAt: $this->occurredAt, aggregateId: $this->aggregateIdType->toColumn(name: $this->aggregateIdName), aggregateType: $this->aggregateType, - aggregateVersion: $this->aggregateVersion + aggregateVersion: $this->aggregateVersion, + correlationId: $this->correlationId ); } @@ -167,4 +169,19 @@ public function withAggregateVersion(string $name): ColumnsBuilder $this->aggregateVersion = $name; return $this; } + + /** + * Enables the correlation id column under the given name. + * + *

No column is written by default. Once enabled, every row carries the correlation id the repository reads + * at write time, or NULL when there is none.

+ * + * @param string $name The correlation id column name. + * @return ColumnsBuilder The builder for chaining. + */ + public function withCorrelationId(string $name): ColumnsBuilder + { + $this->correlationId = $name; + return $this; + } } diff --git a/tests/Integration/DoctrineOutboxRepositoryTest.php b/tests/Integration/DoctrineOutboxRepositoryTest.php index 6b1f54e..5743f89 100644 --- a/tests/Integration/DoctrineOutboxRepositoryTest.php +++ b/tests/Integration/DoctrineOutboxRepositoryTest.php @@ -1252,4 +1252,161 @@ public function testPushWhenSerializerDoesNotSupportIntegrationEventThenPayloadS /** @When pushing the unsupported event */ $repository->push(records: $records); } + + public function testPushWhenCorrelationColumnIsEnabledThenRowCarriesTheCorrelationIdOfTheUnitOfWork(): void + { + /** @Given a layout that enables the correlation id column */ + $tableLayout = self::correlatedLayout(); + + /** @And the outbox table carries that column */ + OutboxTableFactory::addCorrelationIdColumn(connection: self::$connection, tableLayout: $tableLayout); + + /** @And a repository that reads the correlation id of the unit of work in flight */ + $repository = new DoctrineOutboxRepository( + connection: self::$connection, + serializers: PayloadSerializers::createFrom(elements: [new OrderPlacedSerializer()]), + translators: IntegrationEventTranslators::createFrom(elements: [new OrderPlacedTranslator()]), + tableLayout: $tableLayout, + correlationId: static fn(): string => 'req-0001' + ); + + /** @When a record is pushed inside a committed transaction */ + self::$connection->beginTransaction(); + $repository->push(records: EventRecords::createFrom(elements: [ + EventRecordFactory::create(event: new OrderPlaced(), aggregateType: 'Order') + ])); + self::$connection->commit(); + + /** @Then the row carries the correlation id */ + self::assertSame('req-0001', self::$connection->fetchOne('SELECT correlation_id FROM outbox_events')); + } + + public function testPushWhenCorrelationIdIsEmptyThenColumnIsNull(): void + { + /** @Given a layout that enables the correlation id column, on a table that carries it */ + $tableLayout = self::correlatedLayout(); + OutboxTableFactory::addCorrelationIdColumn(connection: self::$connection, tableLayout: $tableLayout); + + /** @And a repository whose unit of work has no correlation id, as a worker outside any request */ + $repository = new DoctrineOutboxRepository( + connection: self::$connection, + serializers: PayloadSerializers::createFrom(elements: [new OrderPlacedSerializer()]), + translators: IntegrationEventTranslators::createFrom(elements: [new OrderPlacedTranslator()]), + tableLayout: $tableLayout, + correlationId: static fn(): string => '' + ); + + /** @When a record is pushed inside a committed transaction */ + self::$connection->beginTransaction(); + $repository->push(records: EventRecords::createFrom(elements: [ + EventRecordFactory::create(event: new OrderPlaced(), aggregateType: 'Order') + ])); + self::$connection->commit(); + + /** @Then the column holds NULL instead of an empty string */ + self::assertNull(self::$connection->fetchOne('SELECT correlation_id FROM outbox_events')); + } + + public function testPushWhenCorrelationColumnIsEnabledWithoutReaderThenColumnIsNull(): void + { + /** @Given a layout that enables the correlation id column, on a table that carries it */ + $tableLayout = self::correlatedLayout(); + OutboxTableFactory::addCorrelationIdColumn(connection: self::$connection, tableLayout: $tableLayout); + + /** @And a repository built without a correlation id reader */ + $repository = new DoctrineOutboxRepository( + connection: self::$connection, + serializers: PayloadSerializers::createFrom(elements: []), + translators: IntegrationEventTranslators::createFrom(elements: []), + tableLayout: $tableLayout + ); + + /** @When a domain event record without a translator is pushed inside a committed transaction */ + self::$connection->beginTransaction(); + $repository->push(records: EventRecords::createFrom(elements: [ + EventRecordFactory::create(event: new InventoryReserved(sku: 'SKU-1', quantity: 3), aggregateType: 'Order') + ])); + self::$connection->commit(); + + /** @Then the column holds NULL */ + self::assertNull(self::$connection->fetchOne('SELECT correlation_id FROM outbox_events')); + } + + public function testPushWhenReaderIsGivenButColumnIsNotEnabledThenBindingsCarryNoCorrelationId(): void + { + /** @Given a mocked connection with an active transaction */ + $connection = $this->createMock(Connection::class); + $connection->method('isTransactionActive')->willReturn(true); + + /** @And a variable to capture the SQL and the parameters passed to executeStatement */ + $captured = ['sql' => '', 'params' => []]; + $connection->expects(self::once()) + ->method('executeStatement') + ->willReturnCallback( + function (string $sql, array $params) use (&$captured): int { + $captured = ['sql' => $sql, 'params' => $params]; + return 1; + } + ); + + /** @When a record is pushed with a correlation id reader but the default layout */ + new DoctrineOutboxRepository( + connection: $connection, + serializers: PayloadSerializers::createFrom(elements: [new OrderPlacedSerializer()]), + translators: IntegrationEventTranslators::createFrom(elements: [new OrderPlacedTranslator()]), + correlationId: static fn(): string => 'req-0001' + )->push( + records: EventRecords::createFrom(elements: [ + EventRecordFactory::create(event: new OrderPlaced(), aggregateType: 'Order') + ]) + ); + + /** @Then neither the statement nor the bindings mention the correlation id */ + self::assertStringNotContainsString('correlation', $captured['sql']); + self::assertArrayNotHasKey('correlationId', $captured['params']); + self::assertCount(8, $captured['params']); + } + + public function testPushWhenCorrelationColumnIsEnabledThenStatementNamesTheColumnAndBindsTheValue(): void + { + /** @Given a mocked connection with an active transaction */ + $connection = $this->createMock(Connection::class); + $connection->method('isTransactionActive')->willReturn(true); + + /** @And a variable to capture the SQL and the parameters passed to executeStatement */ + $captured = ['sql' => '', 'params' => []]; + $connection->expects(self::once()) + ->method('executeStatement') + ->willReturnCallback( + function (string $sql, array $params) use (&$captured): int { + $captured = ['sql' => $sql, 'params' => $params]; + return 1; + } + ); + + /** @When a record is pushed with a layout that enables the correlation id column */ + new DoctrineOutboxRepository( + connection: $connection, + serializers: PayloadSerializers::createFrom(elements: [new OrderPlacedSerializer()]), + translators: IntegrationEventTranslators::createFrom(elements: [new OrderPlacedTranslator()]), + tableLayout: self::correlatedLayout(), + correlationId: static fn(): string => 'req-0001' + )->push( + records: EventRecords::createFrom(elements: [ + EventRecordFactory::create(event: new OrderPlaced(), aggregateType: 'Order') + ]) + ); + + /** @Then the statement names the column last, with its placeholder, and binds the value */ + self::assertStringContainsString(', occurred_at, correlation_id) VALUES (', $captured['sql']); + self::assertStringEndsWith(', :occurredAt, :correlationId)', $captured['sql']); + self::assertSame('req-0001', $captured['params']['correlationId']); + } + + private static function correlatedLayout(): TableLayout + { + return TableLayout::builder() + ->withColumns(columns: Columns::builder()->withCorrelationId(name: 'correlation_id')->build()) + ->build(); + } } diff --git a/tests/Integration/OutboxTableFactory.php b/tests/Integration/OutboxTableFactory.php index 424cb6f..305889b 100644 --- a/tests/Integration/OutboxTableFactory.php +++ b/tests/Integration/OutboxTableFactory.php @@ -98,4 +98,15 @@ public static function createWithStringIdentities(Connection $connection, TableL ) ); } + + public static function addCorrelationIdColumn(Connection $connection, TableLayout $tableLayout): void + { + $connection->executeStatement( + sprintf( + 'ALTER TABLE %s ADD COLUMN %s VARCHAR(255) NULL', + $tableLayout->tableName, + $tableLayout->columns->correlationId + ) + ); + } }