ianfoxdev / inbox
Consumer-side deduplication for PHP: the message key is written in the same database transaction as the change it causes, so a redelivered message is skipped. Works with outbox's CloudEvents headers, Symfony Messenger and Laravel Queue.
Requires (Dev)
- friendsofphp/php-cs-fixer: ^3.64
- illuminate/database: ^11.0 || ^12.0 || ^13.0
- phpstan/phpstan: ^2.1
- phpunit/phpunit: ^12.5 || ^13.0
- symfony/messenger: ^7.3 || ^8.0
Suggests
- ext-pdo_mysql: To keep the inbox in MySQL
- ext-pdo_pgsql: To keep the inbox in PostgreSQL
- illuminate/queue: For the Laravel job middleware
- symfony/messenger: For the Messenger middleware
Provides
None
Conflicts
None
Replaces
None
This package is auto-updated.
Last update: 2026-10-06 08:12:55 UTC
README
Consumer-side deduplication for PHP. The message key is written in the same database transaction as the change the message causes, so a redelivered message is skipped and a failed one is processed again. Reads the CloudEvents ids that outbox puts on every Kafka and RabbitMQ message, and plugs into Symfony Messenger and Laravel Queue.
Status: v0.1. Until 1.0 a minor version may change the API; such changes are marked BREAKING in the CHANGELOG. The table and the meaning of a stored key do not change in a minor version.
Brokers deliver at least once. Kafka hands a record out again when a consumer dies before committing its offset, or after a rebalance; RabbitMQ after a missed ack. The common fixes leave a gap:
- "Seen it?" before the handler, "mark it" after. Two instances that get the same message during a rebalance both see "not yet" and both credit the wallet. A crash between the credit and the mark credits it again on the next delivery.
- Deduplication on the way into the queue (Symfony's
DeduplicateStamp, Laravel'sShouldBeUnique) stops a job from being queued twice. It does not stop a redelivered one from running twice. - A key in Redis commits separately from the database change it protects.
Install
composer require ianfoxdev/inbox
PHP 8.3 or later, PostgreSQL or MySQL 8. Create the table from schema/postgresql.sql or schema/mysql.sql. No dependencies; Symfony Messenger and Laravel are optional.
Example
use IanFoxDev\Inbox\{Inbox, MessageKey}; use IanFoxDev\Inbox\Storage\PdoStore; $inbox = new Inbox(new PdoStore($pdo)); // the connection your handler writes with $outcome = $inbox->handle( 'billing.cashback', // the consumer MessageKey::cloudEvents($record->headers), // "/orders 0199b1d4-..." from outbox function () use ($pdo, $event): void { $pdo->prepare('UPDATE wallets SET cashback = cashback + ? WHERE user_id = ?') ->execute([$event['cashback'], $event['user_id']]); }, ); $outcome->duplicate; // true on every delivery after the first
examples/redelivery replays one event the way Kafka does:
1. OrderPaid arrives credited 1.90 wallet 1.90
2. Consumer died before committing the offset: again duplicate, skipped wallet 1.90
3. Rebalance gives the partition to another instance duplicate, skipped wallet 1.90
4. Next event, the provider times out failed: rolled back wallet 1.90
5. The same event is redelivered credited 1.90 wallet 3.80
6. Another consumer of the first event receipt sent wallet 3.80
What it guarantees
- The key and the change commit together. The key is inserted first, then the handler runs, then one commit. A handler that throws rolls its key back with its writes, so the next delivery processes the message; nothing is lost and nothing is applied twice.
- A repeat is a unique-index conflict. No "select, then insert" window. The insert does
not end the transaction (
ON CONFLICT DO NOTHINGin PostgreSQL, a caught duplicate-key error in MySQL). - Concurrent duplicates wait for each other. When two instances get the same message, the second blocks on the first's index entry. If the first commits, the second skips; if it rolls back, the second processes. Tested with real processes on PostgreSQL 17 and MySQL 8.4, including eight consumers racing over the same 60 messages: each credited exactly once.
- Each consumer processes a message once. The key is the consumer name plus the message id, so billing and notifications both handle the same event.
- It joins your transaction. Give it your PDO connection; it uses an open transaction or opens one at READ COMMITTED.
When the message comes again
Mode |
The handler | handle() returns |
|---|---|---|
Skip (default) |
is not called | Outcome with duplicate: true |
Fail |
is not called | throws Duplicate |
ReturnStored |
is not called | the first result, stored as JSON; for webhooks and requests that must get the same answer |
Message keys
| Source | Call |
|---|---|
Kafka record from outbox (ce_id, ce_source) |
MessageKey::cloudEvents($headers) |
RabbitMQ message from outbox (message_id, cloudEvents_id) |
MessageKey::cloudEvents($headersAndProperties) |
CloudEvents over HTTP (ce-id) or structured JSON |
MessageKey::cloudEvents($headers, $body) |
| a header of your own | MessageKey::header($headers, 'X-Webhook-Id') |
| a payload field | MessageKey::field($payload, 'data.transaction_id') |
CloudEvents guarantees that source plus id is unique, not id alone, so the key is
"source id" when the source is there. A key over 255 bytes is replaced by its SHA-256. A
message without a key throws NoKey: processing it without deduplication would be silent.
Symfony Messenger
framework: messenger: buses: messenger.bus.default: middleware: - doctrine_transaction - IanFoxDev\Inbox\Messenger\InboxMiddleware
Put it after doctrine_transaction, so the key joins the transaction your handlers write
in, and build its PdoStore on $connection->getNativeConnection(). Only received messages
are deduplicated; the consumer is the transport name. The key comes from an InboxKeyStamp
your transport serializer adds, or from the message if it implements HasInboxKey. A
skipped message gets a DuplicateStamp.
Laravel
public function middleware(): array { return [new InboxJobMiddleware(app(Inbox::class), DB::connection())]; }
The job's work and the key commit in one DB::transaction(); a transaction inside the job
becomes a savepoint. The job implements HasInboxKey, or you pass a closure that returns the
key.
Kafka and the offset
Commit the database transaction first, the offset after it. If the process dies in between, the record comes again and is skipped. The other order loses the message. docs/kafka.md has the consumer loop.
Keeping the table small
vendor/bin/inbox purge --older-than=30 --dsn='pgsql:host=db;dbname=app' --user=app
Keep keys longer than a message can come back: Kafka's retention, the replay window of your dead-letter queue. A purged key makes the next delivery of its message new again.
Counters (a Listener) counts processed messages and duplicates per consumer for your
metrics; a rising duplicate count is a consumer that crashes or rebalances too often.
More: docs/usage.md. Why it is built this way: docs/adr/0001-the-key-is-written-with-the-change.md.
Not yet
Handlers whose only effect is outside the database (an API call) get no exactly-once from this: pass the message key on as the API's idempotency key. Redis storage, a ready Kafka consumer, a Symfony bundle with configuration, dead-lettering (the broker's job). Open an issue if you need one of them first.