Search by

ianfoxdev / inbox

IanFoxDev

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.

Package info

github.com/IanFoxDev/inbox

pkg:composer/ianfoxdev/inbox

Statistics

Installs: 2

Dependents: 0

Suggesters: 0

Stars: 0

Open Issues: 0

v0.1.0 2026-10-06 08:09 UTC

This package is auto-updated.

Last update: 2026-10-06 08:12:55 UTC


README

php

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's ShouldBeUnique) 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 NOTHING in 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.

License

MIT