# Implementing Webhook Idempotency with a Transactional Inbox: The Lost Events and Double Processing That Event-ID De-duplication Alone Won't Stop

> How to make webhook handling idempotent with a transactional inbox that commits the receipt record and the business update in one PostgreSQL transaction. Why marking an event processed first loses it forever and marking it afterwards processes it twice, how to handle concurrent duplicates, out-of-order delivery and separate Event objects for the same fact, and how this compares with DynamoDB de-duplication — with Python code tested against PostgreSQL 18.

- Published: 2026-10-04
- Author: 友田 陽大
- Tags: 決済, 信頼性, PostgreSQL, Python, アーキテクチャ設計
- URL: https://tomodahinata.com/en/blog/webhook-idempotency-transactional-inbox-postgresql-guide
- Category: Reliability, async & real-time
- Pillar guide: https://tomodahinata.com/en/blog/transactional-outbox-pattern-reliable-event-publishing-guide

## Key points

- Design webhook handlers on the assumption that the same event arrives more than once. Stripe documents that it retries for up to three days in live mode, that an endpoint may receive the same event more than once, and that delivery order is not guaranteed
- Committing 'processed' in its own transaction before the business update means that if the update fails, the retry is treated as a duplicate and the event is lost forever (the at-most-once failure)
- Committing the business update and only then recording the event means a crash in between makes the retry process it twice (the at-least-once failure)
- The fix is a transactional inbox: the receipt record (INSERT ... ON CONFLICT DO NOTHING) and the business update in the same transaction. On failure both disappear and the retry reprocesses correctly; concurrent duplicates are serialized by waiting on the unique index
- Event-ID de-duplication isn't enough on its own. Separate Event objects for the same business fact and out-of-order delivery are handled by making the handler itself idempotent, with state transitions and unique constraints on business keys

---

**The short answer.** Webhook idempotency isn't finished when you put a UNIQUE constraint on the event ID. You must **commit the receipt record and the business update in the same database transaction.** Record "processed" first and a failed update **loses the event forever**; record it afterwards and a crash **processes it twice**. This receiving-side design is called a **transactional inbox**. On top of that, make **the handler itself idempotent**, because the same business fact can arrive as a different event and events can arrive out of order.

This article reproduces both failure modes and implements a transactional inbox with PostgreSQL's `INSERT ... ON CONFLICT DO NOTHING`, in both synchronous and asynchronous forms. The code was **verified on PostgreSQL 18.6 / psycopg 3.3 / stripe-python 16 with 18 tests, including concurrent duplicate deliveries and worker crashes** (October 2026). It also compares the approach with DynamoDB de-duplication without declaring either one "wrong".

> **Ground rules**: Stripe's behaviour is taken from its [official documentation](https://docs.stripe.com/webhooks) as of October 2026. Check other providers' retry and ordering behaviour in their own documentation; I take care not to present Stripe's behaviour as true of webhooks in general.

---

## 1. Premise: a webhook isn't guaranteed to arrive exactly once

Stripe's documentation states the following about delivery (as of October 2026):

- **Retries**: in live mode, Stripe attempts to deliver an event **for up to three days with exponential backoff**.
- **Duplicates**: webhook endpoints **might occasionally receive the same event more than once**. Stripe recommends logging the event IDs you've processed and not processing already-logged events.
- **Duplicates as separate objects**: in some cases **two separate Event objects are generated and sent**. Stripe suggests identifying these by the ID in `data.object` together with `event.type`.
- **Ordering**: Stripe **does not guarantee delivery in the order events were generated.**

Retries themselves are unavoidable. Stripe treats anything other than a 2xx as a failure. If your handler processed the event correctly but the response was lost to a timeout or a network fault, the same event arrives again. The receiver must be built for **the same event arriving several times, possibly at the same moment.**

---

## 2. Two failure modes

Recording event IDs to reject duplicates is the right idea, but if you get **the order and the transaction boundaries of the record and the business update** wrong, it fails in two opposite ways.

The example is `invoice.paid`: book the payment in a ledger (`ledger_entries`) and mark the invoice `paid`.

### Case A: commit "processed" first → the event is lost forever

```python
# ❌ Broken: the receipt record commits first, in its own transaction
def receive_mark_first_broken(conn, *, endpoint, event):
    with conn.transaction():
        inserted = conn.execute(
            "INSERT INTO webhook_inbox (provider, endpoint, event_id, event_type, payload) "
            "VALUES ('stripe', %s, %s, %s, %s) ON CONFLICT DO NOTHING RETURNING event_id",
            (endpoint, event.id, event.type, Jsonb(event.payload)),
        ).fetchone()
    if inserted is None:
        return Outcome.DUPLICATE        # logged = treated as processed, return 200
    with conn.transaction():
        apply_event(conn, event)        # ← if this fails...
    return Outcome.PROCESSED
```

