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
24 changes: 24 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -350,6 +350,30 @@ $store = new SqlEventStore($pdo);
> MySQL workers use `SELECT … FOR UPDATE SKIP LOCKED` on the `event_outbox` table for safe
> concurrent processing.

#### Committing the event with the business change

`add()` is meant to run inside the caller's transaction, so the event and the change it describes commit
or roll back together. A fixed `PDO` cannot do that when the application opens its transaction on a
different connection: the two are separate sessions, so the event commits immediately and a worker can
pick it up before the business row exists — or keep it after that row rolls back.

Pass a `Closure(): PDO` instead and the store resolves the connection on every call, joining whatever
transaction is currently open:

```php
$store = new SqlEventStore(
fn(): PDO => $currentTransactionConnection() ?? $platformConnection,
);
```

If the outbox tables live in a different schema from the connection doing the insert, name it — every
statement is then qualified (`platform.event_outbox`). All schemas must be on the same server for a
single session to write across them.

```php
$store = new SqlEventStore($connection, schema: 'platform');
```

If you wire up a `RedeliveryStore` (see [Automatic retry & failure tracking](#automatic-retry--failure-tracking)),
a third table `event_outbox_redelivery` holds per-listener retry state — install it from the same `migrations/` directory.

Expand Down
115 changes: 74 additions & 41 deletions src/Infrastructure/SqlEventStore.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

use Carbon\CarbonImmutable;
use Carbon\CarbonInterval;
use Closure;
use Override;
use PDO;
use Throwable;
Expand All @@ -13,14 +14,28 @@

readonly class SqlEventStore implements EventStore
{
public function __construct(private PDO $connection) {}
/**
* A closure is resolved on every call so add() can join whichever transaction the caller
* has open. Handing over a fixed PDO writes the event on a second session, which commits
* the event independently of the business change it describes.
*
* @param PDO|Closure(): PDO $connection
* @param string|null $schema Qualifies the outbox tables when they live outside the
* caller's default schema.
*/
public function __construct(
private PDO|Closure $connection,
private ?string $schema = null,
) {}

#[Override]
public function add(RawEvent $event): void
{
$stmt = $this->connection->prepare(
$connection = $this->connection();

$stmt = $connection->prepare(
<<<SQL
INSERT INTO event_outbox (id, name, status, payload, created_at, publish_at)
INSERT INTO {$this->table('event_outbox')} (id, name, status, payload, created_at, publish_at)
VALUES (:id, :name, :status, :payload, :created_at, :publish_at)
SQL,
);
Expand All @@ -34,19 +49,21 @@ public function add(RawEvent $event): void
'publish_at' => $event->publishAt->format('Y-m-d H:i:s.u'),
]);

$this->insertStatusAudit($event->id, $event->status->value);
$this->insertStatusAudit($connection, $event->id, $event->status->value);
}

#[Override]
public function next(): ?RawEvent
{
$this->connection->beginTransaction();
$connection = $this->connection();

$connection->beginTransaction();

try {
$row = $this->fetchNextPendingRow();
$row = $this->fetchNextPendingRow($connection);

if ($row === null) {
$this->connection->commit();
$connection->commit();
return null;
}

Expand All @@ -62,17 +79,17 @@ public function next(): ?RawEvent
publishAt: new CarbonImmutable($row['publish_at']),
)->claim();

$this->connection->prepare(
"UPDATE event_outbox SET status = :status WHERE id = :id",
$connection->prepare(
"UPDATE {$this->table('event_outbox')} SET status = :status WHERE id = :id",
)->execute(['status' => $event->status->value, 'id' => $event->id]);

$this->insertStatusAudit($event->id, $event->status->value);
$this->insertStatusAudit($connection, $event->id, $event->status->value);

$this->connection->commit();
$connection->commit();

return $event;
} catch (Throwable $e) {
$this->rollBackIfActive();
$this->rollBackIfActive($connection);

throw $e;
}
Expand All @@ -82,22 +99,24 @@ public function next(): ?RawEvent
#[Override]
public function markProcessed(RawEvent $event): void
{
$this->connection->beginTransaction();
$connection = $this->connection();

$connection->beginTransaction();

try {
$this->connection->prepare(
$connection->prepare(
<<<SQL
UPDATE event_outbox
UPDATE {$this->table('event_outbox')}
SET status = :status
WHERE id = :id AND status = 'processing'
SQL,
)->execute(['status' => $event->status->value, 'id' => $event->id]);

$this->insertStatusAudit($event->id, $event->status->value);
$this->insertStatusAudit($connection, $event->id, $event->status->value);

$this->connection->commit();
$connection->commit();
} catch (Throwable $e) {
$this->rollBackIfActive();
$this->rollBackIfActive($connection);

throw $e;
}
Expand All @@ -113,14 +132,16 @@ public function markProcessed(RawEvent $event): void
#[Override]
public function recoverStuckEvents(CarbonInterval $olderThan): int
{
$connection = $this->connection();

$thresholdAt = CarbonImmutable::now()->sub($olderThan)->format('Y-m-d H:i:s.u');

$stmt = $this->connection->prepare(
$stmt = $connection->prepare(
<<<SQL
SELECT e.id FROM event_outbox e
SELECT e.id FROM {$this->table('event_outbox')} e
WHERE e.status = 'processing'
AND NOT EXISTS (
SELECT 1 FROM event_outbox_status s
SELECT 1 FROM {$this->table('event_outbox_status')} s
WHERE s.event_id = e.id
AND s.status = 'processing'
AND s.created_at >= :threshold
Expand All @@ -136,11 +157,11 @@ public function recoverStuckEvents(CarbonInterval $olderThan): int
return 0;
}

$this->connection->beginTransaction();
$connection->beginTransaction();

try {
$updateStmt = $this->connection->prepare(
"UPDATE event_outbox SET status = 'pending' WHERE id = :id AND status = 'processing'",
$updateStmt = $connection->prepare(
"UPDATE {$this->table('event_outbox')} SET status = 'pending' WHERE id = :id AND status = 'processing'",
);

$recovered = 0;
Expand All @@ -151,31 +172,43 @@ public function recoverStuckEvents(CarbonInterval $olderThan): int
continue;
}

$this->insertStatusAudit($id, RawEventStatus::pending->value, 'Recovered from stuck processing state');
$this->insertStatusAudit($connection, $id, RawEventStatus::pending->value, 'Recovered from stuck processing state');
$recovered++;
}

$this->connection->commit();
$connection->commit();

return $recovered;
} catch (Throwable $e) {
$this->rollBackIfActive();
$this->rollBackIfActive($connection);

throw $e;
}
}

private function connection(): PDO
{
return $this->connection instanceof Closure
? ($this->connection)()
: $this->connection;
}

private function table(string $name): string
{
return $this->schema === null ? $name : "{$this->schema}.{$name}";
}

/**
* @return array{id: string, name: string, status: string, payload: string, created_at: string, publish_at: string}|null
*/
private function fetchNextPendingRow(): ?array
private function fetchNextPendingRow(PDO $connection): ?array
{
$lockClause = $this->lockingClause();
$lockClause = $this->lockingClause($connection);

$stmt = $this->connection->prepare(
$stmt = $connection->prepare(
<<<SQL
SELECT id, name, status, payload, created_at, publish_at
FROM event_outbox WHERE
FROM {$this->table('event_outbox')} WHERE
status = 'pending' AND
publish_at <= :now
ORDER BY publish_at
Expand All @@ -193,11 +226,11 @@ private function fetchNextPendingRow(): ?array
return $row === false ? null : $row;
}

private function insertStatusAudit(string $eventId, string $status, ?string $errorMessage = null): void
private function insertStatusAudit(PDO $connection, string $eventId, string $status, ?string $errorMessage = null): void
{
$stmt = $this->connection->prepare(
$stmt = $connection->prepare(
<<<SQL
INSERT INTO event_outbox_status (event_id, status, error_message, created_at)
INSERT INTO {$this->table('event_outbox_status')} (event_id, status, error_message, created_at)
VALUES (:event_id, :status, :error_message, :created_at)
SQL,
);
Expand All @@ -211,24 +244,24 @@ private function insertStatusAudit(string $eventId, string $status, ?string $err
}

/** Guards against the implicit-commit case where a DDL statement closed the transaction before the failure. */
private function rollBackIfActive(): void
private function rollBackIfActive(PDO $connection): void
{
if ($this->connection->inTransaction()) {
$this->connection->rollBack();
if ($connection->inTransaction()) {
$connection->rollBack();
}
}

private function lockingClause(): string
private function lockingClause(PDO $connection): string
{
return match ($this->driverName()) {
return match ($this->driverName($connection)) {
'mysql' => 'FOR UPDATE SKIP LOCKED',
default => '',
};
}

private function driverName(): string
private function driverName(PDO $connection): string
{
/** @var string */
return $this->connection->getAttribute(PDO::ATTR_DRIVER_NAME);
return $connection->getAttribute(PDO::ATTR_DRIVER_NAME);
}
}
36 changes: 36 additions & 0 deletions tests/Feature/SqlEventStoreTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -141,6 +141,42 @@ public function test_schema_create_is_idempotent(): void
self::assertTrue(true);
}

public function test_add_resolves_the_connection_on_every_call(): void
{
$resolved = 0;
$store = new SqlEventStore(function () use (&$resolved): PDO {
$resolved++;

return $this->pdo;
});

$store->add(self::createEvent('order.placed'));
$store->add(self::createEvent('order.shipped'));

self::assertSame(2, $resolved);
}

public function test_add_and_next_use_the_schema_qualified_tables(): void
{
$pdo = new PDO('sqlite::memory:');
$pdo->setAttribute(PDO::ATTR_ERRMODE, PDO::ERRMODE_EXCEPTION);
$pdo->exec("ATTACH DATABASE ':memory:' AS platform");
$pdo->exec('CREATE TABLE platform.event_outbox (
id TEXT NOT NULL PRIMARY KEY, name TEXT NOT NULL, status TEXT NOT NULL,
payload TEXT NOT NULL, created_at TEXT NOT NULL, publish_at TEXT NOT NULL)');
$pdo->exec('CREATE TABLE platform.event_outbox_status (
event_id TEXT NOT NULL, status TEXT NOT NULL, error_message TEXT, created_at TEXT NOT NULL)');

$store = new SqlEventStore($pdo, schema: 'platform');
$event = self::createEvent('order.placed', ['order_id' => 99]);
$store->add($event);

$retrieved = $store->next();

self::assertNotNull($retrieved);
self::assertSame($event->id, $retrieved->id);
}

/**
* @param array<string, mixed> $payload
*/
Expand Down
Loading