From 89c3b3466ca148f9bf0ca48fef6ed5618c0a1c78 Mon Sep 17 00:00:00 2001 From: Kocsis Csongor Date: Tue, 18 Aug 2026 10:59:38 +0200 Subject: [PATCH] fix: resolve the event store connection at call time - SqlEventStore accepts a Closure(): PDO so add() joins whichever transaction the caller has open, instead of committing the event on a second session independently of the business change it describes - Add an optional schema so the outbox tables can be qualified when they live outside the calling connection's default schema - Resolve the connection once per method so a closure cannot hand back a different handle mid-transaction --- README.md | 24 ++++++ src/Infrastructure/SqlEventStore.php | 115 +++++++++++++++++---------- tests/Feature/SqlEventStoreTest.php | 36 +++++++++ 3 files changed, 134 insertions(+), 41 deletions(-) diff --git a/README.md b/README.md index 65bac5e..2b45174 100644 --- a/README.md +++ b/README.md @@ -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. diff --git a/src/Infrastructure/SqlEventStore.php b/src/Infrastructure/SqlEventStore.php index 7661372..eb576f2 100644 --- a/src/Infrastructure/SqlEventStore.php +++ b/src/Infrastructure/SqlEventStore.php @@ -4,6 +4,7 @@ use Carbon\CarbonImmutable; use Carbon\CarbonInterval; +use Closure; use Override; use PDO; use Throwable; @@ -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( <<table('event_outbox')} (id, name, status, payload, created_at, publish_at) VALUES (:id, :name, :status, :payload, :created_at, :publish_at) SQL, ); @@ -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; } @@ -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; } @@ -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( <<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; } @@ -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( <<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 @@ -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; @@ -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( <<table('event_outbox')} WHERE status = 'pending' AND publish_at <= :now ORDER BY publish_at @@ -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( <<table('event_outbox_status')} (event_id, status, error_message, created_at) VALUES (:event_id, :status, :error_message, :created_at) SQL, ); @@ -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); } } diff --git a/tests/Feature/SqlEventStoreTest.php b/tests/Feature/SqlEventStoreTest.php index 71d54b0..93fa3b7 100644 --- a/tests/Feature/SqlEventStoreTest.php +++ b/tests/Feature/SqlEventStoreTest.php @@ -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 $payload */