```text
1st delivery:  receipt record COMMITs → business update hits a DB error → ROLLBACK → return 500
2nd delivery:  receipt INSERT collides → judged a "duplicate" → return 200
→ the payment is never booked, and Stripe stops retrying
```

This is **the at-most-once failure**: no duplicates, but events go missing. In my tests, failing the first business update made the second delivery return `DUPLICATE`, and the ledger stayed empty. For payments this is as serious as a double charge.

### Case B: commit the business update, then record → double processing

```python
# ❌ Broken: the business commit and the receipt record are separate transactions
def receive_mark_after_broken(conn, *, endpoint, event):
    seen = conn.execute(
        "SELECT 1 FROM webhook_inbox WHERE provider='stripe' AND endpoint=%s AND event_id=%s",
        (endpoint, event.id),
    ).fetchone()
    if seen:
        return Outcome.DUPLICATE
    with conn.transaction():
        apply_event(conn, event)        # business update COMMITs
    # ← the process dies here (deploy, OOM, timeout)
    with conn.transaction():
        conn.execute("INSERT INTO webhook_inbox ...")
    return Outcome.PROCESSED
```

```text
1st delivery:  business update COMMITs (payment booked) → crash before the receipt record → no 2xx
2nd delivery:  no receipt record → judged unprocessed → business update COMMITs again
→ the payment is booked twice
```

This is **the at-least-once failure**. In my tests, stopping just before the receipt record and then redelivering left two ledger rows. Because it does `SELECT` then `INSERT`, two simultaneous deliveries both see "unprocessed" and double-process too (TOCTOU).

**Both failures have the same cause**: **two writes — the receipt record and the business update — committed separately.** Whichever goes first, a crash in between leaves only one of them.

---

## 3. The fix: one transaction for the receipt record and the business update (transactional inbox)

Commit the receipt `INSERT` and the business update **in the same transaction.**

```text
BEGIN
  INSERT INTO webhook_inbox (event_id ...) ON CONFLICT DO NOTHING RETURNING ...
    → no row returned means duplicate: COMMIT without doing anything, return 200
  UPDATE invoices / INSERT ledger_entries (business update)
  INSERT INTO outbox (if something downstream must be notified)
  UPDATE webhook_inbox SET processed_at = now()
COMMIT → 200
```

