Skip to content
Closed
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
51 changes: 45 additions & 6 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)):

Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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.

Expand All @@ -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
<?php

declare(strict_types=1);

use TinyBlocks\BuildingBlocks\Event\IntegrationEventTranslators;
use TinyBlocks\Outbox\DoctrineOutboxRepository;
use TinyBlocks\Outbox\Schema\Columns;
use TinyBlocks\Outbox\Schema\TableLayout;
use TinyBlocks\Outbox\Serialization\PayloadSerializerReflection;
use TinyBlocks\Outbox\Serialization\PayloadSerializers;

$repository = new DoctrineOutboxRepository(
connection: $connection,
serializers: PayloadSerializers::createFrom(elements: [new PayloadSerializerReflection()]),
translators: IntegrationEventTranslators::createFrom(elements: [new OrderPlacedTranslator()]),
tableLayout: TableLayout::builder()
->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,
Expand Down
16 changes: 14 additions & 2 deletions src/DoctrineOutboxRepository.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

namespace TinyBlocks\Outbox;

use Closure;
use Doctrine\DBAL\Connection;
use TinyBlocks\BuildingBlocks\Event\EventRecord;
use TinyBlocks\BuildingBlocks\Event\EventRecords;
Expand All @@ -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
);
}

Expand Down
61 changes: 35 additions & 26 deletions src/Internal/OutboxInsert.php
Original file line number Diff line number Diff line change
Expand Up @@ -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 = <<<SQL
INSERT INTO %s (%s, %s, %s, %s, %s, %s, %s, %s)
VALUES (:id, :aggregateId, :aggregateType, :eventType, :revision,
:aggregateVersion, :payload, :occurredAt)
SQL;

$columns = $tableLayout->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
);
}
}
23 changes: 20 additions & 3 deletions src/Internal/OutboxWriter.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

namespace TinyBlocks\Outbox\Internal;

use Closure;
use Doctrine\DBAL\Connection;
use Doctrine\DBAL\Exception\UniqueConstraintViolationException;
use TinyBlocks\BuildingBlocks\Event\EventRecord;
Expand All @@ -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);
Expand All @@ -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);
Expand All @@ -77,4 +89,9 @@ public function write(EventRecord $eventRecord): void
);
}
}

private function correlationId(): string
{
return is_null($this->correlationId) ? '' : ($this->correlationId)();
}
}
11 changes: 8 additions & 3 deletions src/Schema/Columns.php
Original file line number Diff line number Diff line change
Expand Up @@ -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
) {
}

Expand All @@ -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(
Expand All @@ -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,
Expand All @@ -53,7 +57,8 @@ public static function from(
occurredAt: $occurredAt,
aggregateId: $aggregateId,
aggregateType: $aggregateType,
aggregateVersion: $aggregateVersion
aggregateVersion: $aggregateVersion,
correlationId: $correlationId
);
}

Expand Down
19 changes: 18 additions & 1 deletion src/Schema/ColumnsBuilder.php
Original file line number Diff line number Diff line change
Expand Up @@ -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()
{
Expand Down Expand Up @@ -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
);
}

Expand Down Expand Up @@ -167,4 +169,19 @@ public function withAggregateVersion(string $name): ColumnsBuilder
$this->aggregateVersion = $name;
return $this;
}

/**
* Enables the correlation id column under the given name.
*
* <p>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.</p>
*
* @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;
}
}
Loading
Loading