- If the business update fails, **the receipt record rolls back with it.** Nothing is left behind, so the retry is reprocessed as unprocessed (Case A can't happen).
- If COMMIT succeeds, **both the business update and the receipt record are durable.** Even if the process dies before returning 2xx, the retry is rejected as a duplicate (Case B can't happen).

That's a transactional inbox. It is the mirror image of the [transactional outbox](/blog/transactional-outbox-pattern-reliable-event-publishing-guide), which puts "the business update and the event to send" in one transaction on the sending side.

| | Transactional outbox (sending side) | Transactional inbox (receiving side) |
| --- | --- | --- |
| What goes in the same transaction | Business update + event to send | Received event ID + business update |
| Failure prevented | Business updated but no event / event sent but no business change | Received but never processed / processed twice |
| Remaining property | Publishing is at-least-once (duplicates possible) | Absorbs duplicate deliveries so the business effect converges to one |

Note that an inbox has this atomicity **only when there is a single database.** If the receipt record and the business data live in different stores, you are back to two writes (see the DynamoDB comparison in section 7).

---

## 4. Schema

```sql
-- Receipt ledger. The primary key lets the DB decide "accept the same delivery once"
CREATE TABLE webhook_inbox (
    provider     text        NOT NULL,              -- e.g. 'stripe'; separates ID collisions across providers
    endpoint     text        NOT NULL,              -- separates endpoints that receive the same event
    event_id     text        NOT NULL,              -- evt_... for Stripe
    event_type   text        NOT NULL,
    payload      jsonb       NOT NULL,              -- the signature-verified raw JSON
    received_at  timestamptz NOT NULL DEFAULT now(),
    available_at timestamptz NOT NULL DEFAULT now(), -- retry time for async processing
    attempts     integer     NOT NULL DEFAULT 0,
    last_error   text,
    processed_at timestamptz,                       -- NULL = not yet processed
    PRIMARY KEY (provider, endpoint, event_id)
);
CREATE INDEX webhook_inbox_unprocessed
    ON webhook_inbox (available_at) WHERE processed_at IS NULL;

CREATE TABLE invoices (
    id          text        PRIMARY KEY,           -- in_...
    customer_id text        NOT NULL,
    status      text        NOT NULL CHECK (status IN ('draft', 'open', 'paid', 'void', 'uncollectible')),
    amount_paid bigint      NOT NULL DEFAULT 0,
    paid_at     timestamptz,
    updated_at  timestamptz NOT NULL DEFAULT now()
);

-- Downstream notifications (same shape as in the transactional outbox article)
CREATE TABLE outbox (
    id             bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
    aggregate_type text        NOT NULL,
    aggregate_id   text        NOT NULL,
    event_type     text        NOT NULL,
    payload        jsonb       NOT NULL,
    created_at     timestamptz NOT NULL DEFAULT now(),
    published_at   timestamptz
);

-- Business-level idempotency: never book the same invoice's payment twice (even via a different Event object)
CREATE TABLE ledger_entries (
    id         bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
    invoice_id text   NOT NULL,
    kind       text   NOT NULL,
    amount     bigint NOT NULL,
    UNIQUE (invoice_id, kind)
);
```

Three design decisions:

**Why `endpoint` is part of the primary key**: if you register several endpoints on one Stripe account (say, billing and analytics), the same event goes to both. Both should process it, so a primary key on the event ID alone would let one endpoint's processing mark the other's as a "duplicate". In my tests, separate `endpoint` values made both return `PROCESSED`. `provider` is included so that IDs from different services that happen to match never mix.

**Why store `payload`**: so that a worker can process it asynchronously (section 6), and so you can later check "what actually arrived" during an investigation. If events contain personal data, decide the retention period (section 9) and who may read it.

**Why `UNIQUE (invoice_id, kind)` on the ledger**: the receipt ledger only prevents duplicates of the same event ID. When Stripe sends the same fact as a separate Event object, the ID differs and passes straight through the inbox. A unique constraint on the business key (the invoice ID) is the second line of defence (section 5.2).

---

## 5. Synchronous version: receive and process on the spot

If the business update is short and stays in the database, this is the simplest and strongest shape.

> **Assumption in the code**: the psycopg code in this article assumes connections opened with `psycopg.connect(dsn, autocommit=True)`, and expresses transaction boundaries only with `with conn.transaction():`. On a connection without autocommit, running SQL outside `transaction()` starts an implicit transaction, later `transaction()` blocks become SAVEPOINTs, and the commit timing described here no longer holds.

### 5.1 Signature verification: on the raw body

```python
import json
from dataclasses import dataclass
from typing import Any

import stripe


@dataclass(frozen=True)
class VerifiedEvent:
    """An event that passed signature verification. Nothing else enters business logic."""
    id: str
    type: str
    payload: dict[str, Any]


class InvalidWebhook(Exception):
    pass


def verify_stripe_event(raw_body: bytes, sig_header: str | None, secret: str) -> VerifiedEvent:
    """Verify the signature on the bytes before any framework turns the body into a dict."""
    try:
        stripe.Webhook.construct_event(raw_body, sig_header, secret)
    except (ValueError, stripe.SignatureVerificationError) as exc:
        raise InvalidWebhook(str(exc)) from exc
    payload = json.loads(raw_body)
    return VerifiedEvent(id=payload["id"], type=payload["type"], payload=payload)
```

Stripe computes the signature over **the exact bytes received.** Parse the JSON and re-serialize it, and whitespace or key order changes and the signature no longer matches. In my tests, a body with identical content but re-serialized with different separators raised `InvalidWebhook`. Signature verification is the only boundary that decides whether the event's contents can be trusted; don't write anything unverified to the database.

By default Stripe's libraries allow **5 minutes** between the signature timestamp and the current time. The documentation warns against a tolerance of `0`, which disables the recency check entirely. Keep server clocks in sync with NTP.

### 5.2 Make the handler itself idempotent and order-independent

```python
def _apply_invoice_created(conn: psycopg.Connection, inv: dict[str, Any]) -> None:
    # Creation only inserts if absent. Arriving late after 'paid' never rolls the state back
    conn.execute(
        """
        INSERT INTO invoices (id, customer_id, status)
        VALUES (%(id)s, %(customer)s, 'open')
        ON CONFLICT (id) DO NOTHING
        """,
        {"id": inv["id"], "customer": inv["customer"]},
    )


def _apply_invoice_paid(conn: psycopg.Connection, inv: dict[str, Any]) -> None:
    # 'paid' can arrive before 'created' → upsert. Never overwrite a voided invoice
    conn.execute(
        """
        INSERT INTO invoices (id, customer_id, status, amount_paid, paid_at)
        VALUES (%(id)s, %(customer)s, 'paid', %(amount_paid)s, now())
        ON CONFLICT (id) DO UPDATE
           SET status      = 'paid',
               amount_paid = EXCLUDED.amount_paid,
               paid_at     = COALESCE(invoices.paid_at, EXCLUDED.paid_at),
               updated_at  = now()
         WHERE invoices.status NOT IN ('paid', 'void')
        """,
        {"id": inv["id"], "customer": inv["customer"], "amount_paid": inv["amount_paid"]},
    )
    row = conn.execute("SELECT status FROM invoices WHERE id = %s", (inv["id"],)).fetchone()
    if row is None or row[0] != "paid":
        # A payment on a voided invoice is an anomaly: don't book it; reconciliation sends it to a human
        return
    booked = conn.execute(
        """
        INSERT INTO ledger_entries (invoice_id, kind, amount)
        VALUES (%(id)s, 'payment', %(amount_paid)s)
        ON CONFLICT (invoice_id, kind) DO NOTHING
        RETURNING id
        """,
        {"id": inv["id"], "amount_paid": inv["amount_paid"]},
    ).fetchone()
    if booked is not None:
        # Notify downstream (receipt e-mail, etc.) via the outbox, in the same transaction as the business update
        conn.execute(
            """
            INSERT INTO outbox (aggregate_type, aggregate_id, event_type, payload)
            VALUES ('Invoice', %(id)s, 'InvoicePaymentBooked', %(payload)s)
            """,
            {"id": inv["id"], "payload": Jsonb({"invoice_id": inv["id"], "amount": inv["amount_paid"]})},
        )


HANDLERS = {
    "invoice.created": _apply_invoice_created,
    "invoice.paid": _apply_invoice_paid,
}
```

Four design choices are doing the work:

1. **Written as state transitions**: `invoice.created` only creates the row if it's absent and never changes existing state. `invoice.paid` advances to `paid` only from states other than `paid` / `void`. So even if `invoice.paid` arrives before `invoice.created`, the final state is `paid` (verified).
2. **No ordering by `created`**: Stripe's documentation explicitly says that because `created` is recorded in seconds and distinct events can share a timestamp, **you shouldn't use `created` to determine event order or whether an event was already processed** (as of October 2026). Decisions that genuinely need the latest state (such as a subscription's current status) are made by fetching the object from the Stripe API; the documentation suggests the same, using information from an event that arrives first to retrieve the related objects.
3. **A unique constraint on the business key makes the effect happen once**: thanks to `UNIQUE (invoice_id, kind)` on `ledger_entries`, the same invoice's payment arriving under two different Event objects (`evt_A` and `evt_B`) is booked once (verified).
4. **Book only once the invoice is actually `paid`**: if `invoice.paid` arrives for an invoice that is already `void`, the upsert leaves its state unchanged. Booking it anyway would record "a payment on an invalid invoice", so the handler re-reads the status and books nothing unless it is `paid`. That combination is an anomaly, so reconciliation routes it to a human (verified).

### 5.3 The receive function

```python
from enum import Enum


class Outcome(Enum):
    PROCESSED = "processed"
    DUPLICATE = "duplicate"
    IGNORED = "ignored"


def apply_event(conn: psycopg.Connection, event: VerifiedEvent) -> bool:
    handler = HANDLERS.get(event.type)
    if handler is None:
        return False
    handler(conn, event.payload["data"]["object"])
    return True


def receive_sync(
    conn: psycopg.Connection, *, provider: str, endpoint: str, event: VerifiedEvent
) -> Outcome:
    """Commit the receipt record and the business update in one transaction."""
    with conn.transaction():
        inserted = conn.execute(
            """
            INSERT INTO webhook_inbox (provider, endpoint, event_id, event_type, payload)
            VALUES (%s, %s, %s, %s, %s)
            ON CONFLICT (provider, endpoint, event_id) DO NOTHING
            RETURNING event_id
            """,
            (provider, endpoint, event.id, event.type, Jsonb(event.payload)),
        ).fetchone()
        if inserted is None:
            return Outcome.DUPLICATE

        handled = apply_event(conn, event)
        conn.execute(
            """
            UPDATE webhook_inbox SET processed_at = now()
             WHERE provider = %s AND endpoint = %s AND event_id = %s
            """,
            (provider, endpoint, event.id),
        )
    return Outcome.PROCESSED if handled else Outcome.IGNORED
```

If an exception is raised inside the `conn.transaction()` block, psycopg rolls the transaction back and re-raises. The caller turns that into a 5xx and leaves it to Stripe's retries.

### 5.4 What happens when two copies arrive at once

How the same event arriving twice concurrently is handled is the heart of a transactional inbox.

```text
time  Delivery 1 (connection 1)                   Delivery 2 (connection 2)
t1    BEGIN; INSERT inbox(evt_1) → ok (uncommitted)
t2                                                BEGIN; INSERT inbox(evt_1) → waits on the unique index
t3    business update
t4    COMMIT
t5                                                → conflict confirmed; ON CONFLICT DO NOTHING → 0 rows
t6                                                COMMIT → DUPLICATE (200)
```

PostgreSQL's `INSERT ... ON CONFLICT` **waits while another uncommitted transaction is inserting the same key, until its outcome is known.** If the first commits, the second sees a conflict and does nothing; if the first rolls back, the second's insert succeeds and it takes over.

In my tests, sending the same event from eight connections at once produced "one `PROCESSED`, seven `DUPLICATE`, one ledger row". In a case where the first delivery's business update failed after 0.5 seconds, the second waited for the first to roll back and then returned `PROCESSED`, with one ledger row. **You don't write any locking in the application: the unique index serializes concurrent duplicates for you.**

### 5.5 The HTTP boundary (FastAPI example)

```python
import os

from fastapi import FastAPI, Request, Response
from psycopg_pool import ConnectionPool
from starlette.concurrency import run_in_threadpool

app = FastAPI()
pool = ConnectionPool(os.environ["DATABASE_URL"], kwargs={"autocommit": True}, open=True)
WEBHOOK_SECRET = os.environ["STRIPE_WEBHOOK_SECRET"]  # never hard-code this


def _handle(event: VerifiedEvent) -> Outcome:
    with pool.connection() as conn:
        return receive_sync(conn, provider="stripe", endpoint="billing", event=event)


@app.post("/webhooks/stripe")
async def stripe_webhook(request: Request) -> Response:
    raw_body = await request.body()  # the raw bytes, before any parsing
    try:
        event = verify_stripe_event(raw_body, request.headers.get("stripe-signature"), WEBHOOK_SECRET)
    except InvalidWebhook:
        return Response(status_code=400)
    await run_in_threadpool(_handle, event)  # an exception becomes a 500 and Stripe retries
    return Response(status_code=200)
```

`DUPLICATE` and `IGNORED` (event types you don't handle) also return 200: a non-2xx makes Stripe keep retrying. Conversely, whenever you couldn't process something, always return a non-2xx. "Return 200 anyway and log it" produces the same loss as Case A.

---

## 6. Asynchronous version: separate accepting from processing

Stripe's documentation recommends **returning a 2xx quickly, before any complex logic that could cause a timeout.** If the business update is heavy or calls external APIs, separate accepting the event from processing it.

```python
def receive_async(
    conn: psycopg.Connection, *, provider: str, endpoint: str, event: VerifiedEvent
) -> Outcome:
    """Commit only the receipt record and return 2xx. From here, retry responsibility moves from the provider to us."""
    with conn.transaction():
        inserted = conn.execute(
            """
            INSERT INTO webhook_inbox (provider, endpoint, event_id, event_type, payload)
            VALUES (%s, %s, %s, %s, %s)
            ON CONFLICT (provider, endpoint, event_id) DO NOTHING
            RETURNING event_id
            """,
            (provider, endpoint, event.id, event.type, Jsonb(event.payload)),
        ).fetchone()
    return Outcome.PROCESSED if inserted else Outcome.DUPLICATE
```

At first glance this looks like Case A, but there is a decisive difference: **the receipt record is left as "unprocessed" (`processed_at IS NULL`), not "processed".** In Case A the record was treated as proof of processing, so a processing failure became permanent loss. Here, the record is **proof that we have taken responsibility for processing it.** Once it commits and we return 2xx, responsibility for retrying moves from Stripe to our worker.

The worker splits "claiming" and "processing" into separate transactions.

```python
from datetime import timedelta

InboxKey = tuple[str, str, str]  # (provider, endpoint, event_id)

CLAIM_INBOX_SQL = """
UPDATE webhook_inbox AS w
   SET attempts     = w.attempts + 1,
       available_at = now() + %(lease)s       -- lease: hide it from other workers while it's being processed
  FROM (SELECT provider, endpoint, event_id
          FROM webhook_inbox
         WHERE processed_at IS NULL
           AND available_at <= now()
           AND attempts < %(max_attempts)s
         ORDER BY available_at
         LIMIT 1
           FOR UPDATE SKIP LOCKED) AS picked
 WHERE (w.provider, w.endpoint, w.event_id) = (picked.provider, picked.endpoint, picked.event_id)
RETURNING w.provider, w.endpoint, w.event_id
"""


def claim_next(
    conn: psycopg.Connection, *, max_attempts: int = 10, lease: timedelta = timedelta(minutes=2)
) -> InboxKey | None:
    """Pick one unprocessed receipt and commit the attempt count first (counted before processing,
    so it survives even if the whole process dies mid-processing, and the row stops at the limit)."""
    with conn.transaction():
        row = conn.execute(CLAIM_INBOX_SQL, {"lease": lease, "max_attempts": max_attempts}).fetchone()
    return (row[0], row[1], row[2]) if row else None


def process_claimed(conn: psycopg.Connection, key: InboxKey, *, timeout: str = "30s") -> bool:
    """Process one claimed receipt. Lock the row, re-check processed_at, and commit the
    business update and processed_at in the same transaction.

    Even if the lease expires and another worker takes the same row, the lock and the re-check
    mean the business update happens once. (No fencing token is needed because the business effect
    lives in the same DB transaction as the "already processed?" check. That doesn't hold for work
    that calls external APIs — use the outbox dispatcher's lease + fencing there.)
    """
    try:
        with conn.transaction():
            # Cap the time for one item; keep it well below the lease
            conn.execute("SELECT set_config('statement_timeout', %s, true)", (timeout,))
            row = conn.execute(
                """
                SELECT event_type, payload FROM webhook_inbox
                 WHERE provider = %s AND endpoint = %s AND event_id = %s
                   AND processed_at IS NULL
                   FOR UPDATE
                """,
                key,
            ).fetchone()
            # Not SKIP LOCKED: a competing worker's claim statement briefly locking this row as a candidate
            # would make us skip it until the lease expires. statement_timeout bounds the wait
            if row is None:
                return False  # another worker completed it while we waited
            event_type, payload = row
            apply_event(conn, VerifiedEvent(id=key[2], type=event_type, payload=payload))
            conn.execute(
                """
                UPDATE webhook_inbox SET processed_at = now(), last_error = NULL
                 WHERE provider = %s AND endpoint = %s AND event_id = %s
                """,
                key,
            )
        return True
    except Exception as exc:  # worker boundary: one failure must not stop the whole loop
        # The business update has rolled back. Record the failure and reschedule after a backoff instead of the rest of the lease
        # (the attempt was already counted in claim_next)
        with conn.transaction():
            conn.execute(
                """
                UPDATE webhook_inbox
                   SET last_error   = %s,
                       available_at = now() + make_interval(secs => least(300, 2 ^ attempts) * random())
                 WHERE provider = %s AND endpoint = %s AND event_id = %s
                   AND processed_at IS NULL   -- never touch a row another worker has completed
                """,
                (type(exc).__name__, *key),
            )
        return False


def process_next(conn: psycopg.Connection, *, max_attempts: int = 10) -> bool:
    """Claim one receipt and process it. False if nothing could be claimed (the worker waits briefly and retries)."""
    key = claim_next(conn, max_attempts=max_attempts)
    if key is None:
        return False
    process_claimed(conn, key)
    return True
```

The split exists **to count attempts reliably.** If claiming and processing share one transaction, then when **the process itself** dies mid-processing (killed for OOM, say), the attempt increment rolls back with it, and an event that kills the worker every time is re-taken forever without reaching the limit. `claim_next` increments the count in a short transaction, pushes `available_at` forward by the lease, and commits. If processing dies, the count remains; once the lease expires the row is re-claimed, and at the limit it stops being taken (verified).

On the processing side, `process_claimed` locks the row with `FOR UPDATE` and **re-checks `processed_at`** before the business update. If the lease expires and another worker takes the same row, whichever arrives second waits for the lock, sees the row processed and does nothing. In my tests, letting worker A's lease expire after it claimed and stalled, having worker B process the row, then resuming A still left one ledger row.

There's a reason not to use `SKIP LOCKED` here. A competing worker's claim statement can briefly lock this row as a candidate. With `SKIP LOCKED`, processing would give up at that instant and the row would sit unprocessed until the lease expires. I first wrote it with `SKIP LOCKED` and actually observed exactly that miss in the four-worker concurrency test. The wait is bounded by `statement_timeout` (keep it well below the lease), so it can't wait forever; processing that exceeds it is recorded as a `QueryCanceled` failure (verified).

Unlike the dispatcher, no fencing token is needed, because **the business effect (booking the ledger) is inside the same DB transaction as the "already processed?" check.** The lock and the re-check alone guarantee the business update can't commit twice.

If the handler calls external APIs (sending e-mail, writing back to Stripe), that premise breaks: an external call isn't part of the DB transaction, so the lock and re-check can't prevent it happening twice. In that case, put the external call in an outbox and run it from a separate dispatcher. The lease and fencing token that keep several dispatchers from executing the same row twice are covered in [running an outbox dispatcher safely with multiple workers](/blog/outbox-dispatcher-skip-locked-lease-fencing-token-guide).

Rows whose `attempts` reach the limit stop being picked up. They **need a human to look at them.** Monitor the count (section 9), and once the cause is fixed, reset `attempts` so they're reprocessed.

---

## 7. Comparison with DynamoDB de-duplication

On serverless stacks, de-duplicating with a DynamoDB conditional write (`attribute_not_exists`) is the standard approach. I have run both DynamoDB conditional-write de-duplication for webhooks and a design that puts the de-duplication marker in the same transaction as the business update, on separate production projects. The question isn't which is correct; it's **where your business data lives.**

| Aspect | DynamoDB de-dup store (business data elsewhere) | DynamoDB (business data also in DynamoDB) | PostgreSQL transactional inbox |
| --- | --- | --- | --- |
| Atomicity | **None.** Marker and business update are two writes to two stores; Case A or Case B remains | `TransactWriteItems` writes the marker and business update atomically (up to 100 actions) | Atomic in one transaction |
| Concurrent duplicates | Only one passes the conditional write | Same | Serialized by waiting on the unique index; if the first fails, the second takes over |
| Retention | TTL deletes automatically, but officially "typically within a few days" of expiry (not immediately) | Same | You write the cleanup job |
| Failure isolation | A de-dup store outage stops all receiving (one more dependency) | One store | Same failure domain as the business DB (if the DB is down the business is down anyway; no extra dependency) |
| Load | No load on the business DB | — | One extra write to the business DB per receipt |
| Throughput | High; scales horizontally per key | High | Limited by the business DB's write capacity |
| Operations | More IAM, tables and TTL to manage | — | Rides on existing DB operations |

My decision rules:

- **Business data in PostgreSQL**: use a transactional inbox. Putting only de-duplication in DynamoDB gives up atomicity for little in return.
- **Business data also in DynamoDB**: use `TransactWriteItems` to put the marker's conditional Put and the business update in one transaction. You get atomicity equivalent to PostgreSQL's.
- **The de-dup store and the business data must be separate**: decide which failure to accept. To avoid losing events, do "business update → marker" (the Case B side) and absorb double processing with unique constraints on business keys and idempotent handlers.

If you use TTL, **make retention longer than the provider's retry window.** Stripe retries for up to three days, so with a 24-hour TTL a retry arriving on day three isn't recognised as a duplicate and gets processed. TTL only sets a lower bound on deletion, and expired items can linger for days; if your reads filter out expired items, that filter is part of the retention design too.

---

## 8. Writing back to the provider: idempotency keys and their lifetime

A handler sometimes calls back into the Stripe API (for example, updating a customer's metadata after `invoice.paid`). That write-back needs idempotency too.

The Stripe API supports idempotent retries with the `Idempotency-Key` header. According to the documentation, subsequent requests with the same key return **the result of the first request, whether it succeeded or failed** (including `500` errors). Keys can be up to 255 characters, and Stripe advises against putting personal data such as e-mail addresses in them.

Derive the key **deterministically from the event ID and the kind of operation**:

```python
idempotency_key = f"webhook:{event.id}:update-customer-metadata"
```

A random key per request makes every retry a different request, which isn't idempotency at all.

There's an easily overlooked trap here. Stripe's documentation explains that **keys can be removed after they are at least 24 hours old, and a key reused after removal creates a new request.** Webhooks, meanwhile, are retried for up to three days. So **a write-back made with the same key during a day-three retry may be executed by Stripe as a brand-new request.**

The remedy is to **record the result of the write-back (the ID of the created object, for instance) in your own database, and skip the call if a record exists.** An idempotency key protects "retries within hours"; it doesn't guarantee "once, however many days later". Designing external side effects to converge to one over long periods is covered in [What exactly-once actually guarantees](/blog/exactly-once-at-least-once-idempotency-effectively-once-guide).

---

## 9. Operations: monitoring and retention

```sql
-- Unprocessed receipts and how long the oldest has waited
SELECT count(*)                        AS unprocessed,
       max(now() - received_at)        AS oldest_age,
       count(*) FILTER (WHERE attempts >= 10) AS needs_review
  FROM webhook_inbox
 WHERE processed_at IS NULL;

-- Receipts in the last hour by type (count DUPLICATE responses in the application)
SELECT event_type, count(*)
  FROM webhook_inbox
 WHERE received_at > now() - interval '1 hour'
 GROUP BY event_type
 ORDER BY 2 DESC;
```

What to watch:

- **Age of the oldest unprocessed receipt**: the first sign that async processing has stalled.
- **`needs_review` (rows that hit the attempt limit)**: alert on anything other than zero.
- **Number of 5xx responses**: cross-check against failed deliveries in the Stripe dashboard.
- **Share of `DUPLICATE`**: a spike suggests responses are so slow that Stripe treats them as timeouts.

Keep **retention** well beyond the provider's retry window (three days for Stripe). I usually delete processed rows after about 30 days, which also covers incident investigation.

```sql
DELETE FROM webhook_inbox
 WHERE processed_at IS NOT NULL
   AND processed_at < now() - interval '30 days';
```

Even so, relying on webhooks alone misses events Stripe never sent and events that arrived while you were down for longer than the retry window. Periodically compare important state with the Stripe API and repair it; that design is covered in [Designing reconciliation](/blog/reconciliation-architecture-drift-detection-repair-guide).

---

## 10. Tests

To verify PostgreSQL's locking and unique-index behaviour you need **tests that open several connections to a real PostgreSQL.** These are the cases I ran against a PostgreSQL 18 container and confirmed passing:

| Test | Expected result |
| --- | --- |
| Same event twice | 1st `PROCESSED`, 2nd `DUPLICATE`, one ledger row |
| Same event from 8 connections at once | 1 `PROCESSED`, 7 `DUPLICATE`, one ledger row |
| 1st delivery fails mid-processing while the 2nd waits | 2nd returns `PROCESSED` after the 1st rolls back; one ledger row |
| Business update fails → redelivery | No receipt record remains; redelivery returns `PROCESSED` |
| `invoice.paid` arrives before `invoice.created` | Final state is `paid` |
| Same invoice payment under two different event IDs | One ledger row |
| `invoice.paid` arrives for a voided invoice | Nothing booked, no notification queued |
| Same event received on a different endpoint | Both `PROCESSED` |
| Case A (record first) | Redelivery returns `DUPLICATE`, zero ledger rows (failure reproduced) |
| Case B (record last, crash in between) | Two ledger rows (failure reproduced) |
| Async: worker processing fails | `attempts=1`, `last_error` recorded, processed on retry |
| Async: 50 receipts across 4 workers | Exactly 50 ledger rows, none unprocessed, one attempt each |
| Async: crash right after claiming, three times (limit 3) | Stops being claimed at `attempts=3` (to human review) |
| Async: lease expires, another worker processes, the original worker returns | The original worker does nothing; one ledger row |
| Async: a second worker tries to process a row being processed | Waits for the lock, sees it done, backs off |
| Async: processing exceeds `statement_timeout` | Recorded as `QueryCanceled`; stays unprocessed |

Concurrency tests align thread start with `threading.Barrier` and give each thread **its own connection**.

```python
def test_concurrent_same_event(dsn, conn):
    e = ev("evt_1", "invoice.paid")
    barrier = threading.Barrier(8)

    def worker():
        with psycopg.connect(dsn) as c:
            barrier.wait()
            return receive_sync(c, provider="stripe", endpoint="billing", event=e)

    with ThreadPoolExecutor(8) as ex:
        results = [f.result() for f in [ex.submit(worker) for _ in range(8)]]

    assert results.count(Outcome.PROCESSED) == 1
    assert results.count(Outcome.DUPLICATE) == 7
    assert conn.execute("SELECT count(*) FROM ledger_entries").fetchone()[0] == 1
```

---

## Summary

- Build webhook receivers for the same event arriving several times, concurrently, and out of order (Stripe documents all three).
- Committing the receipt record first loses events (at-most-once); committing it afterwards double-processes (at-least-once). Both come from committing two writes separately.
- A transactional inbox puts the receipt record and the business update in one transaction. Concurrent duplicates are serialized by the unique index, and on failure both disappear so the retry reprocesses.
- Event-ID de-duplication isn't enough. Write state transitions, put unique constraints on business keys, and don't order by `created`.
- Business data in PostgreSQL → inbox; in DynamoDB → `TransactWriteItems`. If the stores must differ, choose explicitly which failure to accept.
- Idempotency keys expire. For write-backs that must happen once even on a retry days later, record the result in your own database.

The full Stripe integration (Checkout, the subscription state machine, testing) is in [Implementing Stripe webhooks and idempotency to production quality](/blog/stripe-payments-production-guide-webhooks-idempotency-subscriptions), and consuming at-least-once deliveries idempotently with SQS and Lambda is in [Idempotent async processing with SQS + Lambda + EventBridge](/blog/aws-sqs-lambda-eventbridge-idempotent-async-processing-guide).
